You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

为何额外的log()会影响fromCallable()使用的线程池?

为什么Reactor中log()会改变Mono的线程执行行为?

先看你给出的示例代码:

Mono.fromCallable(() -> calculate())
//.log()
  .publishOn(Schedulers.boundedElastic())
  .doOnNext(i -> System.out.println("thread: " + Thread.currentThread().getName()))
  .block();

核心原因在于Reactor的订阅信号传递逻辑,以及publishOn和log()操作符的实现差异:

  • Reactor是懒加载模型,只有调用block()时才会触发订阅,订阅信号是从下游(block())向上游(fromCallable)传递的。
  • publishOn的特殊行为:它不仅负责切换数据向下游流动的线程,还会在订阅信号向上游传递时,把后续的订阅处理(包括上游生产者的执行)切换到自身指定的调度器线程。
    • 当没有log()时,publishOn是最靠近下游的操作符:订阅信号从main线程(block()所在线程)传到publishOn后,publishOn会把订阅请求提交到boundedElastic线程池,由该线程触发fromCallable的执行,所以calculate()运行在boundedElastic线程。
  • 加上log()后,log()位于publishOn和fromCallable之间:log()是一个无线程切换的操作符,它会在当前线程(main)直接处理订阅信号并传递给上游的fromCallable,因此calculate()在main线程执行。数据产生后,publishOn仍会把数据切换到boundedElastic线程,所以doOnNext的打印还是在该线程。

如果想让calculate()的执行线程不受这类操作符顺序影响,建议用subscribeOn专门指定上游的执行线程,比如:

Mono.fromCallable(() -> calculate())
    .subscribeOn(Schedulers.boundedElastic())
    .doOnNext(i -> System.out.println("thread: " + Thread.currentThread().getName()))
    .block();

内容的提问来源于stack exchange,提问作者Grzegorz Piwowarek

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.10 02:27:18