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

Kotlin Flow中Map操作符能否无需等待前次发射并行执行?

Kotlin Flow 并行处理耗时任务问题

问题描述

我在Flow的map操作中包含耗时任务(如网络请求),当前map会阻塞后续值的发射,必须等前一次map执行完成才会处理下一个值。我不想使用mapLatest,因为它会取消前一次的执行。请问有没有办法让每个值一发射就并行执行map内的代码?

以下代码可说明需求:

suspend fun main() = coroutineScope {

    val flow = flow {
        repeat(4) {
            emit(it) // 快速发射值
        }
    }.map {
        println("Loading")
        delay(1000) // 模拟网络请求
        println("I'm slow")
        it
    }

    flow.collect {
        println("HEY $it")
    }

}

解决方案

可以使用flatMapMerge操作符实现并行处理,它会为每个发射的值启动独立协程执行任务,还能通过参数控制并发数量。

示例代码1:基础并行处理

suspend fun main() = coroutineScope {

    val flow = flow {
        repeat(4) {
            emit(it) // 快速发射值
        }
    }.flatMapMerge {
        flow {
            println("Loading $it")
            delay(1000) // 模拟网络请求
            println("I'm slow $it")
            emit(it)
        }
    }

    flow.collect {
        println("HEY $it")
    }

}

运行后会看到4个"Loading"几乎同时打印,1秒后陆续打印"I'm slow"和"HEY",实现了并行执行。

示例代码2:控制并发度

如果需要限制同时运行的协程数量,给flatMapMerge传入concurrency参数即可:

.flatMapMerge(concurrency = 2) { // 最多同时运行2个协程
    flow {
        // 耗时任务逻辑
    }
}

另一种实现方式:结合map与async

也可以通过map配合async启动异步任务,再在收集时等待结果:

suspend fun main() = coroutineScope {

    val flow = flow {
        repeat(4) {
            emit(it)
        }
    }.map {
        async {
            println("Loading $it")
            delay(1000)
            println("I'm slow $it")
            it
        }
    }

    flow.collect { deferred ->
        println("HEY ${deferred.await()}")
    }

}

这种方式的收集顺序会与任务完成顺序一致,而flatMapMerge默认会按子流的发射顺序(即任务完成顺序)收集结果。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 15:52:27