This is an automated email from the ASF dual-hosted git repository.

wenjin272 pushed a commit to branch main
in repository 
https://gitbox.apache.org/repos/asf/flink-connector-elasticsearch.git


The following commit(s) were added to refs/heads/main by this push:
     new c63a8d7  [FLINK-40381] Fix thread leak in Elasticsearch8AsyncWriter by 
closing HTTP transport on sink close (#162)
c63a8d7 is described below

commit c63a8d7732bbcc3a58187312c698bda47c5742d5
Author: nadavFux <[email protected]>
AuthorDate: Thu Aug 27 11:51:27 2026 +0300

    [FLINK-40381] Fix thread leak in Elasticsearch8AsyncWriter by closing HTTP 
transport on sink close (#162)
---
 .../sink/Elasticsearch8AsyncWriter.java            |  9 ++++++-
 .../sink/Elasticsearch8AsyncWriterITCase.java      | 28 +++++++---------------
 2 files changed, 16 insertions(+), 21 deletions(-)

diff --git 
a/flink-connector-elasticsearch8/src/main/java/org/apache/flink/connector/elasticsearch/sink/Elasticsearch8AsyncWriter.java
 
b/flink-connector-elasticsearch8/src/main/java/org/apache/flink/connector/elasticsearch/sink/Elasticsearch8AsyncWriter.java
index 1d8ceaf..3073ba0 100644
--- 
a/flink-connector-elasticsearch8/src/main/java/org/apache/flink/connector/elasticsearch/sink/Elasticsearch8AsyncWriter.java
+++ 
b/flink-connector-elasticsearch8/src/main/java/org/apache/flink/connector/elasticsearch/sink/Elasticsearch8AsyncWriter.java
@@ -39,6 +39,7 @@ import 
co.elastic.clients.elasticsearch.core.bulk.BulkOperation;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
+import java.io.IOException;
 import java.net.ConnectException;
 import java.net.NoRouteToHostException;
 import java.util.ArrayList;
@@ -204,7 +205,13 @@ public class Elasticsearch8AsyncWriter<InputT> extends 
AsyncSinkWriter<InputT, O
     public void close() {
         if (!close) {
             close = true;
-            esClient.shutdown();
+            try {
+                if (esClient != null && esClient._transport() != null) {
+                    esClient._transport().close();
+                }
+            } catch (IOException e) {
+                LOG.warn("Failed to close Elasticsearch transport during sink 
close", e);
+            }
         }
     }
 }
diff --git 
a/flink-connector-elasticsearch8/src/test/java/org/apache/flink/connector/elasticsearch/sink/Elasticsearch8AsyncWriterITCase.java
 
b/flink-connector-elasticsearch8/src/test/java/org/apache/flink/connector/elasticsearch/sink/Elasticsearch8AsyncWriterITCase.java
index d65ae61..05559c8 100644
--- 
a/flink-connector-elasticsearch8/src/test/java/org/apache/flink/connector/elasticsearch/sink/Elasticsearch8AsyncWriterITCase.java
+++ 
b/flink-connector-elasticsearch8/src/test/java/org/apache/flink/connector/elasticsearch/sink/Elasticsearch8AsyncWriterITCase.java
@@ -37,9 +37,7 @@ import java.io.IOException;
 import java.util.Collections;
 import java.util.List;
 import java.util.Optional;
-import java.util.concurrent.locks.Condition;
-import java.util.concurrent.locks.Lock;
-import java.util.concurrent.locks.ReentrantLock;
+import java.util.concurrent.CountDownLatch;
 
 import static org.assertj.core.api.Assertions.assertThat;
 
@@ -47,13 +45,12 @@ import static org.assertj.core.api.Assertions.assertThat;
 public class Elasticsearch8AsyncWriterITCase extends 
ElasticsearchSinkBaseITCase {
     private TestSinkInitContext context;
 
-    private final Lock lock = new ReentrantLock();
-
-    private final Condition completed = lock.newCondition();
+    private CountDownLatch latch;
 
     @BeforeEach
     void setUp() {
         this.context = new TestSinkInitContext();
+        this.latch = new CountDownLatch(1);
     }
 
     @TestTemplate
@@ -82,6 +79,7 @@ public class Elasticsearch8AsyncWriterITCase extends 
ElasticsearchSinkBaseITCase
     public void testBulkOnBufferTimeFlush() throws Exception {
         String index = "test-bulk-on-time-in-buffer";
         int maxBatchSize = 3;
+        this.latch = new CountDownLatch(2);
 
         try (final Elasticsearch8AsyncWriter<DummyData> writer =
                 createWriter(index, maxBatchSize)) {
@@ -191,9 +189,9 @@ public class Elasticsearch8AsyncWriterITCase extends 
ElasticsearchSinkBaseITCase
                 createWriter(maxBatchSize, elementConverter)) {
             writer.write(new DummyData("test-1", "test-1-updated"), null);
             writer.write(new DummyData("test-2", "test-2-updated"), null);
-        }
 
-        await();
+            await();
+        }
 
         
assertThat(context.metricGroup().getNumRecordsOutErrorsCounter().getCount()).isEqualTo(1);
         assertIdsAreWritten(index, new String[] {"test-2"});
@@ -287,20 +285,10 @@ public class Elasticsearch8AsyncWriterITCase extends 
ElasticsearchSinkBaseITCase
     }
 
     private void signal() {
-        lock.lock();
-        try {
-            completed.signal();
-        } finally {
-            lock.unlock();
-        }
+        latch.countDown();
     }
 
     private void await() throws InterruptedException {
-        lock.lock();
-        try {
-            completed.await();
-        } finally {
-            lock.unlock();
-        }
+        latch.await();
     }
 }

Reply via email to