Kotlin Flow如何实现收集指定数量值且遇值发射超时自动终止
问题根因分析
- 问题1丢包:你当前的
withNullOnTimeout实现中,merge操作会对上游原始Flow发起两次独立订阅。如果上游是单播冷流(大部分业务场景下的Flow都是这类),两次订阅会导致上游数据被拆分到两个消费分支,一部分数据被debounce分支消费走,自然就会出现原始流分支丢包的问题。 - 问题2超时流残留:原始流被
take终止后,你merge的debounce分支流仍然处于活跃状态,会等到最后一次数据到达后的debounce时长结束,发射null之后整个合并流才会终止,这就导致最后多余的超时触发。
修复方案
我们重新实现takeUntilTimeout操作符,全程只对上游做一次订阅,每次收到数据后重置超时计时器,上游结束后直接终止所有残留任务:
import kotlinx.coroutines.coroutineScope import kotlinx.coroutines.delay import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.flow import kotlinx.coroutines.launch import java.io.ByteArrayOutputStream suspend fun receive(packages: Flow<ByteArray>, amount: Int): ByteArray { val buffer = ByteArrayOutputStream(blockSize.toInt()) packages .take(amount) // 修复原代码硬编码take(10)的笔误 .takeUntilTimeout(100) .collect { pck -> buffer.write(pck) // 修复原代码pck.data的笔误,ByteArray本身就是数据载体 } return buffer.toByteArray() } fun <T> Flow<T>.takeUntilTimeout(durationMillis: Long): Flow<T> = flow { coroutineScope { var timeoutJob: Job? = null try { // 仅对上游做一次订阅,不会出现分流丢包 collect { value -> // 收到新值先取消之前的超时任务 timeoutJob?.cancel() // 发射当前有效值 emit(value) // 启动新的超时计时任务 timeoutJob = launch { delay(durationMillis) // 触发超时直接终止整个流收集 this@coroutineScope.cancel() } } } finally { // 上游结束(比如take到指定数量)后主动取消超时任务,避免残留触发 timeoutJob?.cancel() } } }
效果说明
- 所有上游发射的数据包都会唯一进入处理逻辑,不会出现分流丢包
- 当
take(amount)拿到足够的数量后,上游终止时会主动取消还在等待的超时任务,不会再触发多余的超时逻辑 - 如果两次数据包的发射间隔超过指定的durationMillis,超时任务会自动触发,终止整个流的收集
内容的提问来源于stack exchange,提问作者Alektas
相关产品推荐
相关产品推荐

