使用Mockito测试SharedFlow逻辑时流不发射数据致测试超时失败
问题现象
在为返回Flow的登录请求方法编写单元测试时,使用Mockito配置桩返回SharedFlow后,测试运行1分钟后超时失败,抛出异常:
This job has not completed yet java.lang.IllegalStateException: This job has not completed yet
原测试代码
@ExperimentalCoroutinesApi @Test fun `when response body have error in request login`() = runBlockingTest { runCurrent() Mockito.`when`(webSocketClient.isConnect()).thenReturn(true) Mockito.`when`(mapper.createRPC(userLoginObject)).thenReturn(rpc) Mockito.`when`(requestManager.sendRequest(rpc)).thenReturn(userLoginFlow) userLoginFlow.emit(errorObject) loginServiceImpl.requestLogin(userLoginObject).drop(1).collectLatest { assert(it == errorObject) } }
被测方法代码
override fun requestLogin(userLoginObject: BaseDomain): Flow<DataState<BaseDomain>> = flow { emit(DataState.Loading(ProgressBarState.Loading)) if (webSocketClient.isConnect()) { requestManager.sendRequest(mapper.createRPC(userLoginObject)!!)?.filterNotNull()?.collectLatest { if (it is IG_RPC.Error) { emit(DataState.Error(ErrorObject(it.major, it.minor, it.wait))) } else if (it is IG_RPC.Res_User_Register) { val userLoginObject = userLoginObject as UserLoginObject emit( DataState.Data( UserLoginObject( userName = it.userName, phoneNumber = userLoginObject.phoneNumber, userId = it.userId, authorHash = it.authorHash, regex = it.codeRegex, resendCodeDelay = it.resendDelayTime ) ) ) } } } else { emit(DataState.Error(ErrorObject(-1, -1, 0))) } }
问题根因
- 热流无限挂起:
SharedFlow是热流,默认没有结束边界,只要不主动取消,对它的collect操作会一直挂起等待新数据,永远不会自行完成,最终触发测试超时。 - 发射时序错误:原代码在启动流收集之前就提前调用
emit发射数据,如果SharedFlow没有配置足够的replay缓存,这部分数据会直接丢弃,后续收集操作根本拿不到。 - 断言逻辑错误:被测方法收到
IG_RPC.Error类型的数据后,会将其包装为DataState.Error对象再发射,原断言直接对比裸的ErrorObject,就算数据成功接收,断言也无法通过。 - 测试API过时:使用已经废弃的
runBlockingTest,协程调度时序控制不可靠,容易出现执行顺序错乱。
修复方案
- 替换过时的协程测试API:将
runBlockingTest替换为当前稳定版协程测试库的runTest,使用测试调度器精准控制协程执行时序。 - 调整发射收集时序:先启动流收集协程,等收集动作就绪后再往SharedFlow发射测试数据,避免数据丢失。
- 主动取消收集协程:拿到目标断言结果后立刻取消收集操作,避免热流无限挂起导致测试超时。
- 修正断言逻辑:对比被测方法实际返回的
DataState.Error包装对象,而非裸的错误结构体。
修复后测试代码示例
@ExperimentalCoroutinesApi @Test fun `when response body have error in request login`() = runTest { // 初始化测试用的SharedFlow,配置replay=1缓存发射的数据 val testError = IG_RPC.Error.newBuilder() .setMajor(1) .setMinor(2) .setWait(0) .build() val expectedResult = DataState.Error(ErrorObject(1, 2, 0)) val userLoginFlow = MutableSharedFlow<IG_RPC.Base>(replay = 1) // 配置所有Mock桩返回 Mockito.`when`(webSocketClient.isConnect()).thenReturn(true) Mockito.`when`(mapper.createRPC(userLoginObject)).thenReturn(rpc) Mockito.`when`(requestManager.sendRequest(rpc)).thenReturn(userLoginFlow) // 单独启动协程收集流,避免阻塞当前测试协程 val collectJob = launch(UnconfinedTestDispatcher(testScheduler)) { loginServiceImpl.requestLogin(userLoginObject) .drop(1) // 跳过默认发射的Loading状态 .collectLatest { // 校验结果符合预期 assert(it == expectedResult) // 拿到结果立刻取消收集,避免热流无限挂起 cancel() } } // 等收集协程启动就绪后发射测试数据 advanceUntilIdle() userLoginFlow.emit(testError) // 等待所有协程执行完成,确认收集逻辑正常执行 advanceUntilIdle() collectJob.join() }
内容的提问来源于stack exchange,提问作者Parsa Dehghan
相关产品推荐
相关产品推荐

