为何额外的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
相关产品推荐
相关产品推荐

