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-12393-9b4814726cd50d2c72557b8ba6e40b8c4cc16cba
in repository https://gitbox.apache.org/repos/asf/seatunnel.git

commit 2f92ff7559dcc782427926756baf3614581c4b2d
Author: Sepuri Sai Krishna <[email protected]>
AuthorDate: Sat Sep 19 14:42:17 2026 +0000

    [Test][E2E] Remove RocketMqIT.deleteTopicIfExist, which never deleted 
anything (#12393)
---
 .../e2e/connector/rocketmq/RocketMqIT.java         | 38 ----------------------
 1 file changed, 38 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 41b8af5cde..337b6278ba 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
@@ -51,12 +51,9 @@ import org.apache.rocketmq.common.admin.TopicOffset;
 import org.apache.rocketmq.common.message.Message;
 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;
 
 import org.awaitility.Awaitility;
 import org.junit.jupiter.api.AfterAll;
@@ -91,7 +88,6 @@ import java.util.concurrent.CompletableFuture;
 import java.util.concurrent.ExecutionException;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicLong;
-import java.util.stream.Collectors;
 
 @Slf4j
 public class RocketMqIT extends TestSuiteBase implements TestResource {
@@ -224,9 +220,6 @@ public class RocketMqIT extends TestSuiteBase implements 
TestResource {
         String topic = "test_topic_text_tag";
         String tag = "tag_test";
 
-        // delete topic if exist
-        deleteTopicIfExist(topic);
-
         DefaultSeaTunnelRowSerializer serializer =
                 new DefaultSeaTunnelRowSerializer(
                         topic, tag, SEATUNNEL_ROW_TYPE, SchemaFormat.TEXT, 
DEFAULT_FIELD_DELIMITER);
@@ -242,9 +235,6 @@ public class RocketMqIT extends TestSuiteBase implements 
TestResource {
         String topic = "test_topic_text_error_tag";
         String tag = "test_error_tag";
 
-        // delete topic if exist
-        deleteTopicIfExist(topic);
-
         DefaultSeaTunnelRowSerializer serializer =
                 new DefaultSeaTunnelRowSerializer(
                         topic, tag, SEATUNNEL_ROW_TYPE, SchemaFormat.TEXT, 
DEFAULT_FIELD_DELIMITER);
@@ -957,32 +947,4 @@ public class RocketMqIT extends TestSuiteBase implements 
TestResource {
     private String uniqueTestSuffix() {
         return UUID.randomUUID().toString().replace("-", "");
     }
-
-    private void deleteTopicIfExist(String topicName) {
-        DefaultMQAdminExt admin = new DefaultMQAdminExt();
-        admin.setInstanceName(UUID.randomUUID().toString());
-        try {
-            admin.start();
-            TopicRouteData topicRouteData = 
admin.examineTopicRouteInfo(topicName);
-            if (topicRouteData != null
-                    && topicRouteData.getQueueDatas() != null
-                    && !topicRouteData.getQueueDatas().isEmpty()) {
-                Set<String> brokerNames =
-                        topicRouteData.getQueueDatas().stream()
-                                .map(QueueData::getBrokerName)
-                                .collect(Collectors.toSet());
-                admin.deleteTopicInBroker(brokerNames, topicName);
-                admin.deleteTopicInNameServer(brokerNames, topicName, 
"delete_topic");
-                log.info("Deleted topic: {}", topicName);
-            } else {
-                log.info("Topic {} does not exist", topicName);
-            }
-        } catch (Exception e) {
-            log.warn("Failed to delete topic {}: {}", topicName, 
e.getMessage());
-        } finally {
-            if (admin != null) {
-                admin.shutdown();
-            }
-        }
-    }
 }

Reply via email to