如何根据Flow执行结果分支调度不同Coroutine Flow?
Kotlin Flow 分支执行实现方案
现有可正常运行的Flow结构
现有代码实现了flow1执行完成后启动flow2的顺序执行逻辑,代码示例如下:
import kotlinx.coroutines.* import kotlinx.coroutines.flow.* val flow_1 = flowOf("loading_1", "loading_1", "success_1").onEach { delay(100) } val flow_2 = flowOf("loading_2", "loading_2", "success_2").onEach { delay(200) } fun main() = runBlocking<Unit> { flowOf( flow_1 .onEach { if (it == "failure_1") throw Exception("Failed!") } .retry(3), flow_2 ).flattenConcat() .catch { println("Failed with $it") } .collect { println("Result $it") } }
需求说明
需要实现分支执行逻辑:
- 若
flow_1执行成功,则继续执行flow_3 - 若
flow_1执行失败,则转而执行flow_2
场景示例
场景1:flow_1执行成功
此时按顺序执行flow_1和flow_3:
val flow_1 = flowOf("loading_1", "loading_1", "success_1").onEach { delay(100) } val flow_2 = flowOf("loading_2", "loading_2", "success_2").onEach { delay(200) } val flow_3 = flowOf("loading_3", "loading_3", "success_3").onEach { delay(200) }
场景2:flow_1执行失败
此时按顺序执行flow_1和flow_2:
val flow_1 = flowOf("loading_1", "loading_1", "error_1").onEach { delay(100) } val flow_2 = flowOf("loading_2", "loading_2", "success_2").onEach { delay(200) } val flow_3 = flowOf("loading_3", "loading_3", "success_3").onEach { delay(200) }
实现代码
这里提供两种可行的实现方式:
方式一:利用catch与flatMapConcat实现
通过catch捕获flow_1的异常,在异常分支切换到flow_2;正常完成后通过flatMapConcat串联flow_3:
import kotlinx.coroutines.* import kotlinx.coroutines.flow.* fun main() = runBlocking<Unit> { // 测试成功场景用这行,失败场景替换为下一行注释的代码 val flow_1 = flowOf("loading_1", "loading_1", "success_1").onEach { delay(100) } // val flow_1 = flowOf("loading_1", "loading_1", "error_1").onEach { delay(100) }.onEach { if (it == "error_1") throw Exception("flow_1 failed") } val flow_2 = flowOf("loading_2", "loading_2", "success_2").onEach { delay(200) } val flow_3 = flowOf("loading_3", "loading_3", "success_3").onEach { delay(200) } flow_1 .onEach { println("Result $it") } // flow_1成功完成后,切换执行flow_3 .flatMapConcat { flow_3 } // 捕获flow_1的异常,转而执行flow_2 .catch { println("Failed with $it") emitAll(flow_2) } .collect { println("Result $it") } }
方式二:自定义成功判断逻辑(更灵活)
通过手动收集flow_1,自定义判断成功/失败的条件,再选择后续执行的Flow:
import kotlinx.coroutines.* import kotlinx.coroutines.flow.* fun main() = runBlocking<Unit> { val flow_1 = flowOf("loading_1", "loading_1", "success_1").onEach { delay(100) } // val flow_1 = flowOf("loading_1", "loading_1", "error_1").onEach { delay(100) }.onEach { if (it == "error_1") throw Exception("flow_1 failed") } val flow_2 = flowOf("loading_2", "loading_2", "success_2").onEach { delay(200) } val flow_3 = flowOf("loading_3", "loading_3", "success_3").onEach { delay(200) } flow { var isSuccess = false try { var lastValue: String? = null flow_1.collect { println("Result $it") lastValue = it } // 自定义判断成功的条件:最后一个元素为success_1 isSuccess = lastValue == "success_1" } catch (e: Exception) { println("Failed with $e") isSuccess = false } // 根据结果选择后续执行的Flow emitAll(if (isSuccess) flow_3 else flow_2) }.collect { println("Result $it") } }
代码说明
- 方式一:适合仅通过异常判断失败的场景,代码简洁,利用Flow的操作符链式调用完成分支逻辑。
- 方式二:适合需要自定义成功条件的场景(比如根据Flow的最终输出值判断,而非仅异常),逻辑更灵活可控。
内容的提问来源于stack exchange,提问作者Metropol_Tilkisi
相关产品推荐
相关产品推荐

