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

为何在launch内绑定onEach会导致MutableSharedFlow事件丢失?

SharedFlow中onEach位置不同导致测试失败的原因

测试现象

通过的测试用例

@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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 01:05:22