This is an automated email from the ASF dual-hosted git repository.
Abacn pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/master by this push:
new a2df8610996 Deflake MqttIOTest read tests (#40393)
a2df8610996 is described below
commit a2df86109966f3d429ffbde047279fc96a7c6f13
Author: Maksym Tymoshyk <[email protected]>
AuthorDate: Fri Oct 9 00:33:10 2026 +0300
Deflake MqttIOTest read tests (#40393)
testReadWithMetadata and testReadNoClientId waited for any broker connection
with a non-empty id before publishing. The test's own publish client
connects
first, so the wait returned at once. If the reader had not subscribed yet,
the
broker dropped the messages: testReadWithMetadata read nothing, and
testReadNoClientId, which has no max read time, hung until the test timeout.
Wait for a broker subscription that matches the publish topic instead, and
drop the 2-second sleep added in #34133. Re-enable testReadNoClientId and
log
publisher failures in both tests.
Fixes #18723
---
.../org/apache/beam/sdk/io/mqtt/MqttIOTest.java | 47 ++++++++++++++++++----
1 file changed, 39 insertions(+), 8 deletions(-)
diff --git
a/sdks/java/io/mqtt/src/test/java/org/apache/beam/sdk/io/mqtt/MqttIOTest.java
b/sdks/java/io/mqtt/src/test/java/org/apache/beam/sdk/io/mqtt/MqttIOTest.java
index dfadc818a46..b026fd857b6 100644
---
a/sdks/java/io/mqtt/src/test/java/org/apache/beam/sdk/io/mqtt/MqttIOTest.java
+++
b/sdks/java/io/mqtt/src/test/java/org/apache/beam/sdk/io/mqtt/MqttIOTest.java
@@ -43,8 +43,15 @@ import java.util.Set;
import java.util.UUID;
import java.util.concurrent.ConcurrentSkipListMap;
import java.util.concurrent.ConcurrentSkipListSet;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
import org.apache.activemq.broker.BrokerService;
import org.apache.activemq.broker.Connection;
+import org.apache.activemq.broker.region.AbstractRegion;
+import org.apache.activemq.broker.region.RegionBroker;
+import org.apache.activemq.broker.region.Subscription;
+import org.apache.activemq.command.ActiveMQDestination;
+import org.apache.activemq.command.ActiveMQTopic;
import org.apache.beam.sdk.coders.ByteArrayCoder;
import org.apache.beam.sdk.coders.KvCoder;
import org.apache.beam.sdk.coders.StringUtf8Coder;
@@ -61,7 +68,6 @@ import org.joda.time.Duration;
import org.joda.time.Instant;
import org.junit.After;
import org.junit.Before;
-import org.junit.Ignore;
import org.junit.Rule;
import org.junit.Test;
import org.junit.runner.RunWith;
@@ -107,7 +113,6 @@ public class MqttIOTest {
}
@Test(timeout = 60 * 1000)
- @Ignore("https://github.com/apache/beam/issues/18723 Test timeout failure.")
public void testReadNoClientId() throws Exception {
final String topicName = "READ_TOPIC_NO_CLIENT_ID";
Read<byte[]> mqttReader =
@@ -138,7 +143,7 @@ public class MqttIOTest {
new Thread(
() -> {
try {
- doConnect(connection ->
!connection.getConnectionId().isEmpty());
+ waitForSubscription(new ActiveMQTopic(topicName));
for (int i = 0; i < 10; i++) {
publishClient
.publishWith()
@@ -148,7 +153,7 @@ public class MqttIOTest {
.send();
}
} catch (Exception e) {
- // nothing to do
+ LOG.warn("Failed to publish the test messages", e);
}
});
publisherThread.start();
@@ -245,9 +250,9 @@ public class MqttIOTest {
new Thread(
() -> {
try {
- doConnect(connection ->
!connection.getConnectionId().isEmpty());
- // Sleep two seconds, to give enough time for client to be
ready to accept messages
- Thread.sleep(2 * 1000);
+ // The broker drops messages published before the reader
subscribes. ActiveMQ names
+ // the MQTT topic "topic/1" as "topic.1".
+ waitForSubscription(new ActiveMQTopic("topic.1"));
for (int i = 0; i < 5; i++) {
publishClient
.publishWith()
@@ -266,7 +271,7 @@ public class MqttIOTest {
}
} catch (Exception e) {
- // nothing to do
+ LOG.warn("Failed to publish the test messages", e);
}
});
@@ -668,6 +673,32 @@ public class MqttIOTest {
}
}
+ /**
+ * Blocks until a client has subscribed to a topic filter that matches the
given destination.
+ *
+ * <p>The publishing clients in these tests never subscribe, so a matching
subscription belongs to
+ * the pipeline's reader.
+ *
+ * @param destination the ActiveMQ destination the test is about to publish
to
+ * @throws Exception if no client subscribes within 30 seconds
+ */
+ private void waitForSubscription(ActiveMQDestination destination) throws
Exception {
+ LOG.info(
+ "Waiting for the pipeline to subscribe to {} before sending messages
...", destination);
+ AbstractRegion topicRegion =
+ (AbstractRegion) ((RegionBroker)
brokerService.getRegionBroker()).getTopicRegion();
+ long deadlineNanos = System.nanoTime() + TimeUnit.SECONDS.toNanos(30);
+ while (System.nanoTime() < deadlineNanos) {
+ for (Subscription subscription :
topicRegion.getSubscriptions().values()) {
+ if (subscription.matches(destination)) {
+ return;
+ }
+ }
+ Thread.sleep(100);
+ }
+ throw new TimeoutException("No subscription to " + destination + " after
30 seconds");
+ }
+
@After
public void stopBroker() throws Exception {
if (brokerService != null) {