diff --git a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/rsocket/TracingRequesterRSocketProxy.java b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/rsocket/TracingRequesterRSocketProxy.java index 90e3d6ec1..b0b9d0f1f 100644 --- a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/rsocket/TracingRequesterRSocketProxy.java +++ b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/rsocket/TracingRequesterRSocketProxy.java @@ -71,13 +71,19 @@ public class TracingRequesterRSocketProxy extends RSocketProxy { this.isZipkinPropagationEnabled = isZipkinPropagationEnabled; } + private void clearThreadLocal() { + this.tracer.withSpan(null); + } + @Override public Mono fireAndForget(Payload payload) { + clearThreadLocal(); return setSpan(super::fireAndForget, payload, FrameType.REQUEST_FNF); } @Override public Mono requestResponse(Payload payload) { + clearThreadLocal(); return setSpan(super::requestResponse, payload, FrameType.REQUEST_RESPONSE); } @@ -138,11 +144,13 @@ public class TracingRequesterRSocketProxy extends RSocketProxy { @Override public Flux requestStream(Payload payload) { + clearThreadLocal(); return Flux.deferContextual(contextView -> setSpan(super::requestStream, payload, contextView)); } @Override public Flux requestChannel(Publisher inbound) { + clearThreadLocal(); return Flux.from(inbound).switchOnFirst((firstSignal, flux) -> { final Payload firstPayload = firstSignal.get(); if (firstPayload != null) { diff --git a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/rsocket/TracingResponderRSocketProxy.java b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/rsocket/TracingResponderRSocketProxy.java index 7cbc8e7df..122b6cb57 100644 --- a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/rsocket/TracingResponderRSocketProxy.java +++ b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/rsocket/TracingResponderRSocketProxy.java @@ -76,6 +76,7 @@ public class TracingResponderRSocketProxy extends RSocketProxy { @Override public Mono fireAndForget(Payload payload) { + clearThreadLocal(); // called on Netty EventLoop // there can't be trace context in thread local here Span handle = consumerSpanBuilder(payload.sliceMetadata(), FrameType.REQUEST_FNF); @@ -86,8 +87,13 @@ public class TracingResponderRSocketProxy extends RSocketProxy { return ReactorSleuth.tracedMono(this.tracer, handle, () -> super.fireAndForget(newPayload)); } + private void clearThreadLocal() { + this.tracer.withSpan(null); + } + @Override public Mono requestResponse(Payload payload) { + clearThreadLocal(); Span handle = consumerSpanBuilder(payload.sliceMetadata(), FrameType.REQUEST_RESPONSE); if (log.isDebugEnabled()) { log.debug("Created consumer span " + handle); @@ -98,6 +104,7 @@ public class TracingResponderRSocketProxy extends RSocketProxy { @Override public Flux requestStream(Payload payload) { + clearThreadLocal(); Span handle = consumerSpanBuilder(payload.sliceMetadata(), FrameType.REQUEST_STREAM); if (log.isDebugEnabled()) { log.debug("Created consumer span " + handle); @@ -108,6 +115,7 @@ public class TracingResponderRSocketProxy extends RSocketProxy { @Override public Flux requestChannel(Publisher payloads) { + clearThreadLocal(); return Flux.from(payloads).switchOnFirst((firstSignal, flux) -> { final Payload firstPayload = firstSignal.get(); if (firstPayload != null) {