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

如何在收集器挂起时让SharedFlow的emit操作同步挂起?

如何让SharedFlow的emit在收集器挂起时同步挂起?

调用挂起函数时协程会随之挂起,但将逻辑迁移到SharedFlow后,emit操作并未随收集器的挂起而同步挂起:

示例代码与问题现象

基础挂起逻辑(符合预期)

fun main()= runBlocking<Unit> {
    println("before, time is ${System.currentTimeMillis()/1000} s")
    delay5s()
    println("after, time is ${System.currentTimeMillis()/1000} s")
}

private suspend fun delay5s() {
    delay(5000)
}

这段代码中delay5s()会挂起协程,after打印会在5秒后执行。

SharedFlow场景(不符合预期)

fun main()= runBlocking<Unit> {
    val flow = MutableSharedFlow<Int>(
        replay = 0,
        extraBufferCapacity = 0,
        onBufferOverflow = BufferOverflow.SUSPEND
    )
    launch {
        flow.collect {
            println("before on collect, time is ${System.currentTimeMillis()/1000} s")
            delay5s()
            println("after on collect, time is ${System.currentTimeMillis()/1000} s")
        }
    }
    launch {
        println("before emit, time is ${System.currentTimeMillis()/1000} s")
        flow.emit(1)
        println("after emit, time is ${System.currentTimeMillis()/1000} s")
    }
}

实际输出:

before emit, time is 1733335171 s
before on collect, time is 1733335171 s
after emit, time is 1733335171 s
after on collect, time is 1733335176 s

仅收集器被挂起,emit操作的前后打印几乎同时执行。

期望输出:

before emit, time is 1733335171 s
before on collect, time is 1733335171 s
after on collect, time is 1733335176 s
after emit, time is 1733335176 s

问题原因

SharedFlow是热流,emit的核心职责是将事件分发到所有活跃收集器,不会等待收集器完成事件处理。一旦事件被交付至收集器的处理队列,emit就会立即恢复执行,因此无法同步跟随收集器的挂起逻辑。

解决方案

方案1:改用冷流(单收集器场景)

如果仅需要单个收集器,直接使用冷流flow { ... }即可。冷流的代码块会在收集器的协程上下文执行,emit后的代码会等待收集器处理完当前元素后才继续执行:

fun main()= runBlocking<Unit> {
    val flow = flow {
        println("before emit, time is ${System.currentTimeMillis()/1000} s")
        emit(1)
        println("after emit, time is ${System.currentTimeMillis()/1000} s")
    }
    launch {
        flow.collect {
            println("before on collect, time is ${System.currentTimeMillis()/1000} s")
            delay5s()
            println("after on collect, time is ${System.currentTimeMillis()/1000} s")
        }
    }
}

private suspend fun delay5s() {
    delay(5000)
}

输出完全符合期望。

方案2:基于SharedFlow实现确认机制(多收集器场景)

如果必须使用SharedFlow(如需要多播至多个收集器),可通过双向通信让emit等待收集器的处理确认:

data class EventWithAck<T>(val data: T, val ack: () -> Unit)

fun main()= runBlocking<Unit> {
    val flow = MutableSharedFlow<EventWithAck<Int>>(
        replay = 0,
        extraBufferCapacity = 0,
        onBufferOverflow = BufferOverflow.SUSPEND
    )
    
    // 收集器:处理完成后发送确认信号
    launch {
        flow.collect { event ->
            println("before on collect, time is ${System.currentTimeMillis()/1000} s")
            delay5s()
            println("after on collect, time is ${System.currentTimeMillis()/1000} s")
            event.ack()
        }
    }
    
    // 发送方:等待收集器确认后再继续执行
    launch {
        println("before emit, time is ${System.currentTimeMillis()/1000} s")
        val ackSignal = CompletableDeferred<Unit>()
        flow.emit(EventWithAck(1) { ackSignal.complete(Unit) })
        ackSignal.await()
        println("after emit, time is ${System.currentTimeMillis()/1000} s")
    }
}

private suspend fun delay5s() {
    delay(5000)
}

若存在多个收集器,可修改逻辑等待所有收集器的确认信号(例如使用awaitAll)。

方案3:使用Channel(严格顺序场景)

如果不需要多播能力,Channel结合确认机制也能实现需求:

fun main()= runBlocking<Unit> {
    val channel = Channel<Pair<Int, CompletableDeferred<Unit>>>(Channel.RENDEZVOUS)
    
    launch {
        for ((data, ack) in channel) {
            println("before on collect, time is ${System.currentTimeMillis()/1000} s")
            delay5s()
            println("after on collect, time is ${System.currentTimeMillis()/1000} s")
            ack.complete(Unit)
        }
    }
    
    launch {
        println("before emit, time is ${System.currentTimeMillis()/1000} s")
        val ack = CompletableDeferred<Unit>()
        channel.send(1 to ack)
        ack.await()
        println("after emit, time is ${System.currentTimeMillis()/1000} s")
    }
}

private suspend fun delay5s() {
    delay(5000)
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 18:15:17