Apache Http Async客户端超时与连接关闭异常处理咨询
高并发场景下Flink RichSink+HttpAsyncClient的配置优化问题
背景
在AWS托管的Apache Flink项目中,我实现了一个RichSinkFunction,用于通过REST POST请求向HTTP服务发送数据。为提升性能,选用了Apache Http Async Client(版本:org.apache.httpcomponents.client5.httpclient5 v5.3.1)。
初始客户端配置
在RichSinkFunction.open()方法中创建客户端,仅配置了30秒的保活策略,其余保持默认:
public void open(Configuration parameters) { CloseableHttpAsyncClient client = HttpAsyncClients.custom() .setKeepAliveStrategy((httpResponse, httpContext) -> TimeValue.of(30, TimeUnit.SECONDS)) .build(); client.start(); }
请求发送逻辑
RichSinkFunction.invoke()方法中构造并发送POST请求的代码如下:
SimpleHttpRequest simpleHttpRequest = SimpleRequestBuilder.post(endPoint).setHeader(HttpHeaders.AUTHORIZATION, authHeader).build(); simpleHttpRequest.setBody(JSON_content, ContentType.APPLICATION_JSON); simpleHttpRequest.setHeader("Accept", "application/json"); simpleHttpRequest.setHeader("Content-type", "application/json"); httpClient.execute(simpleHttpRequest, new FutureCallback<SimpleHttpResponse>() { // 回调方法已实现 });
初始方案的异常问题
当请求量极大(阈值未明确)且REST端点响应缓慢时,出现大量DeadlineTimeoutException:
org.apache.hc.core5.util.DeadlineTimeoutException: Deadline: 2024-03-09T12:40:32.438+0000, -199656424 MILLISECONDS overdue at org.apache.hc.core5.util.DeadlineTimeoutException.from(DeadlineTimeoutException.java:49) at org.apache.hc.core5.pool.StrictConnPool.processPendingRequest(StrictConnPool.java:323) at org.apache.hc.core5.pool.StrictConnPool.processNextPendingRequest(StrictConnPool.java:304) at org.apache.hc.core5.pool.StrictConnPool.release(StrictConnPool.java:266) at org.apache.hc.client5.http.impl.nio.PoolingAsyncClientConnectionManager.release(PoolingAsyncClientConnectionManager.java:410) at org.apache.hc.client5.http.impl.async.InternalHttpAsyncExecRuntime.discardEndpoint(InternalHttpAsyncExecRuntime.java:147) at org.apache.hc.client5.http.impl.async.InternalHttpAsyncExecRuntime.discardEndpoint(InternalHttpAsyncExecRuntime.java:170) at org.apache.hc.client5.http.impl.async.InternalAbstractHttpAsyncClient$2.failed(InternalAbstractHttpAsyncClient.java:346) at org.apache.hc.client5.http.impl.async.AsyncRedirectExec$1.failed(AsyncRedirectExec.java:248) at org.apache.hc.client5.http.impl.async.AsyncHttpRequestRetryExec$1.failed(AsyncHttpRequestRetryExec.java:197) at org.apache.hc.client5.http.impl.async.AsyncProtocolExec$1.failed(AsyncProtocolExec.java:295) at org.apache.hc.client5.http.impl.async.HttpAsyncMainClientExec$1.failed(HttpAsyncMainClientExec.java:131) at org.apache.hc.core5.http2.impl.nio.ClientH2StreamHandler.failed(ClientH2StreamHandler.java:253) at org.apache.hc.core5.http2.impl.nio.AbstractH2StreamMultiplexer$H2Stream.reset(AbstractH2StreamMultiplexer.java:1668) at org.apache.hc.core5.http2.impl.nio.AbstractH2StreamMultiplexer$H2Stream.cancel(AbstractH2StreamMultiplexer.java:1693) at org.apache.hc.core5.http2.impl.nio.AbstractH2StreamMultiplexer.onDisconnect(AbstractH2StreamMultiplexer.java:575) at org.apache.hc.core5.http2.impl.nio.AbstractH2IOEventHandler.disconnected(AbstractH2IOEventHandler.java:96) at org.apache.hc.core5.http2.impl.nio.ClientH2IOEventHandler.disconnected(ClientH2IOEventHandler.java:39) at org.apache.hc.core5.reactor.ssl.SSLIOSession$1.disconnected(SSLIOSession.java:247) at org.apache.hc.core5.reactor.InternalDataChannel.disconnected(InternalDataChannel.java:204) at org.apache.hc.core5.reactor.SingleCoreIOReactor.processClosedSessions(SingleCoreIOReactor.java:231) at org.apache.hc.core5.reactor.SingleCoreIOReactor.doTerminate(SingleCoreIOReactor.java:106) at org.apache.hc.core5.reactor.AbstractSingleCoreIOReactor.execute(AbstractSingleCoreIOReactor.java:93) at org.apache.hc.core5.reactor.IOReactorWorker.run(IOReactorWorker.java:44) at java.base/java.lang.Thread.run(Thread.java:829)
优化尝试(方案2)
尝试简化客户端配置,移除保活策略:
CloseableHttpAsyncClient client = HttpAsyncClients.custom().build(); client.start();
错误数量有所减少,但仍出现ConnectionClosedException:
org.apache.hc.core5.http.ConnectionClosedException: Connection is closed at org.apache.hc.core5.http2.impl.nio.AbstractH2StreamMultiplexer.onDisconnect(AbstractH2StreamMultiplexer.java:581) at org.apache.hc.core5.http2.impl.nio.AbstractH2IOEventHandler.disconnected(AbstractH2IOEventHandler.java:96) at org.apache.hc.core5.http2.impl.nio.ClientH2IOEventHandler.disconnected(ClientH2IOEventHandler.java:39) at org.apache.hc.core5.reactor.ssl.SSLIOSession$1.disconnected(SSLIOSession.java:247) at org.apache.hc.core5.reactor.InternalDataChannel.disconnected(InternalDataChannel.java:204) at org.apache.hc.core5.reactor.SingleCoreIOReactor.processClosedSessions(SingleCoreIOReactor.java:231) at org.apache.hc.core5.reactor.SingleCoreIOReactor.doExecute(SingleCoreIOReactor.java:133) at org.apache.hc.core5.reactor.AbstractSingleCoreIOReactor.execute(AbstractSingleCoreIOReactor.java:86) at org.apache.hc.core5.reactor.IOReactorWorker.run(IOReactorWorker.java:44) at java.base/java.lang.Thread.run(Thread.java:829)
核心需求
由于该Sink的HTTP更新至关重要,需要配置客户端实现:
- 请求池化
- 确保请求完成不丢失
- 支持约10万条并行消息的高吞吐量
内容的提问来源于stack exchange,提问作者Shankar
相关产品推荐
相关产品推荐

