如何在Kotlin Flows中实现值累加并按指定大小批量发射?
问题描述
我有一个整数类型的Flow,需要将其中的值累加到列表中,当列表达到指定大小(如5)时,将该列表发射出去。目前使用runningFold可实现类似效果,但存在列表满后无法自动重建的问题,只能手动清空列表,这种方式不够优雅。例如,希望将输入[1,2,3,4,5,6,7,8,9,10]转换为[1,2,3,4,5]和[6,7,8,9,10]两个列表发射,请问最优实现方式是什么?
当前实现代码
fun main() = runBlocking { launch { repeat(20) { emitValue(it) } } flow .collect { // it should be of type Result println(it) } } private val sharedFlow = MutableSharedFlow<Int>() val flow = sharedFlow.asSharedFlow() .onEach { println("received: $it") } .runningFold(mutableListOf<Int>()) { list, value -> list.add(value) list } .filter { it.size >= 5 } .map { val result = Result(it.toList()) it.clear() // that isn't very nice result } suspend fun emitValue(value: Int) { sharedFlow.emit(value) } data class Result(val list: List<Int>)
最优实现方式
方式1:基于scan的无依赖实现
这种方式通过scan跟踪当前累积的临时列表,当列表达到指定大小后,自动创建新的空列表作为下一轮的累积容器,完全避免手动清空列表的操作,状态流转更清晰:
fun main() = runBlocking { launch { repeat(20) { emitValue(it) } } flow .collect { println(it) } } private val sharedFlow = MutableSharedFlow<Int>() private const val CHUNK_SIZE = 5 // 指定分块大小 val flow = sharedFlow.asSharedFlow() .onEach { println("received: $it") } // scan维护当前累积的临时列表作为状态 .scan(mutableListOf<Int>()) { currentList, value -> currentList.add(value) if (currentList.size == CHUNK_SIZE) { // 列表已满,返回新的空列表用于下一轮累积 mutableListOf() } else { // 列表未满,继续使用当前列表 currentList } } // 仅筛选出达到指定大小的完整列表 .filter { it.size == CHUNK_SIZE } // 转换为包含不可变列表的Result对象,避免后续修改风险 .map { Result(it.toList()) } suspend fun emitValue(value: Int) { sharedFlow.emit(value) } data class Result(val list: List<Int>)
方式2:使用chunked操作符(简洁优先)
如果你的项目使用Kotlin Coroutines 1.6.0及以上版本,可以直接使用Flow官方提供的chunked扩展操作符,它已经封装了分块逻辑,代码最简洁高效:
fun main() = runBlocking { launch { repeat(20) { emitValue(it) } } flow .collect { println(it) } } private val sharedFlow = MutableSharedFlow<Int>() private const val CHUNK_SIZE = 5 // 指定分块大小 val flow = sharedFlow.asSharedFlow() .onEach { println("received: $it") } .chunked(CHUNK_SIZE) // 自动按指定大小分块,满块即发射 .map { Result(it) } // 转换为Result对象 suspend fun emitValue(value: Int) { sharedFlow.emit(value) } data class Result(val list: List<Int>)
方案对比
- 方式1:不依赖特定版本的Coroutines库,手动控制状态流转,避免了原实现中可变列表共享修改的潜在问题,兼容性更强。
- 方式2:官方封装的操作符,代码量最少、可读性最高,是优先推荐的方案,适合版本允许的项目。
内容的提问来源于stack exchange,提问作者d_eder
相关产品推荐
相关产品推荐

