[ 
https://issues.apache.org/jira/browse/FLINK-40381?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

ASF GitHub Bot updated FLINK-40381:
-----------------------------------
    Labels: pull-request-available  (was: )

> Fix thread leak in Elasticsearch8AsyncWriter by closing HTTP transport on 
> sink close
> ------------------------------------------------------------------------------------
>
>                 Key: FLINK-40381
>                 URL: https://issues.apache.org/jira/browse/FLINK-40381
>             Project: Flink
>          Issue Type: Bug
>          Components: Connectors / ElasticSearch
>    Affects Versions: elasticsearch-3.0.0, elasticsearch-3.1.0, 
> elasticsearch-3.0.1, elasticsearch-4.0.0
>            Reporter: nadav fux
>            Priority: Major
>              Labels: pull-request-available
>
> In `flink-connector-elasticsearch8`, running jobs on a persistent Flink 
> Session Cluster causes a thread leak of `I/O dispatcher` threads per job run. 
> When jobs are submitted and completed repeatedly, zombie reactor worker 
> threads accumulate in the TaskManager JVM until the TaskManager process 
> crashes or exhausts system thread limits.
>  
> *Root Cause*
> In `Elasticsearch8AsyncWriter.java`, the `close()` method contains:
> public void close() {
> if (!this.close) {
> this.close = true;
> this.esClient.shutdown();
> }
> }
> In the official Elastic Java API Client 
> ({{{}co.elastic.clients.elasticsearch.ElasticsearchAsyncClient{}}}), 
> {{esClient.shutdown()}} returns an {{ElasticsearchShutdownAsyncClient}} 
> builder for making requests to Elasticsearch Cluster Shutdown REST API 
> endpoints. It does *NOT* close the underlying {{RestClient}} or HTTP Async 
> Client IO reactor thread pool.
> h3. Solution
> Update {{Elasticsearch8AsyncWriter.close()}} to invoke 
> {{{}esClient._transport().close(){}}}, which closes the underlying 
> {{RestClient}} and terminates the {{I/O dispatcher}} worker thread pool upon 
> sink teardown.
>  
>  
>  



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to