diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/Span.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/Span.java index 445e26b12..6abd557ae 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/Span.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/Span.java @@ -64,6 +64,10 @@ public class Span { public static final String SPAN_ID_NAME = "X-Span-Id"; public static final String SPAN_EXPORT_NAME = "X-Span-Export"; public static final String SPAN_LOCAL_COMPONENT_TAG_NAME = "lc"; + /** + * As in Open Tracing + */ + public static final String SPAN_PEER_SERVICE_TAG_NAME = "peer.service"; private final long begin; private long end = 0; diff --git a/spring-cloud-sleuth-zipkin-stream/src/main/java/org/springframework/cloud/sleuth/zipkin/stream/SamplingZipkinSpanIterator.java b/spring-cloud-sleuth-zipkin-stream/src/main/java/org/springframework/cloud/sleuth/zipkin/stream/SamplingZipkinSpanIterator.java index 409c947b7..99e5c7e3c 100644 --- a/spring-cloud-sleuth-zipkin-stream/src/main/java/org/springframework/cloud/sleuth/zipkin/stream/SamplingZipkinSpanIterator.java +++ b/spring-cloud-sleuth-zipkin-stream/src/main/java/org/springframework/cloud/sleuth/zipkin/stream/SamplingZipkinSpanIterator.java @@ -15,7 +15,6 @@ */ package org.springframework.cloud.sleuth.zipkin.stream; -import java.util.ArrayList; import java.util.Arrays; import java.util.Iterator; import java.util.List; @@ -27,6 +26,7 @@ import org.springframework.cloud.sleuth.stream.Host; import org.springframework.cloud.sleuth.stream.SleuthSink; import org.springframework.cloud.sleuth.stream.Spans; import org.springframework.util.StringUtils; + import zipkin.BinaryAnnotation; import zipkin.Constants; import zipkin.Endpoint; @@ -41,23 +41,9 @@ import zipkin.Span.Builder; * @since 1.0.0 */ final class SamplingZipkinSpanIterator implements Iterator { - private static final List ZIPKIN_CLIENT_ANNOTATIONS = Arrays.asList( - Constants.CLIENT_ADDR, Constants.CLIENT_RECV, Constants.CLIENT_SEND, - Constants.CLIENT_RECV_FRAGMENT, Constants.CLIENT_SEND_FRAGMENT + private static final List ZIPKIN_START_EVENTS = Arrays.asList( + Constants.CLIENT_RECV, Constants.SERVER_RECV ); - private static final List ZIPKIN_ANNOTATIONS = zipkinAnnotations(); - - private static List zipkinAnnotations() { - List annotations = new ArrayList<>(); - annotations.addAll(Arrays.asList( - Constants.SERVER_ADDR, Constants.SERVER_RECV, Constants.SERVER_SEND, - Constants.SERVER_RECV_FRAGMENT, Constants.SERVER_SEND_FRAGMENT, - Constants.LOCAL_COMPONENT, - Constants.WIRE_RECV, Constants.WIRE_SEND - )); - annotations.addAll(ZIPKIN_CLIENT_ANNOTATIONS); - return annotations; - } private static final Log log = org.apache.commons.logging.LogFactory .getLog(SamplingZipkinSpanIterator.class); @@ -119,6 +105,10 @@ final class SamplingZipkinSpanIterator implements Iterator { *
  • Create timeline annotations based on data from Span object. *
  • Create binary annotations based on data from Span object. * + * + * When logging {@link Constants#CLIENT_SEND}, instrumentation should also log the {@link Constants#SERVER_ADDR} + * Check Zipkin code + * for more information */ // VisibleForTesting static zipkin.Span convert(Span span, Host host) { @@ -129,17 +119,13 @@ final class SamplingZipkinSpanIterator implements Iterator { // A zipkin span without any annotations cannot be queried, add special "lc" to // avoid that. - if (span.logs().isEmpty() && span.tags().isEmpty()) { - addLocalComponentAnnotation(span, zipkinSpan, ep); + if (notClientOrServer(span)) { + ensureLocalComponent(span, zipkinSpan, ep); } - else { - ZipkinMessageListener.addZipkinAnnotations(zipkinSpan, span, ep); - ZipkinMessageListener.addZipkinBinaryAnnotations(zipkinSpan, span, ep); - } - if (!spanContainsAnyZipkinConstant(span)) { - addLocalComponentAnnotation(span, zipkinSpan, ep); - } else if (spanContainsAnyZipkinClientConstantAndSaIsNotSet(span)) { - addServerAddressAnnotation(zipkinSpan, ep); + ZipkinMessageListener.addZipkinAnnotations(zipkinSpan, span, ep); + ZipkinMessageListener.addZipkinBinaryAnnotations(zipkinSpan, span, ep); + if (hasClientSend(span)) { + ensureServerAddr(span, zipkinSpan, ep); } zipkinSpan.timestamp(span.getBegin() * 1000); zipkinSpan.duration(span.getAccumulatedMillis() * 1000); @@ -159,8 +145,11 @@ final class SamplingZipkinSpanIterator implements Iterator { return zipkinSpan.build(); } - private static void addLocalComponentAnnotation(Span span, Builder zipkinSpan, + private static void ensureLocalComponent(Span span, Builder zipkinSpan, Endpoint ep) { + if (span.tags().containsKey(Constants.LOCAL_COMPONENT)) { + return; + } String processId = span.getProcessId() != null ? span.getProcessId().toLowerCase() : ZipkinMessageListener.UNKNOWN_PROCESS_ID; @@ -168,35 +157,30 @@ final class SamplingZipkinSpanIterator implements Iterator { BinaryAnnotation.create(Constants.LOCAL_COMPONENT, processId, ep)); } - private static void addServerAddressAnnotation(zipkin.Span.Builder zipkinSpan, + private static void ensureServerAddr(Span span, zipkin.Span.Builder zipkinSpan, Endpoint ep) { - BinaryAnnotation component = new BinaryAnnotation.Builder() - .type(BinaryAnnotation.Type.STRING) - .key(Constants.SERVER_ADDR) - .value(ep.serviceName) - .endpoint(ep).build(); - zipkinSpan.addBinaryAnnotation(component); + String serviceName = span.tags().containsKey(Span.SPAN_PEER_SERVICE_TAG_NAME) ? + span.tags().get(Span.SPAN_PEER_SERVICE_TAG_NAME) : ep.serviceName; + zipkinSpan.addBinaryAnnotation(BinaryAnnotation.address(Constants.SERVER_ADDR, + Endpoint.create(serviceName, ep.ipv4, ep.port))); } - private static boolean spanContainsAnyZipkinConstant(Span span) { + private static boolean notClientOrServer(Span span) { for (org.springframework.cloud.sleuth.Log log : span.logs()) { - if (ZIPKIN_ANNOTATIONS.contains(log.getEvent())) { - return true; + if (ZIPKIN_START_EVENTS.contains(log.getEvent())) { + return false; + } + } + return true; + } + + private static boolean hasClientSend(Span span) { + for (org.springframework.cloud.sleuth.Log log : span.logs()) { + if (Constants.CLIENT_SEND.equals(log.getEvent())) { + return !span.tags().containsKey(Constants.SERVER_ADDR); } } return false; } - private static boolean spanContainsAnyZipkinClientConstantAndSaIsNotSet(Span span) { - boolean containsAnyZipkinClientConstant = false; - for (org.springframework.cloud.sleuth.Log log : span.logs()) { - if (ZIPKIN_CLIENT_ANNOTATIONS.contains(log.getEvent())) { - containsAnyZipkinClientConstant = true; - break; - } - } - return containsAnyZipkinClientConstant && !span.tags().containsKey(Constants.SERVER_ADDR); - } - - } \ No newline at end of file diff --git a/spring-cloud-sleuth-zipkin-stream/src/test/java/org/springframework/cloud/sleuth/zipkin/stream/SamplingZipkinSpanIteratorTests.java b/spring-cloud-sleuth-zipkin-stream/src/test/java/org/springframework/cloud/sleuth/zipkin/stream/SamplingZipkinSpanIteratorTests.java index 766dd3dcd..0976fabdc 100644 --- a/spring-cloud-sleuth-zipkin-stream/src/test/java/org/springframework/cloud/sleuth/zipkin/stream/SamplingZipkinSpanIteratorTests.java +++ b/spring-cloud-sleuth-zipkin-stream/src/test/java/org/springframework/cloud/sleuth/zipkin/stream/SamplingZipkinSpanIteratorTests.java @@ -15,7 +15,6 @@ */ package org.springframework.cloud.sleuth.zipkin.stream; -import java.nio.charset.Charset; import java.util.Arrays; import java.util.Collections; import java.util.Iterator; @@ -95,7 +94,7 @@ public class SamplingZipkinSpanIteratorTests { @Test public void appendsServerAddressTagIfClientLogIsPresent() { Span span = span("foo"); - span.logEvent(Constants.CLIENT_RECV); + span.logEvent(Constants.CLIENT_SEND); Spans spans = new Spans(this.host, Collections.singletonList(span)); Iterator result = new SamplingZipkinSpanIterator( @@ -105,15 +104,15 @@ public class SamplingZipkinSpanIteratorTests { .hasSize(1) .flatExtracting(input1 -> input1.binaryAnnotations) .filteredOn("key", Constants.SERVER_ADDR) - .extracting(input -> input.value) - .contains("myservice".getBytes(Charset.forName("UTF-8"))); + .extracting(input -> input.endpoint.serviceName) + .contains("myservice"); } @Test public void shouldReuseServerAddressTag() { Span span = span("foo"); - span.logEvent(Constants.CLIENT_RECV); - span.tag(Constants.SERVER_ADDR, "barservice"); + span.logEvent(Constants.CLIENT_SEND); + span.tag(Span.SPAN_PEER_SERVICE_TAG_NAME, "barservice"); Spans spans = new Spans(this.host, Collections.singletonList(span)); Iterator result = new SamplingZipkinSpanIterator( @@ -123,8 +122,8 @@ public class SamplingZipkinSpanIteratorTests { .hasSize(1) .flatExtracting(input1 -> input1.binaryAnnotations) .filteredOn("key", Constants.SERVER_ADDR) - .extracting(input -> input.value) - .contains("barservice".getBytes(Charset.forName("UTF-8"))); + .extracting(input -> input.endpoint.serviceName) + .contains("barservice"); } Span span(String name) { diff --git a/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin/ZipkinSpanListener.java b/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin/ZipkinSpanListener.java index b7ce058fe..dd5e610bd 100644 --- a/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin/ZipkinSpanListener.java +++ b/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin/ZipkinSpanListener.java @@ -17,7 +17,6 @@ package org.springframework.cloud.sleuth.zipkin; import java.nio.charset.Charset; -import java.util.ArrayList; import java.util.Arrays; import java.util.List; import java.util.Map; @@ -47,23 +46,9 @@ import zipkin.Endpoint; * @since 1.0.0 */ public class ZipkinSpanListener { - private static final List ZIPKIN_CLIENT_ANNOTATIONS = Arrays.asList( - Constants.CLIENT_ADDR, Constants.CLIENT_RECV, Constants.CLIENT_SEND, - Constants.CLIENT_RECV_FRAGMENT, Constants.CLIENT_SEND_FRAGMENT + private static final List ZIPKIN_START_EVENTS = Arrays.asList( + Constants.CLIENT_RECV, Constants.SERVER_RECV ); - private static final List ZIPKIN_ANNOTATIONS = zipkinAnnotations(); - - private static List zipkinAnnotations() { - List annotations = new ArrayList<>(); - annotations.addAll(Arrays.asList( - Constants.SERVER_ADDR, Constants.SERVER_RECV, Constants.SERVER_SEND, - Constants.SERVER_RECV_FRAGMENT, Constants.SERVER_SEND_FRAGMENT, - Constants.LOCAL_COMPONENT, - Constants.WIRE_RECV, Constants.WIRE_SEND - )); - annotations.addAll(ZIPKIN_CLIENT_ANNOTATIONS); - return annotations; - } private static final org.apache.commons.logging.Log log = org.apache.commons.logging.LogFactory .getLog(ZipkinSpanListener.class); @@ -143,22 +128,23 @@ public class ZipkinSpanListener { *
  • Create timeline annotations based on data from Span object. *
  • Create binary annotations based on data from Span object. * + * + * When logging {@link Constants#CLIENT_SEND}, instrumentation should also log the {@link Constants#SERVER_ADDR} + * Check Zipkin code + * for more information */ // Visible for testing zipkin.Span convert(Span span) { zipkin.Span.Builder zipkinSpan = new zipkin.Span.Builder(); // A zipkin span without any annotations cannot be queried, add special "lc" to avoid that. - if (span.logs().isEmpty() && span.tags().isEmpty()) { - addLocalComponentAnnotation(span, zipkinSpan); - } else { - addZipkinAnnotations(zipkinSpan, span, this.localEndpoint); - addZipkinBinaryAnnotations(zipkinSpan, span, this.localEndpoint); + if (notClientOrServer(span)) { + ensureLocalComponent(span, zipkinSpan); } - if (!spanContainsAnyZipkinConstant(span)) { - addLocalComponentAnnotation(span, zipkinSpan); - } else if (spanContainsAnyZipkinClientConstantAndSaIsNotSet(span)) { - addServerAddressAnnotation(zipkinSpan); + addZipkinAnnotations(zipkinSpan, span, this.localEndpoint); + addZipkinBinaryAnnotations(zipkinSpan, span, this.localEndpoint); + if (hasClientSend(span)) { + ensureServerAddr(span, zipkinSpan); } zipkinSpan.timestamp(span.getBegin() * 1000L); zipkinSpan.duration(span.getAccumulatedMillis() * 1000L); @@ -177,7 +163,10 @@ public class ZipkinSpanListener { return zipkinSpan.build(); } - private void addLocalComponentAnnotation(Span span, zipkin.Span.Builder zipkinSpan) { + private void ensureLocalComponent(Span span, zipkin.Span.Builder zipkinSpan) { + if (span.tags().containsKey(Constants.LOCAL_COMPONENT)) { + return; + } byte[] processId = span.getProcessId() != null ? span.getProcessId().toLowerCase().getBytes(UTF_8) : UNKNOWN_BYTES; @@ -189,35 +178,31 @@ public class ZipkinSpanListener { zipkinSpan.addBinaryAnnotation(component); } - private void addServerAddressAnnotation(zipkin.Span.Builder zipkinSpan) { - BinaryAnnotation component = new BinaryAnnotation.Builder() - .type(BinaryAnnotation.Type.STRING) - .key(Constants.SERVER_ADDR) - .value(this.localEndpoint.serviceName) - .endpoint(this.localEndpoint).build(); - zipkinSpan.addBinaryAnnotation(component); + private void ensureServerAddr(Span span, zipkin.Span.Builder zipkinSpan) { + String serviceName = span.tags().containsKey(Span.SPAN_PEER_SERVICE_TAG_NAME) ? + span.tags().get(Span.SPAN_PEER_SERVICE_TAG_NAME) : this.localEndpoint.serviceName; + zipkinSpan.addBinaryAnnotation(BinaryAnnotation.address(Constants.SERVER_ADDR, + Endpoint.create(serviceName, this.localEndpoint.ipv4, this.localEndpoint.port))); } - private boolean spanContainsAnyZipkinConstant(Span span) { + private boolean notClientOrServer(Span span) { for (Log log : span.logs()) { - if (ZIPKIN_ANNOTATIONS.contains(log.getEvent())) { - return true; + if (ZIPKIN_START_EVENTS.contains(log.getEvent())) { + return false; + } + } + return true; + } + + private boolean hasClientSend(Span span) { + for (org.springframework.cloud.sleuth.Log log : span.logs()) { + if (Constants.CLIENT_SEND.equals(log.getEvent())) { + return !span.tags().containsKey(Constants.SERVER_ADDR); } } return false; } - private boolean spanContainsAnyZipkinClientConstantAndSaIsNotSet(Span span) { - boolean containsAnyZipkinClientConstant = false; - for (org.springframework.cloud.sleuth.Log log : span.logs()) { - if (ZIPKIN_CLIENT_ANNOTATIONS.contains(log.getEvent())) { - containsAnyZipkinClientConstant = true; - break; - } - } - return containsAnyZipkinClientConstant && !span.tags().containsKey(Constants.SERVER_ADDR); - } - /** * Add annotations from the sleuth Span. */ @@ -236,13 +221,13 @@ public class ZipkinSpanListener { * Adds binary annotation from the sleuth Span */ private void addZipkinBinaryAnnotations(zipkin.Span.Builder zipkinSpan, - Span span, Endpoint endpoint) { + Span span, Endpoint ep) { for (Map.Entry e : span.tags().entrySet()) { BinaryAnnotation binaryAnn = new BinaryAnnotation.Builder() .type(BinaryAnnotation.Type.STRING) .key(e.getKey()) .value(e.getValue().getBytes(UTF_8)) - .endpoint(endpoint).build(); + .endpoint(ep).build(); zipkinSpan.addBinaryAnnotation(binaryAnn); } } diff --git a/spring-cloud-sleuth-zipkin/src/test/java/org/springframework/cloud/sleuth/zipkin/ZipkinSpanListenerTests.java b/spring-cloud-sleuth-zipkin/src/test/java/org/springframework/cloud/sleuth/zipkin/ZipkinSpanListenerTests.java index 4e93a061f..fe3055554 100644 --- a/spring-cloud-sleuth-zipkin/src/test/java/org/springframework/cloud/sleuth/zipkin/ZipkinSpanListenerTests.java +++ b/spring-cloud-sleuth-zipkin/src/test/java/org/springframework/cloud/sleuth/zipkin/ZipkinSpanListenerTests.java @@ -16,7 +16,6 @@ package org.springframework.cloud.sleuth.zipkin; -import java.nio.charset.Charset; import java.util.ArrayList; import java.util.List; @@ -160,22 +159,22 @@ public class ZipkinSpanListenerTests { assertThat(result.binaryAnnotations) .filteredOn("key", Constants.SERVER_ADDR) - .extracting(input -> input.value) - .containsOnly("unknown".getBytes(Charset.forName("UTF-8"))); + .extracting(input -> input.endpoint.serviceName) + .containsOnly("unknown"); } @Test public void shouldReuseServerAddressTag() { this.parent.logEvent(Constants.CLIENT_SEND); - this.parent.tag(Constants.SERVER_ADDR, "fooservice"); + this.parent.tag(Span.SPAN_PEER_SERVICE_TAG_NAME, "fooservice"); this.parent.stop(); zipkin.Span result = this.listener.convert(this.parent); assertThat(result.binaryAnnotations) .filteredOn("key", Constants.SERVER_ADDR) - .extracting(input -> input.value) - .containsOnly("fooservice".getBytes(Charset.forName("UTF-8"))); + .extracting(input -> input.endpoint.serviceName) + .containsOnly("fooservice"); } @Configuration