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-12396-36401ab4bbb652c02cd721aa70f53c8ef4fe959a in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit 5f7f5c1dbbb490436605cef37b47a5b30d5acdd1 Author: Sepuri Sai Krishna <[email protected]> AuthorDate: Sun Sep 20 02:18:38 2026 +0000 [Test][E2E] Isolate the RocketMqIT tag tests with unique topic and consumer group (#12396) --- .../seatunnel/e2e/connector/rocketmq/RocketMqIT.java | 16 ++++++++++++---- .../rocketmq-source_text_error_tag_to_console.conf | 3 ++- .../resources/rocketmq-source_text_tag_to_console.conf | 3 ++- 3 files changed, 16 insertions(+), 6 deletions(-) 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 337b6278ba..8bbcd2d986 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 @@ -217,7 +217,9 @@ public class RocketMqIT extends TestSuiteBase implements TestResource { @TestTemplate public void testSourceRocketMqTextTagToConsole(TestContainer container) throws IOException, InterruptedException { - String topic = "test_topic_text_tag"; + final String uniqueSuffix = uniqueTestSuffix(); + final String topic = "test_topic_text_tag_" + uniqueSuffix; + final String consumerGroup = "SeaTunnel-Consumer-Group-" + uniqueSuffix; String tag = "tag_test"; DefaultSeaTunnelRowSerializer serializer = @@ -225,14 +227,18 @@ public class RocketMqIT extends TestSuiteBase implements TestResource { topic, tag, SEATUNNEL_ROW_TYPE, SchemaFormat.TEXT, DEFAULT_FIELD_DELIMITER); generateTestData(serializer::serializeRow, topic, 0, 32); Container.ExecResult execResult = - container.executeJob("/rocketmq-source_text_tag_to_console.conf"); + container.executeJob( + "/rocketmq-source_text_tag_to_console.conf", + Arrays.asList("sourceTopic=" + topic, "consumerGroup=" + consumerGroup)); Assertions.assertEquals(0, execResult.getExitCode(), execResult.getStderr()); } @TestTemplate public void testSourceRocketMqTextErrorTagToConsole(TestContainer container) throws IOException, InterruptedException { - String topic = "test_topic_text_error_tag"; + final String uniqueSuffix = uniqueTestSuffix(); + final String topic = "test_topic_text_error_tag_" + uniqueSuffix; + final String consumerGroup = "SeaTunnel-Consumer-Group-" + uniqueSuffix; String tag = "test_error_tag"; DefaultSeaTunnelRowSerializer serializer = @@ -240,7 +246,9 @@ public class RocketMqIT extends TestSuiteBase implements TestResource { topic, tag, SEATUNNEL_ROW_TYPE, SchemaFormat.TEXT, DEFAULT_FIELD_DELIMITER); generateTestData(serializer::serializeRow, topic, 0, 32); Container.ExecResult execResult = - container.executeJob("/rocketmq-source_text_error_tag_to_console.conf"); + container.executeJob( + "/rocketmq-source_text_error_tag_to_console.conf", + Arrays.asList("sourceTopic=" + topic, "consumerGroup=" + consumerGroup)); Assertions.assertEquals(0, execResult.getExitCode(), execResult.getStderr()); } diff --git a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-rocketmq-e2e/src/test/resources/rocketmq-source_text_error_tag_to_console.conf b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-rocketmq-e2e/src/test/resources/rocketmq-source_text_error_tag_to_console.conf index 71c06a181e..1221629d74 100644 --- a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-rocketmq-e2e/src/test/resources/rocketmq-source_text_error_tag_to_console.conf +++ b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-rocketmq-e2e/src/test/resources/rocketmq-source_text_error_tag_to_console.conf @@ -34,7 +34,8 @@ source { Rocketmq { plugin_output = "rocketmq_table" name.srv.addr = "rocketmq-e2e:9876" - topics = "test_topic_text_error_tag" + topics = "${sourceTopic:test_topic_text_error_tag}" + consumer.group = "${consumerGroup:SeaTunnel-Consumer-Group}" format = text # The default field delimiter is "," field_delimiter = "," diff --git a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-rocketmq-e2e/src/test/resources/rocketmq-source_text_tag_to_console.conf b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-rocketmq-e2e/src/test/resources/rocketmq-source_text_tag_to_console.conf index e69aad1614..a500b4cd80 100644 --- a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-rocketmq-e2e/src/test/resources/rocketmq-source_text_tag_to_console.conf +++ b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-rocketmq-e2e/src/test/resources/rocketmq-source_text_tag_to_console.conf @@ -34,7 +34,8 @@ source { Rocketmq { plugin_output = "rocketmq_table" name.srv.addr = "rocketmq-e2e:9876" - topics = "test_topic_text_tag" + topics = "${sourceTopic:test_topic_text_tag}" + consumer.group = "${consumerGroup:SeaTunnel-Consumer-Group}" format = text # The default field delimiter is "," field_delimiter = ","
