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

Spring Webflux应用WebClient重复调用外部服务问题排查

问题原因分析

1. 外部服务被重复调用4次

Mono是冷发布者,每一次订阅都会重新触发整个流的执行。你的代码里存在4次独立订阅:

  • ClientAPI中手动调用result.subscribe()
  • Service中手动调用response.subscribe()
  • Controller中手动调用response.subscribe()
  • Spring WebFlux框架会自动订阅Controller返回的Mono以处理HTTP响应

这四次订阅直接导致WebClient的POST请求被执行了4次,对应外部服务的4次调用记录。

2. 400错误与成功响应并存的矛盾

WebClient的retrieve()方法默认会在收到4xx/5xx状态码时抛出WebClientResponseException,但你的外部服务日志显示请求成功、数据入库,同时ClientAPI打印了成功的响应日志,核心原因是:

  • 手动订阅的流与框架订阅的流相互独立,手动订阅的流捕获到了成功响应,而框架订阅的流捕获到了错误响应(可能是某一次请求因重试或网络问题返回400),导致日志出现矛盾。
  • 未自定义处理WebClient的错误响应,默认抛出的异常被Spring框架捕获后返回500错误。
修复方案

步骤1:移除所有手动subscribe()调用,改用副作用操作符

业务代码中绝对不要手动调用subscribe(),应使用doOnSuccess、doOnError等操作符处理日志、后续操作等副作用,让Spring框架统一管理订阅生命周期。

修改后的Controller.java

@Slf4j
@RestController
@RequestMapping("/v1")
@CrossOrigin(origins = "http://localhost:4200")
public class Controller {
    private static final String PATH_SERVICE = "/pathservice";

    @Autowired
    Service service;

    @PostMapping(path = PATH_SERVICE, consumes = MediaType.APPLICATION_JSON_VALUE, produces = MediaType.APPLICATION_JSON_VALUE)
    public Mono<TheResponseClass> mySpringbootService(@RequestHeader("requestId") String requestId,
            @RequestBody @Valid InputData inputData) {
        return service.save(requestId, inputData)
                .doOnSuccess(status -> log.debug("Controller :  response status = " + status.getResponse().getStatus()))
                .doOnError(error -> log.error("Controller : The following error happened!", error));
    }
}

修改后的Service.java

@Slf4j
@Service
public class Service {

    @Autowired
    ClientAPI clientAPI;

    @Autowired
    BuilRequestBody builRequestBody;

    @Autowired
    BuildHeaderRequest buildHeaderRequest;
    
    @Autowired
    DoneService doneService;

    @Autowired
    Config config;

    public Mono<TheResponseClass> save(String requestId, @Valid InputData inputData) {
        if (inputData.getClassInfo() != null) {
            return clientAPI.sendRequest(
                    builRequestBody.buildRequest(inputData),
                    buildHeaderRequest.computeHeader(requestId, config.getHeaderID()))
                    .doOnSuccess(responseResult -> {
                        Status resultStatus = responseResult.getResponse().getStatus();
                        log.debug("Service : response status = " + resultStatus);
                        if (resultStatus == Status.SUCCESS) {
                            // 若无需等待通知服务完成,直接subscribe;若需等待则改用flatMap
                            doneService.notify().subscribe();
                        }
                    })
                    .doOnError(error -> log.error("Service : error happened !", error));
        } else {
            TheResponseClass errorResponse = new TheResponseClass();
            errorResponse.setName("name");
            errorResponse.setOperation("operation");
            return Mono.just(errorResponse);
        }
    }
}

修改后的ClientAPI.java

@Slf4j
@Service
public class ClientAPI {
    private WebClient webClient;
    private String serverUrl;

    @Autowired
    public ClientAPI(@Qualifier("CustomWebClient") WebClient webClient) {
        this.serverUrl = "https://localhost:8456/myServer";
        this.webClient = webClient;
    }

    public Mono<TheResponseClass> sendRequest(Request bodyRequest, HeaderParameters headerParameters) {
        return this.webClient.post().uri(this.serverUrl)
                .header(HttpHeaders.CONTENT_TYPE, MediaType.APPLICATION_JSON_VALUE)
                .header(HeaderParameters.REQUEST_ID, headerParameters.getRequestID())
                .body(Mono.just(bodyRequest), Request.class)
                .retrieve()
                // 自定义处理4xx/5xx响应,避免直接抛出异常
                .onStatus(HttpStatus::isError, clientResponse -> 
                    clientResponse.bodyToMono(TheResponseClass.class)
                            .switchIfEmpty(Mono.error(new WebClientResponseException(
                                    clientResponse.statusCode().value(),
                                    clientResponse.statusCode().getReasonPhrase(),
                                    clientResponse.headers().asHttpHeaders(),
                                    null,
                                    null,
                                    null)))
                )
                .bodyToMono(TheResponseClass.class)
                .doOnSuccess(responseResult -> {
                    ObjectMapper mapper = new ObjectMapper();
                    try {
                        String jsonOutput = mapper.writeValueAsString(responseResult);
                        log.info("ClientAPI : response result = "+jsonOutput);
                    } catch (JsonProcessingException e) {
                        log.error("Failed to serialize response", e);
                    }
                })
                .doOnError(error -> log.error("ClientAPI : The following error happened on response getStatus!", error));
    }
}

步骤2:按需调整DoneService的调用逻辑

如果需要等待doneService.notify()完成后再返回响应,应使用flatMap替代doOnSuccess+subscribe(),确保异步操作被正确等待:

// Service.java中替换doOnSuccess为flatMap
.flatMap(responseResult -> {
    Status resultStatus = responseResult.getResponse().getStatus();
    log.debug("Service : response status = " + resultStatus);
    if (resultStatus == Status.SUCCESS) {
        return doneService.notify().thenReturn(responseResult);
    }
    return Mono.just(responseResult);
})
验证修复效果
  1. 外部服务日志只会收到1次调用
  2. 本地日志不再出现矛盾的400错误与成功响应
  3. Controller返回的响应与外部服务实际状态一致

内容的提问来源于stack exchange,提问作者cknelle

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 01:37:03