Spring Reactive管道逻辑重复执行3次问题排查求助
问题原因分析
你的问题核心是Reactor冷流的多次订阅导致重复执行。Reactor中的Flux/Mono默认是冷序列,每一次订阅都会从头触发整个序列的执行逻辑。看你的上游代码:
currentProps = getProps()是一个冷Flux- 它被触发了三次订阅:
propsForEvCalculation = currentProps.flatMap(...)第一次订阅currentPropsfinalProps中的currentProps.collectList()第二次订阅currentPropspropsWithEv依赖propsForEvCalculation,而propsForEvCalculation又依赖currentProps,当finalProps中订阅propsWithEv.collectList()时,会间接触发第三次订阅currentProps
每一次订阅都会重新执行getProps()里的所有逻辑,包括makeAnApiCall()和doSomeOtherStuff(),所以你看到这些方法被执行了3次。
解决方案
要避免重复执行,需要把冷流转为可复用的热流,或者缓存序列的结果,确保只执行一次。以下是两种常用方案:
方案1:使用cache()缓存结果
直接在currentProps上添加cache()操作符,它会缓存序列的结果,后续订阅直接复用缓存的数据:
// 缓存currentProps的结果,避免多次订阅重复执行 Flux<PlayerPropDetail> currentProps = getProps().cache(); Flux<EvCalculationRequest> propsForEvCalculation = currentProps.flatMap(this::compareWithPreviousPropAndCreateRequests); Flux<PlayerPropDetail> propsWithEv = propsForEvCalculation .filter(prop -> Objects.nonNull(prop.getPropId())) .buffer(200) .flatMap(this::calculateValueForProp); Flux<PlayerPropDetail> finalProps = Flux.zip(propsForEvCalculation.collectList(), currentProps.collectList(), propsWithEv.collectList()) .map(zippedPropsData -> combineAllProps(zippedPropsData)) .flatMap(Flux::fromIterable).map(nbaPropsRepository::updateProp); return finalProps;
方案2:提前收集为Mono复用
如果不需要保持Flux的流式特性,可以提前将currentProps收集为Mono<List>,后续所有依赖都复用这个Mono,确保只执行一次:
// 提前收集并缓存结果,仅执行一次getProps() Mono<List<PlayerPropDetail>> currentPropsList = getProps().collectList().cache(); // 从缓存的List生成Flux供后续使用 Flux<PlayerPropDetail> currentProps = currentPropsList.flatMapMany(Flux::fromIterable); Flux<EvCalculationRequest> propsForEvCalculation = currentProps.flatMap(this::compareWithPreviousPropAndCreateRequests); Flux<PlayerPropDetail> propsWithEv = propsForEvCalculation .filter(prop -> Objects.nonNull(prop.getPropId())) .buffer(200) .flatMap(this::calculateValueForProp); // 直接复用缓存的currentPropsList,避免再次collect触发订阅 Flux<PlayerPropDetail> finalProps = Flux.zip(propsForEvCalculation.collectList(), currentPropsList, propsWithEv.collectList()) .map(zippedPropsData -> combineAllProps(zippedPropsData)) .flatMap(Flux::fromIterable).map(nbaPropsRepository::updateProp); return finalProps;
额外优化建议
- 如果你需要更精细的缓存控制(比如缓存过期、限制缓存大小),可以使用
cache(Duration)或cache(int)等重载方法 - 对于
makeAnApiCall()返回的Flux,也可以考虑在内部添加cache(),避免多次订阅重复发起API请求
内容的提问来源于stack exchange,提问作者Ryan Queenan
相关产品推荐
相关产品推荐

