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

使用Mockito测试SharedFlow逻辑时流不发射数据致测试超时失败

单元测试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,协程调度时序控制不可靠,容易出现执行顺序错乱。

修复方案

  1. 替换过时的协程测试API:将runBlockingTest替换为当前稳定版协程测试库的runTest,使用测试调度器精准控制协程执行时序。
  2. 调整发射收集时序:先启动流收集协程,等收集动作就绪后再往SharedFlow发射测试数据,避免数据丢失。
  3. 主动取消收集协程:拿到目标断言结果后立刻取消收集操作,避免热流无限挂起导致测试超时。
  4. 修正断言逻辑:对比被测方法实际返回的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 14:15:32