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 a0f3518d076 Deflake JmsIO tests (#39571)
a0f3518d076 is described below

commit a0f3518d076c7535f3348438d35f77b5a4373c90
Author: Yi Hu <[email protected]>
AuthorDate: Wed Aug 5 10:03:02 2026 -0400

    Deflake JmsIO tests (#39571)
    
    * Split unit tests not requiring a broker out of JmsIO
    
    * Run messaging xlang Python postcommit on single Python version
    
    * Update website about io support status
---
 ...m_PostCommit_Python_Xlang_Messaging_Direct.json |   2 +-
 CHANGES.md                                         |   1 +
 .../java/org/apache/beam/sdk/io/jms/JmsIOTest.java | 176 +--------------
 .../org/apache/beam/sdk/io/jms/JmsLocalTest.java   | 245 +++++++++++++++++++++
 sdks/python/test-suites/direct/build.gradle        |   3 +-
 .../site/content/en/documentation/io/connectors.md |  10 +-
 6 files changed, 260 insertions(+), 177 deletions(-)

diff --git 
a/.github/trigger_files/beam_PostCommit_Python_Xlang_Messaging_Direct.json 
b/.github/trigger_files/beam_PostCommit_Python_Xlang_Messaging_Direct.json
index 455144f02a3..d6a91b7e2e8 100644
--- a/.github/trigger_files/beam_PostCommit_Python_Xlang_Messaging_Direct.json
+++ b/.github/trigger_files/beam_PostCommit_Python_Xlang_Messaging_Direct.json
@@ -1,4 +1,4 @@
 {
   "comment": "Modify this file in a trivial way to cause this test suite to 
run",
-  "modification": 6
+  "modification": 7
 }
diff --git a/CHANGES.md b/CHANGES.md
index 0c4f71aa5a5..fcb011d1489 100644
--- a/CHANGES.md
+++ b/CHANGES.md
@@ -67,6 +67,7 @@
 * Upgraded Iceberg dependency to 1.11.0 (Java) 
([#38925](https://github.com/apache/beam/issues/38925)).
 * Support for X source added (Java/Python) 
([#X](https://github.com/apache/beam/issues/X)).
 * Add ArrowFlight IO (Java) 
([#20116](https://github.com/apache/beam/issues/20116)).
+* (Python) JmsIO (IBM MQ, ActiveMQ, and other providers) is now supported in 
Python via cross-language 
([#30716](https://github.com/apache/beam/issues/30716)).
 
 ## New Features / Improvements
 
diff --git 
a/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/JmsIOTest.java 
b/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/JmsIOTest.java
index eb6fb4faec0..d23b33873e1 100644
--- a/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/JmsIOTest.java
+++ b/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/JmsIOTest.java
@@ -32,27 +32,20 @@ import static org.hamcrest.Matchers.contains;
 import static org.hamcrest.Matchers.greaterThanOrEqualTo;
 import static org.hamcrest.Matchers.hasProperty;
 import static org.hamcrest.Matchers.is;
-import static org.hamcrest.Matchers.isA;
 import static org.hamcrest.Matchers.lessThan;
 import static org.hamcrest.core.StringContains.containsString;
 import static org.hamcrest.object.HasToString.hasToString;
 import static org.junit.Assert.assertEquals;
-import static org.junit.Assert.assertFalse;
 import static org.junit.Assert.assertNotNull;
 import static org.junit.Assert.assertNull;
 import static org.junit.Assert.assertTrue;
 import static org.junit.Assert.fail;
-import static org.mockito.ArgumentMatchers.any;
-import static org.mockito.ArgumentMatchers.anyLong;
-import static org.mockito.ArgumentMatchers.eq;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.times;
 import static org.mockito.Mockito.verify;
-import static org.mockito.Mockito.verifyNoMoreInteractions;
 import static org.mockito.Mockito.when;
 
 import java.io.IOException;
-import java.io.NotSerializableException;
 import java.io.Serializable;
 import java.lang.reflect.Proxy;
 import java.nio.ByteBuffer;
@@ -66,8 +59,6 @@ import java.util.Enumeration;
 import java.util.HashSet;
 import java.util.List;
 import java.util.Set;
-import java.util.concurrent.ScheduledExecutorService;
-import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicInteger;
 import java.util.function.Function;
 import javax.jms.BytesMessage;
@@ -84,7 +75,6 @@ import org.apache.activemq.ActiveMQConnectionFactory;
 import org.apache.activemq.command.ActiveMQMessage;
 import org.apache.activemq.util.Callback;
 import org.apache.beam.sdk.PipelineResult;
-import org.apache.beam.sdk.coders.Coder;
 import org.apache.beam.sdk.coders.SerializableCoder;
 import org.apache.beam.sdk.coders.StringUtf8Coder;
 import org.apache.beam.sdk.io.UnboundedSource;
@@ -94,16 +84,12 @@ import org.apache.beam.sdk.io.jms.JmsIO.UnboundedJmsReader;
 import org.apache.beam.sdk.metrics.MetricNameFilter;
 import org.apache.beam.sdk.metrics.MetricQueryResults;
 import org.apache.beam.sdk.metrics.MetricsFilter;
-import org.apache.beam.sdk.options.ExecutorOptions;
-import org.apache.beam.sdk.options.PipelineOptions;
 import org.apache.beam.sdk.options.PipelineOptionsFactory;
-import org.apache.beam.sdk.testing.CoderProperties;
 import org.apache.beam.sdk.testing.PAssert;
 import org.apache.beam.sdk.testing.TestPipeline;
 import org.apache.beam.sdk.transforms.Count;
 import org.apache.beam.sdk.transforms.Create;
 import org.apache.beam.sdk.transforms.SerializableBiFunction;
-import org.apache.beam.sdk.util.SerializableUtils;
 import org.apache.beam.sdk.values.PCollection;
 import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Throwables;
 import org.apache.qpid.jms.JmsAcknowledgeCallback;
@@ -117,7 +103,6 @@ import org.junit.Rule;
 import org.junit.Test;
 import org.junit.runner.RunWith;
 import org.junit.runners.Parameterized;
-import org.mockito.ArgumentCaptor;
 import org.mockito.Mockito;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -163,9 +148,9 @@ public class JmsIOTest {
   private ConnectionFactory connectionFactory;
   private final Class<? extends ConnectionFactory> connectionFactoryClass;
   private ConnectionFactory connectionFactoryWithSyncAcksAndWithoutPrefetch;
-  private final String brokerUrl;
-  private final Integer brokerPort;
-  private final String forceAsyncAcksParam;
+  private String brokerUrl;
+  private String forceAsyncAcksParam;
+  private int brokerPort;
 
   public JmsIOTest(
       String brokerUrl,
@@ -252,20 +237,6 @@ public class JmsIOTest {
     assertQueueIsEmpty();
   }
 
-  @Test
-  public void testPipelineWithNonSerializableCF() {
-    SerializableUtils.ensureSerializable(
-        JmsIO.read()
-            .withConnectionFactoryProviderFn(__ -> new 
MockNonSerializableConnectionFactory()));
-    try {
-      SerializableUtils.ensureSerializable(
-          JmsIO.read().withConnectionFactory(new 
MockNonSerializableConnectionFactory()));
-      fail();
-    } catch (Exception e) {
-      assertThat(Throwables.getRootCause(e), 
isA(NotSerializableException.class));
-    }
-  }
-
   @Test
   public void testReadMessagesWithCFProviderFn() throws Exception {
     long count = 5;
@@ -522,32 +493,6 @@ public class JmsIOTest {
     assertEquals(100, count);
   }
 
-  @Test
-  public void testSplitForQueue() throws Exception {
-    JmsIO.Read read = JmsIO.read().withQueue(QUEUE);
-    PipelineOptions pipelineOptions = PipelineOptionsFactory.create();
-    int desiredNumSplits = 5;
-    JmsIO.UnboundedJmsSource initialSource = new 
JmsIO.UnboundedJmsSource(read);
-    List<JmsIO.UnboundedJmsSource> splits = 
initialSource.split(desiredNumSplits, pipelineOptions);
-    // in the case of a queue, we have concurrent consumers by default, so the 
initial number
-    // splits is equal to the desired number of splits
-    assertEquals(desiredNumSplits, splits.size());
-  }
-
-  @Test
-  public void testSplitForTopic() throws Exception {
-    JmsIO.Read read = JmsIO.read().withTopic(TOPIC);
-    PipelineOptions pipelineOptions = PipelineOptionsFactory.create();
-    int desiredNumSplits = 5;
-    JmsIO.UnboundedJmsSource initialSource = new 
JmsIO.UnboundedJmsSource(read);
-    List<JmsIO.UnboundedJmsSource> splits = 
initialSource.split(desiredNumSplits, pipelineOptions);
-    // in the case of a topic, we can have only a unique subscriber on the 
topic per pipeline
-    // else it means we can have duplicate messages (all subscribers on the 
topic receive every
-    // message).
-    // So, whatever the desizedNumSplits is, the actual number of splits 
should be 1.
-    assertEquals(1, splits.size());
-  }
-
   private boolean advanceWithRetry(UnboundedSource.UnboundedReader reader) 
throws IOException {
     for (int attempt = 0; attempt < 10; attempt++) {
       if (reader.advance()) {
@@ -685,63 +630,6 @@ public class JmsIOTest {
     assertEquals(5, count(QUEUE));
   }
 
-  @Test
-  public void testJmsCheckpointMarkIndividualAcknowledgeAllMessages() throws 
Exception {
-    Message msg1 = Mockito.mock(Message.class);
-    Message msg2 = Mockito.mock(Message.class);
-    Message msg3 = Mockito.mock(Message.class);
-
-    JmsCheckpointMark.Preparer preparer =
-        
JmsCheckpointMark.newPreparer(JmsIO.AcknowledgeMode.INDIVIDUAL_ACKNOWLEDGE);
-    preparer.add(msg1);
-    preparer.add(msg2);
-    preparer.add(msg3);
-
-    AtomicInteger activeCheckpoints = new AtomicInteger(0);
-    JmsCheckpointMark mark =
-        preparer.newCheckpoint(
-            null, null, JmsIO.AcknowledgeMode.INDIVIDUAL_ACKNOWLEDGE, 
activeCheckpoints);
-    assertNotNull(mark.getMessages());
-    assertEquals(3, mark.getMessages().size());
-    assertNull(mark.getConsumer());
-    assertNull(mark.getSession());
-    assertEquals(1, activeCheckpoints.get());
-
-    mark.finalizeCheckpoint();
-
-    Mockito.verify(msg1, Mockito.times(1)).acknowledge();
-    Mockito.verify(msg2, Mockito.times(1)).acknowledge();
-    Mockito.verify(msg3, Mockito.times(1)).acknowledge();
-    assertEquals(0, activeCheckpoints.get());
-  }
-
-  @Test
-  public void 
testJmsCheckpointMarkClientAcknowledgeUnsafeNoSessionRecreation() throws 
Exception {
-    Message msg1 = Mockito.mock(Message.class);
-    Message msg2 = Mockito.mock(Message.class);
-
-    JmsCheckpointMark.Preparer preparer =
-        
JmsCheckpointMark.newPreparer(JmsIO.AcknowledgeMode.CLIENT_ACKNOWLEDGE_UNSAFE);
-    preparer.add(msg1);
-    preparer.add(msg2);
-
-    AtomicInteger activeCheckpoints = new AtomicInteger(0);
-    JmsCheckpointMark mark =
-        preparer.newCheckpoint(
-            null, null, JmsIO.AcknowledgeMode.CLIENT_ACKNOWLEDGE_UNSAFE, 
activeCheckpoints);
-    assertNotNull(mark.getMessages());
-    assertEquals(1, mark.getMessages().size());
-    assertNull(mark.getConsumer());
-    assertNull(mark.getSession());
-    assertEquals(1, activeCheckpoints.get());
-
-    mark.finalizeCheckpoint();
-
-    Mockito.verify(msg2, Mockito.times(1)).acknowledge();
-    Mockito.verify(msg1, Mockito.never()).acknowledge();
-    assertEquals(0, activeCheckpoints.get());
-  }
-
   private JmsIO.UnboundedJmsReader setupReaderForTest() throws JMSException {
     return setupReaderForTest(null);
   }
@@ -890,17 +778,6 @@ public class JmsIOTest {
     runner.join();
   }
 
-  /** Test the checkpoint mark default coder, which is actually AvroCoder. */
-  @Test
-  public void testCheckpointMarkDefaultCoder() throws Exception {
-    JmsCheckpointMark jmsCheckpointMark =
-        JmsCheckpointMark.newPreparer(JmsIO.AcknowledgeMode.CLIENT_ACKNOWLEDGE)
-            .newCheckpoint(null, null, 
JmsIO.AcknowledgeMode.CLIENT_ACKNOWLEDGE, null);
-    Coder coder = new JmsIO.UnboundedJmsSource(null).getCheckpointMarkCoder();
-    CoderProperties.coderSerializable(coder);
-    CoderProperties.coderDecodeEncodeEqual(coder, jmsCheckpointMark);
-  }
-
   @Test
   public void testDefaultAutoscaler() throws IOException {
     JmsIO.Read spec =
@@ -945,53 +822,6 @@ public class JmsIOTest {
     verify(autoScaler, times(1)).stop();
   }
 
-  @Test
-  public void testCloseWithTimeout() throws IOException, JMSException {
-    Duration closeTimeout = Duration.millis(2000L);
-    JmsIO.Read spec =
-        JmsIO.read()
-            .withConnectionFactory(connectionFactory)
-            .withUsername(USERNAME)
-            .withPassword(PASSWORD)
-            .withQueue(QUEUE)
-            .withCloseTimeout(closeTimeout);
-
-    JmsIO.UnboundedJmsSource source = new JmsIO.UnboundedJmsSource(spec);
-
-    ScheduledExecutorService mockScheduledExecutorService =
-        Mockito.mock(ScheduledExecutorService.class);
-    ExecutorOptions options = PipelineOptionsFactory.as(ExecutorOptions.class);
-    options.setScheduledExecutorService(mockScheduledExecutorService);
-    ArgumentCaptor<Runnable> runnableArgumentCaptor = 
ArgumentCaptor.forClass(Runnable.class);
-    when(mockScheduledExecutorService.schedule(
-            runnableArgumentCaptor.capture(), anyLong(), any(TimeUnit.class)))
-        .thenReturn(null /* unused */);
-
-    JmsIO.UnboundedJmsReader reader = source.createReader(options, null);
-    reader.start();
-    assertFalse(getDiscardedValue(reader));
-    reader.checkpointMarkPreparer.add(Mockito.mock(Message.class));
-    CheckpointMark mark = reader.getCheckpointMark();
-    reader.close();
-    assertTrue(getDiscardedValue(reader));
-    verify(mockScheduledExecutorService)
-        .schedule(any(Runnable.class), eq(1L), eq(TimeUnit.SECONDS));
-    mark.finalizeCheckpoint();
-    runnableArgumentCaptor.getValue().run();
-    assertTrue(getDiscardedValue(reader));
-    verifyNoMoreInteractions(mockScheduledExecutorService);
-  }
-
-  private boolean getDiscardedValue(JmsIO.UnboundedJmsReader reader) {
-    JmsCheckpointMark.Preparer preparer = reader.checkpointMarkPreparer;
-    preparer.lock.readLock().lock();
-    try {
-      return preparer.discarded;
-    } finally {
-      preparer.lock.readLock().unlock();
-    }
-  }
-
   @Test
   public void testDiscardCheckpointMark() throws Exception {
     Connection connection =
diff --git 
a/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/JmsLocalTest.java 
b/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/JmsLocalTest.java
new file mode 100644
index 00000000000..4ea8f6d317a
--- /dev/null
+++ 
b/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/JmsLocalTest.java
@@ -0,0 +1,245 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.beam.sdk.io.jms;
+
+import static org.hamcrest.MatcherAssert.assertThat;
+import static org.hamcrest.Matchers.isA;
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertNotNull;
+import static org.junit.Assert.assertNull;
+import static org.junit.Assert.assertTrue;
+import static org.junit.Assert.fail;
+
+import java.io.IOException;
+import java.io.NotSerializableException;
+import java.util.List;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+import javax.jms.Connection;
+import javax.jms.ConnectionFactory;
+import javax.jms.JMSException;
+import javax.jms.Message;
+import javax.jms.MessageConsumer;
+import javax.jms.Queue;
+import javax.jms.Session;
+import org.apache.activemq.ActiveMQConnectionFactory;
+import org.apache.beam.sdk.coders.Coder;
+import org.apache.beam.sdk.options.ExecutorOptions;
+import org.apache.beam.sdk.options.PipelineOptions;
+import org.apache.beam.sdk.options.PipelineOptionsFactory;
+import org.apache.beam.sdk.testing.CoderProperties;
+import org.apache.beam.sdk.util.SerializableUtils;
+import 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Throwables;
+import org.joda.time.Duration;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.junit.runners.JUnit4;
+import org.mockito.ArgumentCaptor;
+import org.mockito.Mockito;
+
+/** Local unit tests for {@link JmsIO} that do not require an active JMS 
broker. */
+@RunWith(JUnit4.class)
+public class JmsLocalTest {
+
+  private static final String QUEUE = "queue";
+  private static final String TOPIC = "topic";
+
+  @Test
+  public void testPipelineWithNonSerializableCF() {
+    SerializableUtils.ensureSerializable(
+        JmsIO.read()
+            .withConnectionFactoryProviderFn(__ -> new 
MockNonSerializableConnectionFactory()));
+    try {
+      SerializableUtils.ensureSerializable(
+          JmsIO.read().withConnectionFactory(new 
MockNonSerializableConnectionFactory()));
+      fail();
+    } catch (Exception e) {
+      assertThat(Throwables.getRootCause(e), 
isA(NotSerializableException.class));
+    }
+  }
+
+  @Test
+  public void testSplitForQueue() throws Exception {
+    JmsIO.Read<JmsRecord> read = JmsIO.read().withQueue(QUEUE);
+    PipelineOptions pipelineOptions = PipelineOptionsFactory.create();
+    int desiredNumSplits = 5;
+    JmsIO.UnboundedJmsSource<JmsRecord> initialSource = new 
JmsIO.UnboundedJmsSource<>(read);
+    List<JmsIO.UnboundedJmsSource<JmsRecord>> splits =
+        initialSource.split(desiredNumSplits, pipelineOptions);
+    assertEquals(desiredNumSplits, splits.size());
+  }
+
+  @Test
+  public void testSplitForTopic() throws Exception {
+    JmsIO.Read<JmsRecord> read = JmsIO.read().withTopic(TOPIC);
+    PipelineOptions pipelineOptions = PipelineOptionsFactory.create();
+    int desiredNumSplits = 5;
+    JmsIO.UnboundedJmsSource<JmsRecord> initialSource = new 
JmsIO.UnboundedJmsSource<>(read);
+    List<JmsIO.UnboundedJmsSource<JmsRecord>> splits =
+        initialSource.split(desiredNumSplits, pipelineOptions);
+    assertEquals(1, splits.size());
+  }
+
+  @Test
+  public void testPublisherWithRetryConfiguration() {
+    RetryConfiguration retryPolicy =
+        RetryConfiguration.create(5, Duration.standardSeconds(15), null);
+    JmsIO.Write<String> publisher =
+        JmsIO.<String>write()
+            .withConnectionFactory(new 
ActiveMQConnectionFactory("vm://localhost"))
+            .withRetryConfiguration(retryPolicy)
+            .withQueue(QUEUE)
+            .withUsername("user")
+            .withPassword("password");
+    assertEquals(
+        publisher.getRetryConfiguration(),
+        RetryConfiguration.create(5, Duration.standardSeconds(15), null));
+  }
+
+  @Test
+  public void testJmsCheckpointMarkIndividualAcknowledgeAllMessages() throws 
Exception {
+    Message msg1 = Mockito.mock(Message.class);
+    Message msg2 = Mockito.mock(Message.class);
+    Message msg3 = Mockito.mock(Message.class);
+
+    JmsCheckpointMark.Preparer preparer =
+        
JmsCheckpointMark.newPreparer(JmsIO.AcknowledgeMode.INDIVIDUAL_ACKNOWLEDGE);
+    preparer.add(msg1);
+    preparer.add(msg2);
+    preparer.add(msg3);
+
+    AtomicInteger activeCheckpoints = new AtomicInteger(0);
+    JmsCheckpointMark mark =
+        preparer.newCheckpoint(
+            null, null, JmsIO.AcknowledgeMode.INDIVIDUAL_ACKNOWLEDGE, 
activeCheckpoints);
+    assertNotNull(mark.getMessages());
+    assertEquals(3, mark.getMessages().size());
+    assertNull(mark.getConsumer());
+    assertNull(mark.getSession());
+    assertEquals(1, activeCheckpoints.get());
+
+    mark.finalizeCheckpoint();
+
+    Mockito.verify(msg1, Mockito.times(1)).acknowledge();
+    Mockito.verify(msg2, Mockito.times(1)).acknowledge();
+    Mockito.verify(msg3, Mockito.times(1)).acknowledge();
+    assertEquals(0, activeCheckpoints.get());
+  }
+
+  @Test
+  public void 
testJmsCheckpointMarkClientAcknowledgeUnsafeNoSessionRecreation() throws 
Exception {
+    Message msg1 = Mockito.mock(Message.class);
+    Message msg2 = Mockito.mock(Message.class);
+
+    JmsCheckpointMark.Preparer preparer =
+        
JmsCheckpointMark.newPreparer(JmsIO.AcknowledgeMode.CLIENT_ACKNOWLEDGE_UNSAFE);
+    preparer.add(msg1);
+    preparer.add(msg2);
+
+    AtomicInteger activeCheckpoints = new AtomicInteger(0);
+    JmsCheckpointMark mark =
+        preparer.newCheckpoint(
+            null, null, JmsIO.AcknowledgeMode.CLIENT_ACKNOWLEDGE_UNSAFE, 
activeCheckpoints);
+    assertNotNull(mark.getMessages());
+    assertEquals(1, mark.getMessages().size());
+    assertNull(mark.getConsumer());
+    assertNull(mark.getSession());
+    assertEquals(1, activeCheckpoints.get());
+
+    mark.finalizeCheckpoint();
+
+    Mockito.verify(msg2, Mockito.times(1)).acknowledge();
+    Mockito.verify(msg1, Mockito.never()).acknowledge();
+    assertEquals(0, activeCheckpoints.get());
+  }
+
+  /** Test the checkpoint mark default coder, which is actually AvroCoder. */
+  @Test
+  public void testCheckpointMarkDefaultCoder() throws Exception {
+    JmsCheckpointMark jmsCheckpointMark =
+        JmsCheckpointMark.newPreparer(JmsIO.AcknowledgeMode.CLIENT_ACKNOWLEDGE)
+            .newCheckpoint(null, null, 
JmsIO.AcknowledgeMode.CLIENT_ACKNOWLEDGE, null);
+    Coder<JmsCheckpointMark> coder =
+        new JmsIO.UnboundedJmsSource<JmsRecord>(null).getCheckpointMarkCoder();
+    CoderProperties.coderSerializable(coder);
+    CoderProperties.coderDecodeEncodeEqual(coder, jmsCheckpointMark);
+  }
+
+  @Test
+  public void testCloseWithTimeout() throws IOException, JMSException {
+    ConnectionFactory connectionFactory = 
Mockito.mock(ConnectionFactory.class);
+    Connection connection = Mockito.mock(Connection.class);
+    Session session = Mockito.mock(Session.class);
+    MessageConsumer consumer = Mockito.mock(MessageConsumer.class);
+    Queue queue = Mockito.mock(Queue.class);
+
+    Mockito.when(connectionFactory.createConnection(Mockito.any(), 
Mockito.any()))
+        .thenReturn(connection);
+    Mockito.when(connection.createSession(Mockito.anyBoolean(), 
Mockito.anyInt()))
+        .thenReturn(session);
+    Mockito.when(session.createQueue(Mockito.anyString())).thenReturn(queue);
+    Mockito.when(session.createConsumer(Mockito.any())).thenReturn(consumer);
+
+    Duration closeTimeout = Duration.millis(2000L);
+    JmsIO.Read<JmsRecord> spec =
+        JmsIO.read()
+            .withConnectionFactory(connectionFactory)
+            .withUsername("user")
+            .withPassword("password")
+            .withQueue(QUEUE)
+            .withCloseTimeout(closeTimeout);
+
+    JmsIO.UnboundedJmsSource<JmsRecord> source = new 
JmsIO.UnboundedJmsSource<>(spec);
+
+    ScheduledExecutorService mockScheduledExecutorService =
+        Mockito.mock(ScheduledExecutorService.class);
+    ExecutorOptions options = PipelineOptionsFactory.as(ExecutorOptions.class);
+    options.setScheduledExecutorService(mockScheduledExecutorService);
+    ArgumentCaptor<Runnable> runnableArgumentCaptor = 
ArgumentCaptor.forClass(Runnable.class);
+    Mockito.when(
+            mockScheduledExecutorService.schedule(
+                runnableArgumentCaptor.capture(), Mockito.anyLong(), 
Mockito.any(TimeUnit.class)))
+        .thenReturn(null /* unused */);
+
+    JmsIO.UnboundedJmsReader<JmsRecord> reader = source.createReader(options, 
null);
+    reader.start();
+    assertFalse(getDiscardedValue(reader));
+    reader.checkpointMarkPreparer.add(Mockito.mock(Message.class));
+    org.apache.beam.sdk.io.UnboundedSource.CheckpointMark mark = 
reader.getCheckpointMark();
+    reader.close();
+    assertTrue(getDiscardedValue(reader));
+    Mockito.verify(mockScheduledExecutorService)
+        .schedule(Mockito.any(Runnable.class), Mockito.eq(1L), 
Mockito.eq(TimeUnit.SECONDS));
+    mark.finalizeCheckpoint();
+    runnableArgumentCaptor.getValue().run();
+    assertTrue(getDiscardedValue(reader));
+    Mockito.verifyNoMoreInteractions(mockScheduledExecutorService);
+  }
+
+  private boolean getDiscardedValue(JmsIO.UnboundedJmsReader<JmsRecord> 
reader) {
+    JmsCheckpointMark.Preparer preparer = reader.checkpointMarkPreparer;
+    preparer.lock.readLock().lock();
+    try {
+      return preparer.discarded;
+    } finally {
+      preparer.lock.readLock().unlock();
+    }
+  }
+}
diff --git a/sdks/python/test-suites/direct/build.gradle 
b/sdks/python/test-suites/direct/build.gradle
index d1fe45683a8..2c71c81afaa 100644
--- a/sdks/python/test-suites/direct/build.gradle
+++ b/sdks/python/test-suites/direct/build.gradle
@@ -44,7 +44,8 @@ task ioCrossLanguagePostCommit {
 }
 
 task messagingCrossLanguagePostCommit {
-  getVersionsAsList('cross_language_validates_py_versions').each {
+  // Messaging E2E tests has testcontainers overhead. Single Python version 
suffices and reducing CI flakiness
+  getVersionsAsList('cross_language_validates_py_versions').take(1).each {
     
dependsOn.add(":sdks:python:test-suites:direct:py${getVersionSuffix(it)}:messagingCrossLanguagePythonUsingJava")
   }
 }
diff --git a/website/www/site/content/en/documentation/io/connectors.md 
b/website/www/site/content/en/documentation/io/connectors.md
index 242a255ba82..679b6bf1e0e 100644
--- a/website/www/site/content/en/documentation/io/connectors.md
+++ b/website/www/site/content/en/documentation/io/connectors.md
@@ -445,7 +445,10 @@ This table provides a consolidated, at-a-glance overview 
of the available built-
       ✔
       <a 
href="https://beam.apache.org/releases/javadoc/current/org/apache/beam/sdk/io/mqtt/MqttIO.html";>native</a>
     </td>
-    <td>Not available</td>
+    <td class="present">
+      ✔
+      <a 
href="https://beam.apache.org/releases/pydoc/current/apache_beam.transforms.xlang.io.html#apache_beam.transforms.xlang.io.ReadFromMqtt";>via
 X-language</a>
+    </td>
     <td>Not available</td>
     <td>Not available</td>
     <td>Not available</td>
@@ -1269,7 +1272,10 @@ This table provides a consolidated, at-a-glance overview 
of the available built-
       ✔
       <a 
href="https://beam.apache.org/releases/javadoc/current/org/apache/beam/sdk/io/datadog/DatadogIO.html";>native</a>
     </td>
-    <td>Not available</td>
+    <td class="present">
+      ✔
+      <a 
href="https://beam.apache.org/releases/pydoc/current/apache_beam.transforms.xlang.io.html#apache_beam.transforms.xlang.io.DatadogWrite";>via
 X-language</a>
+    </td>
     <td>Not available</td>
     <td>Not available</td>
     <td class="present">

Reply via email to