ApacheHttpClient连接池并发
基于 HttpClient 4.5 的线程池与连接池参数设置指南
HttpClient(通常使用其实现类 CloseableHttpClient)与 PoolingHttpClientConnectionManager 是协作关系,前者负责执行 HTTP 请求,后者负责管理底层的 HTTP 连接池。
连接池管理器(PoolingHttpClientConnectionManager)
连接池管理 HTTP 连接的复用,避免频繁创建 / 销毁连接的开销,线程安全,可在多线程环境中使用。显著提高多线程环境下的 HTTP 客户端性能,特别是需要频繁发送请求到相同服务器的场景。
类UML图展示了PoolingHttpClientConnectionManager类的主要结构和关系:
1、PoolingHttpClientConnectionManager实现了三个接口:HttpClientConnectionManager提供http链接管理的基本功能;ConnPoolControl提供链接池控制功能;Closeable提供资源关闭功能
2、PoolingHttpClientConnectionManager包含两个内部类:InternalConnectionFactory负责创建HTTP连接;ConfigData存储配置数据,Socket配置和连接配置
3、属性:log日志;configData配置;pool具体连接池;connectionOperator连接操作器;isShutdown连接管理器是否关闭。
+---------------------------+ +---------------------------+
| <<interface>> | | <<interface>> |
| HttpClientConnectionManager| | ConnPoolControl<HttpRoute>|
+---------------------------+ +---------------------------+
| + requestConnection() | | + getMaxTotal() |
| + connect() | | + setMaxTotal() |
| + upgrade() | | + getDefaultMaxPerRoute() |
| + routeComplete() | | + setDefaultMaxPerRoute() |
| + releaseConnection() | | + getMaxPerRoute() |
| + closeIdleConnections() | | + setMaxPerRoute() |
| + closeExpiredConnections()| | + getTotalStats() |
+---------------------------+ | + getStats() |
^ +---------------------------+
| ^
| |
| |
+---------------------------+ +---------------------------+
| <<interface>> | | |
| Closeable |<----------------| PoolingHttpClientConnMgr |
+---------------------------+ +---------------------------+
| + close() | | - log: Log |
+---------------------------+ | - configData: ConfigData |
| - pool: CPool |
| - connectionOperator |
| - isShutDown: AtomicBoolean|
+---------------------------+
| + PoolingHttpClientConnMgr()|
| + finalize() |
| + close() |
| + requestConnection() |
| + leaseConnection() |
| + releaseConnection() |
| + connect() |
| + upgrade() |
| + routeComplete() |
| + shutdown() |
| + closeIdleConnections() |
| + closeExpiredConnections()|
| + getMaxTotal() |
| + setMaxTotal() |
| + getDefaultMaxPerRoute() |
| + setDefaultMaxPerRoute() |
| + getMaxPerRoute() |
| + setMaxPerRoute() |
| + getTotalStats() |
| + getStats() |
| + getRoutes() |
| + getDefaultSocketConfig()|
| + setDefaultSocketConfig()|
| + getDefaultConnectionConfig()|
| + setDefaultConnectionConfig()|
| + getSocketConfig() |
| + setSocketConfig() |
| + getConnectionConfig() |
| + setConnectionConfig() |
| + getValidateAfterInactivity()|
| + setValidateAfterInactivity()|
+---------------------------+
|
|
|
+------------------+------------------+
| |
| |
+---------------------------+ +---------------------------+
| InternalConnectionFactory | | ConfigData |
+---------------------------+ +---------------------------+
| - configData: ConfigData | | - socketConfigMap |
| - connFactory | | - connectionConfigMap |
+---------------------------+ | - defaultSocketConfig |
| + create() | | - defaultConnectionConfig|
+---------------------------+ +---------------------------+
| + getDefaultSocketConfig()|
| + setDefaultSocketConfig()|
| + getDefaultConnectionConfig()|
| + setDefaultConnectionConfig()|
| + getSocketConfig() |
| + setSocketConfig() |
| + getConnectionConfig() |
| + setConnectionConfig() |
+---------------------------+
连接池管理器进行HTTP请求操作流程
- 初始化连接池管理器:创建PoolingHttpClientConnectionManager实例,配置连接池参数(
setMaxTotal(int max)最大连接数、setDefaultMaxPerRoute(int max)每路由(域名)最大连接数等)(1)静态方法getDefaultRegistry()方法,其创建一个连接套接字工厂的注册表(Registry<ConnectionSocketFactory>)【Map】,将<http,PlainConnectionSocketFactory>,<https,SSLConnectionSocketFactory>分别注册用来处理普通http连接和https加密连接;()(2)DefaultHttpClientConnectionOperator连接操作器,负责创建、建立和升级HTTP连接,其构造函数有三个参数ConnectionSocketFactory注册表,SchemePortResolver解析协议和端口、DnsResolver(用于DNS解析),核心方法connect:负责建立到目标主机的物理连接根据协议选择合适的ConnectionSocketFactory,使用DnsResolver解析主机,创建并配置socket,建立TCP连接,设置超时参数,将Socket与ManagedHttpClientConnection关联。 - 创建HttpClient:使用HttpClients.custom()创建HttpClientBuilder构建器,设置连接管理器、重试处理器和请求配置,构建CloseableHttpClient实例
- 请求连接:HttpClient(CloseableHttpClient实例)调用连接管理器的requestConnection方法,连接管理器从连接池中租借连接或创建新连接,返回ConnectionRequest对象
- 获取连接:调用ConnectionRequest的get方法获取HttpClientConnection,如果连接池中有可用连接,直接返回,如果没有可用连接且未达到最大连接数,创建新连接,如果达到最大连接数,等待连接释放或超时。
- 建立连接:调用连接管理器的connect方法,连接管理器使用connectionOperator建立物理连接,设置连接超时、Socket超时等参数
- 执行HTTP请求:创建HttpRequestBase对象(如HttpGet、HttpPost等),设置请求头、请求体等,使用HttpClient执行请求,获取HttpResponse响应
- 处理响应:读取响应状态码、响应头和响应体,处理响应数据,关闭响应资源
- 释放连接:调用连接管理器的releaseConnection方法,设置连接保持活动的时间,连接管理器将连接放回连接池或关闭连接
- 连接池维护:定期调用closeExpiredConnections()关闭过期连接,调用closeIdleConnections()关闭空闲连接,监控连接池状态
- 关闭资源:应用程序退出时调用连接管理器的close()或shutdown()方法,关闭所有连接并释放资源
HttpRequestInterceptor
HttpRequestInterceptor 是一个接口,用于在 HTTP 请求发送到服务器之前拦截并处理请求。它允许开发者在请求发出前统一修改请求信息(如添加头信息、设置认证、记录日志等),是实现请求层面横切逻辑的重要机制。
线程池参数(ThreadPoolExecutor)
线程池负责调度并发任务,与连接池配合使用,核心参数:
-
核心线程数(
corePoolSize)
保持存活的最小线程数,默认建议设为CPU核心数 * 2。 -
最大线程数(
maximumPoolSize)
线程池允许的最大线程数,通常不超过连接池的setMaxTotal值(避免线程等待连接)。 -
空闲线程存活时间(
keepAliveTime)
超过核心线程数的空闲线程的存活时间,建议设为60秒左右。 -
任务队列(
workQueue)
用于存放等待执行的任务,建议使用LinkedBlockingQueue并指定容量(如200),避免任务无限制堆积导致 OOM。 -
拒绝策略(
RejectedExecutionHandler)
任务队列满时的处理策略,推荐CallerRunsPolicy(让提交任务的线程执行任务,避免任务丢失)。
参数调优建议
-
连接池与线程池的匹配
- 线程池最大线程数 ≤ 连接池最大总连接数(避免线程等待连接)。
- 单个路由的最大连接数 ≥ 该路由的并发线程数(避免同一域名的请求排队)。
-
根据业务场景调整
- 高并发短请求(如 API 调用):可适当增大
setMaxTotal和线程池最大线程数(如50-100)。 - 低并发长请求(如下载文件):减少连接数和线程数,避免资源浪费。
- 高并发短请求(如 API 调用):可适当增大
-
监控与动态调整
- 通过
cm.getTotalStats()监控连接池状态(已使用 / 空闲连接数): - 根据监控数据动态调整参数(如通过定时任务调整
setMaxTotal)。
- 通过
-
避免连接泄漏
- 确保
CloseableHttpResponse和HttpEntity被正确关闭(使用 try-with-resources)。 - 为连接池设置连接存活时间(
setValidateAfterInactivity),自动清理无效连接:
- 确保
package org.example;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.CompletableFuture;
@Component //Spring框架下使用Component 和 PreDestroy注解实现销毁时资源自动释放
@Slf4j
public class HttpClientUtils{
//链接池
private static final PoolingHttpClientConnectionManager connectionManager;
// HttpClient 实例(线程安全)
private static final CloseableHttpClient httpClient;
// 默认请求配置
private static final RequestConfig defaultRequestConfig;
// 线程池,用于处理并发请求
private static final ExecutorService executorService;
static {
connectionManager = new PoolingHttpClientConnectionManager();
connectionManager.setMaxTotal(20); //最大连接数,默认20
connectionManager.setDefaultMaxPerRoute(10);//同域名最大连接数默认2,防止单域名耗尽
defaultRequestConfig = RequestConfig.custom()
.setConnectTimeout(3000) // 连接超时时间(毫秒)
.setConnectionRequestTimeout(3000) // 从连接池获取连接的超时时间
.setSocketTimeout(6000) // 数据传输超时时间
.build();
// 初始化 HttpClient
httpClient = HttpClients.custom()
.setConnectionManager(connectionManager)
.setDefaultRequestConfig(defaultRequestConfig)
.build();
}
static class HttpRequest{
private String url;
private Map<String, String> params;
private Map<String, String> headers;
private String method;
public HttpRequest(String url, Map<String, String> params, Map<String, String> headers, String method) {
if (url == null) {
throw new RuntimeException("url is null");
}
this.url = url;
this.params = params;
this.headers = headers;
this.method = method;
}
}
//简写
public String doGet(String url){
URI uri = new URIBuilder(url).build();//可以添加参数
HttpGet httpGet = new HttpGet(uri); //可以设置请求头
// 执行GET请求,资源自动关闭
try (CloseableHttpResponse response = httpClient.execute(httpGet)) {
return EntityUtils.toString(response.getEntity(),
StandardCharsets.UTF_8);
}
}
//可以结合线程池 实现并发请求方法
public static Map<String, String> multiHttp(List<HttpRequest> httpRequests) {
Map<String, String> resultMap = new ConcurrentHashMap<>();
if (httpRequests != null && !httpRequests.isEmpty()) {
List<CompletableFuture<Map<String,String>>> futures = new ArrayList<>();
httpRequests.forEach((httpRequest) -> {
if (httpRequest != null && httpRequest.getUrl() != null) {
final String url = httpRequest.getUrl();
if (resultMap.containsKey(url)) {
return;
}
resultMap.put(url, "");
CompletableFuture<Map<String,String>> future = CompletableFuture.supplyAsync(new Supplier<Map<String,String>>() {
public Map<String,String> get() {
Map<String,String> ret = new HashMap<>();
ret.put(url,doGet(url));
return ret;
}
},executorService); //有返回值异步任务,并自定义线程池
futures.add(future);
}
}
if(!futures.isEmpty()){
List<CompletableFuture<Void>> processingFutures = new ArrayList<>();
futures.forEach(f->{
CompletableFuture<Void> processingFuture = f.exceptionally(ex -> {
log.error("multiHttp {}", ex.getMessage());
return null;
}).thenAccept(resultMap::putAll);
processingFutures.add(processingFuture);
});
//List toArray() 转成指定类型数组
//get和join阻塞等待任务完成,get会抛出受检异常InterruptedException和ExecutionException(包原始异常)必须try-catch显示处理,
//join将原线程中异常包装在非受检异常CompletionException中抛出,通过 getCause() 获取
//join可以在lambda 表达式或流式操作中使用(如 Stream.map(CompletableFuture::join))
CompletableFuture.allOf(processingFutures.toArray(new CompletableFuture[0])).join();
}
}
return resultMap;
}
@PreDestroy
public void shutdown() {
try {
httpClient.close();
} catch (IOException e) {
log.error("httpClient Error {}", e.getMessage());
}
executorService.shutdown();
try {
if (!executorService.awaitTermination(5, TimeUnit.SECONDS)) {
executorService.shutdownNow();
}
} catch (InterruptedException e) {
executorService.shutdownNow();
}
log.info("httpClientUtil shutdown");
}
}
常见问题与解决方案
| 问题场景 | 可能原因 | 解决方案 |
|---|---|---|
| 请求排队严重 | 连接池总连接数不足 | 增大 setMaxTotal,检查是否有连接泄漏 |
| 某一域名请求阻塞 | 该路由连接数不足 | 使用 setRouteMaxPerRoute 单独配置 |
| 线程池任务堆积 | 线程数不足或队列满 | 增大最大线程数或队列容量,优化拒绝策略 |
| 连接超时频繁 | 目标服务响应慢或网络差 | 调大 setConnectTimeout,增加重试机制 |
更多推荐
所有评论(0)