如何并行运行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
相关产品推荐
相关产品推荐

