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

如何根据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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 04:27:45