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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 22:12:52