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

Kotlin SharedFlow测试:flowOn与launchIn用法及失败修复

SharedFlow测试失败原因分析与修复

问题描述

我编写了两个SharedFlow测试用例,test shared flow A可正常通过,但test shared flow B执行失败。原以为两种写法等价,想明确两个问题:

  • test shared flow B失败的原因是什么?
  • 能否在保留launchIn方法的前提下让测试通过?

测试代码如下:

import kotlinx.coroutines.ExperimentalCoroutinesApi
import kotlinx.coroutines.flow.MutableSharedFlow
import kotlinx.coroutines.flow.flowOn
import kotlinx.coroutines.flow.launchIn
import kotlinx.coroutines.flow.onEach
import kotlinx.coroutines.launch
import kotlinx.coroutines.test.UnconfinedTestDispatcher
import kotlinx.coroutines.test.runTest
import org.junit.Test

@OptIn(ExperimentalCoroutinesApi::class)
class SomethingTest {

    @Test
    fun `test shared flow A`() = runTest {
        val flow = MutableSharedFlow<Int>()
        val items = mutableListOf<Int>()
        val job = launch(UnconfinedTestDispatcher()) {
            flow.collect {
                items.add(it)
            }
        }
        flow.emit(1)
        assert(items.size == 1)
        job.cancel()
    }

    @Test
    fun `test shared flow B`() = runTest {
        val flow = MutableSharedFlow<Int>()
        val items = mutableListOf<Int>()
        val job = flow.onEach { items.add(it) }
            .flowOn(UnconfinedTestDispatcher())
            .launchIn(this)
        flow.emit(1)
        assert(items.size == 1)
        job.cancel()
    }
}

失败原因解析

两个测试的核心差异在于协程调度器的作用范围:

  • test shared flow A中,收集协程通过launch(UnconfinedTestDispatcher())启动,收集逻辑会立即在当前线程执行——UnconfinedTestDispatcher的特性就是会立即调度执行协程体,所以emit(1)后,收集逻辑同步完成,断言时items已经有值。
  • test shared flow B中,flowOn(UnconfinedTestDispatcher())仅指定了上游流操作的调度器,而launchIn(this)是用runTest默认的测试调度器启动收集协程。这个默认调度器不会立即执行所有逻辑,当我们调用emit(1)后马上执行断言时,收集协程还处于待调度状态,items还没被添加元素,导致断言失败。

保留launchIn的修复方案

方案一:给launchIn指定UnconfinedTestDispatcher

直接让收集协程使用UnconfinedTestDispatcher启动,和Test A的行为对齐:

@Test
fun `test shared flow B fixed`() = runTest {
    val flow = MutableSharedFlow<Int>()
    val items = mutableListOf<Int>()
    val job = flow.onEach { items.add(it) }
        .launchIn(UnconfinedTestDispatcher()) // 直接指定Unconfined调度器
    flow.emit(1)
    assert(items.size == 1)
    job.cancel()
}

这里不需要flowOn,因为我们需要的是收集协程的调度行为,而非上游流的调度。

方案二:让测试调度器执行完所有待处理任务

在断言前调用advanceUntilIdle(),等待测试调度器完成所有pending的协程任务,确保收集逻辑执行完毕:

@Test
fun `test shared flow B fixed`() = runTest {
    val flow = MutableSharedFlow<Int>()
    val items = mutableListOf<Int>()
    val job = flow.onEach { items.add(it) }
        .flowOn(UnconfinedTestDispatcher())
        .launchIn(this)
    flow.emit(1)
    advanceUntilIdle() // 等待所有协程任务执行完成
    assert(items.size == 1)
    job.cancel()
}

advanceUntilIdle()是runTest提供的测试工具方法,会驱动调度器运行到没有待处理任务为止,此时收集逻辑已经执行,断言就能通过。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 07:25:20