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

github-merge-queue[bot] 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 bbfc188e61 [Fix][Zeta] Stabilize CI test isolation (#12050)
bbfc188e61 is described below

commit bbfc188e61b87e79624155043f151b8e84bca6eb
Author: Daniel <[email protected]>
AuthorDate: Fri Sep 4 08:13:33 2026 +0000

    [Fix][Zeta] Stabilize CI test isolation (#12050)
    
    Co-authored-by: DanielLeens <[email protected]>
    Co-authored-by: davidzollo <[email protected]>
    Co-authored-by: Claude Sonnet 5 <[email protected]>
---
 .../python/spawn_stdout_child_then_exit.py         |  4 +-
 .../azurecosmosdb/AbstractAzureCosmosDBIT.java     | 33 +++++++++++++--
 .../seatunnel/cdc/postgres/PostgresCDCIT.java      | 13 +++++-
 .../e2e/connector/couchbase/CouchbaseIT.java       |  4 ++
 .../e2e/connector/databend/DatabendIT.java         |  2 +
 .../elasticsearch/ElasticsearchAuthIT.java         | 48 +++++++++++++++++++---
 .../e2e/connector/hugegraph/HugeGraphSourceIT.java | 12 +++++-
 .../e2e/connector/v2/milvus/MilvusIT.java          | 45 ++++++++++++++++----
 .../e2e/connector/rocketmq/RocketMqContainer.java  |  2 +-
 .../e2e/connector/rocketmq/RocketMqIT.java         | 44 +++++++++++++++++++-
 .../src/test/resources/hazelcast-client.yaml       |  6 +--
 .../src/test/resources/hazelcast.yaml              |  4 +-
 .../engine/server/CoordinatorServiceTest.java      | 16 +++++++-
 .../engine/server/rest/BaseServletTest.java        |  2 +-
 14 files changed, 206 insertions(+), 29 deletions(-)

diff --git 
a/seatunnel-connectors-v2/connector-python/src/test/resources/python/spawn_stdout_child_then_exit.py
 
b/seatunnel-connectors-v2/connector-python/src/test/resources/python/spawn_stdout_child_then_exit.py
index e544a49fdf..fb03de1a38 100644
--- 
a/seatunnel-connectors-v2/connector-python/src/test/resources/python/spawn_stdout_child_then_exit.py
+++ 
b/seatunnel-connectors-v2/connector-python/src/test/resources/python/spawn_stdout_child_then_exit.py
@@ -22,7 +22,9 @@ import tempfile
 def main():
     sys.stdin.readline()
     subprocess.Popen(
-        [sys.executable, "-c", "import time; time.sleep(10)"], 
cwd=tempfile.gettempdir()
+        [sys.executable, "-c", "import time; time.sleep(10)"],
+        close_fds=False,
+        cwd=tempfile.gettempdir(),
     )
     print("1,python_1", flush=True)
 
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-azurecosmosdb-e2e/src/test/java/org/apache/seatunnel/e2e/connector/azurecosmosdb/AbstractAzureCosmosDBIT.java
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-azurecosmosdb-e2e/src/test/java/org/apache/seatunnel/e2e/connector/azurecosmosdb/AbstractAzureCosmosDBIT.java
index 0c89ff8149..325415130e 100644
--- 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-azurecosmosdb-e2e/src/test/java/org/apache/seatunnel/e2e/connector/azurecosmosdb/AbstractAzureCosmosDBIT.java
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-azurecosmosdb-e2e/src/test/java/org/apache/seatunnel/e2e/connector/azurecosmosdb/AbstractAzureCosmosDBIT.java
@@ -32,6 +32,7 @@ import 
org.apache.seatunnel.connectors.seatunnel.azurecosmosdb.source.AzureCosmo
 import org.apache.seatunnel.e2e.common.TestResource;
 import org.apache.seatunnel.e2e.common.TestSuiteBase;
 
+import org.awaitility.Awaitility;
 import org.junit.jupiter.api.AfterAll;
 import org.junit.jupiter.api.BeforeAll;
 import org.slf4j.Logger;
@@ -54,6 +55,7 @@ import java.util.Collections;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
+import java.util.concurrent.TimeUnit;
 
 public abstract class AbstractAzureCosmosDBIT extends TestSuiteBase implements 
TestResource {
 
@@ -84,13 +86,13 @@ public abstract class AbstractAzureCosmosDBIT extends 
TestSuiteBase implements T
                         .endpointDiscoveryEnabled(false)
                         .gatewayMode()
                         .buildClient();
-        seedContainer(BASIC_CONTAINER, item("1", "alpha", 10), item("2", 
"beta", 20));
-        seedContainer(
+        seedContainerWhenReady(BASIC_CONTAINER, item("1", "alpha", 10), 
item("2", "beta", 20));
+        seedContainerWhenReady(
                 FILTER_CONTAINER,
                 item("1", "low-score", 5),
                 item("2", "high-score", 30),
                 item("3", "higher-score", 40));
-        seedContainer(
+        seedContainerWhenReady(
                 PAGINATION_CONTAINER,
                 item("1", "page-one", 10),
                 item("2", "page-two", 20),
@@ -221,6 +223,31 @@ public abstract class AbstractAzureCosmosDBIT extends 
TestSuiteBase implements T
         System.setProperty("javax.net.ssl.trustStorePassword", 
TRUST_STORE_PASSWORD);
     }
 
+    /**
+     * Retries idempotent seeding until the emulator accepts data-plane 
requests.
+     *
+     * <p>The emulator's container health endpoint can be available before its 
database service
+     * finishes accepting collection creation requests. Beyond that initial 
gap, the emulator's
+     * gateway has also been observed (see CI run 33715122800, job 
"all-connectors-it-4") to stall
+     * on an individual request for the client's full {@code 
httpNetworkRequestTimeout} (1 minute)
+     * while it continues initializing in the background, even after an 
earlier request already
+     * succeeded. A 2-minute ceiling only allows one such stall to be absorbed 
before the retry
+     * budget is exhausted, so this is widened to 5 minutes to reliably ride 
out that warm-up
+     * window, matching the more generous ceiling already used for the 
similarly slow-starting
+     * Milvus readiness check in this same test module family.
+     *
+     * @param containerName container to initialize
+     * @param items rows to persist in the container
+     */
+    @SafeVarargs
+    private final void seedContainerWhenReady(String containerName, 
Map<String, Object>... items) {
+        Awaitility.await()
+                .ignoreExceptions()
+                .pollInterval(1, TimeUnit.SECONDS)
+                .atMost(5, TimeUnit.MINUTES)
+                .untilAsserted(() -> seedContainer(containerName, items));
+    }
+
     @SafeVarargs
     private final void seedContainer(String containerName, Map<String, 
Object>... items) {
         client.createDatabaseIfNotExists(DATABASE);
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cdc-postgres-e2e/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/PostgresCDCIT.java
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cdc-postgres-e2e/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/PostgresCDCIT.java
index f2cb2375d3..a8963b7ec5 100644
--- 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cdc-postgres-e2e/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/PostgresCDCIT.java
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cdc-postgres-e2e/src/test/java/org/apache/seatunnel/connectors/seatunnel/cdc/postgres/PostgresCDCIT.java
@@ -592,7 +592,14 @@ public class PostgresCDCIT extends TestSuiteBase 
implements TestResource {
             Assertions.assertEquals(
                     0, 
container.savepointJob(String.valueOf(committedOffsetJobId)).getExitCode());
             committedOffsetJob.get(30, TimeUnit.SECONDS);
-            insertSourceTableRow(POSTGRESQL_SCHEMA, SOURCE_TABLE_1, 15);
+            // The completed job can retain the replication slot briefly after 
savepoint creation.
+            // Wait for that connection to close so the next active state 
belongs to the restored
+            // job.
+            await().atMost(30000, TimeUnit.MILLISECONDS)
+                    .untilAsserted(
+                            () ->
+                                    Assertions.assertFalse(
+                                            
isReplicationSlotActive(committedSlotName)));
 
             committedOffsetJob =
                     CompletableFuture.runAsync(
@@ -610,6 +617,10 @@ public class PostgresCDCIT extends TestSuiteBase 
implements TestResource {
             CompletableFuture<Void> restoredCommittedOffsetJob = 
committedOffsetJob;
             // Restoring the checkpoint and reconnecting the existing 
replication slot can take
             // longer on shared GitHub runners than the initial CDC startup.
+            waitForReplicationSlotActive(committedSlotName);
+            // Insert only after the restored replication connection is active 
so this CDC record
+            // is not written before the restored slot can consume it.
+            insertSourceTableRow(POSTGRESQL_SCHEMA, SOURCE_TABLE_1, 15);
             await().atMost(RESTORE_ASSERT_TIMEOUT_MILLIS, 
TimeUnit.MILLISECONDS)
                     .untilAsserted(
                             () -> {
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-couchbase-e2e/src/test/java/org/apache/seatunnel/e2e/connector/couchbase/CouchbaseIT.java
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-couchbase-e2e/src/test/java/org/apache/seatunnel/e2e/connector/couchbase/CouchbaseIT.java
index bdf5bb6dd8..b9be68fc96 100644
--- 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-couchbase-e2e/src/test/java/org/apache/seatunnel/e2e/connector/couchbase/CouchbaseIT.java
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-couchbase-e2e/src/test/java/org/apache/seatunnel/e2e/connector/couchbase/CouchbaseIT.java
@@ -116,6 +116,10 @@ public class CouchbaseIT extends TestSuiteBase implements 
TestResource {
                         couchbaseContainer.getUsername(),
                         couchbaseContainer.getPassword());
 
+        // The container's HTTP bootstrap may finish before the KV service 
accepts authenticated
+        // client connections. Confirm KV readiness before a SeaTunnel job 
uses the Docker alias.
+        cluster.bucket(COUCHBASE_BUCKET).waitUntilReady(Duration.ofMinutes(2));
+
         // Wait for the query/management service to be ready before issuing 
DDL. The container
         // signals readiness at the KV/bucket level but the query and index 
services can still
         // reject requests for a short window after Cluster.connect() returns. 
Wrapping the DDL
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-databend-e2e/src/test/java/org/apache/seatunnel/e2e/connector/databend/DatabendIT.java
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-databend-e2e/src/test/java/org/apache/seatunnel/e2e/connector/databend/DatabendIT.java
index 9e6cdc198f..04cc0b34d4 100644
--- 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-databend-e2e/src/test/java/org/apache/seatunnel/e2e/connector/databend/DatabendIT.java
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-databend-e2e/src/test/java/org/apache/seatunnel/e2e/connector/databend/DatabendIT.java
@@ -259,6 +259,8 @@ public class DatabendIT extends TestSuiteBase implements 
TestResource {
                 new DatabendContainer(DATABEND_DOCKER_IMAGE)
                         .withNetwork(NETWORK)
                         .withNetworkAliases(DATABEND_CONTAINER_HOST)
+                        // The nightly image can reset JDBC while its object 
storage becomes ready.
+                        .withStartupAttempts(3)
                         .withUsername("root")
                         .withPassword("")
                         .withEnv("STORAGE_TYPE", "s3")
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-elasticsearch-e2e/src/test/java/org/apache/seatunnel/e2e/connector/elasticsearch/ElasticsearchAuthIT.java
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-elasticsearch-e2e/src/test/java/org/apache/seatunnel/e2e/connector/elasticsearch/ElasticsearchAuthIT.java
index 4a8bf2c822..8a575ef945 100644
--- 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-elasticsearch-e2e/src/test/java/org/apache/seatunnel/e2e/connector/elasticsearch/ElasticsearchAuthIT.java
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-elasticsearch-e2e/src/test/java/org/apache/seatunnel/e2e/connector/elasticsearch/ElasticsearchAuthIT.java
@@ -41,6 +41,7 @@ import org.apache.http.impl.client.HttpClients;
 import org.apache.http.ssl.SSLContextBuilder;
 import org.apache.http.util.EntityUtils;
 
+import org.awaitility.Awaitility;
 import org.junit.jupiter.api.AfterEach;
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.BeforeEach;
@@ -72,6 +73,9 @@ public class ElasticsearchAuthIT extends TestSuiteBase 
implements TestResource {
 
     private static final String ELASTICSEARCH_IMAGE = "elasticsearch:8.9.0";
     private static final long INDEX_REFRESH_DELAY = 2000L;
+    // Retry only the transient 503 returned while Elasticsearch initializes 
its security index.
+    private static final int API_KEY_CREATION_MAX_ATTEMPTS = 15;
+    private static final long API_KEY_CREATION_RETRY_DELAY_SECONDS = 2L;
 
     // Test data constants
     private static final String TEST_INDEX = "auth_test_index";
@@ -167,7 +171,7 @@ public class ElasticsearchAuthIT extends TestSuiteBase 
implements TestResource {
     /** Wait for Elasticsearch to be ready */
     private void waitForElasticsearchReady() throws IOException, 
InterruptedException {
         String elasticsearchUrl = "https://"; + 
elasticsearchContainer.getHttpHostAddress();
-        String healthUrl = elasticsearchUrl + "/_cluster/health";
+        String healthUrl = elasticsearchUrl + 
"/_cluster/health?wait_for_status=yellow&timeout=30s";
 
         log.info("Waiting for Elasticsearch to be ready at: {}", healthUrl);
 
@@ -197,7 +201,7 @@ public class ElasticsearchAuthIT extends TestSuiteBase 
implements TestResource {
     }
 
     /** Create real API keys using Elasticsearch API */
-    private void createRealApiKeys() throws IOException {
+    private void createRealApiKeys() {
         String elasticsearchUrl = "https://"; + 
elasticsearchContainer.getHttpHostAddress();
         String apiKeyUrl = elasticsearchUrl + "/_security/api_key";
 
@@ -225,20 +229,53 @@ public class ElasticsearchAuthIT extends TestSuiteBase 
implements TestResource {
                         + "  }\n"
                         + "}";
 
-        HttpPost request = new HttpPost(apiKeyUrl);
         String auth =
                 Base64.getEncoder()
                         .encodeToString(
                                 (VALID_USERNAME + ":" + VALID_PASSWORD)
                                         .getBytes(StandardCharsets.UTF_8));
+
+        // Retry budget matches the previous hand-rolled loop (15 attempts, 2s 
apart, <=30s
+        // total). Using Awaitility here keeps this in line with the 
readiness-wait idiom used
+        // throughout this test module family instead of a bespoke sleep loop.
+        Awaitility.await()
+                
.pollInterval(Duration.ofSeconds(API_KEY_CREATION_RETRY_DELAY_SECONDS))
+                .atMost(
+                        Duration.ofSeconds(
+                                API_KEY_CREATION_MAX_ATTEMPTS
+                                        * 
API_KEY_CREATION_RETRY_DELAY_SECONDS))
+                .untilAsserted(() -> attemptCreateApiKey(apiKeyUrl, 
requestBody, auth));
+    }
+
+    /**
+     * Performs a single API key creation attempt and verifies it on success.
+     *
+     * <p>Elasticsearch answers 503 while its internal security index is still 
initializing after
+     * container startup. That case is raised as an assertion failure so the 
caller's Awaitility
+     * loop treats it as "not ready yet" and retries. Any other non-200 status 
is a genuine,
+     * non-retryable failure and is thrown immediately as a plain {@link 
RuntimeException}, which
+     * Awaitility does not retry on.
+     *
+     * @param apiKeyUrl the Elasticsearch API key creation endpoint
+     * @param requestBody the JSON payload describing the key's role and 
metadata
+     * @param auth the pre-encoded HTTP Basic authorization header value
+     */
+    private void attemptCreateApiKey(String apiKeyUrl, String requestBody, 
String auth)
+            throws IOException {
+        HttpPost request = new HttpPost(apiKeyUrl);
         request.setHeader("Authorization", "Basic " + auth);
         request.setHeader("Content-Type", "application/json");
         request.setEntity(new StringEntity(requestBody, 
StandardCharsets.UTF_8));
 
         HttpResponse response = httpClient.execute(request);
         String responseBody = EntityUtils.toString(response.getEntity());
-
-        if (response.getStatusLine().getStatusCode() != 200) {
+        int statusCode = response.getStatusLine().getStatusCode();
+        if (statusCode == 503) {
+            log.info("Elasticsearch security index is not ready yet; retrying 
API key creation");
+        }
+        Assertions.assertNotEquals(
+                503, statusCode, "Elasticsearch security index is not ready 
yet");
+        if (statusCode != 200) {
             throw new RuntimeException("Failed to create API key: " + 
responseBody);
         }
 
@@ -261,7 +298,6 @@ public class ElasticsearchAuthIT extends TestSuiteBase 
implements TestResource {
 
             // Verify the API key works
             verifyApiKey();
-
         } catch (Exception e) {
             throw new RuntimeException("Failed to parse API key response: " + 
responseBody, e);
         }
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-hugegraph-e2e/src/test/java/org/apache/seatunnel/e2e/connector/hugegraph/HugeGraphSourceIT.java
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-hugegraph-e2e/src/test/java/org/apache/seatunnel/e2e/connector/hugegraph/HugeGraphSourceIT.java
index f960e170ab..9d8d5c58f9 100644
--- 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-hugegraph-e2e/src/test/java/org/apache/seatunnel/e2e/connector/hugegraph/HugeGraphSourceIT.java
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-hugegraph-e2e/src/test/java/org/apache/seatunnel/e2e/connector/hugegraph/HugeGraphSourceIT.java
@@ -27,6 +27,7 @@ import org.apache.hugegraph.structure.constant.IdStrategy;
 import org.apache.hugegraph.structure.graph.Edge;
 import org.apache.hugegraph.structure.graph.Vertex;
 
+import org.awaitility.Awaitility;
 import org.junit.jupiter.api.AfterAll;
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.BeforeAll;
@@ -231,7 +232,16 @@ public class HugeGraphSourceIT extends TestSuiteBase 
implements TestResource {
     }
 
     private void clearGraphWithoutSchema() {
-        hugeClient.graphs().clearGraph(GRAPH_NAME, "I'm sure to delete all 
data");
+        // The server can close its REST connection while completing a 
previous graph mutation.
+        Awaitility.given()
+                .ignoreExceptions()
+                .pollInterval(Duration.ofSeconds(2))
+                .atMost(Duration.ofMinutes(2))
+                .untilAsserted(
+                        () ->
+                                hugeClient
+                                        .graphs()
+                                        .clearGraph(GRAPH_NAME, "I'm sure to 
delete all data"));
     }
 
     private void setupSchema() {
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-milvus-e2e/src/test/java/org/apache/seatunnel/e2e/connector/v2/milvus/MilvusIT.java
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-milvus-e2e/src/test/java/org/apache/seatunnel/e2e/connector/v2/milvus/MilvusIT.java
index dd934c263d..4f72d01ea5 100644
--- 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-milvus-e2e/src/test/java/org/apache/seatunnel/e2e/connector/v2/milvus/MilvusIT.java
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-milvus-e2e/src/test/java/org/apache/seatunnel/e2e/connector/v2/milvus/MilvusIT.java
@@ -135,11 +135,47 @@ public class MilvusIT extends TestSuiteBase implements 
TestResource {
         Startables.deepStart(Stream.of(this.container)).join();
         log.info("Milvus host is {}", container.getHost());
         log.info("Milvus container started");
-        Awaitility.given().ignoreExceptions().await().atMost(720L, 
TimeUnit.SECONDS);
+        waitForMilvusReady();
         this.initMilvus();
         this.initSourceData();
     }
 
+    /**
+     * Waits until the Milvus Proxy can serve a real catalog RPC.
+     *
+     * <p>The container health endpoint becomes available before the Proxy 
accepts gRPC requests.
+     */
+    private void waitForMilvusReady() {
+        Awaitility.await()
+                .ignoreExceptions()
+                .pollInterval(1, TimeUnit.SECONDS)
+                .atMost(720L, TimeUnit.SECONDS)
+                .untilAsserted(
+                        () -> {
+                            MilvusServiceClient readinessClient = 
createMilvusClient();
+                            try {
+                                R<ListDatabasesResponse> response = 
readinessClient.listDatabases();
+                                Assertions.assertEquals(
+                                        R.Status.Success.getCode(), 
response.getStatus());
+                            } finally {
+                                readinessClient.close();
+                            }
+                        });
+    }
+
+    /**
+     * Creates a client bound to the current Milvus test container.
+     *
+     * @return client configured with the test endpoint and token
+     */
+    private MilvusServiceClient createMilvusClient() {
+        return new MilvusServiceClient(
+                ConnectParam.newBuilder()
+                        .withUri(this.container.getEndpoint())
+                        .withToken(TOKEN)
+                        .build());
+    }
+
     private void initMilvus()
             throws SQLException, ClassNotFoundException, 
InstantiationException,
                     IllegalAccessException {
@@ -149,12 +185,7 @@ public class MilvusIT extends TestSuiteBase implements 
TestResource {
         ReadonlyConfig readonlyConfig = ReadonlyConfig.fromMap(config);
         catalog = new MilvusCatalog(COLLECTION_NAME, readonlyConfig);
         catalog.open();
-        milvusClient =
-                new MilvusServiceClient(
-                        ConnectParam.newBuilder()
-                                .withUri(this.container.getEndpoint())
-                                .withToken(TOKEN)
-                                .build());
+        milvusClient = createMilvusClient();
     }
 
     private void initSourceData() {
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-rocketmq-e2e/src/test/java/org/apache/seatunnel/e2e/connector/rocketmq/RocketMqContainer.java
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-rocketmq-e2e/src/test/java/org/apache/seatunnel/e2e/connector/rocketmq/RocketMqContainer.java
index f3929d6551..597a068c4f 100644
--- 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-rocketmq-e2e/src/test/java/org/apache/seatunnel/e2e/connector/rocketmq/RocketMqContainer.java
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-rocketmq-e2e/src/test/java/org/apache/seatunnel/e2e/connector/rocketmq/RocketMqContainer.java
@@ -38,7 +38,7 @@ public class RocketMqContainer extends 
GenericContainer<RocketMqContainer> {
     public static final int BROKER_PORT = 10911;
     public static final String BROKER_NAME = "broker-a";
     private static final int DEFAULT_BROKER_PERMISSION = 6;
-    private static final int DEFAULT_TOPIC_QUEUE_NUMS = 1;
+    static final int DEFAULT_TOPIC_QUEUE_NUMS = 4;
 
     public RocketMqContainer(DockerImageName image) {
         super(image);
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-rocketmq-e2e/src/test/java/org/apache/seatunnel/e2e/connector/rocketmq/RocketMqIT.java
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-rocketmq-e2e/src/test/java/org/apache/seatunnel/e2e/connector/rocketmq/RocketMqIT.java
index ef62a04d68..9eb963417d 100644
--- 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-rocketmq-e2e/src/test/java/org/apache/seatunnel/e2e/connector/rocketmq/RocketMqIT.java
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-rocketmq-e2e/src/test/java/org/apache/seatunnel/e2e/connector/rocketmq/RocketMqIT.java
@@ -53,6 +53,7 @@ import org.apache.rocketmq.common.message.MessageExt;
 import org.apache.rocketmq.common.message.MessageQueue;
 import org.apache.rocketmq.common.protocol.route.QueueData;
 import org.apache.rocketmq.common.protocol.route.TopicRouteData;
+import org.apache.rocketmq.common.topic.TopicValidator;
 import org.apache.rocketmq.remoting.exception.RemotingException;
 import org.apache.rocketmq.remoting.protocol.LanguageCode;
 import org.apache.rocketmq.tools.admin.DefaultMQAdminExt;
@@ -152,6 +153,12 @@ public class RocketMqIT extends TestSuiteBase implements 
TestResource {
         rocketMqContainer.start();
         log.info("RocketMq container started");
         initProducer();
+        // Unlike the other topics in this file, test_topic_source is written 
directly via
+        // producer.send(Message, MessageQueue) in generateTestData(), which 
bypasses the normal
+        // route-resolution path a plain send(Message) would use. Establish 
and confirm the route
+        // up front so the name server has already published it before any 
source job (started by
+        // a later @TestTemplate method, sometimes minutes after this write) 
queries it.
+        waitForTopicRoute("test_topic_source");
         log.info("Write 100 records to topic test_topic_source");
         DefaultSeaTunnelRowSerializer serializer =
                 new DefaultSeaTunnelRowSerializer(
@@ -187,6 +194,7 @@ public class RocketMqIT extends TestSuiteBase implements 
TestResource {
 
     @TestTemplate
     public void testSinkRocketMq(TestContainer container) throws IOException, 
InterruptedException {
+        waitForTopicRoute("test_topic");
 
         Container.ExecResult execResult =
                 container.executeJob("/rocketmq-sink_fake_to_rocketmq.conf");
@@ -205,6 +213,8 @@ public class RocketMqIT extends TestSuiteBase implements 
TestResource {
     @TestTemplate
     public void testTextFormatSinkRocketMq(TestContainer container)
             throws IOException, InterruptedException {
+        waitForTopicRoute("test_text_topic");
+
         Container.ExecResult execResult =
                 
container.executeJob("/rocketmq-text-sink_fake_to_rocketmq.conf");
         Assertions.assertEquals(0, execResult.getExitCode(), 
execResult.getStderr());
@@ -457,6 +467,28 @@ public class RocketMqIT extends TestSuiteBase implements 
TestResource {
         }
     }
 
+    /**
+     * Waits for the name server to expose a topic route to the producer.
+     *
+     * @param topic topic used by the next job submission or restored job
+     */
+    private void waitForTopicRoute(String topic) {
+        Awaitility.await()
+                .ignoreExceptions()
+                .atMost(1, TimeUnit.MINUTES)
+                .pollInterval(1, TimeUnit.SECONDS)
+                .untilAsserted(
+                        () -> {
+                            producer.createTopic(
+                                    TopicValidator.AUTO_CREATE_TOPIC_KEY_TOPIC,
+                                    topic,
+                                    
RocketMqContainer.DEFAULT_TOPIC_QUEUE_NUMS);
+                            Assertions.assertFalse(
+                                    
producer.fetchPublishMessageQueues(topic).isEmpty(),
+                                    "Topic route is not ready: " + topic);
+                        });
+    }
+
     private Map<String, RocketMqConsumerMessage> 
getRocketMqConsumerData(String topicName) {
         Map<String, RocketMqConsumerMessage> data = new HashMap<>();
         Map<MessageQueue, Long> consumedOffsets = new HashMap<>();
@@ -531,9 +563,15 @@ public class RocketMqIT extends TestSuiteBase implements 
TestResource {
         if (consumedOffsets.isEmpty()) {
             return;
         }
+        // consumedOffsets now spans up to DEFAULT_TOPIC_QUEUE_NUMS (4) queues 
instead of the
+        // single queue this topic previously had, so broker-side 
offset-commit visibility for
+        // every queue must converge within the window, not just one. Seen to 
exceed the previous
+        // 30s ceiling under real CI load (fork run 33729027165, job
+        // "rocketmq-connector-it (8, ubuntu-latest)": 
ConditionTimeoutException on this exact
+        // assertion), even though the sibling JDK 11 leg of the same run 
passed cleanly.
         Awaitility.await()
                 .ignoreExceptions()
-                .atMost(30, TimeUnit.SECONDS)
+                .atMost(60, TimeUnit.SECONDS)
                 .pollInterval(1, TimeUnit.SECONDS)
                 .untilAsserted(
                         () -> {
@@ -629,6 +667,7 @@ public class RocketMqIT extends TestSuiteBase implements 
TestResource {
                     "consumerGroup=" + consumerGroup
                 };
 
+        waitForTopicRoute(sourceTopic);
         for (int i = 0; i < 20; i++) {
             Message msg = new Message(sourceTopic, (payload + "_initial_" + 
i).getBytes());
             producer.send(msg, new MessageQueue(sourceTopic, 
RocketMqContainer.BROKER_NAME, 0));
@@ -705,6 +744,9 @@ public class RocketMqIT extends TestSuiteBase implements 
TestResource {
                 "Source end offset should advance by at least 25, actual: "
                         + (srcEndAfterAll - srcEndBeforeStart));
 
+        // The name server can briefly drop an auto-created topic route while 
the job is stopped
+        // for a savepoint. Restore only after the dynamic source topic is 
visible again.
+        waitForTopicRoute(sourceTopic);
         CompletableFuture.runAsync(
                 () -> {
                     try {
diff --git 
a/seatunnel-engine/seatunnel-engine-client/src/test/resources/hazelcast-client.yaml
 
b/seatunnel-engine/seatunnel-engine-client/src/test/resources/hazelcast-client.yaml
index 90631fe056..1e9feb5948 100644
--- 
a/seatunnel-engine/seatunnel-engine-client/src/test/resources/hazelcast-client.yaml
+++ 
b/seatunnel-engine/seatunnel-engine-client/src/test/resources/hazelcast-client.yaml
@@ -20,6 +20,6 @@ hazelcast-client:
 
   network:
     cluster-members:
-      - localhost:5801
-      - localhost:5802
-      - localhost:5803
+      - 127.0.0.1:5801
+      - 127.0.0.1:5802
+      - 127.0.0.1:5803
diff --git 
a/seatunnel-engine/seatunnel-engine-client/src/test/resources/hazelcast.yaml 
b/seatunnel-engine/seatunnel-engine-client/src/test/resources/hazelcast.yaml
index 39eca27d9e..453538bcc4 100644
--- a/seatunnel-engine/seatunnel-engine-client/src/test/resources/hazelcast.yaml
+++ b/seatunnel-engine/seatunnel-engine-client/src/test/resources/hazelcast.yaml
@@ -22,7 +22,7 @@ hazelcast:
       tcp-ip:
         enabled: true
         member-list:
-          - localhost
+          - 127.0.0.1
     port:
       auto-increment: true
       port-count: 10
@@ -40,4 +40,4 @@ hazelcast:
     hazelcast.invocation.max.retry.count: 20
     hazelcast.tcp.join.port.try.count: 30
     hazelcast.logging.type: log4j2
-    hazelcast.operation.generic.thread.count: 200
\ No newline at end of file
+    hazelcast.operation.generic.thread.count: 200
diff --git 
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/CoordinatorServiceTest.java
 
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/CoordinatorServiceTest.java
index da1b2ae928..3be7deafdf 100644
--- 
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/CoordinatorServiceTest.java
+++ 
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/CoordinatorServiceTest.java
@@ -759,10 +759,22 @@ public class CoordinatorServiceTest {
     private void stopCoordinatorSchedulers(CoordinatorService 
coordinatorService) {
         ReflectionUtils.getField(coordinatorService, "masterActiveListener")
                 .map(ScheduledExecutorService.class::cast)
-                .ifPresent(ScheduledExecutorService::shutdownNow);
+                .ifPresent(this::shutdownScheduler);
         ReflectionUtils.getField(coordinatorService, 
"pipelineCleanupScheduler")
                 .map(ScheduledExecutorService.class::cast)
-                .ifPresent(ScheduledExecutorService::shutdownNow);
+                .ifPresent(this::shutdownScheduler);
+    }
+
+    private void shutdownScheduler(ScheduledExecutorService scheduler) {
+        scheduler.shutdownNow();
+        try {
+            if (!scheduler.awaitTermination(5, TimeUnit.SECONDS)) {
+                throw new AssertionError("Scheduler did not stop in time");
+            }
+        } catch (InterruptedException e) {
+            Thread.currentThread().interrupt();
+            throw new AssertionError("Interrupted while stopping scheduler", 
e);
+        }
     }
 
     private void invokePendingJobScheduler(CoordinatorService 
coordinatorService) throws Exception {
diff --git 
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/BaseServletTest.java
 
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/BaseServletTest.java
index 3d3d734cfc..f64258de86 100644
--- 
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/BaseServletTest.java
+++ 
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/BaseServletTest.java
@@ -42,7 +42,7 @@ import java.util.Collections;
 
 class BaseServletTest extends AbstractSeaTunnelServerTest {
 
-    private static final int HTTP_PORT = 18080;
+    private static final int HTTP_PORT = TestUtils.getAvailablePort();
 
     private static final Long JOB_1 = System.currentTimeMillis() + 1L;
 

Reply via email to