WebFlux中如何并行执行获取B1与B2的逻辑?
WebFlux并行查询优化实现方案
实体与业务背景
现有实体类对应数据库表关系如下:
class A { private Long id; private Long rid; // 关联A与B1、B2 // 其他数据字段 } class B1 { private Long rid; // 关联A private Long cid; // 关联C // 其他数据字段 } class B2 { private Long rid; // 关联A private Long cid; // 关联C // 其他数据字段 } class C { private Long id; // 其他数据字段 }
业务流程:
- 获取A表最新1000条数据;
- 根据A的
rid关联查询B1、B2数据; - 根据B1、B2的
cid关联查询C数据; - 对比A与C的特定参数生成报告。
当前实现中,查询B1和B2的逻辑为串行执行,需通过缓存Stage1结果、拆分并行查询再合并的方式优化性能。
核心优化思路
- 缓存Stage1结果:对获取A列表的Flux进行缓存,避免重复查询数据库;
- 并行发起B1/B2查询:基于缓存的A列表,同时发起B1和B2的批量查询,利用WebFlux的异步特性提升效率;
- 合并并行结果:等待B1、B2查询完成后,合并结果集;
- 关联查询C并生成报告:基于合并后的B1/B2数据提取
cid,查询C列表,最后关联A与C生成报告。
具体代码实现
假设已定义对应的Reactive Repository(ARepository、B1Repository、B2Repository、CRepository),实现代码如下:
import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import java.util.List; import java.util.Set; import java.util.stream.Collectors; // 业务服务类 public class ReportService { private final ARepository aRepository; private final B1Repository b1Repository; private final B2Repository b2Repository; private final CRepository cRepository; // 构造函数注入Repository public ReportService(ARepository aRepository, B1Repository b1Repository, B2Repository b2Repository, CRepository cRepository) { this.aRepository = aRepository; this.b1Repository = b1Repository; this.b2Repository = b2Repository; this.cRepository = cRepository; } public Mono<Report> generateReport() { // Stage1: 获取最新1000条A数据并缓存 Flux<A> aFlux = aRepository.findTop1000ByOrderByIdDesc() .cache(); // 缓存结果,避免后续多次查询 // 提取所有A的rid集合 Mono<Set<Long>> ridSetMono = aFlux .map(A::getRid) .collect(Collectors.toSet()); // Stage2&3: 并行查询B1和B2 Mono<List<B1>> b1ListMono = ridSetMono.flatMapMany(b1Repository::findByRidIn) .collectList(); Mono<List<B2>> b2ListMono = ridSetMono.flatMapMany(b2Repository::findByRidIn) .collectList(); // 合并B1、B2结果,提取所有cid集合 Mono<Set<Long>> cidSetMono = Mono.zip(b1ListMono, b2ListMono) .map(tuple -> { Set<Long> cids = tuple.getT1().stream().map(B1::getCid).collect(Collectors.toSet()); cids.addAll(tuple.getT2().stream().map(B2::getCid).collect(Collectors.toSet())); return cids; }); // 查询C数据 Mono<List<C>> cListMono = cidSetMono.flatMapMany(cRepository::findByIdIn) .collectList(); // 关联A、C数据并生成报告 return Mono.zip(aFlux.collectList(), cListMono) .map(tuple -> { List<A> aList = tuple.getT1(); List<C> cList = tuple.getT2(); // 这里实现A与C的特定参数对比逻辑,生成Report return buildReport(aList, cList); }); } // 具体的报告生成逻辑,根据业务需求实现 private Report buildReport(List<A> aList, List<C> cList) { // 业务逻辑实现 return new Report(); } }
关键说明
- 缓存Stage1结果:使用
cache()操作符缓存A的查询结果,确保后续提取rid、关联查询时不会重复访问数据库; - 并行查询实现:通过
Mono.zip将B1、B2的查询Mono合并,Reactor会自动并行执行这两个异步查询任务; - 批量查询优化:使用
findByRidIn、findByIdIn等批量查询方法,减少数据库IO次数,提升整体效率。
内容的提问来源于stack exchange,提问作者doptimusprime
相关产品推荐
相关产品推荐

