Clearing thread locals whenever an RSocket request arrives
This commit is contained in:
@@ -71,13 +71,19 @@ public class TracingRequesterRSocketProxy extends RSocketProxy {
|
||||
this.isZipkinPropagationEnabled = isZipkinPropagationEnabled;
|
||||
}
|
||||
|
||||
private void clearThreadLocal() {
|
||||
this.tracer.withSpan(null);
|
||||
}
|
||||
|
||||
@Override
|
||||
public Mono<Void> fireAndForget(Payload payload) {
|
||||
clearThreadLocal();
|
||||
return setSpan(super::fireAndForget, payload, FrameType.REQUEST_FNF);
|
||||
}
|
||||
|
||||
@Override
|
||||
public Mono<Payload> requestResponse(Payload payload) {
|
||||
clearThreadLocal();
|
||||
return setSpan(super::requestResponse, payload, FrameType.REQUEST_RESPONSE);
|
||||
}
|
||||
|
||||
@@ -138,11 +144,13 @@ public class TracingRequesterRSocketProxy extends RSocketProxy {
|
||||
|
||||
@Override
|
||||
public Flux<Payload> requestStream(Payload payload) {
|
||||
clearThreadLocal();
|
||||
return Flux.deferContextual(contextView -> setSpan(super::requestStream, payload, contextView));
|
||||
}
|
||||
|
||||
@Override
|
||||
public Flux<Payload> requestChannel(Publisher<Payload> inbound) {
|
||||
clearThreadLocal();
|
||||
return Flux.from(inbound).switchOnFirst((firstSignal, flux) -> {
|
||||
final Payload firstPayload = firstSignal.get();
|
||||
if (firstPayload != null) {
|
||||
|
||||
@@ -76,6 +76,7 @@ public class TracingResponderRSocketProxy extends RSocketProxy {
|
||||
|
||||
@Override
|
||||
public Mono<Void> 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<Payload> 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<Payload> 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<Payload> requestChannel(Publisher<Payload> payloads) {
|
||||
clearThreadLocal();
|
||||
return Flux.from(payloads).switchOnFirst((firstSignal, flux) -> {
|
||||
final Payload firstPayload = firstSignal.get();
|
||||
if (firstPayload != null) {
|
||||
|
||||
Reference in New Issue
Block a user