本文详解如何利用 reactor 的 subscribeon 显式指定调度器,使独立的异步任务(如并发获取多个资源)真正并行执行,避免默认串行阻塞行为。
本文详解如何利用 reactor 的 subscribeon 显式指定调度器,使独立的异步任务(如并发获取多个资源)真正并行执行,避免默认串行阻塞行为。
在使用 Project Reactor 构建响应式数据流时,一个常见误区是认为 flatMap 本身就能自动实现“真正并行”。实际上,flatMap 仅负责并发订阅(concurrent subscription),但是否真正并行执行,取决于每个内部 Mono/Flux 所绑定的执行上下文——即它在哪个线程或调度器(Scheduler)上运行。
在你提供的测试代码中,getJelly() 和 getPeanutButter() 均基于 Mono.fromSupplier() 创建,而 Supplier 是同步阻塞调用,且默认在上游线程(通常是 Schedulers.immediate())中执行。因此,尽管 flatMap 启动了 10 个内层流,它们仍被顺序地、在一个线程上依次执行:先为 sandwich 1 获取 jelly → 再获取 peanut butter → 组合 → 然后才轮到 sandwich 2……这导致输出严格有序,完全不符合预期的并行语义。
✅ 正确做法是:为每个独立的异步操作显式指定一个支持并发的调度器,例如 Schedulers.boundedElastic()(专为阻塞 I/O 设计,具备动态线程池与拒绝策略),从而让每个原料获取任务在不同线程上真正并发执行。
以下是修复后的关键代码段:
@Test
void testAsyncSandwich() {
Flux<string> sandwiches = Flux.fromStream(IntStream.rangeClosed(1, 10).boxed())
.flatMap(number ->
getJelly(number)
.subscribeOn(Schedulers.boundedElastic()) // ← 关键:解耦执行线程
.zipWith(
getPeanutButter(number)
.subscribeOn(Schedulers.boundedElastic()) // ← 同样应用
)
)
.map(this::makeSandwich);
StepVerifier.create(sandwiches)
.expectNextCount(10)
.verifyComplete();
}</string>
⚠️ 注意事项:
- 不要使用 publishOn() 替代 subscribeOn():publishOn 仅切换下游操作符的执行线程,无法改变 fromSupplier 这类源头的执行时机;只有 subscribeOn 能影响整个链路的订阅起点。
- 避免滥用 Schedulers.parallel() 处理阻塞调用:它适用于 CPU 密集型非阻塞任务;对 Thread.sleep() 或 JDBC 等真实阻塞操作,应选用 boundedElastic(),防止线程耗尽。
- 若实际场景使用的是真正的异步客户端(如 WebClient、R2DBC),通常其内部已自动调度到 IO 友好线程池,此时无需手动 subscribeOn —— 本例因模拟阻塞而需显式干预。
- 最终结果顺序不保证:并行执行后,sandwiches 的元素将按完成先后发出(即“完成序”,而非输入序)。如需保持原始顺序,可改用 concatMap + parallel().runOn(...) 组合,或在最后用 collectList() + sort() 后处理(但会丧失流式优势)。
总结:Reactor 的并行性不是自动发生的魔法,而是由调度器精确控制的契约行为。理解 subscribeOn 与 publishOn 的语义差异,并为阻塞操作选择合适的 Scheduler,是写出高效、真正并发响应式代码的关键基础。











