From 4d9a7e311fcc86006884efdf31be868c165bb68f Mon Sep 17 00:00:00 2001 From: Marcin Grzejszczak Date: Wed, 16 May 2018 17:25:36 +0200 Subject: [PATCH] Polish --- .../instrument/reactor/ReactorSleuth.java | 38 +++++++++++++------ 1 file changed, 26 insertions(+), 12 deletions(-) diff --git a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/ReactorSleuth.java b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/ReactorSleuth.java index df5440731..b959cd0b6 100644 --- a/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/ReactorSleuth.java +++ b/spring-cloud-sleuth-core/src/main/java/org/springframework/cloud/sleuth/instrument/reactor/ReactorSleuth.java @@ -17,7 +17,6 @@ package org.springframework.cloud.sleuth.instrument.reactor; import java.util.function.Function; -import java.util.function.Predicate; import brave.Tracing; import org.apache.commons.logging.Log; @@ -55,9 +54,11 @@ public abstract class ReactorSleuth { * * @return a new lazy span operator pointcut */ + @SuppressWarnings("unchecked") public static Function, ? extends Publisher> spanOperator( BeanFactory beanFactory) { return sourcePub -> { + // TODO: Remove this once Reactor 3.1.8 is released //do the checks directly on actual original Publisher if (sourcePub instanceof Fuseable //Sleuth can't handle that || sourcePub instanceof Fuseable.ScalarCallable //IIRC Sleuth can't handle that @@ -66,21 +67,36 @@ public abstract class ReactorSleuth { ) { return sourcePub; } - //no more POINTCUT_FILTER since mecanism is broken + //no more POINTCUT_FILTER since mechanism is broken Function, ? extends Publisher> lift = Operators.lift((scannable, sub) -> { + if (contextRefreshed(beanFactory)) { + if (log.isTraceEnabled()) { + log.trace("Spring Context already refreshed. Creating a Sleuth span subscriber with Reactor Context " + "[" + sub.currentContext() + "] and name [" + scannable.name() + "]"); + } + return spanSubscriptionProvider(beanFactory, scannable, sub).get(); + } + if (log.isTraceEnabled()) { + log.trace( + "Spring Context is not yet refreshed, falling back to lazy span subscriber. Reactor Context is [" + sub.currentContext() + "] and name is [" + scannable.name() + "]"); + } //rest of the logic unchanged... return new LazySpanSubscriber( - new SpanSubscriptionProvider( - beanFactory, - sub, - sub.currentContext(), - scannable.name()) + spanSubscriptionProvider(beanFactory, scannable, sub) ); }); return lift.apply(sourcePub); }; } + private static SpanSubscriptionProvider spanSubscriptionProvider( + BeanFactory beanFactory, Scannable scannable, CoreSubscriber sub) { + return new SpanSubscriptionProvider( + beanFactory, + sub, + sub.currentContext(), + scannable.name()); + } + /** * Return a span operator pointcut given a {@link Tracing}. This can be used in reactor * via {@link reactor.core.publisher.Flux#transform(Function)}, {@link @@ -98,6 +114,7 @@ public abstract class ReactorSleuth { public static Function, ? extends Publisher> scopePassingSpanOperator( BeanFactory beanFactory) { return sourcePub -> { + // TODO: Remove this once Reactor 3.1.8 is released //do the checks directly on actual original Publisher if (sourcePub instanceof Fuseable //Sleuth can't handle that || sourcePub instanceof Fuseable.ScalarCallable //IIRC Sleuth can't handle that @@ -106,7 +123,7 @@ public abstract class ReactorSleuth { ) { return sourcePub; } - //no more POINTCUT_FILTER since mecanism is broken + //no more POINTCUT_FILTER since mechanism is broken Function, ? extends Publisher> lift = Operators.lift((scannable, sub) -> { //rest of the logic unchanged... if (contextRefreshed(beanFactory)) { @@ -117,7 +134,7 @@ public abstract class ReactorSleuth { } if (log.isTraceEnabled()) { log.trace( - "Spring Context is not yet refreshed, falling back to lazy span subscriber. " + "Reactor Context is [" + sub.currentContext() + "] and name is [" + scannable.name() + "]"); + "Spring Context is not yet refreshed, falling back to lazy span subscriber. Reactor Context is [" + sub.currentContext() + "] and name is [" + scannable.name() + "]"); } return new LazySpanSubscriber( scopePassingSpanSubscription(beanFactory, scannable, sub) @@ -152,9 +169,6 @@ public abstract class ReactorSleuth { }; } - private static final Predicate POINTCUT_FILTER = - s -> !(s instanceof ConnectableFlux) && !(s instanceof Fuseable.ScalarCallable) && s.isScanAvailable(); - private ReactorSleuth() { } } \ No newline at end of file