Spring Boot Webflux中串联API调用及新增请求整合方案咨询
可以在同一反应式流中实现,无需阻塞现有逻辑
完全不需要阻塞现有执行流程,通过WebFlux的反应式操作符就能把新请求集成到现有流中,全程保持非阻塞特性。核心思路是保留第一步API调用的结果(ResolvedObjectIds),用它构建新请求的EventsRequest,再针对每个对象ID并行获取原有ObjectInfo和新请求的AvailableResponse,最后合并结果。
修改后的完整实现代码
public Mono<List<ObjectInfo>> getEventsForReference(String reference, Postcode postcode) { return connector.resolveForExternalReference(reference) .flatMap(resolvedIds -> { // 用第一步的结果构建EventsRequest(根据你的业务逻辑补充字段) EventsRequest eventsRequest = buildEventsRequest(resolvedIds); // 遍历每个对象ID,并行获取原有数据和新请求结果 return Flux.fromIterable(resolvedIds.getIds()) .flatMap(id -> { // 原有逻辑:获取ObjectInfo Mono<ObjectInfo> objectInfoMono = getEventsFor(id) .map(obj -> redact(obj, id, postcode)); // 新请求:获取可用分流,404时返回空列表避免流中断 Mono<AvailableResponse> diversionsMono = getAvailableDiversionsForObject(id.getId(), eventsRequest) .defaultIfEmpty(AvailableResponse.builder().availableDiversions(Collections.emptyList()).build()); // 合并两个结果,填充availableDiversions字段 return Mono.zip(objectInfoMono, diversionsMono) .map(tuple -> { ObjectInfo objInfo = tuple.getT1(); objInfo.setAvailableDiversions(tuple.getT2().getAvailableDiversions()); return objInfo; }); }) .collectList(); }) .doOnError(throwable -> log.error(SERVICE_ERROR, throwable.getMessage())); } // 补充:根据ResolvedObjectIds构建EventsRequest的方法(需适配你的业务字段) private EventsRequest buildEventsRequest(ResolvedObjectIds resolvedIds) { return EventsRequest.builder() .reference(resolvedIds.getExternalReference()) // 示例字段,替换为实际需要的参数 .objectIds(resolvedIds.getIds().stream().map(ResolvedObjectId::getId).toList()) .build(); } // 修正原代码中的笔误(uid改为id,返回类型对齐ObjectInfo) private Mono<ObjectInfo> getEventsFor(ResolvedObjectId id) { return gateway.getObjectBy(id.getId()).zipWith(gateway.getEventsBy(id.getId())) .map(tuple -> combineEvents(tuple.getT1(), tuple.getT2())); } // 修正原代码中的变量名错误(parcel改为obj) private ObjectInfo redact(ObjectInfo obj, ResolvedObjectId id, Postcode postcode) { boolean fullResponse = ResolvedObjectId.isReturn(obj.getBarcode()) || id.isReturn() || obj.hasPostcode(postcode); return fullResponse ? obj : redact(obj); }
关键实现说明
- 保留第一步结果:通过
flatMap包裹住原有流,把ResolvedObjectIds传递下去,用于构建EventsRequest。 - 并行调用提升性能:用
Mono.zip同时发起原有ObjectInfo请求和新的分流请求,避免串行等待,充分利用反应式的异步特性。 - 容错处理:用
defaultIfEmpty处理新请求返回404的场景,确保流不会中断,返回空列表作为默认值。
内容的提问来源于stack exchange,提问作者JsbDev
相关产品推荐
相关产品推荐

