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

自定义Flow算子中select语句的协程缓冲异常问题排查

问题分析与解决方案

问题根源

原实现的核心问题在于未对上游施加背压控制,且共享缓冲区的修改与信号处理不同步:

  • 上游元素直接写入共享的buffer,produce通道的collect逻辑和上游发射是同步执行的,上游可以在溢出信号被处理前,继续往缓冲区添加新元素。
  • 当缓冲区达到maxBufferSize时发送溢出信号,但此时select还未响应该信号,上游已经把下一个元素(Unit)写入缓冲区,最终导致一次flush了6个元素,而非预期的分两次flush。

修正后的算子实现

fun <T> Flow<T>.bufferedWithTimeout(
    maxBufferSize: Int,
    timeout: Duration
): Flow<List<T>> {
    require(maxBufferSize > 0) {
        "Max buffer size has to be greater than zero (got $maxBufferSize)"
    }
    require(timeout.isPositive()) {
        "Timeout has to be positive (got $timeout)"
    }

    return flow {
        coroutineScope {
            val buffer = ArrayList<T>(maxBufferSize)
            val upstreamChannel = produce(capacity = Channel.RENDEZVOUS) {
                this@bufferedWithTimeout.collect { item ->
                    // 缓冲区满时挂起,直到flush完成再继续接收上游元素,实现背压
                    while (buffer.size >= maxBufferSize) {
                        awaitCancellation()
                    }
                    buffer.add(item)
                    println("Received $item. Size ${buffer.size}")
                    // 缓冲区刚满时发送溢出信号
                    if (buffer.size == maxBufferSize) {
                        println("Overflow sent")
                        send(Unit)
                    }
                }
                // 上游流结束时,发送信号触发剩余数据flush
                send(Unit)
            }

            suspend fun flush() {
                if (buffer.isEmpty()) return
                emit(buffer.toList())
                println("Flushing ${buffer.size}")
                buffer.clear()
                // 唤醒挂起的上游收集逻辑
                upstreamChannel.offer(Unit)
            }

            try {
                while (isActive) {
                    select {
                        upstreamChannel.onReceive {
                            println("Overflow/Completion received")
                            flush()
                            // 上游关闭后退出循环
                            if (!upstreamChannel.isActive) break
                        }

                        onTimeout(timeout) {
                            println("Timeout")
                            flush()
                        }
                    }
                }
            } finally {
                // 兜底flush剩余数据,避免遗漏
                flush()
                upstreamChannel.cancel()
            }
        }
    }
}

关键修改说明

  1. 背压控制:在收集上游元素时,若缓冲区已满则挂起上游发射,直到flush完成后通过upstreamChannel.offer(Unit)唤醒,彻底避免上游在缓冲区满后继续写入元素。
  2. 同步信号逻辑:将上游收集、缓冲区修改和溢出信号发送放到同一个produce通道内,确保缓冲区满时立即触发信号,且上游不会在信号处理前继续写入。
  3. 上游结束处理:上游流终止时主动发送信号,确保剩余缓冲区数据被flush,同时在select中判断通道状态,避免无限循环。
  4. 资源兜底:在finally块中强制flush剩余数据并取消上游通道,避免资源泄漏和数据丢失。

测试验证

运行原测试用例后,会输出符合预期的日志:

Received 0. Size 1
Received 1. Size 2
Received 2. Size 3
Received 3. Size 4
Received 4. Size 5
Overflow sent
Overflow/Completion received
Flushing 5
Received kotlin.Unit. Size 1
Timeout
Flushing 1

测试断言会通过,最终得到[[0,1,2,3,4], [Unit]]的结果。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 15:44:52