Java Reactor Core Flux并行处理元素:为何首元素阻塞全流程?
问题原因与修复方案
核心问题
你的代码里用了subscribeOn(Schedulers.newParallel(...)),但这个操作符的作用是把整个订阅流程(包括源元素的发射)放到指定线程,而非让每个元素的处理并行执行。
具体细节:
Flux.just是同步发射元素的,obj1和obj2会在subscribeOn指定的同一个线程里依次传递给flatMap。- 当obj1的
process方法进入无限循环时,该线程被完全占用,obj2根本没机会被处理,最终导致整个程序停滞。
修复代码
要实现真正的并行处理,需让每个元素的process任务在独立线程执行,同时给flatMap配置并发数:
Flux.just(obj1, obj2) .flatMap(obj -> Mono.fromCallable(() -> Transformator.process(obj)) .subscribeOn(Schedulers.newParallel("parallel", 5)), 5 // 允许同时处理的元素数量,和线程池大小匹配即可 )
修复说明
- 用
Mono.fromCallable把process包装成异步任务,再通过内部的subscribeOn让每个任务跑在并行线程池的独立线程,确保单个任务的阻塞不会影响其他任务。 flatMap的第二个参数指定并发数,明确告知Reactor可以同时处理多个元素,这样obj2的任务会立即在另一个线程启动,不会被obj1的无限循环卡住。
内容的提问来源于stack exchange,提问作者Khetag Abramov
相关产品推荐
相关产品推荐

