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 455144f02a35..d6a91b7e2e86 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 d853314a0ad3..caafac37c745 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 eb6fb4faec04..d23b33873e14 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.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.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.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.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.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 static Collection connectionFactories() { private ConnectionFactory connectionFactory; private final Class 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 void testReadMessages() throws Exception { 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 void testWriteDynamicMessage() throws Exception { 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 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 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 void testCheckpointMarkAndFinalizeSeparatelyClientAcknowledgeUnsafe() thr 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 void testCheckpointMarkSafety() throws Exception { 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 void testCustomAutoscaler() throws IOException { 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 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 000000000000..4ea8f6d317a7 --- /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 read = JmsIO.read().withQueue(QUEUE); + PipelineOptions pipelineOptions = PipelineOptionsFactory.create(); + int desiredNumSplits = 5; + JmsIO.UnboundedJmsSource initialSource = new JmsIO.UnboundedJmsSource<>(read); + List> splits = + initialSource.split(desiredNumSplits, pipelineOptions); + 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> splits = + initialSource.split(desiredNumSplits, pipelineOptions); + assertEquals(1, splits.size()); + } + + @Test + public void testPublisherWithRetryConfiguration() { + RetryConfiguration retryPolicy = + RetryConfiguration.create(5, Duration.standardSeconds(15), null); + JmsIO.Write publisher = + JmsIO.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 coder = + new JmsIO.UnboundedJmsSource(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 spec = + JmsIO.read() + .withConnectionFactory(connectionFactory) + .withUsername("user") + .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 runnableArgumentCaptor = ArgumentCaptor.forClass(Runnable.class); + Mockito.when( + mockScheduledExecutorService.schedule( + runnableArgumentCaptor.capture(), Mockito.anyLong(), Mockito.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)); + 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 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 d1fe45683a83..2c71c81afaa9 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 242a255ba82d..679b6bf1e0e3 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- ✔ native - Not available + + ✔ + via X-language + Not available Not available Not available @@ -1269,7 +1272,10 @@ This table provides a consolidated, at-a-glance overview of the available built- ✔ native - Not available + + ✔ + via X-language + Not available Not available