@@ -27,7 +27,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.Sampler;
|
||||
import org.springframework.cloud.sleuth.metric.SpanReporterService;
|
||||
import org.springframework.cloud.sleuth.metric.SpanMetricReporter;
|
||||
import org.springframework.cloud.sleuth.sampler.PercentageBasedSampler;
|
||||
import org.springframework.cloud.sleuth.sampler.SamplerProperties;
|
||||
import org.springframework.cloud.stream.annotation.EnableBinding;
|
||||
@@ -64,14 +64,14 @@ public class SleuthStreamAutoConfiguration {
|
||||
|
||||
@Bean
|
||||
@GlobalChannelInterceptor(patterns = SleuthSource.OUTPUT, order = Ordered.HIGHEST_PRECEDENCE)
|
||||
public ChannelInterceptor zipkinChannelInterceptor(SpanReporterService spanReporterService) {
|
||||
return new TracerIgnoringChannelInterceptor(spanReporterService);
|
||||
public ChannelInterceptor zipkinChannelInterceptor(SpanMetricReporter spanMetricReporter) {
|
||||
return new TracerIgnoringChannelInterceptor(spanMetricReporter);
|
||||
}
|
||||
|
||||
@Bean
|
||||
public StreamSpanListener sleuthTracer(HostLocator endpointLocator,
|
||||
SpanReporterService spanReporterService) {
|
||||
return new StreamSpanListener(endpointLocator, spanReporterService);
|
||||
SpanMetricReporter spanMetricReporter) {
|
||||
return new StreamSpanListener(endpointLocator, spanMetricReporter);
|
||||
}
|
||||
|
||||
@Configuration
|
||||
|
||||
@@ -23,15 +23,8 @@ 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;
|
||||
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.cloud.sleuth.SpanReporter;
|
||||
import org.springframework.cloud.sleuth.metric.SpanMetricReporter;
|
||||
import org.springframework.integration.annotation.InboundChannelAdapter;
|
||||
import org.springframework.integration.annotation.MessageEndpoint;
|
||||
|
||||
@@ -44,70 +37,21 @@ import org.springframework.integration.annotation.MessageEndpoint;
|
||||
* @since 1.0.0
|
||||
*/
|
||||
@MessageEndpoint
|
||||
public class StreamSpanListener {
|
||||
|
||||
public static final String CLIENT_RECV = "cr";
|
||||
public static final String CLIENT_SEND = "cs";
|
||||
public static final String SERVER_RECV = "sr";
|
||||
public static final String SERVER_SEND = "ss";
|
||||
public class StreamSpanListener implements SpanReporter {
|
||||
|
||||
private Collection<Span> queue = new ConcurrentLinkedQueue<>();
|
||||
private final HostLocator endpointLocator;
|
||||
private final SpanReporterService spanReporterService;
|
||||
private final SpanMetricReporter spanMetricReporter;
|
||||
|
||||
public StreamSpanListener(HostLocator endpointLocator, SpanReporterService spanReporterService) {
|
||||
public StreamSpanListener(HostLocator endpointLocator, SpanMetricReporter spanMetricReporter) {
|
||||
this.endpointLocator = endpointLocator;
|
||||
this.spanReporterService = spanReporterService;
|
||||
this.spanMetricReporter = spanMetricReporter;
|
||||
}
|
||||
|
||||
public void setQueue(Collection<Span> queue) {
|
||||
this.queue = queue;
|
||||
}
|
||||
|
||||
@EventListener
|
||||
@Order(0)
|
||||
public void start(SpanAcquiredEvent event) {
|
||||
event.getSpan().logEvent("acquire");
|
||||
}
|
||||
|
||||
@EventListener
|
||||
@Order(0)
|
||||
public void serverReceived(ServerReceivedEvent event) {
|
||||
if (event.getParent() != null && event.getParent().isRemote()) {
|
||||
event.getParent().logEvent(SERVER_RECV);
|
||||
}
|
||||
}
|
||||
|
||||
@EventListener
|
||||
@Order(0)
|
||||
public void clientSend(ClientSentEvent event) {
|
||||
event.getSpan().logEvent(CLIENT_SEND);
|
||||
}
|
||||
|
||||
@EventListener
|
||||
@Order(0)
|
||||
public void clientReceive(ClientReceivedEvent event) {
|
||||
event.getSpan().logEvent(CLIENT_RECV);
|
||||
}
|
||||
|
||||
@EventListener
|
||||
@Order(0)
|
||||
public void serverSend(ServerSentEvent event) {
|
||||
if (event.getParent() != null && event.getParent().isRemote()) {
|
||||
event.getParent().logEvent(SERVER_SEND);
|
||||
this.queue.add(event.getParent());
|
||||
}
|
||||
}
|
||||
|
||||
@EventListener
|
||||
@Order(0)
|
||||
public void release(SpanReleasedEvent event) {
|
||||
event.getSpan().logEvent("release");
|
||||
if (event.getSpan().isExportable()) {
|
||||
this.queue.add(event.getSpan());
|
||||
}
|
||||
}
|
||||
|
||||
@InboundChannelAdapter(value = SleuthSource.OUTPUT)
|
||||
public Spans poll() {
|
||||
List<Span> result = new ArrayList<>(this.queue);
|
||||
@@ -121,8 +65,14 @@ public class StreamSpanListener {
|
||||
if (result.isEmpty()) {
|
||||
return null;
|
||||
}
|
||||
this.spanReporterService.incrementAcceptedSpans(result.size());
|
||||
this.spanMetricReporter.incrementAcceptedSpans(result.size());
|
||||
return new Spans(this.endpointLocator.locate(result.get(0)), result);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void report(Span span) {
|
||||
if (span.isExportable()) {
|
||||
this.queue.add(span);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -17,7 +17,7 @@
|
||||
package org.springframework.cloud.sleuth.stream;
|
||||
|
||||
import org.springframework.cloud.sleuth.Span;
|
||||
import org.springframework.cloud.sleuth.metric.SpanReporterService;
|
||||
import org.springframework.cloud.sleuth.metric.SpanMetricReporter;
|
||||
import org.springframework.integration.support.MessageBuilder;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.MessageChannel;
|
||||
@@ -31,10 +31,10 @@ import org.springframework.messaging.support.ChannelInterceptorAdapter;
|
||||
*/
|
||||
class TracerIgnoringChannelInterceptor extends ChannelInterceptorAdapter {
|
||||
|
||||
private final SpanReporterService spanReporterService;
|
||||
private final SpanMetricReporter spanMetricReporter;
|
||||
|
||||
public TracerIgnoringChannelInterceptor(SpanReporterService spanReporterService) {
|
||||
this.spanReporterService = spanReporterService;
|
||||
public TracerIgnoringChannelInterceptor(SpanMetricReporter spanMetricReporter) {
|
||||
this.spanMetricReporter = spanMetricReporter;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -55,9 +55,9 @@ class TracerIgnoringChannelInterceptor extends ChannelInterceptorAdapter {
|
||||
Spans spans = (Spans) message.getPayload();
|
||||
int spanNumber = spans.getSpans().size();
|
||||
if (sent) {
|
||||
this.spanReporterService.incrementAcceptedSpans(spanNumber);
|
||||
this.spanMetricReporter.incrementAcceptedSpans(spanNumber);
|
||||
} else {
|
||||
this.spanReporterService.incrementDroppedSpans(spanNumber);
|
||||
this.spanMetricReporter.incrementDroppedSpans(spanNumber);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -35,12 +35,11 @@ import org.springframework.boot.autoconfigure.PropertyPlaceholderAutoConfigurati
|
||||
import org.springframework.boot.test.SpringApplicationConfiguration;
|
||||
import org.springframework.cloud.sleuth.Sampler;
|
||||
import org.springframework.cloud.sleuth.Span;
|
||||
import org.springframework.cloud.sleuth.SpanReporter;
|
||||
import org.springframework.cloud.sleuth.Tracer;
|
||||
import org.springframework.cloud.sleuth.autoconfig.TraceAutoConfiguration;
|
||||
import org.springframework.cloud.sleuth.event.ClientReceivedEvent;
|
||||
import org.springframework.cloud.sleuth.event.ClientSentEvent;
|
||||
import org.springframework.cloud.sleuth.event.ServerReceivedEvent;
|
||||
import org.springframework.cloud.sleuth.event.ServerSentEvent;
|
||||
import org.springframework.cloud.sleuth.log.NoOpSpanLogger;
|
||||
import org.springframework.cloud.sleuth.log.SpanLogger;
|
||||
import org.springframework.cloud.sleuth.metric.TraceMetricsAutoConfiguration;
|
||||
import org.springframework.cloud.sleuth.sampler.AlwaysSampler;
|
||||
import org.springframework.cloud.sleuth.stream.StreamSpanListenerTests.TestConfiguration;
|
||||
@@ -68,6 +67,7 @@ public class StreamSpanListenerTests {
|
||||
@Autowired ZipkinTestConfiguration test;
|
||||
@Autowired StreamSpanListener listener;
|
||||
@Autowired CounterService counterService;
|
||||
@Autowired SpanReporter spanReporter;
|
||||
|
||||
@PostConstruct
|
||||
public void init() {
|
||||
@@ -86,20 +86,30 @@ public class StreamSpanListenerTests {
|
||||
Span parent = Span.builder().traceId(1L).name("http:parent").remote(true)
|
||||
.build();
|
||||
Span context = this.tracer.createSpan("http:child", parent);
|
||||
this.application.publishEvent(new ClientSentEvent(this, context));
|
||||
this.application
|
||||
.publishEvent(new ServerReceivedEvent(this, parent, context));
|
||||
this.application
|
||||
.publishEvent(new ServerSentEvent(this, parent, context));
|
||||
this.application.publishEvent(new ClientReceivedEvent(this, context));
|
||||
context.logEvent(Span.CLIENT_SEND);
|
||||
logServerReceived(parent);
|
||||
logServerSent(this.spanReporter, parent);
|
||||
this.tracer.close(context);
|
||||
assertEquals(2, this.test.spans.size());
|
||||
}
|
||||
|
||||
void logServerReceived(Span parent) {
|
||||
if (parent != null && parent.isRemote()) {
|
||||
parent.logEvent(Span.SERVER_RECV);
|
||||
}
|
||||
}
|
||||
|
||||
void logServerSent(SpanReporter spanReporter, Span parent) {
|
||||
if (parent != null && parent.isRemote()) {
|
||||
parent.logEvent(Span.SERVER_SEND);
|
||||
spanReporter.report(parent);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void nullSpanName() {
|
||||
Span span = this.tracer.createSpan(null);
|
||||
this.application.publishEvent(new ClientSentEvent(this, span));
|
||||
span.logEvent(Span.CLIENT_SEND);
|
||||
this.tracer.close(span);
|
||||
assertEquals(1, this.test.spans.size());
|
||||
this.listener.poll();
|
||||
@@ -135,6 +145,11 @@ public class StreamSpanListenerTests {
|
||||
public void handle(Message<?> msg) {
|
||||
}
|
||||
|
||||
@Bean
|
||||
SpanLogger spanLogger() {
|
||||
return new NoOpSpanLogger();
|
||||
}
|
||||
|
||||
@Bean
|
||||
public Sampler defaultSampler() {
|
||||
return new AlwaysSampler();
|
||||
|
||||
@@ -25,7 +25,7 @@ import org.mockito.InjectMocks;
|
||||
import org.mockito.Mock;
|
||||
import org.mockito.runners.MockitoJUnitRunner;
|
||||
import org.springframework.cloud.sleuth.Span;
|
||||
import org.springframework.cloud.sleuth.metric.SpanReporterService;
|
||||
import org.springframework.cloud.sleuth.metric.SpanMetricReporter;
|
||||
import org.springframework.integration.support.MessageBuilder;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.MessageChannel;
|
||||
@@ -40,7 +40,7 @@ import static org.mockito.Mockito.verifyZeroInteractions;
|
||||
public class TracerIgnoringChannelInterceptorTest {
|
||||
|
||||
@Mock MessageChannel messageChannel;
|
||||
@Mock SpanReporterService spanReporterService;
|
||||
@Mock SpanMetricReporter spanMetricReporter;
|
||||
@InjectMocks TracerIgnoringChannelInterceptor tracerIgnoringChannelInterceptor;
|
||||
|
||||
@Test
|
||||
@@ -59,7 +59,7 @@ public class TracerIgnoringChannelInterceptorTest {
|
||||
|
||||
this.tracerIgnoringChannelInterceptor.afterSendCompletion(message, this.messageChannel, true, null);
|
||||
|
||||
verifyZeroInteractions(this.spanReporterService);
|
||||
verifyZeroInteractions(this.spanMetricReporter);
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -71,7 +71,7 @@ public class TracerIgnoringChannelInterceptorTest {
|
||||
|
||||
this.tracerIgnoringChannelInterceptor.afterSendCompletion(message, this.messageChannel, true, null);
|
||||
|
||||
BDDMockito.then(this.spanReporterService).should().incrementAcceptedSpans(2);
|
||||
BDDMockito.then(this.spanMetricReporter).should().incrementAcceptedSpans(2);
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -83,6 +83,6 @@ public class TracerIgnoringChannelInterceptorTest {
|
||||
|
||||
this.tracerIgnoringChannelInterceptor.afterSendCompletion(message, this.messageChannel, false, null);
|
||||
|
||||
BDDMockito.then(this.spanReporterService).should().incrementDroppedSpans(2);
|
||||
BDDMockito.then(this.spanMetricReporter).should().incrementDroppedSpans(2);
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user