如何构建在终端操作间不重启的Kotlin Resumable Flow
实现可恢复的按需拉取式Flow(多次收集不重置状态)
需求场景
我们需要构建一个支持多次收集但不重置状态的Flow,比如调用first()时能依次获取下一个值,而非每次都从头开始。核心是要实现asResumableFlow扩展函数,用法示例如下:
fun <T> Flow<T>.asResumableFlow(): Flow<T> = ??? // 预期行为 val flow = flowOf(1, 2).asResumableFlow() flow.first() // 返回1 flow.first() // 返回2,而非1 flow.firstOrNull() // 返回null
典型使用场景是在循环中仅当主任务失败时才触发回退逻辑,此时才需要生成昂贵的实例:
val resumableFlow: Flow<Expensive> = ... while (condition) { if (mainTaskSucceded) continue fallback(resumableFlow.first()) // 仅在需要时才生成Expensive实例 }
核心约束
- 按需计算:值仅在实际被消费时才生成
- 无值丢失:消费者切换时,未被消费的值不会丢失
- 单值单消费:每个值只会被一个消费者接收
- 懒加载:消费者断开后,只有新消费者连接时才会计算下一个值
- 结束后空Flow:生产者Flow结束后,新的收集操作会得到空结果
现有方案的问题
Channel、StateFlow/MutableFlow属于推送模型,会提前生成值并缓存,不符合按需计算的要求- 普通冷Flow每次收集都会从头开始,无法保留状态
- 手动维护状态机的临时方案需要自己管理迭代状态,不符合协程的设计理念
优化后的实现
基于草稿实现,我们整理出更简洁、符合协程规范的版本:
import kotlinx.coroutines.* import kotlinx.coroutines.flow.* public fun <T> Flow<T>.asResumableFlow( context: CoroutineContext = EmptyCoroutineContext ): Flow<T> = ResumableFlow(this, context) private class ResumableFlow<T>( private val source: Flow<T>, private val context: CoroutineContext ) : Flow<T> { // 控制生产者协程的续体,null表示生产者已结束 private var producerContinuation: Continuation<Unit>? = null // 暂存当前生成的值,等待消费者获取 private var pendingValue: T? = null // 生产者抛出的异常 private var producerError: Throwable? = null init { // 启动生产者协程,初始挂起等待第一个消费者 startProducerCoroutine() } private fun startProducerCoroutine() { val completion = Continuation<Unit>(context) { result -> producerContinuation = null pendingValue = null producerError = result.exceptionOrNull() } producerContinuation = suspend { runProducer() }.createCoroutine(completion) } private suspend fun runProducer() { source.collect { value -> // 生成值后暂存,然后挂起等待消费者获取 pendingValue = value suspendCoroutine<Unit> { cont -> producerContinuation = cont } } } override suspend fun collect(collector: FlowCollector<T>) { // 如果有异常,直接抛出 producerError?.let { throw it } // 恢复生产者协程,生成下一个值(或确认是否已结束) producerContinuation?.resume(Unit) ?: return // 等待值生成完成(单消费者约束下,pendingValue一定非空) val value = pendingValue ?: return pendingValue = null // 发送值给消费者 collector.emit(value) } }
关键逻辑说明
- 生产者协程:初始化时启动,每次生成值后挂起,等待消费者触发才继续生成下一个值
- 状态管理:用
pendingValue暂存当前值,producerContinuation控制生产者的挂起与恢复 - 按需计算:只有当消费者调用
collect(或first()等终端操作)时,才会恢复生产者生成下一个值 - 结束处理:生产者Flow结束后,
producerContinuation置为null,后续收集操作直接返回空结果
验证效果
使用示例中的代码测试,完全符合预期:
- 第一次
flow.first()触发生产者生成1,返回后生产者挂起 - 第二次
flow.first()恢复生产者生成2,返回后生产者挂起 - 第三次
flow.firstOrNull()恢复生产者,此时source已结束,返回null
内容的提问来源于stack exchange,提问作者self_out_ manoeuvred
相关产品推荐
相关产品推荐

