Spring WebFlux顺序调用API合并Mono结果异常及zipWhen写法咨询
问题背景
我尝试在Mono中实现顺序API调用并合并返回结果,其中第一个API返回的数据中包含调用第二个API所需的ID列表。
初始实现代码
AggregateService 初始写法
@Service @AllArgsConstructor public class AggregateService { private OperationClient operationClient; private OuvrageClient ouvrageClient; public Mono<OperationDetailDto> getDetails(String operationId){ final Mono<OperationFullDto> operationMono = operationClient.getOperation(operationId); return operationMono .flatMap(operation -> { final OperationDetailDto detail = new OperationDetailDto(); final Mono<List<OuvrageDto>> ouvrageMono = ouvrageClient.getOuvrages(operation.getOuvrageIds()); Mono<Tuple2<OperationFullDto, List<OuvrageDto>>> tuple2 = Mono.zip(operationMono,ouvrageMono); tuple2.map(result ->{ detail.setId(ServiceUtils.formatUuid(operationId)); detail.setTitre(result.getT1().getTitre()); detail.setOuvrages(result.getT2()); return detail; }); return Mono.just(detail); } ); } }
WebClient 实现
OperationClient
@Service public class OperationClient { private final WebClient webClient; public OperationClient(WebClient.Builder builder) { this.webClient = builder.baseUrl("lb://operation-microservice/operation/").build(); } public Mono<OperationFullDto> getOperation(String operationId){ return webClient .get() .uri("{id}" , operationId) .retrieve() .bodyToMono(OperationFullDto.class) .subscribeOn(Schedulers.parallel()) .log() .onErrorResume(ex->Mono.empty()); } }
OuvrageClient
public class OuvrageClient { private final WebClient webClient; public OuvrageClient(WebClient.Builder builder) { this.webClient = builder.baseUrl("lb://ouvrage-microservice/ouvrage/").build(); } public Mono<List<OuvrageDto>> getOuvrages(List<Long> ids){ System.out.println("the void getOuvrages was invoked " + ids ); return webClient .get() .uri("multi/{ids}",ids) .retrieve() .bodyToFlux(OuvrageDto.class) .subscribeOn(Schedulers.parallel()) .collectList() .onErrorReturn(Collections.emptyList()) .log(); } }
异常现象
运行时观察到OuvrageClient的方法被触发(控制台打印了方法进入的日志),但没有输出WebClient调用的响应式日志,最终AggregateService返回的结果全为空:
{"id": null,"titre": null,"ouvrages": null}
运行日志
2022-06-21 10:54:37.872 INFO 252 --- [ctor-http-nio-3] reactor.Mono.SubscribeOn.1 : onSubscribe(MonoSubscribeOn.SubscribeOnSubscriber) 2022-06-21 10:54:37.875 INFO 252 --- [ctor-http-nio-3] reactor.Mono.SubscribeOn.1 : request(unbounded) 2022-06-21 10:54:38.724 INFO 252 --- [ctor-http-nio-1] reactor.Mono.SubscribeOn.1 : onNext(OperationFullDto(id=c4f62ea6-0387-4021-b5de-d056105612ce, titre=test3, ouvrageIds=[2])) the void getOuvrages was invoked [2] 2022-06-21 10:54:38.796 INFO 252 --- [ctor-http-nio-1] reactor.Mono.SubscribeOn.1 : onComplete()
问题原因
初始写法存在三个核心错误:
- 响应式流未被订阅:
tuple2.map(...)只是定义了流的处理逻辑,但没有把这个流加入到整个响应式链路上,也没有被订阅,因此OuvrageClient的实际HTTP调用根本不会执行。你看到的getOuvrages was invoked打印只是方法同步执行到了创建Mono的步骤,Reactor中的Mono/Flux只有被订阅才会触发实际的异步请求逻辑。 - 提前返回空对象:在flatMap里直接
return Mono.just(detail),此时detail还没有被任何赋值逻辑处理,就直接返回了空的detail对象,因此最终接口返回的字段全是null。 - 重复调用接口:在flatMap里重新zip了外层的
operationMono,会导致operationClient.getOperation被重复调用第二次,产生不必要的网络开销。
写法说明
你最后给出的zipWhen写法是完全合理的,是这个场景下的标准实现方式:zipWhen会先拿到第一个Mono的结果,用这个结果生成第二个依赖Mono,再把两个结果合并成Tuple,天然适配「第一个接口返回值作为第二个接口入参」的顺序依赖场景。
修正后的完整代码如下:
public Mono<OperationDetailDto> getDetails(String operationId){ return operationClient.getOperation(operationId) .zipWhen(operation -> ouvrageClient.getOuvrages(operation.getOuvrageIds())) .map(tuple2 -> { OperationDetailDto detail = new OperationDetailDto(); detail.setId(ServiceUtils.formatUuid(operationId)); detail.setTitre(tuple2.getT1().getTitre()); detail.setOuvrages(tuple2.getT2()); return detail; }); }
这个写法没有多余调用,流是连续可订阅的,会按顺序先查询operation信息,拿到ouvrageIds后再查询ouvrage列表,最后合并成需要的DTO返回。
如果不用zipWhen,用flatMap嵌套的写法也可以实现相同效果,核心是要把内部的Mono作为返回值串联到链路上,不能中断流:
public Mono<OperationDetailDto> getDetails(String operationId){ return operationClient.getOperation(operationId) .flatMap(operation -> ouvrageClient.getOuvrages(operation.getOuvrageIds()) .map(ouvrageList -> { OperationDetailDto detail = new OperationDetailDto(); detail.setId(ServiceUtils.formatUuid(operationId)); detail.setTitre(operation.getTitre()); detail.setOuvrages(ouvrageList); return detail; }) ); }
内容的提问来源于stack exchange,提问作者Vince62
相关产品推荐
相关产品推荐

