diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceSpringIntegrationAutoConfiguration.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceSpringIntegrationAutoConfiguration.java index 4b1694db8..90de3a9e5 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceSpringIntegrationAutoConfiguration.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceSpringIntegrationAutoConfiguration.java @@ -67,10 +67,11 @@ class TraceSpringIntegrationAutoConfiguration { @Bean TracingChannelInterceptor traceChannelInterceptor(Tracing tracing, + SleuthMessagingProperties properties, Propagation.Setter traceMessagePropagationSetter, Propagation.Getter traceMessagePropagationGetter) { - return new TracingChannelInterceptor(tracing, traceMessagePropagationSetter, - traceMessagePropagationGetter); + return new TracingChannelInterceptor(tracing, properties, + traceMessagePropagationSetter, traceMessagePropagationGetter); } } diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceWebSocketAutoConfiguration.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceWebSocketAutoConfiguration.java index d98959cdf..90be737b3 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceWebSocketAutoConfiguration.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceWebSocketAutoConfiguration.java @@ -47,6 +47,9 @@ class TraceWebSocketAutoConfiguration extends AbstractWebSocketMessageBrokerConf @Autowired Tracing tracing; + @Autowired + SleuthMessagingProperties properties; + @Override public void registerStompEndpoints(StompEndpointRegistry registry) { // The user must register their own endpoints @@ -54,18 +57,20 @@ class TraceWebSocketAutoConfiguration extends AbstractWebSocketMessageBrokerConf @Override public void configureMessageBroker(MessageBrokerRegistry registry) { - registry.configureBrokerChannel() - .setInterceptors(TracingChannelInterceptor.create(this.tracing)); + registry.configureBrokerChannel().setInterceptors( + TracingChannelInterceptor.create(this.tracing, this.properties)); } @Override public void configureClientOutboundChannel(ChannelRegistration registration) { - registration.setInterceptors(TracingChannelInterceptor.create(this.tracing)); + registration.setInterceptors( + TracingChannelInterceptor.create(this.tracing, this.properties)); } @Override public void configureClientInboundChannel(ChannelRegistration registration) { - registration.setInterceptors(TracingChannelInterceptor.create(this.tracing)); + registration.setInterceptors( + TracingChannelInterceptor.create(this.tracing, this.properties)); } } diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TracingChannelInterceptor.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TracingChannelInterceptor.java index e630531b0..9a47743f1 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TracingChannelInterceptor.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TracingChannelInterceptor.java @@ -98,6 +98,8 @@ final class TracingChannelInterceptor extends ChannelInterceptorAdapter final TraceContext.Extractor extractor; + final SleuthMessagingProperties properties; + final boolean integrationObjectSupportPresent; private final boolean hasDirectChannelClass; @@ -106,15 +108,16 @@ final class TracingChannelInterceptor extends ChannelInterceptorAdapter private final Class directWithAttributesChannelClass; @Autowired - TracingChannelInterceptor(Tracing tracing) { - this(tracing, MessageHeaderPropagation.INSTANCE, + TracingChannelInterceptor(Tracing tracing, SleuthMessagingProperties properties) { + this(tracing, properties, MessageHeaderPropagation.INSTANCE, MessageHeaderPropagation.INSTANCE); } - TracingChannelInterceptor(Tracing tracing, + TracingChannelInterceptor(Tracing tracing, SleuthMessagingProperties properties, Propagation.Setter setter, Propagation.Getter getter) { this.tracing = tracing; + this.properties = properties; this.tracer = tracing.tracer(); this.threadLocalSpan = ThreadLocalSpan.create(this.tracer); this.injector = tracing.propagation().injector(setter); @@ -128,8 +131,9 @@ final class TracingChannelInterceptor extends ChannelInterceptorAdapter ? ClassUtils.resolveClassName(STREAM_DIRECT_CHANNEL, null) : null; } - public static TracingChannelInterceptor create(Tracing tracing) { - return new TracingChannelInterceptor(tracing); + public static TracingChannelInterceptor create(Tracing tracing, + SleuthMessagingProperties properties) { + return new TracingChannelInterceptor(tracing, properties); } /** @@ -186,10 +190,10 @@ final class TracingChannelInterceptor extends ChannelInterceptorAdapter private String toRemoteServiceName(MessageHeaderAccessor headers) { for (String key : headers.getMessageHeaders().keySet()) { if (key.startsWith("kafka_")) { - return "kafka"; + return this.properties.getMessaging().getKafka().getRemoteServiceName(); } else if (key.startsWith("amqp_")) { - return "rabbitmq"; + return this.properties.getMessaging().getRabbit().getRemoteServiceName(); } } return REMOTE_SERVICE_NAME; diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TracingChannelInterceptorTest.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TracingChannelInterceptorTest.java index a7c1364fd..c898b032b 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TracingChannelInterceptorTest.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/messaging/TracingChannelInterceptorTest.java @@ -66,7 +66,8 @@ public class TracingChannelInterceptorTest { B3Propagation.newFactoryBuilder().injectFormat(SINGLE).build()) .addSpanHandler(this.spans).build(); - ChannelInterceptor interceptor = TracingChannelInterceptor.create(tracing); + ChannelInterceptor interceptor = TracingChannelInterceptor.create(tracing, + new SleuthMessagingProperties()); QueueChannel channel = new QueueChannel();