You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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;
  // 其他数据字段
}

业务流程:

  1. 获取A表最新1000条数据;
  2. 根据A的rid关联查询B1、B2数据;
  3. 根据B1、B2的cid关联查询C数据;
  4. 对比A与C的特定参数生成报告。

当前实现中,查询B1和B2的逻辑为串行执行,需通过缓存Stage1结果、拆分并行查询再合并的方式优化性能。

核心优化思路

  1. 缓存Stage1结果:对获取A列表的Flux进行缓存,避免重复查询数据库;
  2. 并行发起B1/B2查询:基于缓存的A列表,同时发起B1和B2的批量查询,利用WebFlux的异步特性提升效率;
  3. 合并并行结果:等待B1、B2查询完成后,合并结果集;
  4. 关联查询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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.01 00:04:02