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使用
map操作处理响应,map是同步转换,返回值为普通对象,不能返回Mono类型。当状态码非OK或解析JSON异常时,代码调用了Mono.error()但没有return,这个错误信号根本不会被订阅,异常被静默丢弃。 - 错误信号丢失后,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
相关产品推荐
相关产品推荐

