diff --git a/spring-cloud-sleuth-core/pom.xml b/spring-cloud-sleuth-core/pom.xml
index 33e36f35d..73218cb21 100644
--- a/spring-cloud-sleuth-core/pom.xml
+++ b/spring-cloud-sleuth-core/pom.xml
@@ -209,6 +209,12 @@
spring-security-oauth2-autoconfigure
true
+
+
+ io.zipkin.java
+ zipkin
+ true
+
org.springframework.boot
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 e454cefa5..ac7aa81c1 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
@@ -40,6 +40,7 @@ import org.springframework.messaging.support.ExecutorChannelInterceptor;
import org.springframework.messaging.support.GenericMessage;
import org.springframework.messaging.support.MessageHeaderAccessor;
import org.springframework.util.ClassUtils;
+import zipkin2.Endpoint;
/**
* This starts and propagates {@link Span.Kind#PRODUCER} span for each message sent (via native
@@ -55,6 +56,7 @@ public final class TracingChannelInterceptor extends ChannelInterceptorAdapter
implements ExecutorChannelInterceptor {
private static final Log log = LogFactory.getLog(TracingChannelInterceptor.class);
+ private static final String REMOTE_SERVICE_NAME = "broker";
public static TracingChannelInterceptor create(Tracing tracing) {
return new TracingChannelInterceptor(tracing);
@@ -112,6 +114,7 @@ public final class TracingChannelInterceptor extends ChannelInterceptorAdapter
this.injector.inject(span.context(), headers);
if (!span.isNoop()) {
span.kind(Span.Kind.PRODUCER).name("send").start();
+ span.remoteEndpoint(Endpoint.newBuilder().serviceName(REMOTE_SERVICE_NAME).build());
addTags(message, span, channel);
}
if (log.isDebugEnabled()) {
@@ -160,6 +163,7 @@ public final class TracingChannelInterceptor extends ChannelInterceptorAdapter
this.injector.inject(span.context(), headers);
if (!span.isNoop()) {
span.kind(Span.Kind.CONSUMER).name("receive").start();
+ span.remoteEndpoint(Endpoint.newBuilder().serviceName(REMOTE_SERVICE_NAME).build());
addTags(message, span, channel);
}
if (log.isDebugEnabled()) {
@@ -196,6 +200,7 @@ public final class TracingChannelInterceptor extends ChannelInterceptorAdapter
Span consumerSpan = this.tracer.nextSpan(extracted);
if (!consumerSpan.isNoop()) {
consumerSpan.kind(Span.Kind.CONSUMER).start();
+ consumerSpan.remoteEndpoint(Endpoint.newBuilder().serviceName(REMOTE_SERVICE_NAME).build());
addTags(message, consumerSpan, channel);
consumerSpan.finish();
}