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

求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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 19:43:24