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">