求Kotlin Coroutines Flow类似sample但可发射最后值的操作符
实现带最终元素发射的Flow采样操作符
官方Kotlin协程的sample操作符存在一个特性:当上游流在采样时间窗口结束前完成时,最后一个元素不会被发射,文档明确说明:
Note that the latest element is not emitted if it does not fit into the sampling window.
如果你的业务场景必须确保最后一个元素被发射,下面是一个自定义的扩展操作符,行为和sample一致,但会在流结束时补充发射未被采样到的最后一个值。
自定义操作符实现
import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.collectLatest import kotlinx.coroutines.flow.flow import kotlinx.coroutines.flow.onEach import kotlinx.coroutines.flow.sample fun <T> Flow<T>.sampleWithFinalEmit(periodMillis: Long): Flow<T> = flow { var latestValue: T? = null // 追踪最新元素并生成采样流 val sampledStream = this@sampleWithFinalEmit .onEach { latestValue = it } .sample(periodMillis) // 先发射所有采样到的元素 sampledStream.collectLatest { emit(it) } // 流结束后,发射最后一个未被采样的元素(如果存在) latestValue?.let { emit(it) } }
实现逻辑说明
- 通过
latestValue变量持续记录上游流的最新元素 - 基于原流创建标准的
sample流,同时在每个元素到达时更新latestValue - 先收集并发射采样流的所有结果,保证和官方
sample的行为一致 - 当上游流完成后,检查是否存在未被采样窗口捕获的最后一个元素,若存在则补充发射
使用示例
import kotlinx.coroutines.delay import kotlinx.coroutines.flow.asFlow import kotlinx.coroutines.flow.collectLatest import kotlinx.coroutines.runBlocking fun main() = runBlocking { val k = (0..100).asFlow().onEach { delay(200) } // 使用自定义的采样操作符 val k2 = k.sampleWithFinalEmit(500) k2.collectLatest { println(it) } }
这个实现既保留了sample操作符按固定周期采样的核心行为,又解决了最后元素丢失的问题,完全适配需要确保最终值被处理的业务场景。
内容的提问来源于stack exchange,提问作者Cyber Avater
相关产品推荐
相关产品推荐

