Coroutine Flow对应RxJava mergeDelayError的等效实现选型
RxJava
mergeDelayError 的Coroutine Flow等效实现 首先给出明确结论:该场景下既不适合用原生combine,也不能直接用默认配置的flattenMerge,需要对flattenMerge(或其简便写法merge扩展)做异常捕获包装,才能实现和mergeDelayError完全一致的行为。
两个候选操作符的适配性说明
- 关于
combine- 核心逻辑是任意上游发射新值时,取所有上游的最新值做组合后发射,要求所有上游至少发射一次值才会触发第一次组合,和
merge类操作符「谁先返回就先发射谁的值、不需要等待其他流」的逻辑完全不符,匹配不了你原有逐个接收API返回、即时处理的业务逻辑。 - 原生
combine只要任意一个上游抛出异常,会立刻取消整个流的执行,不会等待剩余请求跑完,完全不满足错误延迟的要求。
- 核心逻辑是任意上游发射新值时,取所有上游的最新值做组合后发射,要求所有上游至少发射一次值才会触发第一次组合,和
- 关于默认配置的
flattenMerge- 合并逻辑和RxJava的普通
merge一致:将多个Flow的发射值按到达顺序合并为单个流发射,不做值的组合,这部分逻辑是匹配需求的。 - 但默认配置下,只要任意一个上游抛出异常,
flattenMerge会立刻取消所有其他未完成的上游流,不会延迟错误处理,直接用达不到mergeDelayError的效果。
- 合并逻辑和RxJava的普通
等效实现方案
实现核心思路是:先给每个单独的请求流加异常捕获,将异常包装为普通发射值避免触发全局流取消,再用flattenMerge并发合并所有流,等所有流全部执行完成后,再统一处理暂存的异常,过程中正常消费所有成功返回的结果。
代码实现(匹配原有业务逻辑)
首先定义一个密封类用来区分成功结果和异常,避免异常直接中断流:
sealed class FlowResult<out T> { data class Success<T>(val data: T) : FlowResult<T>() data class Error(val throwable: Throwable) : FlowResult<Nothing>() }
改造原有请求方法:
fun requestHomeDataAtOnce() { // 假设requestTab1~requestTab4都是挂起的API请求方法,直接包装为Flow即可 val requestList = mutableListOf( flow { emit(requestTab1()) }, flow { emit(requestTab2()) }, flow { emit(requestTab3()) }, flow { emit(requestTab4()) } ) // 用页面对应协程作用域启动,比如lifecycleScope、viewModelScope lifecycleScope.launch { requestHome(requestList) } } private suspend fun requestHome(requestList: List<Flow<Result<Any>>>) { val responseList = mutableListOf<Any?>() val errorList = mutableListOf<Throwable>() // 给每个请求流加异常捕获,将异常包装为普通值发射,不会中断其他流执行 val wrappedRequestFlows = requestList.map { requestFlow -> requestFlow .map<Result<Any>, FlowResult<Any>> { FlowResult.Success(it.getOrThrow()) } .catch { emit(FlowResult.Error(it)) } } wrappedRequestFlows .merge() // 等价于flattenMerge(concurrency = Int.MAX_VALUE),并发执行所有请求 .flowOn(Dispatchers.IO) // 等价于原RxJava的subscribeOn(Schedulers.io()) .collect { result -> // 切主线程处理的话,把collect块包在withContext(Dispatchers.Main)里即可,等价于原observeOn(AndroidSchedulers.mainThread(), true) when (result) { is FlowResult.Success -> { responseList.add(result.data) // 这里写原有onNext里的即时处理逻辑,和原RxJava订阅后的行为完全一致 } is FlowResult.Error -> { errorList.add(result.throwable) } } } // 所有请求全部执行完成后,统一处理收集到的异常,实现错误延迟抛出的效果 if (errorList.isNotEmpty()) { // 按业务需求处理异常即可,比如合并为复合异常抛出、打日志、弹错误提示 } }
更轻量的替代方案
如果你的场景是单次API请求(对应原RxJava的Single,不是持续发射值的流),不需要中途接收值做处理,只需要等所有请求跑完统一拿结果,不需要用Flow,直接用supervisorScope+async就能实现完全一样的效果,代码更简洁:
private suspend fun requestHomeLight(requestList: List<suspend () -> Result<Any>>) = supervisorScope { val responseList = mutableListOf<Any?>() val errorList = mutableListOf<Throwable>() // 并发启动所有请求,一个请求失败不会影响其他请求 val deferredList = requestList.map { request -> async(Dispatchers.IO) { runCatching { request() } } } // 等待所有请求执行完成,收集结果和异常 deferredList.forEach { deferred -> deferred.await() .onSuccess { responseList.add(it.getOrNull()) } .onFailure { errorList.add(it) } } // 统一处理结果和异常即可 }
内容的提问来源于stack exchange,提问作者Youngin Lee
相关产品推荐
相关产品推荐

