如何在Reactor Core中并行运行两个无依赖的Mono异步任务?
让两个无依赖的Mono并行执行的方案
当然可以实现并行!既然这两个Mono彼此没有依赖关系,我们完全可以让它们同时跑起来,避免不必要的串行等待,还能充分利用系统资源提升效率。不过先指出你代码里的几个小问题:
- 拼写错误:
mon.just()应该是Mono.just(),类型参数也建议大写(比如Response、Object,符合Java命名规范) - 第二个
Mono<Object>根本没被订阅:反应式流是冷流,只有调用subscribe()/block()这类订阅方法后,流中的逻辑才会执行,你原来的代码里这个任务其实根本没启动 - 直接用
block()会阻塞线程:虽然能拿到结果,但没发挥出反应式编程的异步非阻塞优势
下面是几种靠谱的实现方式,按需选择:
方法1:用Mono.zip()合并结果(推荐)
Mono.zip()会并行订阅所有传入的Mono,等所有任务都完成后,返回一个包含所有结果的Tuple。这种方式既能保证两个任务并行执行,又能同时获取它们的结果,是最常用的方案:
// 假设Response是你的自定义业务类型 Mono<Response> responseMono = Mono.just(new Response()); Mono<Object> objectMono = Mono.just(new Object()); // 并行执行两个Mono,阻塞等待所有结果返回 Tuple2<Response, Object> resultTuple = Mono.zip(responseMono, objectMono).block(); // 从Tuple中取出各自的结果 Response finalResponse = resultTuple.getT1(); Object finalObject = resultTuple.getT2(); // 根据业务需求返回结果,比如返回response return finalResponse;
方法2:用Flux.merge()触发并行执行
如果你的场景只需要确保两个任务都执行完,不需要同时获取它们的结果,可以用Flux.merge():它会订阅所有传入的流并并行执行,直到最后一个任务完成:
Mono<Response> responseMono = Mono.just(new Response()); Mono<Object> objectMono = Mono.just(new Object()); // 合并两个Mono为Flux,触发并行执行,阻塞等待所有任务完成 Flux.merge(responseMono, objectMono).blockLast(); // 之后再获取response的结果(注意:冷流每次订阅都会重新执行逻辑,所以如果想复用结果,还是推荐用zip的方式) return responseMono.block();
生产环境最佳实践:避免直接用block()
如果是在Spring WebFlux这类反应式框架中,尽量不要直接调用block(),而是返回Mono让框架处理订阅和异步流程,这样能保持非阻塞特性:
public Mono<Response> fetchResponse() { Mono<Response> responseMono = Mono.just(new Response()); Mono<Object> objectMono = Mono.just(new Object()); // 并行执行两个任务,完成后返回response结果 return Mono.zip(responseMono, objectMono) .map(resultTuple -> resultTuple.getT1()); }
这样既实现了并行,又完全符合反应式编程的设计理念,性能和扩展性更好。
内容的提问来源于stack exchange,提问作者Ed Diaz
相关产品推荐
相关产品推荐

