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

Apache Http Async客户端超时与连接关闭异常处理咨询

背景

在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 02:26:06