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

Kotlin协程:SharedFlow单元测试单独通过,批量运行失败

批量运行Flow测试时触发UncompletedCoroutinesError问题

单独运行单个测试均可通过,但批量运行时testA始终通过,testB或testC会抛出错误:

kotlinx.coroutines.test.UncompletedCoroutinesError: After waiting for 60000 ms, the test coroutine is not completing, there were active child jobs

相关类定义

interface Time {
    val currentTime: Flow<Int>
}

data class Data(val value: Int)

interface DataRepository {
    fun getCurrentData(): Data
}

核心类实现

@Singleton
class CurrentDataProvider(time: Time, repository: DataRepository) {
    private val coroutineScope = CoroutineScope(SupervisorJob() + Dispatchers.Default)

    private var lastUpdate = Int.MIN_VALUE

    val currentData: SharedFlow<Data> =
        time.currentTime
            .transform {
                if (lastUpdate + 10 <= it) {
                    lastUpdate = it
                    emit(repository.getCurrentData())
                }
            }
            .onCompletion {
                coroutineScope.cancel() // 是否必要?
            }
            .shareIn(
                scope = coroutineScope,
                started = SharingStarted.Eagerly,
            )
}

测试代码

class Test {
    private val time = mock<Time>()
    private val data = Data(123)
    private val repository = mock<DataRepository>{
        on { getCurrentData() } doReturn data
    }

    @Test
    fun `testA`() = runTest {
        // given
        whenever(time.currentTime).doReturn(flowOf(0))

        // when
        val tested = CurrentDataProvider(time, repository)

        // then
        tested.currentData.test {
            val item = awaitItem()
            item shouldBe data
        }
        verify(repository, times(1)).getCurrentData()
    }

    @Test
    fun `testB`() = runTest {
        // given
        whenever(time.currentTime).doReturn(flowOf(0, 1, 2, 3))

        // when
        val tested = CurrentDataProvider(time, repository)

        // then
        tested.currentData.test {
            val item = awaitItem()
            item shouldBe data
        }
        verify(repository, times(1)).getCurrentData()
    }

    @Test
    fun `testC`() = runTest {
        // given
        whenever(time.currentTime).doReturn(flowOf(1, 2, 3, 11))

        // when
        val tested = CurrentDataProvider(time, repository)

        // then
        tested.currentData.test {
            awaitItem()
            val item = awaitItem()
            item shouldBe data
        }
        verify(repository, times(2)).getCurrentData()
    }
}

已尝试的无效方案

添加MainDispatcherRule并尝试用withContext(Dispatchers.Default),但仍约每5次测试就会失败:

MainDispatcherRule定义

@get:Rule
val mainDispatcherRule = MainDispatcherRule()
@OptIn(ExperimentalCoroutinesApi::class)
class MainDispatcherRule(
    private val testDispatcher: TestDispatcher = UnconfinedTestDispatcher()
) : TestWatcher() {
    override fun starting(description: Description) {
        Dispatchers.setMain(testDispatcher)
    }

    override fun finished(description: Description) {
        Dispatchers.resetMain()
    }
}

问题根源与解决方案

核心问题

CurrentDataProvider内的自定义CoroutineScope未被正确取消:

  1. onCompletion中的coroutineScope.cancel()仅在上游流完成时触发,而shareIn使用SharingStarted.Eagerly会让Scope持续活跃;
  2. 批量测试时,前一个测试的未取消Scope会残留活跃协程,与后续测试产生冲突。

具体修复步骤

1. 重构CurrentDataProvider,实现CoroutineScope并提供取消方法

让类托管自身的Scope,并暴露取消接口,确保测试可主动清理:

@Singleton
class CurrentDataProvider(
    time: Time, 
    repository: DataRepository,
    dispatcher: CoroutineDispatcher = Dispatchers.Default
) : CoroutineScope by CoroutineScope(SupervisorJob() + dispatcher) {

    private var lastUpdate = Int.MIN_VALUE

    val currentData: SharedFlow<Data> =
        time.currentTime
            .transform {
                if (lastUpdate + 10 <= it) {
                    lastUpdate = it
                    emit(repository.getCurrentData())
                }
            }
            .shareIn(
                scope = this,
                started = SharingStarted.Eagerly,
            )

    // 供测试调用的取消方法
    fun cleanup() {
        coroutineContext.cancel()
    }
}

注:移除onCompletion中的coroutineScope.cancel(),避免流完成时提前终止Scope

2. 测试结束时主动清理Scope

在每个测试的finally块中调用cleanup(),确保测试结束后销毁协程:

@Test
fun `testA`() = runTest {
    // given
    whenever(time.currentTime).doReturn(flowOf(0))
    val testDispatcher = UnconfinedTestDispatcher(testScheduler)

    // when
    val tested = CurrentDataProvider(time, repository, testDispatcher)
    try {
        // then
        tested.currentData.test {
            val item = awaitItem()
            item shouldBe data
        }
        verify(repository, times(1)).getCurrentData()
    } finally {
        tested.cleanup()
    }
}

testB和testC做相同修改

3. 可选优化:调整shareIn的启动策略

如果业务允许,将started改为SharingStarted.WhileSubscribed(),无订阅者时自动停止流,减少后台活跃协程:

.shareIn(
    scope = this,
    started = SharingStarted.WhileSubscribed(stopTimeoutMillis = 0),
)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 19:20:01