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

Spring WebFlux如何缓存Flux结果并复用?避免重复调用

解决Flux响应复用避免重复外部调用的问题

你的核心问题是data这个Flux被两次订阅(一次在Flux.concat中,一次在aggregate的reduce操作中),每次订阅都会重新执行上游的invokeEndpoint调用。之前尝试的share()没生效,是因为Flux.concat是串行执行的——第一个订阅完成后,share()的引用计数降到0,上游被取消,第二个订阅触发时会重新执行上游逻辑。

方案1:使用cache()缓存流结果

cache()会订阅上游一次并缓存所有元素和终止信号,后续任何订阅都直接复用缓存内容,不会重复触发上游调用,完美适配你的串行订阅场景。

修改后的代码:

private Flux<Integer> getResponse(List<String> requests){
   Flux<Integer> data = Flux.fromIterable(requests)
      .flatMap(this::invokeEndpoint, 10)
      .cache(); // 缓存所有元素,后续订阅直接复用
   return Flux.concat(data, aggregate(data));
}

private Mono<Integer> aggregate(Flux<Integer> res){
   return res.reduce((a, b) -> a + b);
}

方案2:使用ConnectableFlux手动控制订阅

如果不需要永久缓存,只是在当前方法调用内复用流,可以用publish()创建ConnectableFlux,手动触发上游执行,确保所有订阅共享同一个流实例。

修改后的代码:

private Flux<Integer> getResponse(List<String> requests){
   ConnectableFlux<Integer> data = Flux.fromIterable(requests)
      .flatMap(this::invokeEndpoint, 10)
      .publish(); // 转换为可连接的Flux
   data.connect(); // 手动触发上游执行,所有订阅共享此流

   return Flux.concat(data, aggregate(data));
}

private Mono<Integer> aggregate(Flux<Integer> res){
   return res.reduce((a, b) -> a + b);
}

为什么share()不生效?

share()等价于publish().refCount(1),它会在订阅者数量从0变为1时触发上游订阅,从1变为0时取消上游订阅。而Flux.concat是先订阅第一个data并等待其完成,此时订阅者数量回到0,上游被取消;当开始订阅aggregate(data)时,订阅者数量再次从0变为1,share()会重新触发上游调用,导致重复请求。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 05:43:13