From 88ba09467cbe1fce19fbb515ca69be6a5d9e10e9 Mon Sep 17 00:00:00 2001 From: Marcin Grzejszczak Date: Mon, 7 Jan 2019 16:46:43 +0100 Subject: [PATCH] Ignoring the special Stream DirectChannel which is not a DirectChannel; fixes gh-1155 --- .../messaging/TracingChannelInterceptor.java | 24 +++++++++++++++++-- 1 file changed, 22 insertions(+), 2 deletions(-) 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 31452b698..e8c720a0b 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 @@ -81,6 +81,8 @@ public final class TracingChannelInterceptor extends ChannelInterceptorAdapter */ private static final String REMOTE_SERVICE_NAME = "broker"; + public static final String STREAM_DIRECT_CHANNEL = "org.springframework.cloud.stream.messaging.DirectWithAttributesChannel"; + final Tracing tracing; final Tracer tracer; @@ -95,6 +97,9 @@ public final class TracingChannelInterceptor extends ChannelInterceptorAdapter private final boolean hasDirectChannelClass; + // special case of a Stream + private final Class directWithAttributesChannelClass; + @Autowired TracingChannelInterceptor(Tracing tracing) { this(tracing, MessageHeaderPropagation.INSTANCE, @@ -113,6 +118,9 @@ public final class TracingChannelInterceptor extends ChannelInterceptorAdapter "org.springframework.integration.context.IntegrationObjectSupport", null); this.hasDirectChannelClass = ClassUtils .isPresent("org.springframework.integration.channel.DirectChannel", null); + this.directWithAttributesChannelClass = ClassUtils + .isPresent(STREAM_DIRECT_CHANNEL, null) + ? ClassUtils.resolveClassName(STREAM_DIRECT_CHANNEL, null) : null; } public static TracingChannelInterceptor create(Tracing tracing) { @@ -195,8 +203,20 @@ public final class TracingChannelInterceptor extends ChannelInterceptorAdapter } private boolean isDirectChannel(MessageChannel channel) { - return this.hasDirectChannelClass - && DirectChannel.class.isAssignableFrom(AopUtils.getTargetClass(channel)); + Class targetClass = AopUtils.getTargetClass(channel); + boolean directChannel = this.hasDirectChannelClass + && DirectChannel.class.isAssignableFrom(targetClass); + if (!directChannel) { + return false; + } + if (this.directWithAttributesChannelClass == null) { + return true; + } + return !isStreamSpecialDirectChannel(targetClass); + } + + private boolean isStreamSpecialDirectChannel(Class targetClass) { + return this.directWithAttributesChannelClass.isAssignableFrom(targetClass); } @Override