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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 07:15:05