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

Reactor GroupedFlux性能依赖副作用?反直觉测试现象问询

为什么取消注释空的doOnNext会让Flux groupBy+flatMap的执行速度提升3倍?

问题背景

我在测试Reactor聚合逻辑时遇到一个反直觉的现象:取消注释count()后空的doOnNext(n -> {})操作,测试执行速度直接快了3倍。测试代码如下:

@Test
void groupByFluxAggregateLoadTest() {
    int numberOfGroups = 30_000;
    Flux
            .range(0, numberOfGroups * 10)
            .groupBy(i -> i % numberOfGroups)
            .flatMap(
                integerIntegerGroupedFlux -> integerIntegerGroupedFlux
                    .count()
                    //.doOnNext(n -> {})
                ,
                numberOfGroups + 1
            )
            .as(StepVerifier::create)
            .expectNextCount(numberOfGroups)
            .verifyComplete();
}

原因解析

这个现象的核心是Reactor对同步结果Mono的特殊优化,以及flatMap针对不同Mono类型的调度行为差异:

  • 无doOnNext时:串行执行的同步计算
    因为Flux.range是同步发射元素的源,GroupedFlux.count()会直接同步遍历分组内的10个元素,计算出数量后返回一个ScalarMono——这是Reactor专门用于包装已知确定结果的Mono实现。
    虽然你给flatMap设置了numberOfGroups + 1的并发数(理论上可以并行处理所有分组),但ScalarMono的订阅逻辑是同步阻塞的:flatMap会逐个订阅这些Mono,每个订阅都会在当前测试线程上立即完成count计算,完成一个再处理下一个。30000个分组的计算完全串行,自然速度慢。

  • 添加doOnNext后:触发多线程并行处理
    空的doOnNext会把ScalarMono包装成普通的MonoDoOnNext,打破了它的同步Scalar特性。
    对于这种普通Mono,flatMap会切换为异步调度模式:它会把每个分组的count计算任务放入内部队列,Reactor的默认调度器可以利用多线程同时处理多个任务,真正实现30000个分组的并行计算。
    虽然doOnNext本身没做任何实际操作,但它改变了Mono的类型,触发了flatMap的并行调度逻辑,从而让执行效率提升数倍。

补充说明

Reactor对ScalarMono的同步优化是为了减少调度开销,但在高并发场景下,这种优化反而会导致无法利用多核CPU的并行能力。添加无意义的doOnNext相当于“绕过”了这个优化,让任务进入并行调度流程,这才出现了看似反直觉的速度提升。

内容的提问来源于stack exchange,提问作者PeMa

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 06:50:22