Flux.groupBy相比列表partition的优势差异及适用场景对比
Flux
groupBy 与全量收集后分组的差异及适用场景 两种方案最终都能实现按规则拆分数据的效果,但底层执行逻辑完全不同,适用场景边界非常清晰:
Flux groupBy 的核心特性
groupBy 是响应式流原生的流式分组算子,核心优势来自于不需要等待全量数据到达:
- 内存效率更高:元素从上游发出后会直接按分组key路由到对应的
GroupedFlux,不需要提前把所有数据加载到内存暂存,哪怕是无限流(比如实时消费的MQ消息、日志流、传感器上报数据)也能正常运行,不会因为数据总量过大触发OOM。 - 延迟更低:分组后不需要等所有数据收齐就能处理,比如你可以直接对每个
GroupedFlux做逐元素消费、滑动窗口聚合、限流等操作,第一个元素到达对应分组后就可以开始处理,不需要等整个流结束。 - 完整保留响应式特性:整个分组流程都在响应式链路内,背压可以从下游一直传递到最上游,自动匹配上下游的生产消费速度,不会出现消费不及时压垮服务的问题。
- 注意:示例中
groupBy后接collectList的写法其实是把算子的流式优势抵消了——这种写法最终还是要等每个组收齐所有元素才会输出,和全量分组的最终输出结果一致,只有在分组后做流式处理、或者对接大体量/无限数据源时,groupBy的价值才会完全体现。另外使用时要注意分组key的基数不要无限制增长,每个活跃分组都会占用一定资源,key过多的场景要配置分组过期逻辑。
先全量收集再用partition/Stream groupingBy的核心特性
这种方案是典型的同步批量处理逻辑:
- 逻辑简单直观:对于已经全部加载到内存的小体量数据集,代码写起来更直白,调试成本低,没有响应式API的心智负担。
- 小数据量下开销更低:如果数据源本身就是内存里的固定集合(比如示例中0-9的数字列表),不需要走响应式流的事件路由、分组订阅等额外逻辑,同步分组的执行速度更快。
- 局限性非常明显:只能处理能全部放进内存的有限数据集,必须等所有数据加载完成才能开始分组处理,内存占用和数据总量正相关,无法对接实时无限流,也没有原生背压能力。
场景选择参考
优先选
Flux.groupBy的场景:
- 数据源是无限的实时流
- 数据体量太大,无法全部加载到内存
- 分组后不需要依赖全量数据,要做实时处理
- 全链路需要保持响应式特性,依赖背压做流量控制
优先选全量收集后分组的场景:
- 数据源本身就是已经全量驻留内存的小体量有限集合
- 分组后的处理逻辑必须依赖组内全量数据(比如全局排序、全量聚合计算)
- 同步非响应式代码场景,不需要考虑异步流、背压能力
示例代码
@Test fun groupByTwoLists() { val numbers = (0..9).toList() Flux.fromIterable(numbers) .groupBy<Boolean> { it % 2 == 0 } .flatMap { group: GroupedFlux<Boolean, Int> -> if (group.key() == true) { group.collectList().doOnNext { println("We are even") } } else { group.collectList().doOnNext { println("We are odd") } } } .test() .expectNextCount(2) .verifyComplete() } @Test fun partitionLists() { val numbers = (0..9).toList() val (even, odd) = numbers.partition { it % 2 == 0 } val monoEven = Mono.just(even).doOnNext { println("We are even") } val monoOdd = Mono.just(odd).doOnNext { println("We are odd") } Flux.merge(monoEven, monoOdd) .test() .expectNextCount(2) .verifyComplete() }
内容的提问来源于stack exchange,提问作者Dawn
相关产品推荐
相关产品推荐

