Flux并行map处理后sequential仍乱序 如何实现多线程并行且输出保序
问题原因
你目前使用ParallelFlux + sequential()组合出现乱序,是因为sequential()仅会按照并行轨(rail)的顺序合并结果,而非按原始流的元素提交顺序合并,只要不同轨的元素处理耗时不同,输出顺序就会和原始顺序不一致。
解决方案
使用flatMapSequential算子即可满足需求,该算子原生支持并行异步处理元素、按原始输入顺序输出结果,完美匹配你并行解密后按序转发的业务场景。
修改后核心代码如下
final Disposable disposable = workersFlux // 第一个参数为最大并行处理数,可根据你的线程池配置调整,默认值为256 .flatMapSequential(workers -> Mono.just(workers) // 单个元素的处理调度到并行线程池 .subscribeOn(Schedulers.parallel()) .map(w -> { System.out.println("Thread id: " + Thread.currentThread().getId()); w[0].calculate(); return w; }) , 10) // 示例中一共10个元素,这里设为10实现全并行处理 .doOnNext(workers -> System.out.println(workers[0].getId())) .subscribe();
方案说明
- 所有元素的处理逻辑会提交到独立线程执行,你依然可以通过输出的
Thread id验证多线程并行的效果 - 无论单个元素的处理耗时多长,
flatMapSequential都会严格按照原始流的元素顺序输出结果,前面的元素未处理完成时,后续已处理完成的元素会自动缓存,等前面的元素发出后再按序输出 - 直接套用该方案到你的业务场景中,把
calculate()替换为你的字节数组解密逻辑即可。
内容的提问来源于stack exchange,提问作者Anna Klein
相关产品推荐
相关产品推荐

