You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.28 22:12:18