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();
}
}