Spring WebFlux:如何等待多个Mono并行执行完成
优雅实现并行REST请求并等待全部完成
核心方案:用Reactor的ParallelFlux实现并行执行
替代CountDownLatch的关键是利用Reactor的异步编排能力,通过ParallelFlux实现多线程并行请求,同时借助内置操作符轻松等待所有请求完成。
完整代码示例
假设我们用Spring WebClient作为REST客户端,以下是具体实现步骤:
- 配置WebClient
@Bean public WebClient webClient() { return WebClient.create(); }
- 封装REST请求方法
private Mono<MyResponse> callRestEndpoint(String url, WebClient webClient) { return webClient.get() .uri(url) .retrieve() .bodyToMono(MyResponse.class); }
- 并行执行并等待所有请求完成
import reactor.core.publisher.Flux; import reactor.core.scheduler.Schedulers; import java.util.List; import java.util.stream.Collectors; public List<MyResponse> executeParallelRequests(List<String> targetUrls, WebClient webClient) { // 生成所有待执行的请求Mono实例 List<Mono<MyResponse>> requestMonos = targetUrls.stream() .map(url -> callRestEndpoint(url, webClient)) .collect(Collectors.toList()); // 并行执行并收集所有响应结果 return Flux.fromIterable(requestMonos) // 指定并行度,默认是CPU核心数,可根据目标服务并发限制调整 .parallel(Runtime.getRuntime().availableProcessors()) // IO密集型任务用boundedElastic调度器,自动管理线程资源避免爆炸 .runOn(Schedulers.boundedElastic()) // 展开每个Mono并执行请求 .flatMap(mono -> mono) // 阻塞等待所有请求完成并收集结果(非响应式环境适用) .collectList() .block(); }
关键细节说明
- parallel():将普通Flux转为ParallelFlux开启并行模式,可手动指定并行度(比如
parallel(targetUrls.size()),但注意不要过度并发导致资源耗尽)。 - runOn(Schedulers.boundedElastic()):为并行任务指定线程池,
boundedElastic专门适配IO密集型操作(如HTTP请求),会动态创建线程但限制总数;CPU密集型任务改用Schedulers.parallel()。 - 等待完成的不同方式:
- 需要收集所有结果:用
collectList().block(),同步等待并获取响应列表。 - 只需确认所有请求完成无需结果:用
.then().block(),等待所有任务执行完毕即可。
- 需要收集所有结果:用
- 异常处理:可在
flatMap中添加异常兜底逻辑,避免单个请求失败中断整个流程:
.flatMap(mono -> mono.onErrorResume(e -> { log.error("请求执行失败", e); return Mono.just(new MyResponse("请求异常")); }))
响应式环境适配
如果是在Spring WebFlux等响应式框架中,不要使用block(),直接返回Mono<List<MyResponse>>让框架处理订阅逻辑:
public Mono<List<MyResponse>> executeParallelRequestsReactive(List<String> targetUrls, WebClient webClient) { List<Mono<MyResponse>> requestMonos = targetUrls.stream() .map(url -> callRestEndpoint(url, webClient)) .collect(Collectors.toList()); return Flux.fromIterable(requestMonos) .parallel() .runOn(Schedulers.boundedElastic()) .flatMap(mono -> mono) .collectList(); }
内容的提问来源于stack exchange,提问作者Ashika Umanga Umagiliya
相关产品推荐
相关产品推荐

