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 {