WebClient结合Mono zip并行调用时如何定位失败请求并获取请求体
WebClient并行API调用故障定位与请求体追加方案
精准识别故障请求的实现方式
Mono.zip并行调用时默认抛出的异常不携带请求来源信息,可通过两种方式实现精准定位:
- 为每个独立请求打唯一标识:组装WebClient请求时,通过
attribute()方法给每个请求注入自定义标识,比如接口名、业务场景标记。同时给每个单请求的Mono增加onErrorMap逻辑,异常触发时第一时间将请求标识写入自定义异常的字段中,经Mono.zip合并抛出后,直接读取异常字段即可定位故障源。
示例代码:
// 3个并行请求分别注入标识、绑定专属异常映射 Mono<UserInfoResp> userReq = webClient.get() .uri("/api/user/{uid}", uid) .attribute("apiMark", "用户基础信息查询接口") .retrieve() .bodyToMono(UserInfoResp.class) .onErrorMap(err -> new ApiInvokeException("用户基础信息查询接口调用失败", err)); Mono<OrderInfoResp> orderReq = webClient.get() .uri("/api/order/{orderId}", orderId) .attribute("apiMark", "订单详情查询接口") .retrieve() .bodyToMono(OrderInfoResp.class) .onErrorMap(err -> new ApiInvokeException("订单详情查询接口调用失败", err)); Mono<StockInfoResp> stockReq = webClient.get() .uri("/api/stock/{skuId}", skuId) .attribute("apiMark", "库存信息查询接口") .retrieve() .bodyToMono(StockInfoResp.class) .onErrorMap(err -> new ApiInvokeException("库存信息查询接口调用失败", err)); // 执行并行调用 Mono.zip(userReq, orderReq, stockReq) .map(tuple -> assembleBizData(tuple.getT1(), tuple.getT2(), tuple.getT3()));
- 异常处理器中读取请求属性:在已配置的
ExchangeFilterFunction异常处理逻辑中,直接从ClientRequest对象的attributes集合里读取预先写入的apiMark标识,不需要额外修改zip的逻辑即可识别故障请求。
ExchangeFilterFunction中获取失败请求体的方案
默认无法直接从过滤器中读取原始请求体:WebClient的请求体是订阅时才序列化输出的一次性响应式流,过滤器执行时请求体已经开始写出,没有默认缓存,直接读取流会导致流被消费,请求发送失败。
可通过两种方案实现请求体追加:
- 低侵入推荐方案:组装请求时,直接将待发送的请求体对象存入request attribute,比如POST请求传入
UserCreateReq对象时,追加.attribute("reqPayload", reqObj),过滤器异常处理逻辑中直接从attributes读取该对象即可,不需要处理流消费问题,性能损耗极低,是生产环境首选方案。 - 统一过滤器拦截方案:如果需要全局统一处理、不想每个请求手动传参,可自定义请求装饰器,在请求体写出前缓存内容副本存入请求属性,注意要增加请求类型白名单和大小限制,避免大文件上传、超大请求导致内存溢出。
示例过滤器代码:
public ExchangeFilterFunction cachedRequestBodyFilter() { return ExchangeFilterFunction.ofRequestProcessor(req -> { MediaType contentType = req.headers().getContentType(); // 仅缓存JSON类型、小于10KB的请求,跳过multipart等大流量请求 if (contentType != null && contentType.includes(MediaType.APPLICATION_JSON)) { ClientHttpRequestDecorator reqDecorator = new ClientHttpRequestDecorator(req) { @Override public Mono<Void> writeWith(Publisher<? extends DataBuffer> body) { return DataBufferUtils.join(body, 10240) // 限制最大缓存10KB .doOnNext(buffer -> { byte[] content = new byte[buffer.readableByteCount()]; buffer.read(content); DataBufferUtils.release(buffer); req.attributes().put("cachedReqBody", new String(content, StandardCharsets.UTF_8)); // 重新包装缓存的内容写出,避免流被消费 super.writeWith(Mono.just(req.strategy().bufferFactory().wrap(content))).subscribe(); }).then(); } }; return req.body().insert(new BodyInserterContext(), reqDecorator).thenReturn(req); } return Mono.just(req); }); }
在原有通用异常处理器中,直接从请求attributes读取请求标识、缓存的请求体,拼接到异常响应结构中即可。
内容的提问来源于stack exchange,提问作者Abhishek Wadhwa
相关产品推荐
相关产品推荐

