From 6c2c87ab3c4679be437b998bb5d1844b8a62f244 Mon Sep 17 00:00:00 2001 From: Marcin Grzejszczak Date: Tue, 1 Mar 2016 13:45:58 +0100 Subject: [PATCH 1/3] [#191] Added local component if a span doesn't have any Zipkin constants fixes #191 --- .../stream/SamplingZipkinSpanIterator.java | 38 +++++++++++++--- .../sleuth/zipkin/ZipkinSpanListener.java | 45 ++++++++++++++----- 2 files changed, 67 insertions(+), 16 deletions(-) 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 7e6ff577a..c128aec54 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,9 @@ */ package org.springframework.cloud.sleuth.zipkin.stream; +import java.util.Arrays; import java.util.Iterator; +import java.util.List; import java.util.NoSuchElementException; import org.apache.commons.logging.Log; @@ -38,6 +40,14 @@ import zipkin.Span.Builder; * @since 1.0.0 */ final class SamplingZipkinSpanIterator implements Iterator { + private static final List ZIPKIN_ANNOTATIONS = Arrays.asList( + Constants.CLIENT_ADDR, Constants.CLIENT_RECV, Constants.CLIENT_SEND, + Constants.CLIENT_RECV_FRAGMENT, Constants.CLIENT_SEND_FRAGMENT, + 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 + ); private static final Log log = org.apache.commons.logging.LogFactory .getLog(SamplingZipkinSpanIterator.class); @@ -110,17 +120,15 @@ 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()) { - String processId = span.getProcessId() != null - ? span.getProcessId().toLowerCase() - : ZipkinMessageListener.UNKNOWN_PROCESS_ID; - zipkinSpan.addBinaryAnnotation( - BinaryAnnotation.create(Constants.LOCAL_COMPONENT, processId, ep)); + addLocalComponentAnnotation(span, zipkinSpan, ep); } else { ZipkinMessageListener.addZipkinAnnotations(zipkinSpan, span, ep); ZipkinMessageListener.addZipkinBinaryAnnotations(zipkinSpan, span, ep); } - + if (!spanContainsAnyZipkinConstant(span)) { + addLocalComponentAnnotation(span, zipkinSpan, ep); + } zipkinSpan.timestamp(span.getBegin() * 1000); zipkinSpan.duration(span.getAccumulatedMillis() * 1000); zipkinSpan.traceId(span.getTraceId()); @@ -138,4 +146,22 @@ final class SamplingZipkinSpanIterator implements Iterator { } return zipkinSpan.build(); } + + private static void addLocalComponentAnnotation(Span span, Builder zipkinSpan, + Endpoint ep) { + String processId = span.getProcessId() != null + ? span.getProcessId().toLowerCase() + : ZipkinMessageListener.UNKNOWN_PROCESS_ID; + zipkinSpan.addBinaryAnnotation( + BinaryAnnotation.create(Constants.LOCAL_COMPONENT, processId, ep)); + } + + private static boolean spanContainsAnyZipkinConstant(Span span) { + for (org.springframework.cloud.sleuth.Log log : span.logs()) { + if (ZIPKIN_ANNOTATIONS.contains(log.getEvent())) { + return true; + } + } + return false; + } } \ No newline at end of file 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 9b7aeef00..52c6740fa 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,6 +17,8 @@ package org.springframework.cloud.sleuth.zipkin; import java.nio.charset.Charset; +import java.util.Arrays; +import java.util.List; import java.util.Map; import org.springframework.cloud.sleuth.Log; @@ -44,6 +46,14 @@ import zipkin.Endpoint; * @since 1.0.0 */ public class ZipkinSpanListener { + private static final List ZIPKIN_ANNOTATIONS = Arrays.asList( + Constants.CLIENT_ADDR, Constants.CLIENT_RECV, Constants.CLIENT_SEND, + Constants.CLIENT_RECV_FRAGMENT, Constants.CLIENT_SEND_FRAGMENT, + 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 + ); private static final org.apache.commons.logging.Log log = org.apache.commons.logging.LogFactory .getLog(ZipkinSpanListener.class); private static final Charset UTF_8 = Charset.forName("UTF-8"); @@ -129,20 +139,14 @@ public class ZipkinSpanListener { // A zipkin span without any annotations cannot be queried, add special "lc" to avoid that. if (span.logs().isEmpty() && span.tags().isEmpty()) { - byte[] processId = span.getProcessId() != null - ? span.getProcessId().toLowerCase().getBytes(UTF_8) - : UNKNOWN_BYTES; - BinaryAnnotation component = new BinaryAnnotation.Builder() - .type(BinaryAnnotation.Type.STRING) - .key("lc") // LOCAL_COMPONENT - .value(processId) - .endpoint(this.localEndpoint).build(); - zipkinSpan.addBinaryAnnotation(component); + addLocalComponentAnnotation(span, zipkinSpan); } else { addZipkinAnnotations(zipkinSpan, span, this.localEndpoint); addZipkinBinaryAnnotations(zipkinSpan, span, this.localEndpoint); } - + if (!spanContainsAnyZipkinConstant(span)) { + addLocalComponentAnnotation(span, zipkinSpan); + } zipkinSpan.timestamp(span.getBegin() * 1000L); zipkinSpan.duration(span.getAccumulatedMillis() * 1000L); zipkinSpan.traceId(span.getTraceId()); @@ -160,6 +164,27 @@ public class ZipkinSpanListener { return zipkinSpan.build(); } + private void addLocalComponentAnnotation(Span span, zipkin.Span.Builder zipkinSpan) { + byte[] processId = span.getProcessId() != null + ? span.getProcessId().toLowerCase().getBytes(UTF_8) + : UNKNOWN_BYTES; + BinaryAnnotation component = new BinaryAnnotation.Builder() + .type(BinaryAnnotation.Type.STRING) + .key("lc") // LOCAL_COMPONENT + .value(processId) + .endpoint(this.localEndpoint).build(); + zipkinSpan.addBinaryAnnotation(component); + } + + private boolean spanContainsAnyZipkinConstant(Span span) { + for (Log log : span.logs()) { + if (ZIPKIN_ANNOTATIONS.contains(log.getEvent())) { + return true; + } + } + return false; + } + /** * Add annotations from the sleuth Span. */ From 754c8ce96d9151102f62387ce96f6712b693a882 Mon Sep 17 00:00:00 2001 From: Marcin Grzejszczak Date: Tue, 1 Mar 2016 16:07:35 +0100 Subject: [PATCH 2/3] [#191] Added ServerAddress annotation * if a span contains a CS or CR (or any other client related annotation) then SA is passed from the Endpoint service name fixes #191 --- .../stream/SamplingZipkinSpanIterator.java | 47 ++++++++++++++--- .../SamplingZipkinSpanIteratorTests.java | 51 +++++++++++++++++++ .../sleuth/zipkin/ZipkinSpanListener.java | 47 ++++++++++++++--- .../zipkin/ZipkinSpanListenerTests.java | 45 +++++++++++++++- 4 files changed, 176 insertions(+), 14 deletions(-) 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 c128aec54..409c947b7 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,6 +15,7 @@ */ package org.springframework.cloud.sleuth.zipkin.stream; +import java.util.ArrayList; import java.util.Arrays; import java.util.Iterator; import java.util.List; @@ -40,14 +41,23 @@ import zipkin.Span.Builder; * @since 1.0.0 */ final class SamplingZipkinSpanIterator implements Iterator { - private static final List ZIPKIN_ANNOTATIONS = Arrays.asList( + 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, - 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 + Constants.CLIENT_RECV_FRAGMENT, Constants.CLIENT_SEND_FRAGMENT ); + 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); @@ -128,6 +138,8 @@ final class SamplingZipkinSpanIterator implements Iterator { } if (!spanContainsAnyZipkinConstant(span)) { addLocalComponentAnnotation(span, zipkinSpan, ep); + } else if (spanContainsAnyZipkinClientConstantAndSaIsNotSet(span)) { + addServerAddressAnnotation(zipkinSpan, ep); } zipkinSpan.timestamp(span.getBegin() * 1000); zipkinSpan.duration(span.getAccumulatedMillis() * 1000); @@ -156,6 +168,16 @@ final class SamplingZipkinSpanIterator implements Iterator { BinaryAnnotation.create(Constants.LOCAL_COMPONENT, processId, ep)); } + private static void addServerAddressAnnotation(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); + } + private static boolean spanContainsAnyZipkinConstant(Span span) { for (org.springframework.cloud.sleuth.Log log : span.logs()) { if (ZIPKIN_ANNOTATIONS.contains(log.getEvent())) { @@ -164,4 +186,17 @@ final class SamplingZipkinSpanIterator implements Iterator { } 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 624c01e9a..766dd3dcd 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,6 +15,7 @@ */ package org.springframework.cloud.sleuth.zipkin.stream; +import java.nio.charset.Charset; import java.util.Arrays; import java.util.Collections; import java.util.Iterator; @@ -25,6 +26,8 @@ import org.junit.Test; import org.springframework.cloud.sleuth.Span; import org.springframework.cloud.sleuth.stream.Host; import org.springframework.cloud.sleuth.stream.Spans; + +import zipkin.Constants; import zipkin.Sampler; import static org.assertj.core.api.Assertions.assertThat; @@ -76,6 +79,54 @@ public class SamplingZipkinSpanIteratorTests { "message:foo", "message:baz"); } + @Test + public void appendsLocalComponentTagIfNoZipkinLogIsPresent() { + Spans spans = new Spans(this.host, Collections.singletonList(span("foo"))); + + Iterator result = new SamplingZipkinSpanIterator( + Sampler.create(1.0f), spans); + + assertThat(result) + .flatExtracting(s -> s.binaryAnnotations) + .extracting(input -> input.key) + .contains(Constants.LOCAL_COMPONENT); + } + + @Test + public void appendsServerAddressTagIfClientLogIsPresent() { + Span span = span("foo"); + span.logEvent(Constants.CLIENT_RECV); + Spans spans = new Spans(this.host, Collections.singletonList(span)); + + Iterator result = new SamplingZipkinSpanIterator( + Sampler.create(1.0f), spans); + + assertThat(result) + .hasSize(1) + .flatExtracting(input1 -> input1.binaryAnnotations) + .filteredOn("key", Constants.SERVER_ADDR) + .extracting(input -> input.value) + .contains("myservice".getBytes(Charset.forName("UTF-8"))); + } + + @Test + public void shouldReuseServerAddressTag() { + Span span = span("foo"); + span.logEvent(Constants.CLIENT_RECV); + span.tag(Constants.SERVER_ADDR, "barservice"); + Spans spans = new Spans(this.host, Collections.singletonList(span)); + + Iterator result = new SamplingZipkinSpanIterator( + Sampler.create(1.0f), spans); + + assertThat(result) + .hasSize(1) + .flatExtracting(input1 -> input1.binaryAnnotations) + .filteredOn("key", Constants.SERVER_ADDR) + .extracting(input -> input.value) + .contains("barservice".getBytes(Charset.forName("UTF-8"))); + } + Span span(String name) { Long id = new Random().nextLong(); return new Span(1, 3, "message:" + name, id, Collections.emptyList(), id, true, true, 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 52c6740fa..b7ce058fe 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,6 +17,7 @@ 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; @@ -46,14 +47,24 @@ import zipkin.Endpoint; * @since 1.0.0 */ public class ZipkinSpanListener { - private static final List ZIPKIN_ANNOTATIONS = Arrays.asList( + 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, - 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 - ); + Constants.CLIENT_RECV_FRAGMENT, Constants.CLIENT_SEND_FRAGMENT + ); + 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); private static final Charset UTF_8 = Charset.forName("UTF-8"); @@ -146,6 +157,8 @@ public class ZipkinSpanListener { } if (!spanContainsAnyZipkinConstant(span)) { addLocalComponentAnnotation(span, zipkinSpan); + } else if (spanContainsAnyZipkinClientConstantAndSaIsNotSet(span)) { + addServerAddressAnnotation(zipkinSpan); } zipkinSpan.timestamp(span.getBegin() * 1000L); zipkinSpan.duration(span.getAccumulatedMillis() * 1000L); @@ -176,6 +189,15 @@ 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 boolean spanContainsAnyZipkinConstant(Span span) { for (Log log : span.logs()) { if (ZIPKIN_ANNOTATIONS.contains(log.getEvent())) { @@ -185,6 +207,17 @@ public class ZipkinSpanListener { 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. */ 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 d61bbce99..4e93a061f 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,10 +16,12 @@ package org.springframework.cloud.sleuth.zipkin; -import javax.annotation.PostConstruct; +import java.nio.charset.Charset; import java.util.ArrayList; import java.util.List; +import javax.annotation.PostConstruct; + import org.junit.Test; import org.junit.runner.RunWith; import org.springframework.beans.factory.annotation.Autowired; @@ -41,6 +43,8 @@ import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Import; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; +import zipkin.Constants; + import static org.assertj.core.api.Assertions.assertThat; import static org.junit.Assert.assertEquals; @@ -135,6 +139,45 @@ public class ZipkinSpanListenerTests { assertEquals(2, this.test.spans.size()); } + @Test + public void appendsLocalComponentTagIfNoZipkinLogIsPresent() { + this.parent.logEvent("hystrix/retry"); + this.parent.stop(); + + zipkin.Span result = this.listener.convert(this.parent); + + assertThat(result.binaryAnnotations) + .extracting(input -> input.key) + .contains(Constants.LOCAL_COMPONENT); + } + + @Test + public void appendsServerAddressTagIfClientLogIsPresent() { + this.parent.logEvent(Constants.CLIENT_SEND); + this.parent.stop(); + + zipkin.Span result = this.listener.convert(this.parent); + + assertThat(result.binaryAnnotations) + .filteredOn("key", Constants.SERVER_ADDR) + .extracting(input -> input.value) + .containsOnly("unknown".getBytes(Charset.forName("UTF-8"))); + } + + @Test + public void shouldReuseServerAddressTag() { + this.parent.logEvent(Constants.CLIENT_SEND); + this.parent.tag(Constants.SERVER_ADDR, "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"))); + } + @Configuration @Import({ ZipkinTestConfiguration.class, ZipkinAutoConfiguration.class, TraceAutoConfiguration.class, PropertyPlaceholderAutoConfiguration.class }) From 17bc239536eb551a3cbdada7c4a53e030919de74 Mon Sep 17 00:00:00 2001 From: Marcin Grzejszczak Date: Wed, 2 Mar 2016 12:24:33 +0100 Subject: [PATCH 3/3] [#191] Changes following the review --- .../springframework/cloud/sleuth/Span.java | 4 + .../stream/SamplingZipkinSpanIterator.java | 84 ++++++++---------- .../SamplingZipkinSpanIteratorTests.java | 15 ++-- .../sleuth/zipkin/ZipkinSpanListener.java | 85 ++++++++----------- .../zipkin/ZipkinSpanListenerTests.java | 11 ++- 5 files changed, 85 insertions(+), 114 deletions(-) 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