如何基于Flux A的最后值动态创建Flux B并拼接两者?
基于Flux最后一个值拼接新Flux的实现方案
问题描述
存在两个Flux:fluxA和fluxB,其中fluxB必须依赖fluxA的最后一个值才能创建。需要生成一个拼接后的Flux(fluxC),包含fluxA的所有元素以及fluxB的所有元素,且fluxB的创建必须延迟到获取到fluxA的最后一个值之后。用户给出的示例尝试代码为:
fluxC = fluxA.concatWith(lastA -> createFluxB(lastA))
解决方案
核心思路
要实现需求,需满足三个关键点:
- 先完整输出
fluxA的所有元素 - 在
fluxA执行完成后,获取其最后一个值并以此创建fluxB - 将
fluxB的元素续接在fluxA元素之后输出
代码实现
情况1:fluxA为热流(已共享或无需避免重复订阅)
直接通过concatWith结合last()和flatMapMany()实现:
Flux<T> fluxC = fluxA.concatWith(fluxA.last().flatMapMany(this::createFluxB));
fluxA.last()会等待fluxA执行完毕后取出其最后一个值flatMapMany()将该值传入createFluxB()生成目标fluxBconcatWith()保证fluxA的所有元素输出完成后,才开始输出fluxB的元素
情况2:fluxA为冷流(避免重复执行源逻辑)
冷流多次订阅会重复执行源逻辑,因此需要先对fluxA做共享或缓存处理:
方式1:用publish()共享流
ConnectableFlux<T> sharedFluxA = fluxA.publish(); Flux<T> fluxC = sharedFluxA.concatWith(sharedFluxA.last().flatMapMany(this::createFluxB)); sharedFluxA.connect();
publish()将冷流转为可连接的热流,确保fluxA仅被订阅一次connect()触发流的执行
方式2:用cache()缓存元素
Flux<T> cachedFluxA = fluxA.cache(); Flux<T> fluxC = cachedFluxA.concatWith(cachedFluxA.last().flatMapMany(this::createFluxB));
cache()会缓存fluxA的所有元素,后续订阅直接读取缓存内容,避免重复执行源逻辑
内容的提问来源于stack exchange,提问作者George Hawkins
相关产品推荐
相关产品推荐

