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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 09:35:33