From 80f51f76fecff71c3eceb64295c4b6dc06979c1d Mon Sep 17 00:00:00 2001 From: Marcin Grzejszczak Date: Fri, 29 Apr 2016 13:52:11 +0200 Subject: [PATCH] Draining from queue instead of clearing it (#262) now instead of first passing all the spans from the queue to the list and then clearing it we're now draining it contents. If in the meantime any spans will arrive they will be rained at next passing fixes #259 --- .../cloud/sleuth/stream/StreamSpanReporter.java | 12 ++++++------ .../cloud/sleuth/stream/StreamSpanListenerTests.java | 5 ++++- 2 files changed, 10 insertions(+), 7 deletions(-) diff --git a/spring-cloud-sleuth-stream/src/main/java/org/springframework/cloud/sleuth/stream/StreamSpanReporter.java b/spring-cloud-sleuth-stream/src/main/java/org/springframework/cloud/sleuth/stream/StreamSpanReporter.java index 81409eb0f..1d41cd6f1 100644 --- a/spring-cloud-sleuth-stream/src/main/java/org/springframework/cloud/sleuth/stream/StreamSpanReporter.java +++ b/spring-cloud-sleuth-stream/src/main/java/org/springframework/cloud/sleuth/stream/StreamSpanReporter.java @@ -17,10 +17,10 @@ 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 java.util.concurrent.BlockingQueue; +import java.util.concurrent.LinkedBlockingQueue; import org.springframework.cloud.sleuth.Span; import org.springframework.cloud.sleuth.SpanReporter; @@ -38,7 +38,7 @@ import org.springframework.integration.annotation.MessageEndpoint; @MessageEndpoint public class StreamSpanReporter implements SpanReporter { - private Collection queue = new ConcurrentLinkedQueue<>(); + private BlockingQueue queue = new LinkedBlockingQueue<>(); private final HostLocator endpointLocator; private final SpanMetricReporter spanMetricReporter; @@ -47,14 +47,14 @@ public class StreamSpanReporter implements SpanReporter { this.spanMetricReporter = spanMetricReporter; } - public void setQueue(Collection queue) { + public void setQueue(BlockingQueue queue) { this.queue = queue; } @InboundChannelAdapter(value = SleuthSource.OUTPUT) public Spans poll() { - List result = new ArrayList<>(this.queue); - this.queue.clear(); + List result = new ArrayList<>(); + this.queue.drainTo(result); for (Iterator iterator = result.iterator(); iterator.hasNext();) { Span span = iterator.next(); if (span.getName() != null && span.getName().equals("message/" + SleuthSource.OUTPUT)) { 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 d9294b135..2a7ccb5c3 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 @@ -24,6 +24,9 @@ import static org.mockito.Mockito.verify; import java.util.ArrayList; import java.util.List; +import java.util.concurrent.BlockingQueue; +import java.util.concurrent.LinkedBlockingDeque; +import java.util.concurrent.LinkedBlockingQueue; import javax.annotation.PostConstruct; @@ -155,7 +158,7 @@ public class StreamSpanListenerTests { @MessageEndpoint protected static class ZipkinTestConfiguration { - private List spans = new ArrayList<>(); + private BlockingQueue spans = new LinkedBlockingQueue<>(); @Autowired StreamSpanReporter listener;