diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/annotation/ReactorSleuthMethodInvocationProcessor.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/annotation/ReactorSleuthMethodInvocationProcessor.java index b084a95e8..ebff3d4f0 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/annotation/ReactorSleuthMethodInvocationProcessor.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/annotation/ReactorSleuthMethodInvocationProcessor.java @@ -146,7 +146,9 @@ class ReactorSleuthMethodInvocationProcessor Span span; Tracer tracer = this.processor.tracer(); if (this.span == null) { - span = tracer.newTrace(); + // If we aren't continuing a trace from this flow, use nextSpan so that it + // can consider the "current span" (typically, backed by a thread-local) + span = tracer.nextSpan(); this.processor.newSpanParser().parse(this.invocation, this.newSpan, span); span.start(); } diff --git a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/annotation/SleuthSpanCreatorAspectFluxTests.java b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/annotation/SleuthSpanCreatorAspectFluxTests.java index 5f4a6da24..087c6c106 100644 --- a/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/annotation/SleuthSpanCreatorAspectFluxTests.java +++ b/spring-cloud-sleuth-core/src/test/java/org/springframework/cloud/sleuth/annotation/SleuthSpanCreatorAspectFluxTests.java @@ -25,6 +25,8 @@ import java.util.stream.Collectors; import brave.Span; import brave.Tracer; import brave.handler.SpanHandler; +import brave.propagation.CurrentTraceContext; +import brave.propagation.CurrentTraceContext.Scope; import brave.propagation.TraceContext; import brave.sampler.Sampler; import brave.test.TestSpanHandler; @@ -53,12 +55,18 @@ public class SleuthSpanCreatorAspectFluxTests { @Autowired TestBeanInterface testBean; + @Autowired + CurrentTraceContext currentTraceContext; + @Autowired Tracer tracer; @Autowired TestSpanHandler spans; + TraceContext context = TraceContext.newBuilder().traceId(1).spanId(1).sampled(true) + .build(); + private static String toHexString(Long value) { then(value).isNotNull(); return StringUtils.leftPad(Long.toHexString(value), 16, '0'); @@ -84,6 +92,20 @@ public class SleuthSpanCreatorAspectFluxTests { this.testBean.reset(); } + @Test + public void newSpan_shouldContinueExistingTrace() { + try (Scope scope = this.currentTraceContext.newScope(context)) { + Flux flux = this.testBean.testMethod(); + verifyNoSpansUntilFluxComplete(flux); + } + + Awaitility.await().untilAsserted(() -> { + then(this.spans).hasSize(1); + then(this.spans.get(0).traceId()).isEqualTo(context.traceIdString()); + then(this.spans.get(0).parentId()).isEqualTo(context.spanIdString()); + }); + } + @Test public void shouldCreateSpanWhenAnnotationOnInterfaceMethod() { Flux flux = this.testBean.testMethod();