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

Reactor Netty使用Flux::groupBy出现永久冻结问题排查与优化

基于ID合并两个Reactor列表的冻结问题分析与优化

问题场景

我有两个包含不同类型对象但均带有相同ID属性的列表:

list1 => [{ id: "123", x: "xxx" }] | list2 => [{ id: "123", y: "yyy" }, {id: "456", y: "yyy"}]

已知list1的所有ID都存在于list2中,且两个列表均未排序。尝试用Reactor的groupBy方法将同ID的Object1和Object2合并为Object3,实现代码如下:

// list_1 data class Object1(val id: String, val x: String)
val flux1 = Flux.fromIterable(list_1)
                .map { obj -> obj.id to obj }
// list_2 data class Object2(val id: String, val y: String)
val flux2 = Flux.fromIterable(list_2) // 修正原代码笔误fromIteablle
                .map { obj -> obj.id to obj }

Flux.merge(flux1, flux2)
    .groupBy { (id, obj) -> id }
    .flatMap { gFlux ->
      gFlux
        .map { (id, obj) -> obj }
        .collectList()
        .filter { it.size == 2 }
        .map { (obj1, obj2) -> Object3(obj1, obj2) }
    }
    .collectList()

当list2数据量较大时,程序出现永久冻结。临时通过分组前过滤list1中不存在的ID解决:

Flux.merge(flux1, flux2)
    .filter { (id, obj) -> id in list_1.map { it.id } }

冻结原因

  • groupBy的资源持有机制:Reactor的groupBy会为每个唯一ID创建对应的GroupedFlux,且会持续持有该子流直到上游所有数据发射完毕。当list2包含大量list1中没有的ID时,这些ID对应的GroupedFlux会一直处于等待状态(永远凑不齐2个元素),无法被销毁。
  • 并发资源耗尽:flatMap默认并发数为256,大量无效的GroupedFlux会占满并发槽位,导致有效分组的处理被阻塞,最终整个流无法完成,出现永久冻结。

解决方案与优化方案

方案1:优化提前过滤逻辑

临时方案思路正确,但list_1.map { it.id }每次过滤都会生成新列表,效率低下。先将list1的ID存入HashSet,提升判断效率:

val validIds = list_1.map { it.id }.toHashSet()

Flux.merge(flux1, flux2)
    .filter { (id, obj) -> validIds.contains(id) }
    .groupBy { (id, obj) -> id }
    .flatMap { gFlux ->
      gFlux
        .map { (id, obj) -> obj }
        .collectList()
        .map { list ->
          // 已知list1的ID都存在于list2,因此list必然包含2个元素
          Object3(list[0], list[1])
        }
    }
    .collectList()

方案2:使用Reactor原生join操作符

join专门用于基于关联条件合并两个流,无需处理无效分组,更贴合一对一匹配的需求:

val flux1 = Flux.fromIterable(list_1)
val flux2 = Flux.fromIterable(list_2)

flux1.join(
    flux2,
    { Flux.never() }, // 为每个Object1永久等待匹配(符合场景需求)
    { Flux.never() }, // 为每个Object2永久等待匹配
    { obj1, obj2 -> Object3(obj1, obj2) } // 匹配条件:ID相同(注:需确保流中ID唯一,或根据实际调整)
)
.collectList()

若担心内存占用,可设置超时时间(如Flux.delayMillis(1000)),但根据场景list1的ID都存在于list2,用Flux.never()即可。

方案3:直接使用集合操作(内存数据最优解)

如果数据已全部在内存中,无需使用Reactor流操作,直接通过集合映射更高效:

val obj2Map = list_2.associateBy { it.id }
val result = list_1.map { obj1 ->
    Object3(obj1, obj2Map[obj1.id]!!)
}

这种方式避免了Reactor流的额外开销,性能最优。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 13:53:22