WebFlux技术问询:如何传递Flux完成链式调用与对象更新
解决WebClient调用链中传递原始Flux到更新步骤的问题
要实现你需要的调用链,核心是在响应式流中保留原始的Flux<ObjA>数据,同时完成属性提取、详情获取和对象更新的操作。下面提供两种符合Reactor响应式编程思想的实现方案:
方案一:收集原始列表后统一处理
适合数据量不大的场景,把原始Flux收集为List后,一次性处理属性提取、详情获取和对象更新:
import java.util.List; import java.util.Map; import java.util.stream.Collectors; public Flux<ObjA> getUpdatedObjects() { return getListOfObjects() // 将Flux转为Mono<List<ObjA>>,保留所有原始对象 .collectList() .flatMapMany(objAList -> { // 提取所有prop并去重 List<Long> uniqueProps = objAList.stream() .map(ObjA::getProp) .distinct() .collect(Collectors.toList()); // 调用详情API,将结果转为以prop为键的Map,方便后续查找 return thenGetListofPropDetails(uniqueProps) .collectMap(Props::getProp) .flatMapIterable(propDetailsMap -> { // 遍历原始列表,更新每个ObjA的propDescription return objAList.stream() .map(objA -> { Props props = propDetailsMap.get(objA.getProp()); if (props != null) { objA.setPropDescription(props.getDescription()); } return objA; }) .collect(Collectors.toList()); }); }); }
方案二:缓存原始流实现流式处理
适合数据量较大的场景,通过缓存原始Flux避免重复调用第一个API,同时保持流式处理的特性:
import java.util.Map; public Flux<ObjA> getUpdatedObjects() { // 缓存原始Flux,确保第一个API只被调用一次 Flux<ObjA> originalObjAFlux = getListOfObjects().cache(); // 异步获取去重后的prop详情,并转为Map Mono<Map<Long, Props>> propDetailsMapMono = originalObjAFlux .map(ObjA::getProp) .distinct() .collectList() .flatMap(uniqueProps -> thenGetListofPropDetails(uniqueProps) .collectMap(Props::getProp)); // 将每个原始ObjA和详情Map结合,完成更新 return originalObjAFlux .flatMap(objA -> propDetailsMapMono .map(detailsMap -> { Props props = detailsMap.get(objA.getProp()); if (props != null) { objA.setPropDescription(props.getDescription()); } return objA; })); }
关键说明
- 两种方案都遵循响应式编程的非阻塞原则,避免了同步阻塞操作。
- 方案一通过
collectList()将流式数据转为集合,适合小批量数据;方案二用cache()缓存原始流,适合大数据量场景,内存占用更低。 - 详情API调用仅执行一次(基于去重后的prop列表),避免了重复请求。
内容的提问来源于stack exchange,提问作者Nagarjun
相关产品推荐
相关产品推荐

