diff --git a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceFunctionAroundWrapper.java b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceFunctionAroundWrapper.java index 2a831a919..493086d6f 100644 --- a/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceFunctionAroundWrapper.java +++ b/spring-cloud-sleuth-instrumentation/src/main/java/org/springframework/cloud/sleuth/instrument/messaging/TraceFunctionAroundWrapper.java @@ -170,14 +170,14 @@ public class TraceFunctionAroundWrapper extends FunctionAroundWrapper return Mono.deferContextual(contextView -> { MessageAndSpansAndScope msg = contextView.get(MessageAndSpansAndScope.class); return function.doOnNext(message -> { - msg.end(); - msg.handle(); - }).map(msgResult -> { - MessageAndSpan messageAndSpan = traceMessageHandler.wrapOutputMessage(msgResult, - msg.messageAndSpans.parentSpan, outputDestination(targetFunction.getFunctionDefinition())); - traceMessageHandler.afterMessageHandled(messageAndSpan.span, null); - return messageAndSpan.msg; - }) + msg.end(); + msg.handle(); + }).map(msgResult -> { + MessageAndSpan messageAndSpan = traceMessageHandler.wrapOutputMessage(msgResult, + msg.messageAndSpans.parentSpan, outputDestination(targetFunction.getFunctionDefinition())); + traceMessageHandler.afterMessageHandled(messageAndSpan.span, null); + return messageAndSpan.msg; + }) // TODO: Fix me when this is resolved in Reactor // .doOnSubscribe(__ -> scope.close()) .doOnError(msg::error).doFinally(signalType -> { @@ -221,14 +221,14 @@ public class TraceFunctionAroundWrapper extends FunctionAroundWrapper return Flux.deferContextual(contextView -> { MessageAndSpansAndScope msg = contextView.get(MessageAndSpansAndScope.class); return function.doOnNext(message -> { - msg.end(); - msg.handle(); - }).map(msgResult -> { - MessageAndSpan messageAndSpan = traceMessageHandler.wrapOutputMessage(msgResult, - msg.messageAndSpans.parentSpan, outputDestination(targetFunction.getFunctionDefinition())); - traceMessageHandler.afterMessageHandled(messageAndSpan.span, null); - return messageAndSpan.msg; - }) + msg.end(); + msg.handle(); + }).map(msgResult -> { + MessageAndSpan messageAndSpan = traceMessageHandler.wrapOutputMessage(msgResult, + msg.messageAndSpans.parentSpan, outputDestination(targetFunction.getFunctionDefinition())); + traceMessageHandler.afterMessageHandled(messageAndSpan.span, null); + return messageAndSpan.msg; + }) // TODO: Fix me when this is resolved in Reactor // .doOnSubscribe(__ -> scope.close()) .doOnError(msg::error).doFinally(signalType -> { @@ -272,9 +272,9 @@ public class TraceFunctionAroundWrapper extends FunctionAroundWrapper } Mono mono = (Mono) publisher; publisher = ReactorSleuth.tracedMono(tracer, tracer.currentTraceContext(), - targetFunction.getFunctionDefinition(), () -> mono, (msg, s) -> { - customizedInputMessageSpan(s, msg instanceof Message ? (Message) msg : null); - }).map(object -> toMessage(object)) + targetFunction.getFunctionDefinition(), () -> mono, (msg, s) -> { + customizedInputMessageSpan(s, msg instanceof Message ? (Message) msg : null); + }).map(object -> toMessage(object)) .map(object -> this.getMessageAndSpans((Message) object, targetFunction.getFunctionDefinition(), setNameAndTag(targetFunction, tracer.currentSpan()))) .doOnNext(wrappedOutputMessage -> customizedOutputMessageSpan( @@ -289,9 +289,9 @@ public class TraceFunctionAroundWrapper extends FunctionAroundWrapper } Flux flux = (Flux) publisher; publisher = ReactorSleuth.tracedFlux(tracer, tracer.currentTraceContext(), - targetFunction.getFunctionDefinition(), () -> flux, (msg, s) -> { - customizedInputMessageSpan(s, msg instanceof Message ? (Message) msg : null); - }).map(object -> toMessage(object)) + targetFunction.getFunctionDefinition(), () -> flux, (msg, s) -> { + customizedInputMessageSpan(s, msg instanceof Message ? (Message) msg : null); + }).map(object -> toMessage(object)) .map(object -> this.getMessageAndSpans((Message) object, targetFunction.getFunctionDefinition(), setNameAndTag(targetFunction, tracer.currentSpan()))) .doOnNext(wrappedOutputMessage -> customizedOutputMessageSpan(