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未被正确取消:
onCompletion中的coroutineScope.cancel()仅在上游流完成时触发,而shareIn使用SharingStarted.Eagerly会让Scope持续活跃;- 批量测试时,前一个测试的未取消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
相关产品推荐
相关产品推荐

