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

如何构建在终端操作间不重启的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)
    }
}

关键逻辑说明

  1. 生产者协程:初始化时启动,每次生成值后挂起,等待消费者触发才继续生成下一个值
  2. 状态管理:用pendingValue暂存当前值,producerContinuation控制生产者的挂起与恢复
  3. 按需计算:只有当消费者调用collect(或first()等终端操作)时,才会恢复生产者生成下一个值
  4. 结束处理:生产者Flow结束后,producerContinuation置为null,后续收集操作直接返回空结果

验证效果

使用示例中的代码测试,完全符合预期:

  • 第一次flow.first()触发生产者生成1,返回后生产者挂起
  • 第二次flow.first()恢复生产者生成2,返回后生产者挂起
  • 第三次flow.firstOrNull()恢复生产者,此时source已结束,返回null

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 10:13:18