This is an automated email from the ASF dual-hosted git repository.
zhangshenghang pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git
The following commit(s) were added to refs/heads/dev by this push:
new 310d5643dd [Fix][Connector-V2] Share Elasticsearch REST clients in
multi-table sink (#11567)
310d5643dd is described below
commit 310d5643dd1b053d695a6e8e9c2d7e6c34e6fccd
Author: David Zollo <[email protected]>
AuthorDate: Thu Aug 20 23:34:04 2026 +0800
[Fix][Connector-V2] Share Elasticsearch REST clients in multi-table sink
(#11567)
Co-authored-by: DanielLeens <[email protected]>
Co-authored-by: DanielLeens <[email protected]>
---
.../ElasticsearchMultiTableResourceManager.java | 235 +++++++++++++++++
.../sink/ElasticsearchSinkWriter.java | 166 +++++++++---
.../sink/ElasticsearchSinkWriterTest.java | 287 ++++++++++++++++++++-
3 files changed, 656 insertions(+), 32 deletions(-)
diff --git
a/seatunnel-connectors-v2/connector-elasticsearch/src/main/java/org/apache/seatunnel/connectors/seatunnel/elasticsearch/sink/ElasticsearchMultiTableResourceManager.java
b/seatunnel-connectors-v2/connector-elasticsearch/src/main/java/org/apache/seatunnel/connectors/seatunnel/elasticsearch/sink/ElasticsearchMultiTableResourceManager.java
new file mode 100644
index 0000000000..c04c885001
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-elasticsearch/src/main/java/org/apache/seatunnel/connectors/seatunnel/elasticsearch/sink/ElasticsearchMultiTableResourceManager.java
@@ -0,0 +1,235 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.elasticsearch.sink;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.sink.MultiTableResourceManager;
+import
org.apache.seatunnel.connectors.seatunnel.elasticsearch.client.EsRestClient;
+import
org.apache.seatunnel.connectors.seatunnel.elasticsearch.config.ElasticsearchBaseOptions;
+import
org.apache.seatunnel.connectors.seatunnel.elasticsearch.dto.ElasticsearchClusterInfo;
+
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+
+/**
+ * Owns Elasticsearch REST clients shared by table writers with identical
connection settings in one
+ * multi-table sink subtask.
+ *
+ * <p>Table placeholders may produce different hosts, credentials, or TLS
settings. Such writers
+ * must use separate clients, while table-only settings such as index names do
not prevent sharing.
+ */
+class ElasticsearchMultiTableResourceManager implements
MultiTableResourceManager<EsRestClient> {
+
+ // Clients grouped by the complete set of options consumed during client
construction.
+ private final Map<List<Object>, ClientResource> clientResources = new
HashMap<>();
+
+ /**
+ * Creates an empty manager; clients are initialized lazily from each
table writer's config.
+ *
+ * <p>Lazy initialization prevents transient creation of one REST client
per table.
+ */
+ ElasticsearchMultiTableResourceManager() {}
+
+ /**
+ * Returns the client and cluster metadata for one table's connection
settings.
+ *
+ * <p>Initialization is synchronized so writers assigned to different
queues cannot create
+ * duplicate clients for the same connection.
+ *
+ * @param config table-specific Elasticsearch configuration
+ * @return client resource shared by writers with identical connection
settings
+ */
+ synchronized ClientResource getOrCreateClientResource(ReadonlyConfig
config) {
+ List<Object> connectionKey = connectionKey(config);
+ ClientResource existingResource = clientResources.get(connectionKey);
+ if (existingResource != null) {
+ existingResource.retain();
+ return existingResource;
+ }
+ EsRestClient esRestClient = null;
+ try {
+ esRestClient = EsRestClient.createInstance(config);
+ ClientResource newResource =
+ new ClientResource(connectionKey, esRestClient,
esRestClient.getClusterInfo());
+ newResource.retain();
+ clientResources.put(connectionKey, newResource);
+ return newResource;
+ } catch (RuntimeException | Error e) {
+ if (esRestClient != null) {
+ esRestClient.close();
+ }
+ // MultiTableSinkWriter does not close its manager when
construction fails. Release all
+ // connection groups here so task startup retries cannot
accumulate client threads.
+ close();
+ throw e;
+ }
+ }
+
+ /**
+ * Releases one writer's reference to a connection group and closes the
client at the last user.
+ *
+ * <p>Multi-table writers can be removed independently when a table fails
with
+ * CONTINUE_OTHER_TABLES. Releasing per writer prevents unique connection
groups from living
+ * until the whole sink stops.
+ *
+ * @param clientResource resource previously retained by {@link
#getOrCreateClientResource}
+ */
+ synchronized void releaseClientResource(ClientResource clientResource) {
+ if (!clientResources.containsKey(clientResource.getConnectionKey())) {
+ return;
+ }
+ if (clientResource.release()) {
+ clientResources.remove(clientResource.getConnectionKey());
+ clientResource.getEsRestClient().close();
+ }
+ }
+
+ /**
+ * Returns the single shared client when all initialized writers use one
connection group.
+ *
+ * <p>The multi-table writer uses {@link
#getOrCreateClientResource(ReadonlyConfig)} because a
+ * task may legitimately contain multiple connection groups.
+ *
+ * @return the only initialized client, or empty when zero or multiple
groups exist
+ */
+ @Override
+ public synchronized Optional<EsRestClient> getSharedResource() {
+ if (clientResources.size() != 1) {
+ return Optional.empty();
+ }
+ return
Optional.of(clientResources.values().iterator().next().getEsRestClient());
+ }
+
+ /**
+ * Closes every connection group exactly once and makes repeated close
calls harmless.
+ *
+ * <p>The map is cleared only after every owned client has been closed.
+ */
+ @Override
+ public synchronized void close() {
+ clientResources.values().forEach(resource ->
resource.getEsRestClient().close());
+ clientResources.clear();
+ }
+
+ /**
+ * Builds a stable key from every option consumed by {@link
EsRestClient#createInstance}.
+ *
+ * <p>Do not add table-specific sink options here: doing so would recreate
one client per table.
+ *
+ * @param config table-specific Elasticsearch configuration
+ * @return immutable-by-construction connection key
+ */
+ private List<Object> connectionKey(ReadonlyConfig config) {
+ return Arrays.asList(
+ new ArrayList<>(config.get(ElasticsearchBaseOptions.HOSTS)),
+
config.getOptional(ElasticsearchBaseOptions.USERNAME).orElse(null),
+
config.getOptional(ElasticsearchBaseOptions.PASSWORD).orElse(null),
+ config.get(ElasticsearchBaseOptions.TLS_VERIFY_CERTIFICATE),
+ config.get(ElasticsearchBaseOptions.TLS_VERIFY_HOSTNAME),
+
config.getOptional(ElasticsearchBaseOptions.TLS_KEY_STORE_PATH).orElse(null),
+
config.getOptional(ElasticsearchBaseOptions.TLS_KEY_STORE_PASSWORD).orElse(null),
+
config.getOptional(ElasticsearchBaseOptions.TLS_TRUST_STORE_PATH).orElse(null),
+
config.getOptional(ElasticsearchBaseOptions.TLS_TRUST_STORE_PASSWORD).orElse(null),
+ config.get(ElasticsearchBaseOptions.AUTH_TYPE),
+
config.getOptional(ElasticsearchBaseOptions.API_KEY_ID).orElse(null),
+
config.getOptional(ElasticsearchBaseOptions.API_KEY).orElse(null),
+
config.getOptional(ElasticsearchBaseOptions.API_KEY_ENCODED).orElse(null));
+ }
+
+ /**
+ * Holds one connection group's client and immutable cluster metadata.
+ *
+ * <p>Both values have the same lifecycle and are reused by every writer
in the group.
+ */
+ static final class ClientResource {
+
+ // Connection key used to remove this resource when the last writer
releases it.
+ private final List<Object> connectionKey;
+
+ // REST client shared by all writers in this connection group.
+ private final EsRestClient esRestClient;
+
+ // Cluster metadata cached once for all serializers in this connection
group.
+ private final ElasticsearchClusterInfo clusterInfo;
+
+ // Number of active table writers currently using this connection
group.
+ private int referenceCount;
+
+ /**
+ * Creates one connection-group resource.
+ *
+ * @param connectionKey stable connection-group key
+ * @param esRestClient shared REST client
+ * @param clusterInfo cluster metadata loaded through the client
+ */
+ private ClientResource(
+ List<Object> connectionKey,
+ EsRestClient esRestClient,
+ ElasticsearchClusterInfo clusterInfo) {
+ this.connectionKey = connectionKey;
+ this.esRestClient = esRestClient;
+ this.clusterInfo = clusterInfo;
+ }
+
+ /**
+ * Returns the stable connection-group key.
+ *
+ * @return connection key
+ */
+ List<Object> getConnectionKey() {
+ return connectionKey;
+ }
+
+ /**
+ * Returns the shared REST client.
+ *
+ * @return connection-group client
+ */
+ EsRestClient getEsRestClient() {
+ return esRestClient;
+ }
+
+ /**
+ * Returns cached cluster metadata.
+ *
+ * @return connection-group cluster metadata
+ */
+ ElasticsearchClusterInfo getClusterInfo() {
+ return clusterInfo;
+ }
+
+ /** Records one active table writer using this connection group. */
+ void retain() {
+ referenceCount++;
+ }
+
+ /**
+ * Releases one active table writer.
+ *
+ * @return true when the connection group has no remaining users
+ */
+ boolean release() {
+ referenceCount--;
+ return referenceCount == 0;
+ }
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-elasticsearch/src/main/java/org/apache/seatunnel/connectors/seatunnel/elasticsearch/sink/ElasticsearchSinkWriter.java
b/seatunnel-connectors-v2/connector-elasticsearch/src/main/java/org/apache/seatunnel/connectors/seatunnel/elasticsearch/sink/ElasticsearchSinkWriter.java
index 1b39d7708b..13b80ff3f5 100644
---
a/seatunnel-connectors-v2/connector-elasticsearch/src/main/java/org/apache/seatunnel/connectors/seatunnel/elasticsearch/sink/ElasticsearchSinkWriter.java
+++
b/seatunnel-connectors-v2/connector-elasticsearch/src/main/java/org/apache/seatunnel/connectors/seatunnel/elasticsearch/sink/ElasticsearchSinkWriter.java
@@ -18,9 +18,11 @@
package org.apache.seatunnel.connectors.seatunnel.elasticsearch.sink;
import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.sink.MultiTableResourceManager;
import org.apache.seatunnel.api.sink.SinkWriter;
import org.apache.seatunnel.api.sink.SupportMultiTableSinkWriter;
import org.apache.seatunnel.api.sink.SupportSchemaEvolutionSinkWriter;
+import org.apache.seatunnel.api.sink.multitablesink.SinkContextProxy;
import org.apache.seatunnel.api.table.catalog.CatalogTable;
import org.apache.seatunnel.api.table.catalog.Column;
import org.apache.seatunnel.api.table.catalog.TableSchema;
@@ -35,6 +37,7 @@ import
org.apache.seatunnel.api.table.schema.handler.TableSchemaChangeEventDispa
import
org.apache.seatunnel.api.table.schema.handler.TableSchemaChangeEventHandler;
import org.apache.seatunnel.api.table.type.RowKind;
import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
import org.apache.seatunnel.common.exception.CommonErrorCodeDeprecated;
import org.apache.seatunnel.common.utils.RetryUtils;
import org.apache.seatunnel.common.utils.RetryUtils.RetryMaterial;
@@ -44,6 +47,7 @@ import
org.apache.seatunnel.connectors.seatunnel.elasticsearch.client.EsRestClie
import org.apache.seatunnel.connectors.seatunnel.elasticsearch.client.EsType;
import
org.apache.seatunnel.connectors.seatunnel.elasticsearch.config.ElasticsearchSinkOptions;
import
org.apache.seatunnel.connectors.seatunnel.elasticsearch.dto.BulkResponse;
+import
org.apache.seatunnel.connectors.seatunnel.elasticsearch.dto.ElasticsearchClusterInfo;
import org.apache.seatunnel.connectors.seatunnel.elasticsearch.dto.IndexInfo;
import
org.apache.seatunnel.connectors.seatunnel.elasticsearch.exception.ElasticsearchConnectorErrorCode;
import
org.apache.seatunnel.connectors.seatunnel.elasticsearch.exception.ElasticsearchConnectorException;
@@ -66,7 +70,7 @@ import java.util.Optional;
@Slf4j
public class ElasticsearchSinkWriter
implements SinkWriter<SeaTunnelRow, ElasticsearchCommitInfo,
ElasticsearchSinkState>,
- SupportMultiTableSinkWriter<Void>,
+ SupportMultiTableSinkWriter<EsRestClient>,
SupportSchemaEvolutionSinkWriter {
private final Context context;
@@ -76,9 +80,26 @@ public class ElasticsearchSinkWriter
private SeaTunnelRowSerializer seaTunnelRowSerializer;
private final List<String> requestEsList;
private EsRestClient esRestClient;
+
+ // Cluster metadata cached by the owner of esRestClient.
+ private ElasticsearchClusterInfo clusterInfo;
+
+ // Whether this writer, instead of the multi-table resource manager, owns
the REST client.
+ private boolean ownsEsRestClient;
+
+ // Resource manager that owns the shared client injected into this
multi-table writer.
+ private ElasticsearchMultiTableResourceManager multiTableResourceManager;
+
+ // Retained connection-group resource released exactly once when this
writer closes.
+ private ElasticsearchMultiTableResourceManager.ClientResource
multiTableClientResource;
+
private RetryMaterial retryMaterial;
private static final long DEFAULT_SLEEP_TIME_MS = 200L;
private final IndexInfo indexInfo;
+
+ // Initial physical row type used to build the serializer after
shared-client injection.
+ private final SeaTunnelRowType initialRowType;
+
private TableSchema tableSchema;
private final TableSchemaChangeEventHandler tableSchemaChangeEventHandler;
private final ReadonlyConfig config;
@@ -95,30 +116,59 @@ public class ElasticsearchSinkWriter
this.indexInfo =
new
IndexInfo(catalogTable.getTableId().getTableName().toLowerCase(), config);
- esRestClient = EsRestClient.createInstance(config);
-
- // Get vectorization fields and dimension from config
- List<String> vectorizationFields =
-
config.getOptional(ElasticsearchSinkOptions.VECTORIZATION_FIELDS)
- .orElse(Collections.emptyList());
- int vectorDimension =
config.get(ElasticsearchSinkOptions.VECTOR_DIMENSIONS);
-
- this.seaTunnelRowSerializer =
- new ElasticsearchRowSerializer(
- esRestClient.getClusterInfo(),
- indexInfo,
- catalogTable.getSeaTunnelRowType(),
- vectorizationFields,
- vectorDimension);
-
+ this.initialRowType = catalogTable.getSeaTunnelRowType();
this.requestEsList = new ArrayList<>(maxBatchSize);
this.retryMaterial =
new RetryMaterial(maxRetryCount, true, exception -> true,
DEFAULT_SLEEP_TIME_MS);
this.tableSchema = catalogTable.getTableSchema();
this.tableSchemaChangeEventHandler = new
TableSchemaChangeEventDispatcher();
+
+ // MultiTableSinkWriter injects one shared client after constructing
all table writers.
+ // Standalone writers keep the existing fail-fast connection
initialization behavior.
+ if (!(context instanceof SinkContextProxy)) {
+ initializeStandaloneClient();
+ }
context.registerFlushAction(this::timerFlush);
}
+ @Override
+ public MultiTableResourceManager<EsRestClient>
initMultiTableResourceManager(
+ int tableSize, int queueSize) {
+ return new ElasticsearchMultiTableResourceManager();
+ }
+
+ @Override
+ public void setMultiTableResourceManager(
+ MultiTableResourceManager<EsRestClient> multiTableResourceManager,
int queueIndex) {
+ if (!(multiTableResourceManager instanceof
ElasticsearchMultiTableResourceManager)) {
+ throw new IllegalArgumentException(
+ "Elasticsearch multi-table writer requires
ElasticsearchMultiTableResourceManager");
+ }
+ releaseSharedClientResource();
+ closeOwnedClient();
+ ElasticsearchMultiTableResourceManager resourceManager =
+ (ElasticsearchMultiTableResourceManager)
multiTableResourceManager;
+ try {
+ ElasticsearchMultiTableResourceManager.ClientResource
clientResource =
+ resourceManager.getOrCreateClientResource(config);
+ this.esRestClient = clientResource.getEsRestClient();
+ this.clusterInfo = clientResource.getClusterInfo();
+ this.multiTableResourceManager = resourceManager;
+ this.multiTableClientResource = clientResource;
+ this.ownsEsRestClient = false;
+ initializeSerializer(initialRowType);
+ } catch (RuntimeException | Error e) {
+ // A failed injection aborts MultiTableSinkWriter construction,
whose close lifecycle
+ // will not run. Close every cached connection group before
propagating the failure.
+ resourceManager.close();
+ this.multiTableResourceManager = null;
+ this.multiTableClientResource = null;
+ this.esRestClient = null;
+ this.clusterInfo = null;
+ throw e;
+ }
+ }
+
@Override
public void write(SeaTunnelRow element) {
if (RowKind.UPDATE_BEFORE.equals(element.getRowKind())) {
@@ -148,19 +198,7 @@ public class ElasticsearchSinkWriter
this.tableSchema =
tableSchemaChangeEventHandler.reset(tableSchema).apply(event);
- // Get vectorization fields and dimension from config
- List<String> vectorizationFields =
-
config.getOptional(ElasticsearchSinkOptions.VECTORIZATION_FIELDS)
- .orElse(Collections.emptyList());
- int vectorDimension =
config.get(ElasticsearchSinkOptions.VECTOR_DIMENSIONS);
-
- this.seaTunnelRowSerializer =
- new ElasticsearchRowSerializer(
- esRestClient.getClusterInfo(),
- indexInfo,
- tableSchema.toPhysicalRowDataType(),
- vectorizationFields,
- vectorDimension);
+ initializeSerializer(tableSchema.toPhysicalRowDataType());
}
static boolean isCommentOnlyEvent(SchemaChangeEvent event) {
@@ -188,7 +226,11 @@ public class ElasticsearchSinkWriter
return Optional.empty();
}
- /** Flushes pending bulk requests when the Zeta engine delivers a timer
flush signal. */
+ /**
+ * Flushes pending bulk requests when the Zeta engine delivers a timer
flush signal.
+ *
+ * <p>The action is registered before multi-table resource injection but
invoked after startup.
+ */
private void timerFlush() {
bulkEsWithRetry(this.esRestClient, this.requestEsList);
}
@@ -225,10 +267,72 @@ public class ElasticsearchSinkWriter
@Override
public void close() {
+ if (esRestClient == null) {
+ return;
+ }
try {
bulkEsWithRetry(this.esRestClient, this.requestEsList);
} finally {
+ releaseSharedClientResource();
+ closeOwnedClient();
+ }
+ }
+
+ /**
+ * Initializes the client owned by a standalone writer and preserves
fail-fast startup
+ * validation.
+ */
+ private void initializeStandaloneClient() {
+ EsRestClient client = EsRestClient.createInstance(config);
+ try {
+ this.clusterInfo = client.getClusterInfo();
+ this.esRestClient = client;
+ this.ownsEsRestClient = true;
+ initializeSerializer(initialRowType);
+ } catch (RuntimeException e) {
+ client.close();
+ throw e;
+ }
+ }
+
+ /**
+ * Rebuilds the table-specific serializer using cluster metadata cached by
the client owner.
+ *
+ * @param rowType current physical row type
+ */
+ private void initializeSerializer(SeaTunnelRowType rowType) {
+ List<String> vectorizationFields =
+
config.getOptional(ElasticsearchSinkOptions.VECTORIZATION_FIELDS)
+ .orElse(Collections.emptyList());
+ int vectorDimension =
config.get(ElasticsearchSinkOptions.VECTOR_DIMENSIONS);
+ this.seaTunnelRowSerializer =
+ new ElasticsearchRowSerializer(
+ clusterInfo, indexInfo, rowType, vectorizationFields,
vectorDimension);
+ }
+
+ /**
+ * Closes only a client owned by this writer; shared clients are closed by
their manager.
+ *
+ * <p>Clearing the ownership flag makes repeated writer close calls safe.
+ */
+ private void closeOwnedClient() {
+ if (ownsEsRestClient && esRestClient != null) {
esRestClient.close();
+ ownsEsRestClient = false;
+ }
+ }
+
+ /**
+ * Releases the shared multi-table client reference retained during
resource injection.
+ *
+ * <p>The manager closes the underlying REST client only when this writer
was the last active
+ * user of the connection group.
+ */
+ private void releaseSharedClientResource() {
+ if (multiTableResourceManager != null && multiTableClientResource !=
null) {
+
multiTableResourceManager.releaseClientResource(multiTableClientResource);
+ multiTableResourceManager = null;
+ multiTableClientResource = null;
}
}
}
diff --git
a/seatunnel-connectors-v2/connector-elasticsearch/src/test/java/org/apache/seatunnel/connectors/seatunnel/elasticsearch/sink/ElasticsearchSinkWriterTest.java
b/seatunnel-connectors-v2/connector-elasticsearch/src/test/java/org/apache/seatunnel/connectors/seatunnel/elasticsearch/sink/ElasticsearchSinkWriterTest.java
index 672343eb22..5dc18bed73 100644
---
a/seatunnel-connectors-v2/connector-elasticsearch/src/test/java/org/apache/seatunnel/connectors/seatunnel/elasticsearch/sink/ElasticsearchSinkWriterTest.java
+++
b/seatunnel-connectors-v2/connector-elasticsearch/src/test/java/org/apache/seatunnel/connectors/seatunnel/elasticsearch/sink/ElasticsearchSinkWriterTest.java
@@ -17,22 +17,244 @@
package org.apache.seatunnel.connectors.seatunnel.elasticsearch.sink;
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.sink.MultiTableResourceManager;
+import org.apache.seatunnel.api.sink.SinkWriter;
+import org.apache.seatunnel.api.sink.multitablesink.SinkContextProxy;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
import org.apache.seatunnel.api.table.catalog.PhysicalColumn;
import org.apache.seatunnel.api.table.catalog.TableIdentifier;
import org.apache.seatunnel.api.table.catalog.TablePath;
+import org.apache.seatunnel.api.table.catalog.TableSchema;
import org.apache.seatunnel.api.table.schema.event.AlterColumnCommentEvent;
import org.apache.seatunnel.api.table.schema.event.AlterTableAddColumnEvent;
import org.apache.seatunnel.api.table.schema.event.AlterTableCommentEvent;
import org.apache.seatunnel.api.table.type.BasicType;
+import org.apache.seatunnel.api.table.type.SeaTunnelDataType;
+import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
+import
org.apache.seatunnel.connectors.seatunnel.elasticsearch.client.EsRestClient;
+import
org.apache.seatunnel.connectors.seatunnel.elasticsearch.dto.ElasticsearchClusterInfo;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
+import org.mockito.MockedStatic;
-public class ElasticsearchSinkWriterTest {
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.Map;
+
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * Verifies Elasticsearch REST client ownership for standalone and multi-table
sink writers.
+ *
+ * <p>The tests guard resource reduction, table-specific connection isolation,
and the comment-only
+ * schema-change no-op behavior of the writer.
+ */
+class ElasticsearchSinkWriterTest {
private static final TableIdentifier TABLE_IDENTIFIER =
TableIdentifier.of("", TablePath.DEFAULT);
+ /**
+ * Multi-table writers must share one client and release it only after the
last writer closes.
+ */
+ @Test
+ void testMultiTableWritersShareRestClient() {
+ EsRestClient sharedClient = mock(EsRestClient.class);
+ when(sharedClient.getClusterInfo()).thenReturn(clusterInfo());
+
+ try (MockedStatic<EsRestClient> clientFactory =
mockStatic(EsRestClient.class)) {
+ clientFactory.when(() ->
EsRestClient.createInstance(any())).thenReturn(sharedClient);
+
+ ElasticsearchSinkWriter firstWriter =
+ createMultiTableWriter(
+ "first", config("http://localhost:9200",
"first_index", null));
+ ElasticsearchSinkWriter secondWriter =
+ createMultiTableWriter(
+ "second", config("http://localhost:9200",
"second_index", null));
+ clientFactory.verify(() -> EsRestClient.createInstance(any()),
never());
+
+ MultiTableResourceManager<EsRestClient> resourceManager =
+ firstWriter.initMultiTableResourceManager(2, 1);
+ firstWriter.setMultiTableResourceManager(resourceManager, 0);
+ secondWriter.setMultiTableResourceManager(resourceManager, 0);
+
+ clientFactory.verify(() -> EsRestClient.createInstance(any()),
times(1));
+ verify(sharedClient, times(1)).getClusterInfo();
+ assertSame(
+ sharedClient,
+
resourceManager.getSharedResource().orElseThrow(AssertionError::new));
+
+ firstWriter.close();
+ verify(sharedClient, never()).close();
+
+ secondWriter.close();
+ verify(sharedClient, times(1)).close();
+
+ resourceManager.close();
+ verify(sharedClient, times(1)).close();
+ }
+ }
+
+ /**
+ * Writers with different connection settings must not share credentials,
TLS state, or hosts.
+ */
+ @Test
+ void testMultiTableWritersKeepDifferentConnectionsIsolated() {
+ EsRestClient firstClient = mock(EsRestClient.class);
+ EsRestClient secondClient = mock(EsRestClient.class);
+ when(firstClient.getClusterInfo()).thenReturn(clusterInfo());
+ when(secondClient.getClusterInfo()).thenReturn(clusterInfo());
+
+ try (MockedStatic<EsRestClient> clientFactory =
mockStatic(EsRestClient.class)) {
+ clientFactory
+ .when(() -> EsRestClient.createInstance(any()))
+ .thenReturn(firstClient, secondClient);
+
+ ElasticsearchSinkWriter firstWriter =
+ createMultiTableWriter(
+ "first", config("http://localhost:9200",
"first_index", "first_user"));
+ ElasticsearchSinkWriter secondWriter =
+ createMultiTableWriter(
+ "second",
+ config("http://localhost:9200", "second_index",
"second_user"));
+ MultiTableResourceManager<EsRestClient> resourceManager =
+ firstWriter.initMultiTableResourceManager(2, 1);
+
+ firstWriter.setMultiTableResourceManager(resourceManager, 0);
+ secondWriter.setMultiTableResourceManager(resourceManager, 0);
+
+ clientFactory.verify(() -> EsRestClient.createInstance(any()),
times(2));
+ verify(firstClient, times(1)).getClusterInfo();
+ verify(secondClient, times(1)).getClusterInfo();
+ assertFalse(resourceManager.getSharedResource().isPresent());
+
+ firstWriter.close();
+ verify(firstClient, times(1)).close();
+ verify(secondClient, never()).close();
+
+ secondWriter.close();
+ verify(firstClient, times(1)).close();
+ verify(secondClient, times(1)).close();
+
+ resourceManager.close();
+ resourceManager.close();
+ verify(firstClient, times(1)).close();
+ verify(secondClient, times(1)).close();
+ }
+ }
+
+ /**
+ * A later connection-group initialization failure must close earlier
groups immediately.
+ *
+ * <p>This prevents task startup retries from accumulating clients from
partial initialization.
+ */
+ @Test
+ void testConnectionGroupInitializationFailureClosesAllClients() {
+ EsRestClient firstClient = mock(EsRestClient.class);
+ EsRestClient failingClient = mock(EsRestClient.class);
+ when(firstClient.getClusterInfo()).thenReturn(clusterInfo());
+ when(failingClient.getClusterInfo())
+ .thenThrow(new IllegalStateException("cluster info
unavailable"));
+
+ try (MockedStatic<EsRestClient> clientFactory =
mockStatic(EsRestClient.class)) {
+ clientFactory
+ .when(() -> EsRestClient.createInstance(any()))
+ .thenReturn(firstClient, failingClient);
+
+ ElasticsearchSinkWriter firstWriter =
+ createMultiTableWriter(
+ "first", config("http://localhost:9200",
"first_index", "first_user"));
+ ElasticsearchSinkWriter failingWriter =
+ createMultiTableWriter(
+ "second",
+ config("http://localhost:9200", "second_index",
"second_user"));
+ MultiTableResourceManager<EsRestClient> resourceManager =
+ firstWriter.initMultiTableResourceManager(2, 1);
+
+ firstWriter.setMultiTableResourceManager(resourceManager, 0);
+ assertThrows(
+ IllegalStateException.class,
+ () ->
failingWriter.setMultiTableResourceManager(resourceManager, 0));
+
+ verify(firstClient, times(1)).close();
+ verify(failingClient, times(1)).close();
+ resourceManager.close();
+ verify(firstClient, times(1)).close();
+ verify(failingClient, times(1)).close();
+ }
+ }
+
+ /**
+ * A serializer initialization failure must close the client cached during
resource injection.
+ */
+ @Test
+ void testSerializerInitializationFailureClosesSharedClient() {
+ EsRestClient sharedClient = mock(EsRestClient.class);
+ when(sharedClient.getClusterInfo())
+ .thenReturn(
+ ElasticsearchClusterInfo.builder()
+ .clusterVersion("invalid-version")
+ .build());
+
+ try (MockedStatic<EsRestClient> clientFactory =
mockStatic(EsRestClient.class)) {
+ clientFactory.when(() ->
EsRestClient.createInstance(any())).thenReturn(sharedClient);
+
+ ElasticsearchSinkWriter writer =
+ createMultiTableWriter("table",
config("http://localhost:9200", "index", null));
+ MultiTableResourceManager<EsRestClient> resourceManager =
+ writer.initMultiTableResourceManager(1, 1);
+
+ assertThrows(
+ NumberFormatException.class,
+ () -> writer.setMultiTableResourceManager(resourceManager,
0));
+
+ verify(sharedClient, times(1)).close();
+ resourceManager.close();
+ verify(sharedClient, times(1)).close();
+ }
+ }
+
+ /**
+ * Standalone writers must retain the existing eager initialization and
ownership behavior.
+ *
+ * <p>This prevents the multi-table optimization from weakening startup
validation.
+ */
+ @Test
+ void testStandaloneWriterOwnsRestClient() {
+ EsRestClient standaloneClient = mock(EsRestClient.class);
+ when(standaloneClient.getClusterInfo()).thenReturn(clusterInfo());
+
+ try (MockedStatic<EsRestClient> clientFactory =
mockStatic(EsRestClient.class)) {
+ clientFactory
+ .when(() -> EsRestClient.createInstance(any()))
+ .thenReturn(standaloneClient);
+
+ ElasticsearchSinkWriter writer =
+ createWriter(
+ mock(SinkWriter.Context.class),
+ "standalone",
+ config("http://localhost:9200",
"standalone_index", null));
+
+ clientFactory.verify(() -> EsRestClient.createInstance(any()),
times(1));
+ verify(standaloneClient, times(1)).getClusterInfo();
+
+ writer.close();
+ verify(standaloneClient, times(1)).close();
+ }
+ }
+
+ /** Comment-only schema changes must not trigger any Elasticsearch mapping
update. */
@Test
void commentOnlySchemaChangeEventsAreNoOpForElasticsearch() {
Assertions.assertTrue(
@@ -43,6 +265,7 @@ public class ElasticsearchSinkWriterTest {
AlterColumnCommentEvent.of(TABLE_IDENTIFIER, "name",
"old", "new")));
}
+ /** Physical schema changes still require an Elasticsearch mapping change.
*/
@Test
void physicalSchemaChangeEventsStillRequireElasticsearchMappingChanges() {
Assertions.assertFalse(
@@ -54,4 +277,66 @@ public class ElasticsearchSinkWriterTest {
.dataType(BasicType.STRING_TYPE)
.build())));
}
+
+ /**
+ * Creates a writer with the context used by the multi-table sink runtime.
+ *
+ * @param tableName target table name
+ * @param config table-specific Elasticsearch configuration
+ * @return uninitialized multi-table writer
+ */
+ private ElasticsearchSinkWriter createMultiTableWriter(
+ String tableName, ReadonlyConfig config) {
+ SinkWriter.Context context = mock(SinkWriter.Context.class);
+ return createWriter(new SinkContextProxy(0, 1, context), tableName,
config);
+ }
+
+ /**
+ * Creates a writer with deterministic table metadata and connector
options.
+ *
+ * @param context sink writer context
+ * @param tableName target table name
+ * @param config table-specific Elasticsearch configuration
+ * @return Elasticsearch sink writer
+ */
+ private ElasticsearchSinkWriter createWriter(
+ SinkWriter.Context context, String tableName, ReadonlyConfig
config) {
+ SeaTunnelRowType rowType =
+ new SeaTunnelRowType(
+ new String[] {"id"}, new SeaTunnelDataType<?>[]
{BasicType.INT_TYPE});
+ CatalogTable catalogTable = mock(CatalogTable.class);
+ when(catalogTable.getTableId())
+ .thenReturn(TableIdentifier.of("catalog", "database",
tableName));
+ when(catalogTable.getSeaTunnelRowType()).thenReturn(rowType);
+
when(catalogTable.getTableSchema()).thenReturn(mock(TableSchema.class));
+ return new ElasticsearchSinkWriter(context, catalogTable, config,
1000, 3);
+ }
+
+ /**
+ * Returns the minimal connection configuration required by the writer.
+ *
+ * @param host Elasticsearch HTTP endpoint
+ * @param index table-specific target index
+ * @param username optional basic-auth username
+ * @return Elasticsearch connector configuration
+ */
+ private ReadonlyConfig config(String host, String index, String username) {
+ Map<String, Object> options = new HashMap<>();
+ options.put("hosts", Collections.singletonList(host));
+ options.put("index", index);
+ if (username != null) {
+ options.put("username", username);
+ options.put("password", "password");
+ }
+ return ReadonlyConfig.fromMap(options);
+ }
+
+ /**
+ * Returns deterministic cluster metadata for serializer construction.
+ *
+ * @return Elasticsearch cluster metadata
+ */
+ private ElasticsearchClusterInfo clusterInfo() {
+ return
ElasticsearchClusterInfo.builder().clusterVersion("7.17.0").build();
+ }
}