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

如何并行运行Flux但按原始序列顺序收集返回结果

问题原因

你当前使用的flatMap操作符本身不保证下游发射顺序和上游原始序列顺序一致,只会按照异步任务的完成顺序发射结果,因此最终收集到的列表顺序会和预期不符。

解决方案

直接将flatMap替换为Reactor原生提供的flatMapSequential操作符即可。这个操作符既会保留异步任务并行执行的特性,同时内部会维护原始序列的顺序,缓存先完成的后续位置的任务结果,等前序位置的异步任务执行完成后,再按原始顺序向下游发射结果,最终收集到的列表就能符合你的预期。

修改后的代码如下:

val mono = Flux.fromIterable(3 downTo 1)
      .map { it.toString() }
      // 替换为flatMapSequential,也可以传入第二个Int参数自定义并行度,默认和flatMap的并行度一致
      .flatMapSequential { number -> task(number) }
      .doOnNext { println("Number of $it") }
      // 此处仍会按任务实际完成顺序输出1、3、2,符合并行执行的预期
      .collectList()
      .doOnNext { println(it) }
      // 最终得到的列表为预期的[3,2,1]
      .block()

注意:因为需要缓存先完成的后续位置的任务结果,flatMapSequential会产生一定的内存开销,元素量级极大的场景下需要提前评估内存占用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 00:45:04