Project Reactor并行流加sequential后processMap多次触发问题咨询
问题成因
- 你对
sequential()的作用存在误解:sequential()仅负责将ParallelFlux转回普通串行Flux,不会合并流中的元素,只会按顺序传递上游发射的所有元素,并不具备等待所有并行任务完成后仅返回一个汇总结果的能力。 - 你的并行流逻辑中,每处理完一个
cId就会通过return Flux.just(tempHashMap);发射一次tempHashMap的引用,上游有多少个待处理的cId,sequential()之后就会向下游传递多少个元素,因此doOnNext会被触发对应次数。 - 额外风险提示:你直接在并行流中修改外部非线程安全的
HashMap,高并发场景下可能出现数据丢失、扩容死循环等异常,不符合响应式编程无状态的最佳实践。
正确实现方案
你可以使用collect类操作符汇总所有并行处理的结果,最终只生成一个Map向下游发射,即可保证processMap只被调用一次,同时规避外部变量的线程安全问题:
Flux.fromIterable(cList) .parallel(20) .runOn(Schedulers.boundedElastic()) // 每个cId处理后返回有效键值对,空结果直接返回空 .flatMap(cId -> { List<String> lb = api.getlb(p,cId); return lb.isEmpty() ? Mono.empty() : Mono.just(new AbstractMap.SimpleEntry<>(cId, lb)); }) .sequential() // 收集所有有效键值对汇总为Map,所有并行任务完成后仅向下游发射一次汇总结果 .collectMap(Map.Entry::getKey, Map.Entry::getValue) .publishOn(Schedulers.single()) .doOnNext(collectMap-> lb = processMap(collectMap)) .doOnError(throwable -> { log.error("Error while getting list of lb : {} ", throwable.getMessage()); }) .subscribe();
内容的提问来源于stack exchange,提问作者perplexedDev
相关产品推荐
相关产品推荐

