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

