如何条件性使用.publishOn(boundedElastic())?验证Reactor实现方案
我有一段响应式代码链:
callUpstream() // 返回Mono .publishOn(boundedElastic()) .doOnNext(this::sendJmsBlocking) // 该操作是阻塞的,必须在boundedElastic调度器的工作线程执行
99%的场景下,callUpstream()返回的信号已经在boundedElastic线程上执行,此时额外切换到同调度器的线程会在高负载下导致任务入队,完全没必要。但偶尔上游线程可能是WebClient的事件循环线程或并行线程,这会导致非阻塞线程被阻塞。
我想要替换掉.publishOn(boundedElastic()),实现逻辑:如果当前线程不是boundedElastic线程,则切换到该线程;否则继续使用当前线程。
尝试的无效方案
一开始用transformDeferred实现了判断逻辑:
private <T> Mono<T> ensureBoundedElasticThread(Mono<T> mono) { if (!isBoundedElastic()) { return mono.publishOn(Schedulers.boundedElastic()); } return mono; } private static boolean isBoundedElastic() { return Thread.currentThread().getName().startsWith("boundedElastic-"); }
然后改写链:
callUpstream() // 返回Mono .transformDeferred(this::ensureBoundedElasticThread) .doOnNext(this::sendJmsBlocking)
但这个方案无效,因为transformDeferred是在订阅阶段执行的,不是信号流动的运行阶段,导致判断逻辑总是在订阅线程(比如示例里的main线程)执行,每次都会触发线程切换。
对应的示例代码(仅修改postProcess方法):
import java.util.function.Consumer; import reactor.core.publisher.Mono; import reactor.core.scheduler.Schedulers; public class TransferredNotWorkingMain { static Consumer<Mono<Integer>> outboundMessageHandler = TransferredNotWorkingMain::postProcess; private static Mono<Integer> route(String InboundTextMessage) { return Mono.just(Integer.valueOf(InboundTextMessage)) .flatMap(i -> process(i).subscribeOn(Schedulers.boundedElastic())); } private static Mono<Integer> process(Integer in) { return Mono.just(in) .map(integer -> { System.out.println("start process: " + Thread.currentThread().getName()); try { Thread.sleep(2000); } catch (InterruptedException e) { throw new RuntimeException(e); } System.out.println("end process: " + Thread.currentThread().getName()); return integer; }); } static void send(Integer message) { System.out.println("start sending: " + Thread.currentThread().getName()); try { Thread.sleep(1000); } catch (InterruptedException e) { throw new RuntimeException(e); } System.out.println("end sending: " + Thread.currentThread().getName()); } static void postProcess(Mono<Integer> output) { output .transformDeferred(TransferredNotWorkingMain::ensureBoundedElasticThread)// 订阅阶段执行,总是切换线程 .subscribe(TransferredWorkingMain::send, e -> System.out.println("Unexpected error occurred during message sending" + e.getMessage())); } private static <T> Mono<T> ensureBoundedElasticThread(Mono<T> mono) { System.out.println("ensureBoundedElasticThread: " + Thread.currentThread().getName()); if (!isBoundedElastic()) { return mono.publishOn(Schedulers.boundedElastic()); } return mono; } private static boolean isBoundedElastic() { return Thread.currentThread().getName().startsWith("boundedElastic-"); } public static void main(String[] args) throws InterruptedException { Mono<Integer> outboundTextMessage = route("1"); outboundMessageHandler.accept(outboundTextMessage); Thread.sleep(5000); } }
日志输出:
ensureBoundedElasticThread: main start process: boundedElastic-1 end process: boundedElastic-1 start sending: boundedElastic-2 end sending: boundedElastic-2
可以看到,即使上游已经在boundedElastic线程,还是切换到了另一个boundedElastic线程,不符合预期。
可行的第二个方案
改用flatMap包裹阻塞操作,结合transformDeferred和subscribeOn实现:
import java.util.function.Consumer; import reactor.core.publisher.Mono; import reactor.core.scheduler.Schedulers; public class TransferredWorkingMain { static Consumer<Mono<Integer>> outboundMessageHandler = TransferredWorkingMain::postProcess; private static Mono<Integer> route(String message) { return Mono.just(Integer.valueOf(message)) .flatMap(i -> process(i).subscribeOn(Schedulers.boundedElastic())); } private static Mono<Integer> process(Integer in) { return Mono.just(in) .map(integer -> { System.out.println("start process: " + Thread.currentThread().getName()); try { Thread.sleep(2000); } catch (InterruptedException e) { throw new RuntimeException(e); } System.out.println("end process: " + Thread.currentThread().getName()); return integer; }); } static void send(Integer message) { System.out.println("start sending: " + Thread.currentThread().getName()); try { Thread.sleep(1000); } catch (InterruptedException e) { throw new RuntimeException(e); } System.out.println("end sending: " + Thread.currentThread().getName()); } static void postProcess(Mono<Integer> output) { output .flatMap(o -> Mono.fromRunnable(() -> send(o)) .transformDeferred(TransferredWorkingMain::ensureBoundedElasticThread)) .subscribe(null, e -> System.out.println("Unexpected error occurred during message sending" + e.getMessage())); } private static <T> Mono<T> ensureBoundedElasticThread(Mono<T> mono) { System.out.println("ensureBoundedElasticThread: " + Thread.currentThread().getName()); if (!isBoundedElastic()) { return mono.subscribeOn(Schedulers.boundedElastic()); } return mono; } private static boolean isBoundedElastic() { return Thread.currentThread().getName().startsWith("boundedElastic-"); } public static void main(String[] args) throws InterruptedException { Mono<Integer> outboundTextMessage = route("1"); outboundMessageHandler.accept(outboundTextMessage); Thread.sleep(5000); } }
日志输出:
start process: boundedElastic-1 end process: boundedElastic-1 ensureBoundedElasticThread: boundedElastic-1 start sending: boundedElastic-1 end sending: boundedElastic-1
符合预期:上游已经在boundedElastic线程时,直接复用该线程执行阻塞的send操作,没有额外切换。
问题
这个第二个方案是否合理?
这个方案是合理且符合需求的,原因如下:
- 线程判断时机正确:
flatMap的lambda是在信号流动的运行阶段执行的,也就是上游信号所在的线程,此时调用ensureBoundedElasticThread能准确判断当前线程是否为boundedElastic线程,避免了订阅阶段判断的错误。 - 避免不必要的线程切换:当上游已经在boundedElastic线程时,直接复用该线程执行阻塞操作,不会触发额外的调度,减少了高负载下的任务排队风险。
- 保证阻塞操作的线程安全:当上游在非阻塞线程(如事件循环线程)时,通过
subscribeOn(Schedulers.boundedElastic())将阻塞操作切换到boundedElastic线程,避免阻塞非阻塞线程池,符合Reactor的最佳实践。
可优化点
- 替换线程名判断为更可靠的方式:通过线程名前缀判断不够健壮(比如自定义调度器线程名可能修改),可以改用Reactor提供的
Scheduler.isInCurrentThread()方法:
private static boolean isBoundedElastic() { return Schedulers.boundedElastic().isInCurrentThread(); }
这个方法会直接检查当前线程是否属于boundedElastic调度器的工作线程,更准确可靠。
2. 简化逻辑:可以去掉transformDeferred,直接在flatMap里判断,代码更简洁:
static void postProcess(Mono<Integer> output) { output .flatMap(o -> { Runnable sendTask = () -> send(o); if (Schedulers.boundedElastic().isInCurrentThread()) { return Mono.fromRunnable(sendTask); } else { return Mono.fromRunnable(sendTask).subscribeOn(Schedulers.boundedElastic()); } }) .subscribe(null, e -> System.out.println("Unexpected error occurred during message sending" + e.getMessage())); }
这样可以减少一层操作符的嵌套,逻辑更直观。
内容的提问来源于stack exchange,提问作者metalpalo

