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

如何让多个Kotlin Flow在独立协程/线程中并行生产数据?

解决Flow并行生产的问题

你的代码里merge操作默认会顺序收集每个Flow的值,每个Flow的生产者代码(网络请求、CPU操作)都运行在collect所在的协程上下文中。当其中一个Flow执行阻塞或CPU密集任务时,会占用当前协程,导致其他Flow只能等待,所以同一时间只有一个Flow在生产数据。

要让多个Flow并行生产,核心是让每个Flow的生产逻辑在独立的协程/线程中执行,你可以通过flowOn操作符实现:

方案1:给每个输入Flow指定调度器

修改runFlows函数,为每个Flow添加flowOn(Dispatchers.IO)(适合阻塞IO操作,比如网络请求)或Dispatchers.Default(适合CPU密集型操作):

import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.flow.flowOn
import kotlinx.coroutines.flow.merge

fun makeInputFlow() = flow {
    while (shouldMakeRequest()) {
        // 阻塞网络请求
        // CPU密集型操作
        results.forEach { emit(it) }
        yield()
    }
}

suspend fun runFlows() {
    val flow1 = makeInputFlow().flowOn(Dispatchers.IO)
    val flow2 = makeInputFlow().flowOn(Dispatchers.IO)
    val flow3 = makeInputFlow().flowOn(Dispatchers.IO)

    val merged = listOf(flow1, flow2, flow3).merge()
    merged.collect { println(it) }
}

方案2:在Flow定义时直接指定调度器

如果所有输入Flow都需要并行执行,可以直接在makeInputFlow里添加flowOn,避免重复代码:

import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.flow.flowOn
import kotlinx.coroutines.flow.merge

fun makeInputFlow() = flow {
    while (shouldMakeRequest()) {
        // 阻塞网络请求
        // CPU密集型操作
        results.forEach { emit(it) }
        yield()
    }
}.flowOn(Dispatchers.IO) // 直接指定调度器

suspend fun runFlows() {
    val flow1 = makeInputFlow()
    val flow2 = makeInputFlow()
    val flow3 = makeInputFlow()

    val merged = listOf(flow1, flow2, flow3).merge()
    merged.collect { println(it) }
}

原理说明

flowOn操作符会将Flow上游生产逻辑(即flow构建器内的代码)切换到指定的调度器线程执行,而下游的collect仍然保持在原来的上下文。每个Flow的生产过程会在独立的协程/线程中运行,这样多个Flow可以同时处理阻塞或CPU任务,不会互相抢占资源,CPU利用率也会提升。

如果你的代码中CPU密集型操作占比更高,建议替换为Dispatchers.Default,它专门针对CPU密集任务优化了线程池大小(默认是CPU核心数)。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 22:28:27