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
相关产品推荐
相关产品推荐

