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

Spring WebClient连接池耗尽问题排查求助(Spring Boot 2.7.3)

WebClient连接池耗尽问题排查

问题描述

使用Spring Boot 2.7.3版本,通过WebClient发送POST请求,运行一段时间后出现连接池耗尽,报错如下:

2022-09-22 15:07:52,701 ERROR ? [parallel-1] Operator called default onErrorDropped
reactor.core.Exceptions$ErrorCallbackNotImplemented: org.springframework.web.reactive.function.client.WebClientRequestException: Pool#acquire(Duration) has been pending for more than the configured timeout of 30000ms; nested exception is reactor.netty.internal.shaded.reactor.pool.PoolAcquireTimeoutException: Pool#acquire(Duration) has been pending for more than the configured timeout of 30000ms
Caused by: org.springframework.web.reactive.function.client.WebClientRequestException: Pool#acquire(Duration) has been pending for more than the configured timeout of 30000ms; nested exception is reactor.netty.internal.shaded.reactor.pool.PoolAcquireTimeoutException: Pool#acquire(Duration) has been pending for more than the configured timeout of 30000ms
        at org.springframework.web.reactive.function.client.ExchangeFunctions$DefaultExchangeFunction.lambda$wrapException$9(ExchangeFunctions.java:141)
        Suppressed: reactor.core.publisher.FluxOnAssembly$OnAssemblyException: 
Error has been observed at the following site(s):
        *__checkpoint ⇢ Request to POST http://api.call/track [DefaultWebClient]
Original Stack Trace:
                at org.springframework.web.reactive.function.client.ExchangeFunctions$DefaultExchangeFunction.lambda$wrapException$9(ExchangeFunctions.java:141)
                at reactor.core.publisher.MonoErrorSupplied.subscribe(MonoErrorSupplied.java:55)
                at reactor.core.publisher.Mono.subscribe(Mono.java:4397)
                ...(省略后续栈帧)
Caused by: reactor.netty.internal.shaded.reactor.pool.PoolAcquireTimeoutException: Pool#acquire(Duration) has been pending for more than the configured timeout of 30000ms
        at reactor.netty.internal.shaded.reactor.pool.AbstractPool$Borrower.run(AbstractPool.java:413)
        at reactor.core.scheduler.SchedulerTask.call(SchedulerTask.java:68)

WebClient初始化代码

@PostConstruct
public void getWebClient() {
    logger.debug("Initializing webClient..");

    ConnectionProvider provider = ConnectionProvider.builder("fixed")
            .metrics(true)
            .maxConnections(50)
            .maxIdleTime(Duration.ofSeconds(20))
            .maxLifeTime(Duration.ofSeconds(60))
            .pendingAcquireTimeout(Duration.ofSeconds(30))
            .evictInBackground(Duration.ofSeconds(30)).build();

    HttpClient httpClient = HttpClient.create(provider)
            .wiretap(this.getClass().getCanonicalName(), LogLevel.DEBUG, AdvancedByteBufFormat.TEXTUAL)
            .option(ChannelOption.CONNECT_TIMEOUT_MILLIS, 10000)
            .responseTimeout(Duration.ofSeconds(10));

    this.webClient = WebClient
            .builder()
            .clientConnector(new ReactorClientHttpConnector(httpClient))
            .defaultHeader(HttpHeaders.CONTENT_TYPE, org.springframework.http.MediaType.APPLICATION_JSON_VALUE)
            .build();
}

请求发送代码

代码1:返回JsonNode的POST请求

public Mono<JsonNode> sendPostR(String url, JsonNode request) {
    String uuid = UUID.randomUUID().toString();
    Instant start = Instant.now();
    logger.debug("Sending post request {} {} {}",url,request,uuid);
    return webClient.post().uri(url)
            .body(Mono.just(request),Map.class)
            .retrieve()
            .toEntity(String.class)
            .map(entity -> {
                Instant end = Instant.now();
                logger.debug("Time taken for request success{} {} ",Duration.between(start,end).toMillis(),uuid);
                if(entity.getStatusCode() != HttpStatus.OK) {
                    logger.error("Error: API returned with in-valid status code {} {}",entity.getStatusCode(),entity.getHeaders());
                    Mono.error(new RuntimeException("Error in sending post " + url));
                }
                String responseBody = entity.getBody();
                try {
                    return objectMapper.readValue(responseBody, JsonNode.class);
                } catch (JsonProcessingException e) {
                    logger.error("Error in parsing response body",e);
                    Mono.error(new RuntimeException("Error in parsing response body " + url));
                }
                return Utils.getNewObjectNode();
            }).switchIfEmpty(Mono.defer(() -> {
                logger.debug("Request completed with no body response {}",uuid);
                return Mono.just(Utils.getNewObjectNode());
            })).log();
}

