Reactor中subscribeOn未按指定调度器执行的问题求助
问题分析与解决方案
核心原因:subscribeOn 和 publishOn 的本质区别
先把这两个操作符的作用讲透:
- subscribeOn:只负责指定「订阅发起阶段」的线程,也就是调用
subscribe()时,从下游向上游建立订阅链路的线程。整个流里只有第一个subscribeOn会生效,后续的subscribeOn无法改变已经确定的订阅线程。 - publishOn:负责指定「数据向下游流动时」的线程,从
publishOn之后的所有操作(map、buffer、subscribe的onNext等)都会切换到指定线程执行。
结合你的代码分析
你的代码里,couchbaseCollection.reactive().get(docId)是核心数据源,而Couchbase的Reactive API本身已经把这个操作绑定到了它自己的cb-io-kv线程池。这时候:
- 流启动订阅时,订阅链路的线程已经被Couchbase的get操作确定为
cb-io-kv线程,你在bufferTimeout之后加的subscribeOn根本不会生效——因为subscribeOn只认第一个生效的实例,且你的subscribeOn位置太靠下,影响不到上游的数据产生线程。 - 换成
publishOn后,它会在数据流过这个节点时,把线程切换到你指定的processorThread,所以后续的bufferTimeout和subscribe的onNext自然就跑在预期线程里了。
正确写法建议
如果你希望processDocument以及后续的buffer、订阅逻辑都跑在指定线程,应该把publishOn放在map(processDocument)之前,确保整个下游操作都切换线程:
Flux<String> documentIds = getDocumentIds(); documentIds .windowTimeout(WINDOW, Duration.ofMillis(300)) .onBackpressureBuffer() .flatMap( obs -> obs.concatMap( docId -> couchbaseCollection.reactive().get(docId, getOptions).retryWhen(retryFunction)), 512) .takeWhile(d -> !interrupted.get()) .publishOn(Schedulers.fromExecutor(processorThread)) // 在这里切换线程,后续操作统一使用该线程 .map(doc -> processDocument(doc)) .bufferTimeout(this.commitBatch, Duration.ofSeconds(11101)) .subscribe( new Subscriber<>() { // 此处的onNext会在processorThread线程执行 } );
额外提醒
- 永远记住:subscribeOn管「订阅发起」的线程,publishOn管「数据流动」的线程。
- 如果想切换某一段操作的执行线程,用publishOn;如果想指定自定义数据源的执行线程,才用subscribeOn,且要放在流的最上游。
内容的提问来源于stack exchange,提问作者Zizou
相关产品推荐
相关产品推荐

