From 301c6b230ae771d512aa2436ce5b396ac6e5065f Mon Sep 17 00:00:00 2001 From: Marcin Grzejszczak Date: Thu, 21 Jan 2016 15:55:26 +0100 Subject: [PATCH] [#99] Basic Health Metrics Changes following review - Reusing CounterService to automatically profit from Dropwizard if present - NoOp is the default impl for SpanReporterService - SpanReporterService has configurable metric names (it's enough to change the name to 'meter.a.b.c' to profit from Dropwizard's meters) Fixes gh-99 --- .../autoconfig/TraceAutoConfiguration.java | 43 ++++++++++++++++ ...ounterServiceBasedSpanReporterService.java | 35 +++++++++++++ .../metric/NoOpSpanReporterService.java | 17 +++++++ .../sleuth/metric/SleuthMetricProperties.java | 45 +++++++++++++++++ .../sleuth/metric/SpanReporterService.java | 21 ++++++++ .../web/TraceAsyncIntegrationTests.java | 3 +- .../spring-cloud-sleuth-sample-zipkin/pom.xml | 6 +++ .../test/java/integration/ZipkinTests.java | 15 ++++-- .../spring-cloud-sleuth-sample/pom.xml | 4 -- .../stream/SleuthStreamAutoConfiguration.java | 23 +++++++-- .../sleuth/stream/StreamSpanListener.java | 27 ++++++---- .../stream/StreamSpanListenerTests.java | 50 +++++++++++-------- .../sleuth/zipkin/HttpZipkinSpanReporter.java | 30 +++++++---- .../zipkin/ZipkinAutoConfiguration.java | 6 ++- .../zipkin/HttpZipkinSpanReporterTest.java | 19 ++++--- .../sleuth/zipkin/InMemorySpanCounter.java | 48 ++++++++++++++++++ 16 files changed, 328 insertions(+), 64 deletions(-) create mode 100644 spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/metric/CounterServiceBasedSpanReporterService.java create mode 100644 spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/metric/NoOpSpanReporterService.java create mode 100644 spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/metric/SleuthMetricProperties.java create mode 100644 spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/metric/SpanReporterService.java create mode 100644 spring-cloud-sleuth-zipkin/src/test/java/org/springframework/cloud/sleuth/zipkin/InMemorySpanCounter.java diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/autoconfig/TraceAutoConfiguration.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/autoconfig/TraceAutoConfiguration.java index f83a1bb8e..e3b1694b1 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/autoconfig/TraceAutoConfiguration.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/autoconfig/TraceAutoConfiguration.java @@ -18,10 +18,19 @@ package org.springframework.cloud.sleuth.autoconfig; import java.util.Random; +import org.springframework.boot.actuate.metrics.CounterService; +import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; +import org.springframework.boot.autoconfigure.condition.ConditionalOnClass; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; +import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingClass; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.cloud.sleuth.Sampler; import org.springframework.cloud.sleuth.Tracer; +import org.springframework.cloud.sleuth.metric.CounterServiceBasedSpanReporterService; +import org.springframework.cloud.sleuth.metric.NoOpSpanReporterService; +import org.springframework.cloud.sleuth.metric.SleuthMetricProperties; +import org.springframework.cloud.sleuth.metric.SpanReporterService; import org.springframework.cloud.sleuth.sampler.IsTracingSampler; import org.springframework.cloud.sleuth.trace.DefaultTracer; import org.springframework.context.ApplicationEventPublisher; @@ -33,6 +42,7 @@ import org.springframework.context.annotation.Configuration; */ @Configuration @ConditionalOnProperty(value="spring.sleuth.enabled", matchIfMissing=true) +@EnableConfigurationProperties public class TraceAutoConfiguration { @Bean @@ -53,4 +63,37 @@ public class TraceAutoConfiguration { ApplicationEventPublisher publisher) { return new DefaultTracer(sampler, random, publisher); } + + @Bean + @ConditionalOnMissingBean + public SleuthMetricProperties sleuthMetricProperties() { + return new SleuthMetricProperties(); + } + + @Configuration + @ConditionalOnClass(CounterService.class) + @ConditionalOnMissingBean(SpanReporterService.class) + protected static class CounterServiceSpanReporterConfig { + @Bean + @ConditionalOnBean(CounterService.class) + public SpanReporterService spanReporterCounterService(CounterService counterService, + SleuthMetricProperties sleuthMetricProperties) { + return new CounterServiceBasedSpanReporterService(sleuthMetricProperties.getSpan().getAcceptedName(), + sleuthMetricProperties.getSpan().getDroppedName(), counterService); + } + + @Bean + @ConditionalOnMissingBean(CounterService.class) + public SpanReporterService noOpSpanReporterCounterService() { + return new NoOpSpanReporterService(); + } + } + + @Bean + @ConditionalOnMissingClass("org.springframework.boot.actuate.metrics.CounterService") + @ConditionalOnMissingBean(SpanReporterService.class) + public SpanReporterService noOpSpanReporterCounterService() { + return new NoOpSpanReporterService(); + } + } diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/metric/CounterServiceBasedSpanReporterService.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/metric/CounterServiceBasedSpanReporterService.java new file mode 100644 index 000000000..0d89bc48a --- /dev/null +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/metric/CounterServiceBasedSpanReporterService.java @@ -0,0 +1,35 @@ +package org.springframework.cloud.sleuth.metric; + +import org.springframework.boot.actuate.metrics.CounterService; + +/** + * Service to operate on accepted and dropped spans statistics. + * + * @author Marcin Grzejszczak + */ +public class CounterServiceBasedSpanReporterService implements SpanReporterService { + private final String acceptedSpansMetricName; + private final String droppedSpansMetricName; + private final CounterService counterService; + + public CounterServiceBasedSpanReporterService(String acceptedSpansMetricName, + String droppedSpansMetricName, CounterService counterService) { + this.acceptedSpansMetricName = acceptedSpansMetricName; + this.droppedSpansMetricName = droppedSpansMetricName; + this.counterService = counterService; + } + + @Override + public void incrementAcceptedSpans(long quantity) { + for (int i = 0; i < quantity; i++) { + this.counterService.increment(this.acceptedSpansMetricName); + } + } + + @Override + public void incrementDroppedSpans(long quantity) { + for (int i = 0; i < quantity; i++) { + this.counterService.increment(this.droppedSpansMetricName); + } + } +} \ No newline at end of file diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/metric/NoOpSpanReporterService.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/metric/NoOpSpanReporterService.java new file mode 100644 index 000000000..a01bb2655 --- /dev/null +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/metric/NoOpSpanReporterService.java @@ -0,0 +1,17 @@ +package org.springframework.cloud.sleuth.metric; + +/** + * Span reporting service that does nothing + * + * @author Marcin Grzejszczak + */ +public class NoOpSpanReporterService implements SpanReporterService { + + public void incrementAcceptedSpans(long quantity) { + + } + + public void incrementDroppedSpans(long quantity) { + + } +} \ No newline at end of file diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/metric/SleuthMetricProperties.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/metric/SleuthMetricProperties.java new file mode 100644 index 000000000..56dcac1d7 --- /dev/null +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/metric/SleuthMetricProperties.java @@ -0,0 +1,45 @@ +package org.springframework.cloud.sleuth.metric; + +import org.springframework.boot.context.properties.ConfigurationProperties; + +/** + * Configuration properties for Sleuth related metrics + * + * @author Marcin Grzejszczak + */ +@ConfigurationProperties("spring.sleuth.metric") +public class SleuthMetricProperties { + + private Span span = new Span(); + + public Span getSpan() { + return this.span; + } + + public void setSpan(Span span) { + this.span = span; + } + + public static class Span { + + private String acceptedName = "counter.span.accepted"; + + private String droppedName = "counter.span.dropped"; + + public String getAcceptedName() { + return this.acceptedName; + } + + public void setAcceptedName(String acceptedName) { + this.acceptedName = acceptedName; + } + + public String getDroppedName() { + return this.droppedName; + } + + public void setDroppedName(String droppedName) { + this.droppedName = droppedName; + } + } +} diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/metric/SpanReporterService.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/metric/SpanReporterService.java new file mode 100644 index 000000000..8c2a2bb5c --- /dev/null +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/metric/SpanReporterService.java @@ -0,0 +1,21 @@ +package org.springframework.cloud.sleuth.metric; + +/** + * @author Marcin Grzejszczak + */ +public interface SpanReporterService { + + /** + * Called when spans are submitted to SpanCollector for processing. + * + * @param quantity the number of spans accepted. + */ + void incrementAcceptedSpans(long quantity); + + /** + * Called when spans become lost for any reason and won't be delivered to the target collector. + * + * @param quantity the number of spans dropped. + */ + void incrementDroppedSpans(long quantity); +} diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/TraceAsyncIntegrationTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/TraceAsyncIntegrationTests.java index 5923c0acf..b5cb21b62 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/TraceAsyncIntegrationTests.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/instrument/web/TraceAsyncIntegrationTests.java @@ -27,8 +27,7 @@ import static org.springframework.cloud.sleuth.assertions.SleuthAssertions.then; public class TraceAsyncIntegrationTests { @Autowired ClassPerformingAsyncLogic classPerformingAsyncLogic; - @Autowired - Tracer tracer; + @Autowired Tracer tracer; @Test public void should_set_span_on_an_async_annotated_method() { diff --git a/spring-cloud-sleuth-samples/spring-cloud-sleuth-sample-zipkin/pom.xml b/spring-cloud-sleuth-samples/spring-cloud-sleuth-sample-zipkin/pom.xml index c3c8681b4..f6d12fd58 100644 --- a/spring-cloud-sleuth-samples/spring-cloud-sleuth-sample-zipkin/pom.xml +++ b/spring-cloud-sleuth-samples/spring-cloud-sleuth-sample-zipkin/pom.xml @@ -91,6 +91,12 @@ brave-core test + + io.dropwizard.metrics + metrics-core + ${dropwizard-metrics.version} + test + diff --git a/spring-cloud-sleuth-samples/spring-cloud-sleuth-sample-zipkin/src/test/java/integration/ZipkinTests.java b/spring-cloud-sleuth-samples/spring-cloud-sleuth-sample-zipkin/src/test/java/integration/ZipkinTests.java index e1e2cd2b1..64e24f94d 100644 --- a/spring-cloud-sleuth-samples/spring-cloud-sleuth-sample-zipkin/src/test/java/integration/ZipkinTests.java +++ b/spring-cloud-sleuth-samples/spring-cloud-sleuth-sample-zipkin/src/test/java/integration/ZipkinTests.java @@ -23,6 +23,7 @@ import org.junit.Test; import org.junit.runner.RunWith; import org.springframework.boot.test.SpringApplicationConfiguration; import org.springframework.boot.test.WebIntegrationTest; +import org.springframework.cloud.sleuth.metric.SpanReporterService; import org.springframework.cloud.sleuth.zipkin.HttpZipkinSpanReporter; import org.springframework.cloud.sleuth.zipkin.ZipkinProperties; import org.springframework.cloud.sleuth.zipkin.ZipkinSpanReporter; @@ -73,12 +74,14 @@ public class ZipkinTests extends AbstractIntegrationTest { public static class WaitUntilZipkinIsUpConfig { @Bean @SneakyThrows - public ZipkinSpanReporter spanCollector(final ZipkinProperties zipkin) { + public ZipkinSpanReporter spanCollector(final ZipkinProperties zipkin, + final SpanReporterService spanReporterService) { await().until(new Runnable() { @Override public void run() { try { - WaitUntilZipkinIsUpConfig.this.getSpanCollector(zipkin); + WaitUntilZipkinIsUpConfig.this.getSpanCollector(zipkin, + spanReporterService); } catch (Exception e) { log.error("Exception occurred while trying to connect to zipkin [" @@ -87,11 +90,13 @@ public class ZipkinTests extends AbstractIntegrationTest { } } }); - return getSpanCollector(zipkin); + return getSpanCollector(zipkin, spanReporterService); } - private ZipkinSpanReporter getSpanCollector(ZipkinProperties zipkin) { - return new HttpZipkinSpanReporter(zipkin.getBaseUrl(), zipkin.getFlushInterval()); + private ZipkinSpanReporter getSpanCollector(ZipkinProperties zipkin, + SpanReporterService spanReporterService) { + return new HttpZipkinSpanReporter(zipkin.getBaseUrl(), zipkin.getFlushInterval(), + spanReporterService); } } } diff --git a/spring-cloud-sleuth-samples/spring-cloud-sleuth-sample/pom.xml b/spring-cloud-sleuth-samples/spring-cloud-sleuth-sample/pom.xml index 81ab18288..413eac3b4 100644 --- a/spring-cloud-sleuth-samples/spring-cloud-sleuth-sample/pom.xml +++ b/spring-cloud-sleuth-samples/spring-cloud-sleuth-sample/pom.xml @@ -52,10 +52,6 @@ org.springframework.boot spring-boot-starter-aop - - org.springframework.boot - spring-boot-starter-actuator - org.projectlombok lombok diff --git a/spring-cloud-sleuth-stream/src/main/java/org/springframework/cloud/sleuth/stream/SleuthStreamAutoConfiguration.java b/spring-cloud-sleuth-stream/src/main/java/org/springframework/cloud/sleuth/stream/SleuthStreamAutoConfiguration.java index 9d3b9e3b3..0c6b1a65c 100644 --- a/spring-cloud-sleuth-stream/src/main/java/org/springframework/cloud/sleuth/stream/SleuthStreamAutoConfiguration.java +++ b/spring-cloud-sleuth-stream/src/main/java/org/springframework/cloud/sleuth/stream/SleuthStreamAutoConfiguration.java @@ -26,6 +26,7 @@ import org.springframework.boot.autoconfigure.web.ServerProperties; import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.cloud.client.discovery.DiscoveryClient; import org.springframework.cloud.sleuth.Span; +import org.springframework.cloud.sleuth.metric.SpanReporterService; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.config.ChannelBindingAutoConfiguration; import org.springframework.context.annotation.Bean; @@ -55,7 +56,7 @@ public class SleuthStreamAutoConfiguration { @Bean @GlobalChannelInterceptor(patterns = SleuthSource.OUTPUT, order = Ordered.HIGHEST_PRECEDENCE) - public ChannelInterceptor zipkinChannelInterceptor() { + public ChannelInterceptor zipkinChannelInterceptor(final SpanReporterService spanReporterService) { // don't trace the tracer (suppress spans originating from our own source) return new ChannelInterceptorAdapter() { @Override @@ -63,12 +64,28 @@ public class SleuthStreamAutoConfiguration { return MessageBuilder.fromMessage(message) .setHeader(Span.NOT_SAMPLED_NAME, "").build(); } + + @Override + public void afterSendCompletion(Message message, MessageChannel channel, + boolean sent, Exception ex) { + if (!(message.getPayload() instanceof Spans)) { + return; + } + Spans spans = (Spans) message.getPayload(); + int spanNumber = spans.getSpans().size(); + if (sent) { + spanReporterService.incrementAcceptedSpans(spanNumber); + } else { + spanReporterService.incrementDroppedSpans(spanNumber); + } + } }; } @Bean - public StreamSpanListener sleuthTracer(HostLocator endpointLocator) { - return new StreamSpanListener(endpointLocator); + public StreamSpanListener sleuthTracer(HostLocator endpointLocator, + SpanReporterService spanReporterService) { + return new StreamSpanListener(endpointLocator, spanReporterService); } @Configuration 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 5f298f5e7..b9f629f1c 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,12 +16,6 @@ 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,11 +23,19 @@ import org.springframework.cloud.sleuth.event.ServerReceivedEvent; import org.springframework.cloud.sleuth.event.ServerSentEvent; import org.springframework.cloud.sleuth.event.SpanAcquiredEvent; import org.springframework.cloud.sleuth.event.SpanReleasedEvent; +import org.springframework.cloud.sleuth.metric.SpanReporterService; import org.springframework.context.event.EventListener; 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. * @@ -48,10 +50,12 @@ public class StreamSpanListener { public static final String SERVER_SEND = "ss"; private Collection queue = new ConcurrentLinkedQueue<>(); - private HostLocator endpointLocator; + private final HostLocator endpointLocator; + private final SpanReporterService spanReporterService; - public StreamSpanListener(HostLocator endpointLocator) { + public StreamSpanListener(HostLocator endpointLocator, SpanReporterService spanReporterService) { this.endpointLocator = endpointLocator; + this.spanReporterService = spanReporterService; } public void setQueue(Collection queue) { @@ -112,8 +116,11 @@ public class StreamSpanListener { iterator.remove(); } } - return result.isEmpty() ? null - : new Spans(this.endpointLocator.locate(result.get(0)), result); + if (result.isEmpty()) { + return null; + } + this.spanReporterService.incrementAcceptedSpans(result.size()); + return new Spans(this.endpointLocator.locate(result.get(0)), result); } } diff --git a/spring-cloud-sleuth-stream/src/test/java/org/springframework/cloud/sleuth/stream/StreamSpanListenerTests.java b/spring-cloud-sleuth-stream/src/test/java/org/springframework/cloud/sleuth/stream/StreamSpanListenerTests.java index e0a44ed2d..07da4936b 100644 --- a/spring-cloud-sleuth-stream/src/test/java/org/springframework/cloud/sleuth/stream/StreamSpanListenerTests.java +++ b/spring-cloud-sleuth-stream/src/test/java/org/springframework/cloud/sleuth/stream/StreamSpanListenerTests.java @@ -16,16 +16,11 @@ package org.springframework.cloud.sleuth.stream; -import static org.junit.Assert.assertEquals; - -import java.util.ArrayList; -import java.util.List; - -import javax.annotation.PostConstruct; - import org.junit.Test; import org.junit.runner.RunWith; +import org.mockito.Mockito; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.actuate.metrics.CounterService; import org.springframework.boot.autoconfigure.PropertyPlaceholderAutoConfiguration; import org.springframework.boot.test.SpringApplicationConfiguration; import org.springframework.cloud.sleuth.Sampler; @@ -49,6 +44,15 @@ import org.springframework.integration.annotation.ServiceActivator; import org.springframework.messaging.Message; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; +import javax.annotation.PostConstruct; +import java.util.ArrayList; +import java.util.List; + +import static org.junit.Assert.assertEquals; +import static org.mockito.Matchers.anyString; +import static org.mockito.Mockito.atLeastOnce; +import static org.mockito.Mockito.verify; + /** * @author Dave Syer * @@ -57,17 +61,11 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; @RunWith(SpringJUnit4ClassRunner.class) public class StreamSpanListenerTests { - @Autowired - private Tracer tracer; - - @Autowired - private ApplicationContext application; - - @Autowired - private ZipkinTestConfiguration test; - - @Autowired - StreamSpanListener listener; + @Autowired Tracer tracer; + @Autowired ApplicationContext application; + @Autowired ZipkinTestConfiguration test; + @Autowired StreamSpanListener listener; + @Autowired CounterService counterService; @PostConstruct public void init() { @@ -106,12 +104,20 @@ public class StreamSpanListenerTests { assertEquals(0, this.test.spans.size()); } + @Test + public void shouldIncreaseNumberOfAcceptedSpans() { + Span context = this.tracer.startTrace("foo"); + this.tracer.close(context); + this.listener.poll(); + + verify(this.counterService, atLeastOnce()).increment(anyString()); + } + @Configuration @Import({ ZipkinTestConfiguration.class, SleuthStreamAutoConfiguration.class, TestSupportBinderAutoConfiguration.class, ChannelBindingAutoConfiguration.class, TraceAutoConfiguration.class, PropertyPlaceholderAutoConfiguration.class }) - protected static class TestConfiguration { - } + protected static class TestConfiguration {} @Configuration @MessageEndpoint @@ -131,6 +137,10 @@ public class StreamSpanListenerTests { return new AlwaysSampler(); } + @Bean CounterService counterService() { + return Mockito.mock(CounterService.class); + } + @PostConstruct public void init() { this.listener.setQueue(this.spans); diff --git a/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin/HttpZipkinSpanReporter.java b/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin/HttpZipkinSpanReporter.java index fb66a795a..0b9306421 100644 --- a/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin/HttpZipkinSpanReporter.java +++ b/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin/HttpZipkinSpanReporter.java @@ -1,6 +1,7 @@ package org.springframework.cloud.sleuth.zipkin; import lombok.extern.apachecommons.CommonsLog; +import org.springframework.cloud.sleuth.metric.SpanReporterService; import zipkin.Codec; import zipkin.Span; @@ -32,14 +33,18 @@ public final class HttpZipkinSpanReporter private final String url; private final BlockingQueue pending = new LinkedBlockingQueue<>(1000); private final Flusher flusher; // Nullable for testing + private final SpanReporterService spanReporterService; /** * @param baseUrl URL of the zipkin query server instance. Like: http://localhost:9411/ * @param flushInterval in seconds. 0 implies spans are {@link #flush() flushed} externally. + * @param spanReporterService service to count number of accepted / dropped spans */ - public HttpZipkinSpanReporter(String baseUrl, int flushInterval) { + public HttpZipkinSpanReporter(String baseUrl, int flushInterval, + SpanReporterService spanReporterService) { this.url = baseUrl + (baseUrl.endsWith("/") ? "" : "/") + "api/v1/spans"; this.flusher = flushInterval > 0 ? new Flusher(this, flushInterval) : null; + this.spanReporterService = spanReporterService; } /** @@ -47,17 +52,19 @@ public final class HttpZipkinSpanReporter * * @param span Span, should not be null. */ - @Override public void report(Span span) { - // TODO: metrics.incrementAcceptedSpans(1); + @Override + public void report(Span span) { + this.spanReporterService.incrementAcceptedSpans(1); if (!this.pending.offer(span)) { - // TODO: metrics.incrementDroppedSpans(1); + this.spanReporterService.incrementDroppedSpans(1); } } /** * Calling this will flush any pending spans to the http transport on the current thread. */ - @Override public void flush() { + @Override + public void flush() { if (this.pending.isEmpty()) return; List drained = new ArrayList<>(this.pending.size()); @@ -70,7 +77,7 @@ public final class HttpZipkinSpanReporter // NOTE: https://github.com/openzipkin/zipkin-java/issues/66 will throw instead of return null. if (json == null) { log.debug("failed to encode spans, dropping them: " + drained); - // TODO: metrics.incrementDroppedSpans(spanCount); + this.spanReporterService.incrementDroppedSpans(drained.size()); return; } @@ -85,8 +92,7 @@ public final class HttpZipkinSpanReporter "error POSTing spans to " + this.url + ": as json: " + new String(json, UTF_8), e); } - // TODO: metrics.incrementDroppedSpans(spanCount); - return; + this.spanReporterService.incrementDroppedSpans(drained.size()); } } @@ -102,7 +108,8 @@ public final class HttpZipkinSpanReporter this.scheduler.scheduleWithFixedDelay(this, 0, flushInterval, SECONDS); } - @Override public void run() { + @Override + public void run() { try { this.flushable.flush(); } @@ -139,11 +146,12 @@ public final class HttpZipkinSpanReporter * Requests a cease of delivery. There will be at most one in-flight request processing after this * call returns. */ - @Override public void close() { + @Override + public void close() { if (this.flusher != null) this.flusher.scheduler.shutdown(); // throw any outstanding spans on the floor int dropped = this.pending.drainTo(new LinkedList<>()); - // TODO: metrics.incrementDroppedSpans(dropped); + this.spanReporterService.incrementDroppedSpans(dropped); } } diff --git a/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin/ZipkinAutoConfiguration.java b/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin/ZipkinAutoConfiguration.java index cb1fecdbc..4fec778d2 100644 --- a/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin/ZipkinAutoConfiguration.java +++ b/spring-cloud-sleuth-zipkin/src/main/java/org/springframework/cloud/sleuth/zipkin/ZipkinAutoConfiguration.java @@ -25,6 +25,7 @@ import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.boot.autoconfigure.web.ServerProperties; import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.cloud.client.discovery.DiscoveryClient; +import org.springframework.cloud.sleuth.metric.SpanReporterService; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @@ -39,9 +40,10 @@ public class ZipkinAutoConfiguration { @Bean @ConditionalOnMissingBean(ZipkinSpanReporter.class) - public ZipkinSpanReporter reporter() { + public ZipkinSpanReporter reporter(SpanReporterService spanReporterService) { ZipkinProperties zipkin = zipkinProperties(); - return new HttpZipkinSpanReporter(zipkin.getBaseUrl(), zipkin.getFlushInterval()); + return new HttpZipkinSpanReporter(zipkin.getBaseUrl(), zipkin.getFlushInterval(), + spanReporterService); } @Bean diff --git a/spring-cloud-sleuth-zipkin/src/test/java/org/springframework/cloud/sleuth/zipkin/HttpZipkinSpanReporterTest.java b/spring-cloud-sleuth-zipkin/src/test/java/org/springframework/cloud/sleuth/zipkin/HttpZipkinSpanReporterTest.java index 8b8405153..b42bd18c1 100644 --- a/spring-cloud-sleuth-zipkin/src/test/java/org/springframework/cloud/sleuth/zipkin/HttpZipkinSpanReporterTest.java +++ b/spring-cloud-sleuth-zipkin/src/test/java/org/springframework/cloud/sleuth/zipkin/HttpZipkinSpanReporterTest.java @@ -6,6 +6,8 @@ import okhttp3.mockwebserver.RecordedRequest; import okhttp3.mockwebserver.SocketPolicy; import org.junit.Rule; import org.junit.Test; +import org.springframework.cloud.sleuth.metric.CounterServiceBasedSpanReporterService; +import org.springframework.cloud.sleuth.metric.SpanReporterService; import zipkin.Codec; import zipkin.Span; @@ -16,10 +18,13 @@ import static org.assertj.core.api.Assertions.assertThat; public class HttpZipkinSpanReporterTest { @Rule public final MockWebServer server = new MockWebServer(); + InMemorySpanCounter inMemorySpanCounter = new InMemorySpanCounter(); + SpanReporterService spanReporterService = new CounterServiceBasedSpanReporterService("accepted", "dropped", + this.inMemorySpanCounter); // set flush interval to 0 so that tests can drive flushing explicitly HttpZipkinSpanReporter reporter = new HttpZipkinSpanReporter( - this.server.url("").toString(), 0); + this.server.url("").toString(), 0, this.spanReporterService); @Test public void reportDoesntDoIO() throws Exception { @@ -32,8 +37,8 @@ public class HttpZipkinSpanReporterTest { public void reportIncrementsAcceptedMetrics() throws Exception { this.reporter.report(span(1L, "foo")); - // TODO: assertThat(metrics.acceptedSpans.get()).isEqualTo(1); - // TODO: assertThat(metrics.droppedSpans.get()).isZero(); + assertThat(this.inMemorySpanCounter.getAcceptedSpans()).isEqualTo(1); + assertThat(this.inMemorySpanCounter.getDroppedSpans()).isZero(); } @Test @@ -41,8 +46,8 @@ public class HttpZipkinSpanReporterTest { for (int i = 0; i < 1001; i++) this.reporter.report(span(1L, "foo")); - // TODO: assertThat(metrics.acceptedSpans.get()).isEqualTo(1001); - // TODO: assertThat(metrics.droppedSpans.get()).isEqualTo(1); + assertThat(this.inMemorySpanCounter.getAcceptedSpans()).isEqualTo(1001); + assertThat(this.inMemorySpanCounter.getDroppedSpans()).isEqualTo(1); } @Test @@ -73,7 +78,7 @@ public class HttpZipkinSpanReporterTest { this.reporter.flush(); // manually flush the spans - // TODO: assertThat(metrics.droppedSpans.get()).isEqualTo(2); + assertThat(this.inMemorySpanCounter.getDroppedSpans()).isEqualTo(2); } @Test @@ -86,7 +91,7 @@ public class HttpZipkinSpanReporterTest { this.reporter.flush(); // manually flush the spans - // TODO: assertThat(metrics.droppedSpans.get()).isEqualTo(2); + assertThat(this.inMemorySpanCounter.getDroppedSpans()).isEqualTo(2); } static Span span(long traceId, String spanName) { diff --git a/spring-cloud-sleuth-zipkin/src/test/java/org/springframework/cloud/sleuth/zipkin/InMemorySpanCounter.java b/spring-cloud-sleuth-zipkin/src/test/java/org/springframework/cloud/sleuth/zipkin/InMemorySpanCounter.java new file mode 100644 index 000000000..1cb91dcf8 --- /dev/null +++ b/spring-cloud-sleuth-zipkin/src/test/java/org/springframework/cloud/sleuth/zipkin/InMemorySpanCounter.java @@ -0,0 +1,48 @@ +package org.springframework.cloud.sleuth.zipkin; + +import org.springframework.boot.actuate.metrics.CounterService; + +import java.util.concurrent.atomic.AtomicLong; + +/** + * Implementation of the metrics statistics held in-memory. + * + * @author Marcin Grzejszczak + */ +public class InMemorySpanCounter implements CounterService { + + private final AtomicLong acceptedSpans = new AtomicLong(0); + private final AtomicLong droppedSpans = new AtomicLong(0); + + public long getAcceptedSpans() { + return this.acceptedSpans.get(); + } + + public long getDroppedSpans() { + return this.droppedSpans.get(); + } + + @Override + public void increment(String metricName) { + if (metricName.contains("accepted")) { + this.acceptedSpans.incrementAndGet(); + } else { + this.droppedSpans.incrementAndGet(); + } + } + + @Override + public void decrement(String metricName) { + if (metricName.contains("accepted")) { + this.acceptedSpans.decrementAndGet(); + } else { + this.droppedSpans.decrementAndGet(); + } + } + + @Override + public void reset(String metricName) { + this.acceptedSpans.set(0); + this.droppedSpans.set(0); + } +}