Reactor中subscribeOn配置10线程池却仅用单线程的原因?
关于subscribeOn仅使用单线程的疑问与解析
问题场景
在使用Reactor的subscribeOn时,发现即使指定了大小为10的线程池,整个数据流的处理依然只使用单个线程,与预期不符:
测试代码(subscribeOn示例)
@Test void fluxWithSubscribeOnTest() { final Scheduler s = Schedulers.newParallel("parallel", 10); final Publisher<Integer> integerFlux = Flux .range(1, 1000) .doOnNext(integer -> log.info("executed")) .subscribeOn(s); StepVerifier.create(integerFlux) .expectNextCount(1000) .verifyComplete(); }
输出日志
18:39:47.612 [parallel-1] INFO com.performance.TestClass - executed 18:39:47.613 [parallel-1] INFO com.performance.TestClass - executed 18:39:47.613 [parallel-1] INFO com.performance.TestClass - executed 18:39:47.613 [parallel-1] INFO com.performance.TestClass - executed 18:39:47.613 [parallel-1] INFO com.performance.TestClass - executed 18:39:47.613 [parallel-1] INFO com.performance.TestClass - executed 18:39:47.613 [parallel-1] INFO com.performance.TestClass - executed 18:39:47.613 [parallel-1] INFO com.performance.TestClass - executed
而使用ParallelFlux时,却能正常使用线程池中的多个线程:
测试代码(ParallelFlux示例)
@Test void parallelFluxTest() { final Scheduler s = Schedulers.newParallel("parallel", 10); final Publisher<Integer> integerFlux = Flux .range(1, 1000) .parallel(10) .runOn(s) .doOnNext(integer -> log.info("executed")); StepVerifier.create(integerFlux) .expectNextCount(1000) .verifyComplete(); }
输出日志
18:43:42.377 [parallel-1] INFO com.performance.TestClass - executed 18:43:42.377 [parallel-2] INFO com.performance.TestClass - executed 18:43:42.377 [parallel-3] INFO com.performance.TestClass - executed 18:43:42.377 [parallel-4] INFO com.performance.TestClass - executed 18:43:42.378 [parallel-5] INFO com.performance.TestClass - executed
疑问:为何subscribeOn示例中线程池大小为10,却仅使用单个线程?
原因解析
1. subscribeOn的本质作用
subscribeOn的核心是指定订阅动作启动的线程,以及后续整个上游数据流的执行线程,但它不会改变数据流本身的执行模式:
Flux.range是一个同步、连续的数据源,它会在单个线程中依次生成1到1000的所有元素,整个数据流的处理是一个连续的任务。- 线程池的作用只是提供一个线程来承载这个连续任务,而非将每个元素拆分到不同线程执行。即使线程池有10个线程,也只会占用其中一个来完成全部元素的生成与处理。
2. ParallelFlux的多线程逻辑
ParallelFlux实现多线程处理的核心是数据流拆分:
.parallel(10)会将原Flux拆分为10个并行的子数据流,每个子数据流处理一部分元素。.runOn(s)则指定每个子数据流在调度器的不同线程上执行,因此每个子数据流的元素处理会占用线程池中的不同线程,最终呈现多线程输出的效果。
3. 何时subscribeOn会用到多线程
只有当数据流本身包含异步拆分逻辑(比如使用flatMap开启多个异步任务)时,subscribeOn指定的线程池才会被多个线程占用,因为此时会有多个独立的任务需要线程池调度执行。
内容的提问来源于stack exchange,提问作者Piotr Cierpich
相关产品推荐
相关产品推荐

