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

如何使用async/await且不中断Kotlin Flow的执行流程?

解决方案

当然可以实现你要的效果,核心思路是把每个输入元素转换成一个**先发射带null字段的Something,等异步请求完成后再发射填充好数据的Something**的子Flow,然后通过flatMapMerge(或flatMapConcat)将这些子Flow合并到主Flow中,这样就能保证Flow持续执行,不会被单个异步请求阻塞。

修改后的完整代码如下:

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*

val dispatcher = Dispatchers.IO
val scope = CoroutineScope(SupervisorJob() + dispatcher)

val numbersFlow = flowOf(1, 2, 3, 4, 5, 6, 7, 8, 9, 10)

fun main(args: Array<String>) = runBlocking {
    numbersFlow
        .flatMapMerge { num ->
            flow {
                // 先发射初始的null值
                emit(Something(id = null, path = null))
                
                // 异步执行两个请求
                val idDeferred = scope.async(dispatcher) {
                    delay(1000L)
                    num * num
                }
                val pathDeferred = scope.async(dispatcher) {
                    delay(500L)
                    "result $num"
                }
                
                // 等待请求完成后发射填充好的数据
                emit(Something(id = idDeferred.await(), path = pathDeferred.await()))
            }
        }
        .onEach { println("Received: $it") }
        .launchIn(scope)
    
    // 等待所有协程完成后退出,替代死循环
    scope.coroutineContext.job.join()
}

data class Something(
    val id: Int?,
    val path: String?
)

关键说明:

  • flatMapMerge的作用:它会把每个输入元素转换成一个子Flow,然后将所有子Flow的发射结果合并到主Flow中。相比原来的map,它不会等待子Flow完成再处理下一个元素,能让多个元素的异步请求并行执行;如果需要串行处理元素(上一个完成再处理下一个),可以换成flatMapConcat。
  • 子Flow的逻辑:每个子Flow先立即发射初始的Something(null, null),然后启动两个异步任务,等任务完成后再发射填充好数据的实例。
  • 简化异步请求:原代码中suspend fun包裹async的写法没必要,直接在协程作用域内用async即可,避免多余的上下文切换。
  • 替代死循环:用scope.coroutineContext.job.join()等待所有协程完成后退出,比死循环更优雅。

内容的提问来源于stack exchange,提问作者xephosbot

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 12:47:33