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

davidzollo 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 f81e1360e1 [Test][E2E] Start the RocketMQ broker with a single 
identity and wait for name-server routes (#12115)
f81e1360e1 is described below

commit f81e1360e13cb8c2b9fea20d6d803004abe77c11
Author: David Zollo <[email protected]>
AuthorDate: Fri Sep 18 10:53:10 2026 -0400

    [Test][E2E] Start the RocketMQ broker with a single identity and wait for 
name-server routes (#12115)
    
    Co-authored-by: Claude Fable 5.1 <[email protected]>
---
 .../e2e/connector/rocketmq/RocketMqContainer.java  | 123 +++++++++++++++------
 .../e2e/connector/rocketmq/RocketMqIT.java         |  19 +++-
 2 files changed, 107 insertions(+), 35 deletions(-)

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 597a068c4f..d98952da58 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
@@ -18,67 +18,89 @@
 package org.apache.seatunnel.e2e.connector.rocketmq;
 
 import org.testcontainers.containers.GenericContainer;
+import org.testcontainers.images.builder.Transferable;
 import org.testcontainers.utility.DockerImageName;
 
-import com.github.dockerjava.api.command.InspectContainerResponse;
-import lombok.SneakyThrows;
-
+import java.io.IOException;
 import java.net.Inet4Address;
 import java.net.InetAddress;
 import java.net.NetworkInterface;
+import java.net.ServerSocket;
 import java.net.SocketException;
-import java.util.ArrayList;
 import java.util.Enumeration;
-import java.util.List;
 
-/** rocketmq container */
+/**
+ * RocketMQ container for the connector E2E suite.
+ *
+ * <p>The broker identity (broker name, advertised address and listen port) is 
fixed through a
+ * broker.conf that is copied into the container before the broker process 
starts, so the broker
+ * registers with the name server exactly once and under a single identity. 
The previous approach
+ * started the broker with the image defaults and renamed it afterwards 
through {@code mqadmin
+ * updateBrokerConfig}. That left the name server with two broker entries for 
the same process (the
+ * container hostname with the container-internal address plus {@code 
broker-a} with the host
+ * address). Until the stale entry was purged the default-topic route only 
referenced the hostname
+ * broker, so {@code producer.send(msg, new MessageQueue(topic, "broker-a", 
0))} failed with "The
+ * broker[broker-a] not exist" and the admin client could not see consume 
offsets of {@code
+ * broker-a}. In CI that window covered most of the test run (apache/seatunnel 
PRs #12115 and
+ * #12099, job rocketmq-connector-it).
+ *
+ * <p>RocketMQ binds and advertises the same {@code listenPort}, so a single 
identity requires the
+ * host port to equal the container port. The broker port is therefore 
published as a fixed binding
+ * on a free host port chosen at construction time, while the name server 
keeps a dynamic mapping.
+ */
 public class RocketMqContainer extends GenericContainer<RocketMqContainer> {
 
     public static final int NAMESRV_PORT = 9876;
-    public static final int BROKER_PORT = 10911;
     public static final String BROKER_NAME = "broker-a";
     private static final int DEFAULT_BROKER_PERMISSION = 6;
     static final int DEFAULT_TOPIC_QUEUE_NUMS = 4;
+    private static final String BROKER_CONF_PATH = 
"/home/rocketmq/broker.conf";
+    private static final int FREE_PORT_ATTEMPTS = 20;
+
+    /** Broker listen port inside the container, published 1:1 on the host. */
+    private final int brokerPort;
 
     public RocketMqContainer(DockerImageName image) {
         super(image);
-        withExposedPorts(NAMESRV_PORT, BROKER_PORT, BROKER_PORT - 2);
+        this.brokerPort = findFreeBrokerPort();
+        withExposedPorts(NAMESRV_PORT);
+        // listenPort is both the bind port and the advertised port, so it 
must be published 1:1.
+        // The VIP channel listens on listenPort - 2 and is published the same 
way.
+        addFixedExposedPort(brokerPort, brokerPort);
+        addFixedExposedPort(brokerPort - 2, brokerPort - 2);
         this.withEnv("JAVA_OPT_EXT", "-Xms512m -Xmx512m");
     }
 
     @Override
     protected void configure() {
+        // Written before the broker starts so its very first registration 
already carries the
+        // final name, advertised address and permissions; nothing is renamed 
at runtime.
+        String brokerConf =
+                "brokerClusterName=DefaultCluster\n"
+                        + "brokerName="
+                        + BROKER_NAME
+                        + "\n"
+                        + "brokerId=0\n"
+                        + "brokerIP1="
+                        + getLinuxLocalIp()
+                        + "\n"
+                        + "listenPort="
+                        + brokerPort
+                        + "\n"
+                        + "autoCreateTopicEnable=true\n"
+                        + "defaultTopicQueueNums="
+                        + DEFAULT_TOPIC_QUEUE_NUMS
+                        + "\n"
+                        + "brokerPermission="
+                        + DEFAULT_BROKER_PERMISSION
+                        + "\n";
+        withCopyToContainer(Transferable.of(brokerConf), BROKER_CONF_PATH);
         String command = "#!/bin/bash\n";
         command += "./mqnamesrv &\n";
-        command += "./mqbroker -n localhost:" + NAMESRV_PORT;
+        command += "./mqbroker -n localhost:" + NAMESRV_PORT + " -c " + 
BROKER_CONF_PATH;
         withCommand("sh", "-c", command);
     }
 
-    @Override
-    @SneakyThrows
-    protected void containerIsStarted(InspectContainerResponse containerInfo) {
-        List<String> updateBrokerConfigCommands = new ArrayList<>();
-        
updateBrokerConfigCommands.add(updateBrokerConfig("autoCreateTopicEnable", 
true));
-        updateBrokerConfigCommands.add(
-                updateBrokerConfig("defaultTopicQueueNums", 
DEFAULT_TOPIC_QUEUE_NUMS));
-        updateBrokerConfigCommands.add(updateBrokerConfig("brokerName", 
BROKER_NAME));
-        updateBrokerConfigCommands.add(updateBrokerConfig("brokerIP1", 
getLinuxLocalIp()));
-        updateBrokerConfigCommands.add(
-                updateBrokerConfig("listenPort", getMappedPort(BROKER_PORT)));
-        updateBrokerConfigCommands.add(
-                updateBrokerConfig("brokerPermission", 
DEFAULT_BROKER_PERMISSION));
-        final String command = String.join(" && ", updateBrokerConfigCommands);
-        ExecResult result = execInContainer("/bin/sh", "-c", command);
-        if (result != null && result.getExitCode() != 0) {
-            throw new IllegalStateException(result.toString());
-        }
-    }
-
-    private String updateBrokerConfig(final String key, final Object val) {
-        final String brokerAddr = "localhost:" + BROKER_PORT;
-        return "./mqadmin updateBrokerConfig -b " + brokerAddr + " -k " + key 
+ " -v " + val;
-    }
-
     public String getNameSrvAddr() {
         return String.format("%s:%s", getHost(), getMappedPort(NAMESRV_PORT));
     }
@@ -103,4 +125,37 @@ public class RocketMqContainer extends 
GenericContainer<RocketMqContainer> {
         }
         return ip;
     }
+
+    /**
+     * Picks a free host port whose VIP channel port (port - 2) is free as 
well, so both fixed
+     * bindings can be published.
+     */
+    private static int findFreeBrokerPort() {
+        for (int i = 0; i < FREE_PORT_ATTEMPTS; i++) {
+            int port = findFreePort();
+            if (port > 2 && isPortFree(port - 2)) {
+                return port;
+            }
+        }
+        throw new IllegalStateException(
+                "Could not find a free host port pair for the RocketMQ 
broker");
+    }
+
+    private static int findFreePort() {
+        try (ServerSocket socket = new ServerSocket(0)) {
+            socket.setReuseAddress(true);
+            return socket.getLocalPort();
+        } catch (IOException e) {
+            throw new IllegalStateException("Could not allocate a free host 
port", e);
+        }
+    }
+
+    private static boolean isPortFree(int port) {
+        try (ServerSocket socket = new ServerSocket(port)) {
+            socket.setReuseAddress(true);
+            return true;
+        } catch (IOException e) {
+            return false;
+        }
+    }
 }
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 86a3c9c167..943e23a491 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
@@ -487,6 +487,20 @@ public class RocketMqIT extends TestSuiteBase implements 
TestResource {
                             Assertions.assertFalse(
                                     
producer.fetchPublishMessageQueues(topic).isEmpty(),
                                     "Topic route is not ready: " + topic);
+                            // The producer above may answer from its own 
route cache, primed by
+                            // createTopic() against the broker, before the 
name server has
+                            // published the route. The connector resolves 
routes through a
+                            // fresh admin client 
(RocketMqAdminUtil#offsetTopics ->
+                            // examineTopicStats), which is exactly what 
failed with
+                            // "No topic route info in name server" in CI, so 
require the route
+                            // to be visible on that path too before returning.
+                            Assertions.assertFalse(
+                                    RocketMqAdminUtil.offsetTopics(
+                                                    newConfiguration(),
+                                                    
Collections.singletonList(topic))
+                                            .get(0)
+                                            .isEmpty(),
+                                    "Topic route is not visible in the name 
server: " + topic);
                         });
     }
 
@@ -746,8 +760,11 @@ public class RocketMqIT extends TestSuiteBase implements 
TestResource {
                         + (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.
+        // for a savepoint. Restore only after both dynamic topics are visible 
again: the restored
+        // job's sink publishes to sinkTopic first, and a dropped sink route 
fails every send with
+        // "No topic route info in name server" until the route is 
re-published.
         waitForTopicRoute(sourceTopic);
+        waitForTopicRoute(sinkTopic);
         CompletableFuture.runAsync(
                 () -> {
                     try {

Reply via email to