diff --git a/spring-cloud-sleuth-core/pom.xml b/spring-cloud-sleuth-core/pom.xml
index 9f0b269f4..dd5daa76e 100644
--- a/spring-cloud-sleuth-core/pom.xml
+++ b/spring-cloud-sleuth-core/pom.xml
@@ -85,6 +85,11 @@
org.springframework.cloud
spring-cloud-commons
+
+ org.springframework.cloud
+ spring-cloud-stream
+ true
+
org.springframework.cloud
spring-cloud-aws-messaging
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 3f95197b0..2eecfa9ea 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
@@ -16,6 +16,9 @@
package org.springframework.cloud.sleuth.instrument.messaging;
+import java.util.Iterator;
+import java.util.Map;
+
import brave.Span;
import brave.SpanCustomizer;
import brave.Tracer;
@@ -28,8 +31,13 @@ import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.springframework.aop.support.AopUtils;
+import org.springframework.beans.BeansException;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.cloud.sleuth.util.SpanNameUtil;
+import org.springframework.cloud.stream.binder.BinderType;
+import org.springframework.cloud.stream.binder.BinderTypeRegistry;
+import org.springframework.context.ApplicationContext;
+import org.springframework.context.ApplicationContextAware;
import org.springframework.integration.channel.AbstractMessageChannel;
import org.springframework.integration.channel.DirectChannel;
import org.springframework.integration.context.IntegrationObjectSupport;
@@ -63,7 +71,7 @@ import org.springframework.util.ClassUtils;
*/
@Deprecated
public final class TracingChannelInterceptor extends ChannelInterceptorAdapter
- implements ExecutorChannelInterceptor {
+ implements ExecutorChannelInterceptor, ApplicationContextAware {
/**
* Name of the class in Spring Cloud Stream that is a direct channel.
@@ -109,9 +117,13 @@ public final class TracingChannelInterceptor extends ChannelInterceptorAdapter
private final boolean hasDirectChannelClass;
+ private final boolean hasBinderTypeRegistry;
+
// special case of a Stream
private final Class> directWithAttributesChannelClass;
+ private ApplicationContext applicationContext;
+
@Autowired
TracingChannelInterceptor(Tracing tracing, SleuthMessagingProperties properties) {
this(tracing, properties, MessageHeaderPropagation.INSTANCE,
@@ -131,6 +143,8 @@ public final class TracingChannelInterceptor extends ChannelInterceptorAdapter
"org.springframework.integration.context.IntegrationObjectSupport", null);
this.hasDirectChannelClass = ClassUtils
.isPresent("org.springframework.integration.channel.DirectChannel", null);
+ this.hasBinderTypeRegistry = ClassUtils.isPresent(
+ "org.springframework.cloud.stream.binder.BinderTypeRegistry", null);
this.directWithAttributesChannelClass = ClassUtils
.isPresent(STREAM_DIRECT_CHANNEL, null)
? ClassUtils.resolveClassName(STREAM_DIRECT_CHANNEL, null) : null;
@@ -204,6 +218,23 @@ public final class TracingChannelInterceptor extends ChannelInterceptorAdapter
return this.properties.getMessaging().getRabbit().getRemoteServiceName();
}
}
+ if (this.hasBinderTypeRegistry && this.applicationContext != null) {
+ BinderTypeRegistry typeRegistry = this.applicationContext
+ .getBean(BinderTypeRegistry.class);
+ Iterator> iterator = typeRegistry.getAll()
+ .entrySet().iterator();
+ if (iterator.hasNext()) {
+ String binderName = iterator.next().getKey();
+ if (binderName.equals("kafka")) {
+ return this.properties.getMessaging().getKafka()
+ .getRemoteServiceName();
+ }
+ else if (binderName.equals("rabbit")) {
+ return this.properties.getMessaging().getRabbit()
+ .getRemoteServiceName();
+ }
+ }
+ }
return REMOTE_SERVICE_NAME;
}
@@ -435,4 +466,10 @@ public final class TracingChannelInterceptor extends ChannelInterceptorAdapter
return message == null;
}
+ @Override
+ public void setApplicationContext(ApplicationContext applicationContext)
+ throws BeansException {
+ this.applicationContext = applicationContext;
+ }
+
}
diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/SpanAdjusterTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/SpanAdjusterTests.java
index 282f85dba..9cce9ca84 100644
--- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/SpanAdjusterTests.java
+++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/SpanAdjusterTests.java
@@ -29,6 +29,7 @@ import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
import org.springframework.boot.autoconfigure.integration.IntegrationAutoConfiguration;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.cloud.sleuth.util.ArrayListSpanReporter;
+import org.springframework.cloud.stream.function.FunctionConfiguration;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.test.context.junit4.SpringRunner;
@@ -58,7 +59,8 @@ public class SpanAdjusterTests {
}
@Configuration
- @EnableAutoConfiguration(exclude = IntegrationAutoConfiguration.class)
+ @EnableAutoConfiguration(
+ exclude = { IntegrationAutoConfiguration.class, FunctionConfiguration.class })
static class SpanAdjusterAspectTestsConfig {
@Bean
diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/SpanHandlerTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/SpanHandlerTests.java
index 7818dd3e3..51ecaeece 100644
--- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/SpanHandlerTests.java
+++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/SpanHandlerTests.java
@@ -31,6 +31,7 @@ import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
import org.springframework.boot.autoconfigure.integration.IntegrationAutoConfiguration;
import org.springframework.boot.test.context.SpringBootTest;
+import org.springframework.cloud.stream.function.FunctionConfiguration;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.test.context.junit4.SpringRunner;
@@ -62,7 +63,8 @@ public class SpanHandlerTests {
}
@Configuration
- @EnableAutoConfiguration(exclude = IntegrationAutoConfiguration.class)
+ @EnableAutoConfiguration(
+ exclude = { IntegrationAutoConfiguration.class, FunctionConfiguration.class })
static class SpanHandlerAspectTestsConfig {
@Bean