如何并行执行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完成,但核心问题是阻塞任务没有被分配到独立线程池执行。
修改步骤
- 异步化阻塞操作:给每个获取原料的Mono添加
subscribeOn(Schedulers.boundedElastic()),将阻塞任务提交到弹性线程池,避免阻塞主线程,同时实现并行执行。 - 利用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
相关产品推荐
相关产品推荐

