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

Reply via email to