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

如何并行执行Reactor Flux操作以提升数据获取效率?

问题:Reactor并行获取多数据源数据并乱序返回

需要从多个数据源获取独立数据项,合并后返回,希望并行执行数据获取操作(因为这部分耗时最长)。尝试用Flux.fromStream()和Mono.zipWith()实现,但实际是串行执行(需等待前一项完成才开始下一项)。

模拟场景:制作10个三明治,随机延迟获取花生酱和果冻,预期结果因延迟随机而乱序,但实际按顺序完成每个三明治。

测试代码:

class FluxOrderingTest {
  private final Random random = new Random();

  @Test
  void testAsyncSandwich() {
    Flux<String> sandwiches = Flux.fromStream(IntStream.rangeClosed(1, 10).boxed())
            .flatMap(number -> getJelly(number).zipWith(getPeanutButter(number)))
            .map(this::makeSandwich);

    StepVerifier.create(sandwiches)
            .expectNextCount(10)
            .verifyComplete();
  }

  private String makeSandwich(Tuple2<String, String> ingredients) {
    System.out.printf("Combined %s with %s%n", ingredients.getT1(), ingredients.getT2());
    return "Sandwich";
  }

  private Mono<String> getPeanutButter(Integer number) {
    return Mono.fromSupplier(() -> withDelay("the peanut butter for " + number));
  }

  private Mono<String> getJelly(Integer number) {
    return Mono.fromSupplier(() -> withDelay("the jelly for " + number));
  }

  private String withDelay(String s) {
    try {
      Thread.sleep(random.nextInt(1000));
    } catch (InterruptedException e) {
      e.printStackTrace();
    }
    System.out.println("Got " + s);
    return s;
  }
}

实际输出:

Got the jelly for 1
Got the peanut butter for 1
Combined the jelly for 1 with the peanut butter for 1
Got the jelly for 2
Got the peanut butter for 2
Combined the jelly for 2 with the peanut butter for 2
Got the jelly for 3
Got the peanut butter for 3
Combined the jelly for 3 with the peanut butter for 3
Got the jelly for 4
Got the peanut butter for 4
Combined the jelly for 4 with the peanut butter for 4
Got the jelly for 5
Got the peanut butter for 5
Combined the jelly for 5 with the peanut butter for 5
Got the jelly for 6
Got the peanut butter for 6
Combined the jelly for 6 with the peanut butter for 6
Got the jelly for 7
Got the peanut butter for 7
Combined the jelly for 7 with the peanut butter for 7
Got the jelly for 8
Got the peanut butter for 8
Combined the jelly for 8 with the peanut butter for 8
Got the jelly for 9
Got the peanut butter for 9
Combined the jelly for 9 with the peanut butter for 9
Got the jelly for 10
Got the peanut butter for 10
Combined the jelly for 10 with the peanut butter for 10

请问如何实现所有原料的并行获取,且不关心最终结果的顺序?


解决方案

问题根源

代码中getJelly和getPeanutButter里的Mono.fromSupplier是在调用线程上执行同步阻塞操作(Thread.sleep),Reactor默认不会自动将阻塞操作异步化,导致整个流串行执行。zipWith本身只是等待两个Mono完成,但核心问题是阻塞任务没有被分配到独立线程池执行。

修改步骤

  1. 异步化阻塞操作:给每个获取原料的Mono添加subscribeOn(Schedulers.boundedElastic()),将阻塞任务提交到弹性线程池,避免阻塞主线程,同时实现并行执行。
  2. 利用flatMap的并发能力:flatMap默认支持最高256个并发任务,足够覆盖10个三明治的场景,无需额外配置(如果需要调整并发数,可通过flatMap(..., concurrency)指定)。

修改后的代码

class FluxOrderingTest {
  private final Random random = new Random();

  @Test
  void testAsyncSandwich() {
    Flux<String> sandwiches = Flux.fromStream(IntStream.rangeClosed(1, 10).boxed())
            .flatMap(number -> getJelly(number).zipWith(getPeanutButter(number)))
            .map(this::makeSandwich);

    StepVerifier.create(sandwiches)
            .expectNextCount(10)
            .verifyComplete();
  }

  private String makeSandwich(Tuple2<String, String> ingredients) {
    System.out.printf("Combined %s with %s%n", ingredients.getT1(), ingredients.getT2());
    return "Sandwich";
  }

  private Mono<String> getPeanutButter(Integer number) {
    // 将阻塞操作切换到弹性线程池
    return Mono.fromSupplier(() -> withDelay("the peanut butter for " + number))
            .subscribeOn(Schedulers.boundedElastic());
  }

  private Mono<String> getJelly(Integer number) {
    // 将阻塞操作切换到弹性线程池
    return Mono.fromSupplier(() -> withDelay("the jelly for " + number))
            .subscribeOn(Schedulers.boundedElastic());
  }

  private String withDelay(String s) {
    try {
      Thread.sleep(random.nextInt(1000));
    } catch (InterruptedException e) {
      e.printStackTrace();
    }
    System.out.println("Got " + s);
    return s;
  }
}

预期输出示例(乱序)

Got the jelly for 3
Got the peanut butter for 1
Got the jelly for 2
Got the peanut butter for 3
Combined the jelly for 3 with the peanut butter for 3
Got the jelly for 5
Got the peanut butter for 2
Combined the jelly for 2 with the peanut butter for 2
Got the peanut butter for 5
Combined the jelly for 5 with the peanut butter for 5
...(后续输出因随机延迟无固定顺序)

补充说明

  • Schedulers.boundedElastic()专门用于处理阻塞IO操作,会根据需求动态创建线程,同时限制线程总数避免资源耗尽。
  • 如果数据源本身是异步非阻塞调用(比如WebClient请求),不需要subscribeOn,因为这类操作本身就是异步的;但对于JDBC、本地文件读取等阻塞操作,必须通过subscribeOn切换线程池。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 23:35:08