From 0c5c984e4af48d71ab755354872efdecc5cf9102 Mon Sep 17 00:00:00 2001 From: Dave Syer Date: Tue, 26 Jan 2016 13:53:44 +0000 Subject: [PATCH] Use accumulated millisecs instead of computing span duration --- .../sleuth/stream/StreamSpanListener.java | 14 ++--- .../stream/SamplingZipkinSpanIterator.java | 63 +++++++++++++++++-- .../zipkin/stream/ZipkinMessageListener.java | 54 ++-------------- .../stream/ZipkinMessageListenerTests.java | 8 +-- .../sleuth/zipkin/ZipkinSpanListener.java | 11 ++-- 5 files changed, 79 insertions(+), 71 deletions(-) diff --git a/spring-cloud-sleuth-stream/src/main/java/org/springframework/cloud/sleuth/stream/StreamSpanListener.java b/spring-cloud-sleuth-stream/src/main/java/org/springframework/cloud/sleuth/stream/StreamSpanListener.java index 70d1dd6db..ab972c52f 100644 --- a/spring-cloud-sleuth-stream/src/main/java/org/springframework/cloud/sleuth/stream/StreamSpanListener.java +++ b/spring-cloud-sleuth-stream/src/main/java/org/springframework/cloud/sleuth/stream/StreamSpanListener.java @@ -16,6 +16,12 @@ package org.springframework.cloud.sleuth.stream; +import java.util.ArrayList; +import java.util.Collection; +import java.util.Iterator; +import java.util.List; +import java.util.concurrent.ConcurrentLinkedQueue; + import org.springframework.cloud.sleuth.Span; import org.springframework.cloud.sleuth.event.ClientReceivedEvent; import org.springframework.cloud.sleuth.event.ClientSentEvent; @@ -29,12 +35,6 @@ import org.springframework.core.annotation.Order; import org.springframework.integration.annotation.InboundChannelAdapter; import org.springframework.integration.annotation.MessageEndpoint; -import java.util.ArrayList; -import java.util.Collection; -import java.util.Iterator; -import java.util.List; -import java.util.concurrent.ConcurrentLinkedQueue; - /** * A message source for spans. Also handles RPC flavoured annotations. @@ -112,7 +112,7 @@ public class StreamSpanListener { this.queue.clear(); for (Iterator iterator = result.iterator(); iterator.hasNext();) { Span span = iterator.next(); - if (span.getName() != null && span.getName().equals("message/zipkin")) { + if (span.getName() != null && span.getName().equals("message/" + SleuthSource.OUTPUT)) { iterator.remove(); } } 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 051433b56..a2efae68e 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,15 +15,21 @@ */ package org.springframework.cloud.sleuth.zipkin.stream; -import lombok.extern.apachecommons.CommonsLog; +import java.util.Iterator; +import java.util.NoSuchElementException; + import org.springframework.cloud.sleuth.Span; import org.springframework.cloud.sleuth.stream.Host; import org.springframework.cloud.sleuth.stream.SleuthSink; import org.springframework.cloud.sleuth.stream.Spans; -import zipkin.Sampler; +import org.springframework.util.StringUtils; -import java.util.Iterator; -import java.util.NoSuchElementException; +import lombok.extern.apachecommons.CommonsLog; +import zipkin.BinaryAnnotation; +import zipkin.Constants; +import zipkin.Endpoint; +import zipkin.Sampler; +import zipkin.Span.Builder; /** * This converts sleuth spans to zipkin ones, skipping invalid or unsampled. @@ -70,7 +76,7 @@ final class SamplingZipkinSpanIterator implements Iterator { */ zipkin.Span convertAndSample(Span input, Host host) { if (!input.getName().equals("message/" + SleuthSink.INPUT)) { - zipkin.Span result = ZipkinMessageListener.convert(input, host); + zipkin.Span result = SamplingZipkinSpanIterator.convert(input, host); if (this.sampler.isSampled(result.traceId)) { return result; } @@ -80,4 +86,51 @@ final class SamplingZipkinSpanIterator implements Iterator { } return null; } + + /** + * Converts a given Sleuth span to a Zipkin Span. + *
    + *
  • Set ids, etc + *
  • Create timeline annotations based on data from Span object. + *
  • Create binary annotations based on data from Span object. + *
+ */ + // VisibleForTesting + static zipkin.Span convert(Span span, Host host) { + Builder zipkinSpan = new zipkin.Span.Builder(); + + Endpoint ep = Endpoint.create(host.getServiceName(), host.getIpv4(), + host.getPort().shortValue()); + + // 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)); + } + else { + ZipkinMessageListener.addZipkinAnnotations(zipkinSpan, span, ep); + ZipkinMessageListener.addZipkinBinaryAnnotations(zipkinSpan, span, ep); + } + + zipkinSpan.timestamp(span.getBegin() * 1000); + zipkinSpan.duration(span.getAccumulatedMillis() * 1000); + zipkinSpan.traceId(span.getTraceId()); + if (span.getParents().size() > 0) { + if (span.getParents().size() > 1) { + SamplingZipkinSpanIterator.log + .debug("zipkin doesn't support spans with multiple parents. Omitting " + + "other parents for " + span); + } + zipkinSpan.parentId(span.getParents().get(0)); + } + zipkinSpan.id(span.getSpanId()); + if (StringUtils.hasText(span.getName())) { + zipkinSpan.name(span.getName()); + } + return zipkinSpan.build(); + } } \ No newline at end of file diff --git a/spring-cloud-sleuth-zipkin-stream/src/main/java/org/springframework/cloud/sleuth/zipkin/stream/ZipkinMessageListener.java b/spring-cloud-sleuth-zipkin-stream/src/main/java/org/springframework/cloud/sleuth/zipkin/stream/ZipkinMessageListener.java index 20328cd57..e42c15705 100644 --- a/spring-cloud-sleuth-zipkin-stream/src/main/java/org/springframework/cloud/sleuth/zipkin/stream/ZipkinMessageListener.java +++ b/spring-cloud-sleuth-zipkin-stream/src/main/java/org/springframework/cloud/sleuth/zipkin/stream/ZipkinMessageListener.java @@ -11,7 +11,6 @@ import org.springframework.cloud.Cloud; import org.springframework.cloud.CloudFactory; import org.springframework.cloud.sleuth.Log; import org.springframework.cloud.sleuth.Span; -import org.springframework.cloud.sleuth.stream.Host; import org.springframework.cloud.sleuth.stream.SleuthSink; import org.springframework.cloud.sleuth.stream.Spans; import org.springframework.cloud.sleuth.zipkin.stream.ZipkinMessageListener.NotSleuthStreamClient; @@ -22,7 +21,7 @@ import org.springframework.core.env.MapPropertySource; import org.springframework.core.type.AnnotatedTypeMetadata; import org.springframework.integration.annotation.MessageEndpoint; import org.springframework.integration.annotation.ServiceActivator; -import org.springframework.util.StringUtils; + import zipkin.*; import zipkin.BinaryAnnotation.Type; import zipkin.Span.Builder; @@ -38,7 +37,7 @@ import java.util.Map; @Conditional(NotSleuthStreamClient.class) public class ZipkinMessageListener { - private static final String UNKNOWN_PROCESS_ID = "unknown"; + static final String UNKNOWN_PROCESS_ID = "unknown"; @Autowired SpanStore spanStore; @@ -54,55 +53,10 @@ public class ZipkinMessageListener { } } - /** - * Converts a given Sleuth span to a Zipkin Span. - *
    - *
  • Set ids, etc - *
  • Create timeline annotations based on data from Span object. - *
  • Create binary annotations based on data from Span object. - *
- */ - // VisibleForTesting - static zipkin.Span convert(Span span, Host host) { - Builder zipkinSpan = new zipkin.Span.Builder(); - - Endpoint ep = Endpoint.create(host.getServiceName(), host.getIpv4(), - host.getPort().shortValue()); - - // 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() - : UNKNOWN_PROCESS_ID; - zipkinSpan.addBinaryAnnotation( - BinaryAnnotation.create(Constants.LOCAL_COMPONENT, processId, ep) - ); - } else { - addZipkinAnnotations(zipkinSpan, span, ep); - addZipkinBinaryAnnotations(zipkinSpan, span, ep); - } - - zipkinSpan.timestamp(span.getBegin() * 1000); - zipkinSpan.duration((span.getEnd() - span.getBegin()) * 1000); - zipkinSpan.traceId(span.getTraceId()); - if (span.getParents().size() > 0) { - if (span.getParents().size() > 1) { - log.error("zipkin doesn't support spans with multiple parents. Omitting " - + "other parents for " + span); - } - zipkinSpan.parentId(span.getParents().get(0)); - } - zipkinSpan.id(span.getSpanId()); - if (StringUtils.hasText(span.getName())) { - zipkinSpan.name(span.getName()); - } - return zipkinSpan.build(); - } - /** * Add annotations from the sleuth Span. */ - private static void addZipkinAnnotations(Builder zipkinSpan, Span span, Endpoint endpoint) { + static void addZipkinAnnotations(Builder zipkinSpan, Span span, Endpoint endpoint) { for (Log ta : span.logs()) { Annotation zipkinAnnotation = new Annotation.Builder() .endpoint(endpoint) @@ -118,7 +72,7 @@ public class ZipkinMessageListener { * * @return list of Annotations that could be added to Zipkin Span. */ - private static void addZipkinBinaryAnnotations(Builder zipkinSpan, Span span, + static void addZipkinBinaryAnnotations(Builder zipkinSpan, Span span, Endpoint endpoint) { for (Map.Entry e : span.tags().entrySet()) { BinaryAnnotation.Builder binaryAnn = new BinaryAnnotation.Builder(); diff --git a/spring-cloud-sleuth-zipkin-stream/src/test/java/org/springframework/cloud/sleuth/zipkin/stream/ZipkinMessageListenerTests.java b/spring-cloud-sleuth-zipkin-stream/src/test/java/org/springframework/cloud/sleuth/zipkin/stream/ZipkinMessageListenerTests.java index 33e45c53a..47c900b7a 100644 --- a/spring-cloud-sleuth-zipkin-stream/src/test/java/org/springframework/cloud/sleuth/zipkin/stream/ZipkinMessageListenerTests.java +++ b/spring-cloud-sleuth-zipkin-stream/src/test/java/org/springframework/cloud/sleuth/zipkin/stream/ZipkinMessageListenerTests.java @@ -38,7 +38,7 @@ public class ZipkinMessageListenerTests { long start = System.currentTimeMillis(); this.span.logEvent("hystrix/retry"); // System.currentTimeMillis - zipkin.Span result = ZipkinMessageListener.convert(this.span, this.host); + zipkin.Span result = SamplingZipkinSpanIterator.convert(this.span, this.host); assertThat(result.timestamp) .isEqualTo(this.span.getBegin() * 1000); @@ -55,7 +55,7 @@ public class ZipkinMessageListenerTests { this.span.logEvent("hystrix/retry"); this.span.tag("spring-boot/version", "1.3.1.RELEASE"); - zipkin.Span result = ZipkinMessageListener.convert(this.span, this.host); + zipkin.Span result = SamplingZipkinSpanIterator.convert(this.span, this.host); assertThat(result.annotations.get(0).endpoint) .isEqualTo(this.endpoint); @@ -70,7 +70,7 @@ public class ZipkinMessageListenerTests { */ @Test public void spanWithoutAnnotationsLogsComponent() { - zipkin.Span result = ZipkinMessageListener.convert(this.span, this.host); + zipkin.Span result = SamplingZipkinSpanIterator.convert(this.span, this.host); assertThat(result.binaryAnnotations).hasSize(1); assertThat(result.binaryAnnotations.get(0)).isEqualToComparingFieldByField( @@ -82,7 +82,7 @@ public class ZipkinMessageListenerTests { public void nullProcessIdCoercesToUnknownServiceName() { Span noProcessId = Span.builder().traceId(1L).name("parent").remote(true).build(); - zipkin.Span result = ZipkinMessageListener.convert(noProcessId, this.host); + zipkin.Span result = SamplingZipkinSpanIterator.convert(noProcessId, this.host); assertThat(result.binaryAnnotations) .containsOnly(BinaryAnnotation.create("lc", "unknown", this.endpoint)); 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 0f1baa4d0..65d19e78b 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 @@ -16,7 +16,9 @@ package org.springframework.cloud.sleuth.zipkin; -import lombok.extern.apachecommons.CommonsLog; +import java.nio.charset.Charset; +import java.util.Map; + import org.springframework.cloud.sleuth.Log; import org.springframework.cloud.sleuth.Span; import org.springframework.cloud.sleuth.event.ClientReceivedEvent; @@ -28,14 +30,13 @@ import org.springframework.cloud.sleuth.event.SpanReleasedEvent; import org.springframework.context.event.EventListener; import org.springframework.core.annotation.Order; import org.springframework.util.StringUtils; + +import lombok.extern.apachecommons.CommonsLog; import zipkin.Annotation; import zipkin.BinaryAnnotation; import zipkin.Constants; import zipkin.Endpoint; -import java.nio.charset.Charset; -import java.util.Map; - /** * @author Spencer Gibb */ @@ -139,7 +140,7 @@ public class ZipkinSpanListener { } zipkinSpan.timestamp(span.getBegin() * 1000L); - zipkinSpan.duration((span.getEnd() - span.getBegin()) * 1000L); + zipkinSpan.duration(span.getAccumulatedMillis() * 1000L); zipkinSpan.traceId(span.getTraceId()); if (span.getParents().size() > 0) { if (span.getParents().size() > 1) {