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

RxJava3中groupBy未按key分组,逐个发射GroupedFlowable问题

关于RxJava3中groupBy操作符的行为变化解析

问题原因

RxJava3对groupBy的实现逻辑做了关键调整:当下游没有及时订阅生成的GroupedFlowable时,上游会为每个同key的元素创建新的GroupedFlowable实例。而RxJava2中,groupBy会缓存同key的元素,直到对应的GroupedFlowable被订阅,不会重复创建实例。

你的测试代码中,cache().toList().blockingGet()的组合触发了这个变化:cache()订阅了groupBy的输出流,但toList()只是逐个收集GroupedFlowable实例,并没有立即订阅它们。groupBy的上游会认为之前的GroupedFlowable已经被丢弃,于是为下一个同key元素创建新的实例,最终得到6个同key的分组流。

解决方案

方案1:直接收集分组内的元素(推荐)

如果你的目标是获取分组后的元素列表,无需保留GroupedFlowable实例,可以用flatMapSingle先订阅每个分组流并收集元素:

fun runTest() {
    val eventDTOFlowable = Flowable.just(
        "item1", "item2", "item3", "item4", "item5", "item6"
    )
    val groupedResult = eventDTOFlowable.groupBy { _ -> 1 }
        .flatMapSingle { grouped -> 
            grouped.toList().map { grouped.key to it } 
        }
        .toList()
        .blockingGet()
    output = groupedResult.toString()
}

执行后会得到包含单个条目[(1, [item1, item2, item3, item4, item5, item6])]的列表,符合预期。

方案2:保留GroupedFlowable实例并确保订阅

如果必须保留GroupedFlowable实例,需要确保每个分组流被及时订阅,让上游知道该key的流仍在活跃:

fun runTest() {
    val eventDTOFlowable = Flowable.just(
        "item1", "item2", "item3", "item4", "item5", "item6"
    )
    val groupedList = mutableListOf<GroupedFlowable<Int, String>>()
    eventDTOFlowable.groupBy { _ -> 1 }
        .concatMap { grouped ->
            groupedList.add(grouped)
            grouped.ignoreElements() // 订阅分组流,维持其活跃状态
        }
        .blockingAwait()
    output = groupedList.toString() // 此时列表仅含1个GroupedFlowable实例
}

核心变化总结

RxJava3的groupBy调整是为了强化背压管理和避免内存泄漏:不再为未被订阅的分组流缓存元素,而是通过创建新实例的方式,强制下游及时处理每个分组流,避免因长期未订阅的分组流占用内存。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 18:22:12