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
相关产品推荐
相关产品推荐

