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) {

Reply via email to