代码2:使用exchangeToMono的POST请求

public Mono<JsonNode> sendPostR(String url, Map<String,Object> params) {
    String uuid = UUID.randomUUID().toString();
    Instant start = Instant.now();
    logger.debug("Sending post request {} {} {}",url,params,uuid);
    return webClient.post().uri(url)
            .body(Mono.just(params),Map.class)
            .exchangeToMono(response -> {
                logger.debug("response is {}",response.statusCode());
                Instant end = Instant.now();
                logger.debug("Time taken for request success{} {} ",Duration.between(start,end).toMillis(),uuid);
                if(response.statusCode().equals(HttpStatus.OK)) {
                    return response.bodyToMono(JsonNode.class);
                }
                response.releaseBody();
                logger.debug("Error in hitting url {}",url);
                return Mono.error(new Exception("Error in sending post " + url));
            }).onErrorResume(WebClientResponseException.class, ex -> {
                Instant end = Instant.now();
                logger.debug("Time taken for request error {} {} ",Duration.between(start,end).toMillis(),uuid);
                logger.debug("Error in hitting url .... ",ex);
                return Mono.error(ex);
            }).doOnError(error -> {
                logger.error("Error in sending post request ",error);
            }).switchIfEmpty(Mono.defer(() -> {
                logger.debug("Request completed with no body response {}",uuid);
                return Mono.just(Utils.getNewObjectNode());
            }))
            .log();
}

代码3:发后即忘的POST请求

public Mono<ResponseEntity<Void>> raiseEvent(String url, JsonNode request) {
    String uuid = UUID.randomUUID().toString();
    logger.debug("Sending post request {} {} {}",url,request,uuid);
    return webClient.post().uri(url)
            .body(Mono.just(request),JsonNode.class)
            .retrieve()
            .toBodilessEntity() ;
}

问题原因分析

连接池耗尽的核心原因是代码1中错误信号被丢弃,导致WebClient无法正确释放连接:

  1. 代码1使用map操作处理响应,map是同步转换,返回值为普通对象,不能返回Mono类型。当状态码非OK或解析JSON异常时,代码调用了Mono.error()但没有return,这个错误信号根本不会被订阅,异常被静默丢弃。
  2. 错误信号丢失后,WebClient认为请求正常完成,但实际发生了错误,对应的连接没有被回收回连接池,随着请求累积,最终导致连接池耗尽。

修复方案

将代码1中的map替换为flatMap,因为flatMap支持返回Mono类型,能正确传递错误信号,触发连接释放:

public Mono<JsonNode> sendPostR(String url, JsonNode request) {
    String uuid = UUID.randomUUID().toString();
    Instant start = Instant.now();
    logger.debug("Sending post request {} {} {}", url, request, uuid);
    return webClient.post().uri(url)
            .body(Mono.just(request), Map.class)
            .retrieve()
            .toEntity(String.class)
            .flatMap(entity -> { // 替换map为flatMap
                Instant end = Instant.now();
                logger.debug("Time taken for request success{} {} ", Duration.between(start, end).toMillis(), uuid);
                if (entity.getStatusCode() != HttpStatus.OK) {
                    logger.error("Error: API returned with in-valid status code {} {}", entity.getStatusCode(), entity.getHeaders());
                    return Mono.error(new RuntimeException("Error in sending post " + url)); // 返回错误Mono
                }
                String responseBody = entity.getBody();
                try {
                    return Mono.just(objectMapper.readValue(responseBody, JsonNode.class)); // 用Mono.just包装正常结果
                } catch (JsonProcessingException e) {
                    logger.error("Error in parsing response body", e);
                    return Mono.error(new RuntimeException("Error in parsing response body " + url)); // 返回错误Mono
                }
            }).switchIfEmpty(Mono.defer(() -> {
                logger.debug("Request completed with no body response {}", uuid);
                return Mono.just(Utils.getNewObjectNode());
            })).log();
}

额外注意事项

  • 确保所有WebClient返回的Mono都被正确订阅:如果Mono未被订阅,请求不会执行,但如果订阅后错误信号被丢弃,仍会导致连接泄漏。
  • 利用已开启的连接池指标(metrics(true)),结合监控工具查看连接的获取、释放、等待情况,验证修复效果。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 14:15:41