Project Reactor命令式代码中Flux.zip未并行执行的原因与解决办法
未并行执行的原因
- 你代码里的
Thread.sleep和打印逻辑直接写在了blah1/blah2的方法执行流程中,没有封装进Mono的异步处理链路。调用test1.blah1()时会立刻在当前线程执行sleep、打印逻辑,执行完成后才会返回Mono对象。所以代码会先同步执行完blah1的全部逻辑,再执行blah2的逻辑,自然是blah1的输出在前。 - 你的测试代码仅完成了Flux流水线的组装,没有触发订阅。Reactor的异步执行逻辑是在订阅时才会启动,未订阅的情况下所有代码都是方法调用阶段的同步阻塞执行。
- 即便完成订阅,如果没有手动指定调度器,Mono默认会在订阅线程上同步执行,也不会自动并行处理。
修复方案
首先调整方法实现,将阻塞逻辑和打印逻辑封装到Mono的处理链路中:
public class Test1 { public Mono<String> blah1() { return Mono.fromCallable(() -> { try { Thread.sleep(5000); } catch (InterruptedException e) { e.printStackTrace(); } System.out.println("blah1 done"); return "blah1"; }); } public Mono<String> blah2() { return Mono.fromCallable(() -> { try { Thread.sleep(1000); } catch (InterruptedException e) { e.printStackTrace(); } System.out.println("blah2 done"); return "blah2"; }); } }
之后修改测试代码,给两个Mono指定异步调度器,并触发订阅等待执行完成:
@Test public void blah1Test() { Flux<Tuple2<String, String>> s = Flux.zip( test1.blah1().subscribeOn(Schedulers.boundedElastic()), test1.blah2().subscribeOn(Schedulers.boundedElastic()) ); // 等待异步执行完成,避免测试提前结束 s.blockLast(); }
修改后运行即可得到预期的输出顺序:blah2 done在前,blah1 done在后。
内容的提问来源于stack exchange,提问作者Brian
相关产品推荐
相关产品推荐

