为何在launch内绑定onEach会导致MutableSharedFlow事件丢失?
测试现象
通过的测试用例
@Test fun passingTest() { val emitter = MutableSharedFlow<Int>(); var total = 0; val handler = emitter.onEach { // <--- onEach 定义在 launch 外部 total += it; println("-> $it : total: $total") } GlobalScope.launch { handler.collect() } runBlocking { for (i in 1..11) { emitter.emit(1); } } assert(total == 11); }
失败的测试用例
运行后报错:
> expected:<11> but was:<0> > Expected :11 > Actual :0
对应的代码:
@Test fun failingTest() { val emitter = MutableSharedFlow<Int>(); var total = 0; GlobalScope.launch { emitter.onEach { // <--- onEach 定义在 launch 内部 total += it; println("-> $it : total: $total") }.collect() } runBlocking { for (i in 1..11) { emitter.emit(1); } } assertEquals(11, total); }
原因分析
本质是SharedFlow的默认特性+协程调度时机的竞争导致的:
1. SharedFlow的默认规则
MutableSharedFlow默认有三个关键配置:
replay = 0:新订阅者无法收到订阅前发送的任何事件extraBufferCapacity = 0:没有额外缓冲区存储未被接收的事件onBufferOverflow = BufferOverflow.SUSPEND:当没有活跃订阅者且缓冲区满时,emit()会挂起,直到有订阅者接收事件
2. 两个测试的核心差异
通过的测试
handler提前创建后,GlobalScope.launch内的collect()启动速度更快,在runBlocking的emit()挂起前完成了订阅。emit()被唤醒后逐个发送事件,total被正常累加,最终断言通过。
失败的测试
onEach和collect()都放在launch内部,协程需要先构建Flow链再执行订阅,启动延迟略高。runBlocking里的emit()已经发完所有事件,launch的协程才开始订阅——而SharedFlow默认不保留历史事件,所以订阅者什么都收不到,total始终为0,断言失败。
修复思路
如果想确保订阅者能收到所有事件,可以采用以下方式:
- 给
MutableSharedFlow设置足够大的replay或extraBufferCapacity - 让订阅逻辑在发送前完成(比如用
runBlocking启动订阅) - 监听
subscriptionCount,等订阅者就绪后再发送事件
内容的提问来源于stack exchange,提问作者Doug
相关产品推荐
相关产品推荐

