如何在收集器挂起时让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
相关产品推荐
相关产品推荐

