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

Spring Project Reactor实现分组内顺序、组间并行处理方案问询

解决方案:实现组内顺序、组间并行的Reactor处理逻辑

你的核心问题是错误地拆分了分组单元:原代码先通过flatMap(f -> f)把所有分组的元素打散成单个元素的Flux,再开启并行,导致同一组的元素可能被分配到不同线程,完全破坏了组内顺序。正确的思路是将每个分组作为并行处理的单元,保证组内元素顺序执行,同时让不同分组在不同线程并行运行。

修正后的代码示例

List<CurrencyWithValue> pairWithValueList = Arrays.asList(
        new CurrencyWithValue("EUR", 1.0),
        new CurrencyWithValue("USD", 1.0),
        new CurrencyWithValue("CHF", 1.0),
        new CurrencyWithValue("SEK", 1.0),
        new CurrencyWithValue("SEK", 2.0),
        new CurrencyWithValue("EUR", 2.0),
        new CurrencyWithValue("USD", 2.0),
        new CurrencyWithValue("CHF", 2.0),
        new CurrencyWithValue("SEK", 3.0),
        new CurrencyWithValue("EUR", 3.0),
        new CurrencyWithValue("USD", 3.0),
        new CurrencyWithValue("CHF", 3.0),
        new CurrencyWithValue("EUR", 4.0),
        new CurrencyWithValue("USD", 4.0),
        new CurrencyWithValue("CHF", 4.0),
        new CurrencyWithValue("SEK", 4.0),
        new CurrencyWithValue("SEK", 5.0)
);

Flux.fromIterable(pairWithValueList)
    // 1. 按货币属性分组,得到每个分组的独立Flux
    .groupBy(CurrencyWithValue::getCurrency)
    // 2. 开启并行处理,最多同时处理3个分组
    .parallel(3)
    // 3. 指定并行执行的线程池
    .runOn(Schedulers.newBoundedElastic(3, 1000, "k-task"))
    // 4. 对每个分组内的元素做顺序处理:GroupedFlux本身保留原顺序,concatMap确保组内元素依次执行
    .concatMap(groupedFlux -> groupedFlux
        .doOnNext(item -> System.out.println(
            "Thread:" + Thread.currentThread().getName() + 
            "::Data:" + item.getCurrency() + "::" + item.getAmount()
        ))
    )
    // 可选:将并行流转回串行,方便后续统一处理
    .sequential()
    // 等待所有任务执行完成(测试场景用)
    .blockLast();

关键逻辑说明

  1. 分组作为并行单元:groupBy后直接调用parallel(),让每个分组(而非单个元素)成为并行处理的最小单元,确保同一组的所有元素在同一个线程上下文里执行。
  2. 组内顺序保证:GroupedFlux会严格保留原Flux中该分组元素的顺序,配合concatMap(或直接使用map/doOnNext),可以保证组内元素按输入顺序依次处理。
  3. 组间并行控制:parallel(3)限制同时运行的分组数量,避免线程资源耗尽;runOn指定的线程池会为每个并行分组分配独立线程,实现组间并行。

原代码错误分析

你之前的代码先通过flatMap(f -> f)将所有分组的元素合并成一个无分组的Flux,再开启并行,导致:

  • 同一组的元素被拆分成独立的并行任务,可能被分配到不同线程
  • 线程调度的随机性会打乱组内元素的执行顺序,完全不符合需求

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 02:52:47