This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/dev/pr-12050-98182ca59c93c10390717af4a49b8806f5c555f1 in repository https://gitbox.apache.org/repos/asf/seatunnel.git
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;
