使用ParallelFlux实现多ID并行处理及每个ID的双下游API并行调用的可行性问询
使用ParallelFlux实现多ID并行处理及每个ID的双下游API并行调用的可行性问询
当然可以实现!你的思路方向是对的,但目前的代码里有个关键问题——getInfoByID里的block()调用拖了后腿,它会阻塞当前线程,导致原本能并行的两个下游API变成串行执行,还浪费了ParallelFlux分配的线程资源。
我们只需要把阻塞的逻辑改成响应式非阻塞的写法,就能同时实现多ID并行处理和每个ID下双API并行调用的需求,具体调整如下:
第一步:修改主流程代码
原来的flatMap里用Mono.fromCallable包裹getInfoByID是多余的,因为我们要把getInfoByID改成返回Mono<Info>,直接在flatMap里调用即可:
Flux.fromIterable(idList) .parallel() .runOn(Schedulers.boundedElastic()) .flatMap(this::getInfoByID) // 直接调用返回Mono的方法 .sequential() .collectList() // 用Reactor内置的collectList更简洁 .block();
第二步:改造getInfoByID方法
去掉阻塞的block()调用,让方法返回Mono<Info>,这样Mono.zip()会自动并行触发两个API调用,不会阻塞线程:
private Mono<Info> getInfoByID(String id) { // 假设api1.getDetails()和api2.getDetails()本身返回Mono Mono<Details1> mono1 = api1.getDetails(); Mono<Details2> mono2 = api2.getDetails(); return Mono.zip(mono1, mono2) .map(tuple2 -> { // 在这里处理两个API的返回结果,组装成Info对象 Details1 details1 = tuple2.getT1(); Details2 details2 = tuple2.getT2(); return new Info(id, details1, details2); // 根据你的实际构造逻辑调整 }); }
特殊情况处理:如果API调用是阻塞同步的
如果你的api1.getDetails()和api2.getDetails()是传统的阻塞同步方法(比如用了非响应式的HTTP客户端),那需要给它们单独指定调度器,避免阻塞ParallelFlux的线程池:
private Mono<Info> getInfoByID(String id) { // 把阻塞的API调用放到boundedElastic调度器上执行 Mono<Details1> mono1 = Mono.fromCallable(api1::getDetails) .subscribeOn(Schedulers.boundedElastic()); Mono<Details2> mono2 = Mono.fromCallable(api2::getDetails) .subscribeOn(Schedulers.boundedElastic()); return Mono.zip(mono1, mono2) .map(tuple2 -> { // 组装Info对象 return new Info(id, tuple2.getT1(), tuple2.getT2()); }); }
效果说明
- 外层的
parallel()配合runOn(Schedulers.boundedElastic())会把多个ID分配到不同线程并行处理 - 每个ID对应的
mono1和mono2会通过Mono.zip()并行触发,Reactor会自动调度线程执行这两个异步任务 - 全程保持响应式非阻塞,只有最后
block()会同步获取最终结果(如果业务允许,也可以去掉block()保持全响应式链)
备注:内容来源于stack exchange,提问作者Indunil Rathnayake
相关产品推荐
相关产品推荐

