如何在Flow collect中获取当前与前序值并灵活控制执行上下文
需求说明
需要在Flow收集流程中同时获取当前发射值与上一个发射值,预期运行时序如下:
----A----------B-------C-----|---> ---(null+A)---(A+B)---(B+C)--|--->
最初的自定义withPrevious扩展实现代码如下:
fun <T: Any> Flow<T>.withPrevious(): Flow<Pair<T?, T>> = flow { var prev: T? = null this@withPrevious.collect { emit(prev to it) prev = it } }
这个实现在当前简单逻辑下看似可以运行,但手动在flow构建器中直接收集上游的写法,没有强制遵循Flow的上下文透明规范,后续迭代中很容易因为插入withContext切换、异常捕获逻辑不当等问题破坏flowOn等上下文控制操作符的生效逻辑,也容易出现背压、取消传播不符合预期的问题,灵活性不足。
推荐实现方案
优先使用Kotlin Flow标准库提供的变换操作符实现,不需要手动编写上游收集逻辑,天然支持所有Flow标准特性(上下文控制、异常传播、背压处理等),行为和官方操作符完全一致。
方案1:基于scan(最简洁推荐)
scan(别名runningFold)是标准库提供的带状态累积操作符,天生符合上下文透明要求,实现代码非常简洁:
fun <T: Any> Flow<T>.withPrevious(): Flow<Pair<T?, T>> = scan(null as Pair<T?, T>?) { lastPair, currentValue -> // 上一次发射值的second字段就是前一个上游值,和当前值配对 lastPair?.second to currentValue }.filterNotNull() // 过滤掉初始的null值
这个实现的逻辑和预期完全一致:第一个元素发射时,因为还没有历史值,会输出null to A,后续每个元素都和上一个元素配对输出。
方案2:基于transform(适合扩展定制)
如果需要在配对逻辑之外添加额外判断(比如满足特定条件才发射、需要附加其他状态),可以用transform操作符实现,同样是上下文安全的:
fun <T: Any> Flow<T>.withPrevious(): Flow<Pair<T?, T>> = flow { var prev: T? = null // 通过emitAll委托给transform的变换逻辑,保证上下文透明 emitAll( transform { current -> emit(prev to current) prev = current } ) }
注意这里的状态变量prev定义在flow块内部,每次流被重新收集时都会重置初始值,不会出现多收集场景下的状态污染问题。
效果验证
两种实现都可以正常响应flowOn的上下文控制,例如下方测试代码:
import kotlinx.coroutines.* import kotlinx.coroutines.flow.* suspend fun main() { flowOf(1, 2, 3) .onEach { println("上游发射值$it,上下文:${currentCoroutineContext()[CoroutineName]?.name}") } .flowOn(CoroutineName("上游调度上下文")) .withPrevious() .onEach { println("下游收到值$it,上下文:${currentCoroutineContext()[CoroutineName]?.name}") } .flowOn(CoroutineName("下游调度上下文")) .collect() }
运行后可以看到上游逻辑全部运行在上游调度上下文,下游逻辑全部运行在下游调度上下文,和使用标准库操作符的行为完全一致,不会出现上下文错乱的问题。
内容的提问来源于stack exchange,提问作者Vlad Kudoyar

