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

Spring WebFlux如何避免重复查库,两次遍历同一Flux元素?

解决Spring Reactor重复数据库查询的方案

核心思路是利用Flux.cache()将冷序列转换为带缓存的序列,避免多次订阅触发重复查询。以下是具体实现步骤:

1. 缓存查询结果

首先获取数据库查询的Flux并调用cache(),这样后续所有订阅都会复用第一次查询的结果,不会触发新的数据库请求:

Flux<MyEntity> cachedEntities = repo.findAll().cache();

冷序列(比如直接从仓库查询的Flux)每次订阅都会重新执行上游操作(这里就是DB查询),cache()会缓存第一次订阅时发射的所有数据,后续订阅直接读取缓存。

2. 实现方法A:计算总和并保存SumEntity

基于缓存的Flux处理总和计算逻辑:

Mono<Void> processSum() {
    return cachedEntities
        // 提取金额字段
        .map(MyEntity::getAmount)
        // 累加计算总金额
        .reduce(BigDecimal.ZERO, BigDecimal::add)
        // 转换为SumEntity
        .map(totalAmount -> new SumEntity(totalAmount))
        // 保存到sumRepo
        .flatMap(sumRepo::save)
        // 转换为Mono<Void>表示任务完成
        .then();
}

3. 实现方法B:更新实体usable状态并保存

同样基于缓存的Flux执行更新操作:

Mono<Void> updateUsableStatus() {
    return cachedEntities
        // 修改usable为false
        .map(entity -> {
            entity.setUsable(false);
            return entity;
        })
        // 批量保存回repo
        .flatMap(repo::save)
        // 转换为Mono<Void>表示任务完成
        .then();
}

4. 合并执行两个任务

使用Mono.when()合并两个任务,确保它们都基于同一批缓存数据执行:

// 合并两个异步任务,等待全部完成
Mono<Void> combinedTasks = Mono.when(processSum(), updateUsableStatus());

// 订阅触发执行(在Spring环境中如果是WebFlux或Reactive Bean,可直接返回Mono让框架处理订阅)
combinedTasks.subscribe();

为什么不用转List的方式?

将Flux转成List再转回Flux确实能避免重复查询,但这种方式会一次性把所有数据加载到内存中,当数据量较大时会消耗更多内存,也无法利用Reactor的背压机制。而cache()是流式处理,数据按需流动,更符合响应式编程的非阻塞、高效特性。

内容的提问来源于stack exchange,提问作者nadavgam

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 15:45:25