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

