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 40852a71e00 Update JmsIO to ActiveMQ 6.2.5 and jakarta.jms (#38729)
40852a71e00 is described below
commit 40852a71e007b6a16a3e62b77cf03cce102998f9
Author: Jean-Baptiste Onofré <[email protected]>
AuthorDate: Thu Oct 8 23:15:39 2026 +0200
Update JmsIO to ActiveMQ 6.2.5 and jakarta.jms (#38729)
* Update JmsIO to ActiveMQ 6.2.5 and jakarta.jms
Bump activemq to 6.2.5 and qpid-jms-client to 2.10.0, both of which
target jakarta.jms. Replace the geronimo-jms_2.0_spec dependency with
jakarta.jms-api 3.1.0 and migrate all javax.jms imports and javadoc
references in JmsIO main and test sources to jakarta.jms.
---------
Co-authored-by: Derrick Williams <[email protected]>
---
CHANGES.md | 4 ++-
.../org/apache/beam/gradle/BeamModulePlugin.groovy | 4 +--
sdks/java/io/amqp/build.gradle | 16 +++++++++
sdks/java/io/jms/build.gradle | 8 ++---
.../io/jms/BeamGenericJmsConnectionFactory.java | 2 +-
.../beam/sdk/io/jms/ConnectionConfiguration.java | 6 ++--
.../apache/beam/sdk/io/jms/JmsCheckpointMark.java | 8 ++---
.../java/org/apache/beam/sdk/io/jms/JmsIO.java | 38 +++++++++++-----------
.../java/org/apache/beam/sdk/io/jms/JmsRecord.java | 2 +-
.../apache/beam/sdk/io/jms/TextMessageMapper.java | 12 +++----
.../java/org/apache/beam/sdk/io/jms/CommonJms.java | 8 ++---
.../sdk/io/jms/ConnectionConfigurationTest.java | 19 ++++++-----
.../java/org/apache/beam/sdk/io/jms/JmsIOIT.java | 14 ++++----
.../java/org/apache/beam/sdk/io/jms/JmsIOTest.java | 32 +++++++++---------
.../org/apache/beam/sdk/io/jms/JmsLocalTest.java | 14 ++++----
.../jms/MockNonSerializableConnectionFactory.java | 8 ++---
.../io/messaging-expansion-service/build.gradle | 1 +
sdks/java/io/mqtt/build.gradle | 15 +++++++++
.../apache_beam/io/external/xlang_jmsio_it_test.py | 4 +--
sdks/python/apache_beam/yaml/integration_tests.py | 2 +-
sdks/python/apache_beam/yaml/standard_io.yaml | 2 +-
21 files changed, 126 insertions(+), 93 deletions(-)
diff --git a/CHANGES.md b/CHANGES.md
index eb374150cad..cb5eced42f7 100644
--- a/CHANGES.md
+++ b/CHANGES.md
@@ -76,8 +76,10 @@
## Breaking Changes
-* X behavior was changed ([#X](https://github.com/apache/beam/issues/X)).
* (Go) The row coder now encodes `int16` and `uint16` struct fields as 2 byte
big endian INT16 values, matching the Java and Python SDKs. This is an update
incompatible change for streaming pipelines that use rows with `int16` or
`uint16` fields ([#40151](https://github.com/apache/beam/issues/40151)).
+* (Java) JmsIO migrated to `jakarta.jms` (JMS 3.1) and ActiveMQ 6.2.5. User
code implementing `JmsIO.MessageMapper`,
+ `valueMapper`, `topicNameMapper`, or providing a `ConnectionFactory` must
update imports from `javax.jms.*` to
+ `jakarta.jms.*`. The module now requires Java 17 at runtime
([#38729](https://github.com/apache/beam/issues/38729)).
## Deprecations
diff --git
a/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy
b/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy
index 3b93abfcc9e..263585f60b1 100644
--- a/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy
+++ b/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy
@@ -608,7 +608,7 @@ class BeamModulePlugin implements Plugin<Project> {
//
// There are a few versions are determined by the BOMs by running
scripts/tools/bomupgrader.py
// marked as [bomupgrader]. See the documentation of that script for
detail.
- def activemq_version = "5.19.5"
+ def activemq_version = "6.2.5"
def autovalue_version = "1.9"
def autoservice_version = "1.0.1"
def aws_java_sdk2_version = "2.20.162"
@@ -650,7 +650,7 @@ class BeamModulePlugin implements Plugin<Project> {
def postgres_version = "42.6.2"
// [bomupgrader] determined by: com.google.protobuf:protobuf-java,
consistent with: google_cloud_platform_libraries_bom
def protobuf_version = "4.33.6"
- def qpid_jms_client_version = "0.61.0"
+ def qpid_jms_client_version = "2.10.0"
def quickcheck_version = "1.0"
def sbe_tool_version = "1.25.1"
def singlestore_jdbc_version = "1.1.4"
diff --git a/sdks/java/io/amqp/build.gradle b/sdks/java/io/amqp/build.gradle
index 6f2899eeb05..440ff895f73 100644
--- a/sdks/java/io/amqp/build.gradle
+++ b/sdks/java/io/amqp/build.gradle
@@ -19,6 +19,22 @@
plugins { id 'org.apache.beam.module' }
applyJavaNature( automaticModuleName: 'org.apache.beam.sdk.io.amqp')
+// The JMS IO module pulls ActiveMQ 6.2.5 (Java 17) into the shared dependency
+// map. Because the global resolution strategy forces every mapped version,
+// that 6.2.5 would otherwise be forced onto this Java 11 module via the shared
+// org.apache.activemq coordinates. AMQP IO only needs ActiveMQ as an embedded
+// test broker, so pin it back to the Java 11-compatible 5.x line.
+def activemq5_version = '5.19.5'
+configurations.all {
+ resolutionStrategy.dependencySubstitution {
+ substitute module('org.apache.activemq:activemq-broker') using
module("org.apache.activemq:activemq-broker:${activemq5_version}")
+ substitute module('org.apache.activemq:activemq-amqp') using
module("org.apache.activemq:activemq-amqp:${activemq5_version}")
+ substitute module('org.apache.activemq:activemq-client') using
module("org.apache.activemq:activemq-client:${activemq5_version}")
+ substitute module('org.apache.activemq:activemq-kahadb-store') using
module("org.apache.activemq:activemq-kahadb-store:${activemq5_version}")
+ substitute module('org.apache.activemq.tooling:activemq-junit') using
module("org.apache.activemq.tooling:activemq-junit:${activemq5_version}")
+ }
+}
+
description = "Apache Beam :: SDKs :: Java :: IO :: AMQP"
ext.summary = "IO to read and write using AMQP 1.0 protocol
(http://www.amqp.org)."
diff --git a/sdks/java/io/jms/build.gradle b/sdks/java/io/jms/build.gradle
index 369d02b37d3..e801701caa0 100644
--- a/sdks/java/io/jms/build.gradle
+++ b/sdks/java/io/jms/build.gradle
@@ -19,6 +19,7 @@
plugins { id 'org.apache.beam.module' }
applyJavaNature(
automaticModuleName: 'org.apache.beam.sdk.io.jms',
+ requireJavaVersion: JavaVersion.VERSION_17,
)
provideIntegrationTestingDependencies()
enableJavaPerformanceTesting()
@@ -32,13 +33,10 @@ dependencies {
implementation project(path: ":sdks:java:core", configuration: "shadow")
implementation library.java.slf4j_api
implementation library.java.joda_time
- implementation "org.apache.geronimo.specs:geronimo-jms_2.0_spec:1.0-alpha-2"
// Don't put proprietary licensed ibm mq client into runtimeClasspath
// (affects expansion service shadow jar)
- compileOnly("com.ibm.mq:com.ibm.mq.allclient:9.3.0.25") {
- // duplicating geronimo-jms_2.0_spec
- exclude group: "javax.jms", module: "javax.jms-api"
- }
+ compileOnly "com.ibm.mq:com.ibm.mq.jakarta.client:9.3.0.25"
+ implementation "jakarta.jms:jakarta.jms-api:3.1.0"
testImplementation library.java.activemq_amqp
testImplementation library.java.activemq_broker
testImplementation library.java.activemq_jaas
diff --git
a/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/BeamGenericJmsConnectionFactory.java
b/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/BeamGenericJmsConnectionFactory.java
index d14c8c8addc..695fefdc1ac 100644
---
a/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/BeamGenericJmsConnectionFactory.java
+++
b/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/BeamGenericJmsConnectionFactory.java
@@ -17,8 +17,8 @@
*/
package org.apache.beam.sdk.io.jms;
+import jakarta.jms.ConnectionFactory;
import java.io.Serializable;
-import javax.jms.ConnectionFactory;
/**
* An interface for creating custom JMS {@link ConnectionFactory} instances.
diff --git
a/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/ConnectionConfiguration.java
b/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/ConnectionConfiguration.java
index f534f8a0731..cfa94a739b2 100644
---
a/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/ConnectionConfiguration.java
+++
b/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/ConnectionConfiguration.java
@@ -20,14 +20,14 @@ package org.apache.beam.sdk.io.jms;
import static
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkArgument;
import com.google.auto.value.AutoValue;
-import com.ibm.mq.jms.MQConnectionFactory;
-import com.ibm.msg.client.wmq.WMQConstants;
+import com.ibm.mq.jakarta.jms.MQConnectionFactory;
+import com.ibm.msg.client.jakarta.wmq.WMQConstants;
+import jakarta.jms.ConnectionFactory;
import java.io.Serializable;
import java.lang.reflect.InvocationTargetException;
import java.lang.reflect.Method;
import java.net.URI;
import java.util.List;
-import javax.jms.ConnectionFactory;
import org.apache.beam.sdk.schemas.AutoValueSchema;
import org.apache.beam.sdk.schemas.annotations.DefaultSchema;
import org.apache.beam.sdk.schemas.annotations.SchemaFieldDescription;
diff --git
a/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/JmsCheckpointMark.java
b/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/JmsCheckpointMark.java
index 1a00216564a..fdfd2c0c7bc 100644
---
a/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/JmsCheckpointMark.java
+++
b/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/JmsCheckpointMark.java
@@ -17,6 +17,10 @@
*/
package org.apache.beam.sdk.io.jms;
+import jakarta.jms.JMSException;
+import jakarta.jms.Message;
+import jakarta.jms.MessageConsumer;
+import jakarta.jms.Session;
import java.io.IOException;
import java.io.Serializable;
import java.util.ArrayList;
@@ -24,10 +28,6 @@ import java.util.List;
import java.util.Objects;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.locks.ReentrantReadWriteLock;
-import javax.jms.JMSException;
-import javax.jms.Message;
-import javax.jms.MessageConsumer;
-import javax.jms.Session;
import org.apache.beam.sdk.io.UnboundedSource;
import org.apache.beam.sdk.io.jms.JmsIO.AcknowledgeMode;
import
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.annotations.VisibleForTesting;
diff --git
a/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/JmsIO.java
b/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/JmsIO.java
index 7fafe074a9f..8de0ef52698 100644
--- a/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/JmsIO.java
+++ b/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/JmsIO.java
@@ -22,6 +22,15 @@ import static
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Pr
import com.google.auto.value.AutoValue;
import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
+import jakarta.jms.Connection;
+import jakarta.jms.ConnectionFactory;
+import jakarta.jms.Destination;
+import jakarta.jms.JMSException;
+import jakarta.jms.Message;
+import jakarta.jms.MessageConsumer;
+import jakarta.jms.MessageProducer;
+import jakarta.jms.Session;
+import jakarta.jms.TextMessage;
import java.io.IOException;
import java.io.Serializable;
import java.nio.charset.StandardCharsets;
@@ -37,15 +46,6 @@ import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.stream.Stream;
-import javax.jms.Connection;
-import javax.jms.ConnectionFactory;
-import javax.jms.Destination;
-import javax.jms.JMSException;
-import javax.jms.Message;
-import javax.jms.MessageConsumer;
-import javax.jms.MessageProducer;
-import javax.jms.Session;
-import javax.jms.TextMessage;
import org.apache.beam.sdk.coders.Coder;
import org.apache.beam.sdk.coders.SerializableCoder;
import org.apache.beam.sdk.io.Read.Unbounded;
@@ -87,9 +87,9 @@ import org.slf4j.LoggerFactory;
*
* <p>JmsIO source returns unbounded collection of JMS records as {@code
PCollection<JmsRecord>}. A
* {@link JmsRecord} includes JMS headers and properties, along with the JMS
{@link
- * javax.jms.TextMessage} payload.
+ * jakarta.jms.TextMessage} payload.
*
- * <p>To configure a JMS source, you have to provide a {@link
javax.jms.ConnectionFactory} and the
+ * <p>To configure a JMS source, you have to provide a {@link
jakarta.jms.ConnectionFactory} and the
* destination (queue or topic) where to consume. The following example
illustrates various options
* for configuring the source:
*
@@ -103,8 +103,8 @@ import org.slf4j.LoggerFactory;
*
* }</pre>
*
- * <p>It is possible to read any type of JMS {@link javax.jms.Message} into a
custom POJO using the
- * following configuration:
+ * <p>It is possible to read any type of JMS {@link jakarta.jms.Message} into
a custom POJO using
+ * the following configuration:
*
* <pre>{@code
* pipeline.apply(JmsIO.<T>readMessage()
@@ -122,10 +122,10 @@ import org.slf4j.LoggerFactory;
* <h4>Acknowledgment Modes and Client Prefetch Configuration</h4>
*
* <p>By default, {@link JmsIO} consumes messages using {@link
AcknowledgeMode#CLIENT_ACKNOWLEDGE}
- * where a new {@link javax.jms.Session} is created for each checkpoint to
prevent premature
+ * where a new {@link jakarta.jms.Session} is created for each checkpoint to
prevent premature
* acknowledgments across bundles. When using {@link
AcknowledgeMode#CLIENT_ACKNOWLEDGE}, if your
* JMS broker or client library utilizes client-side message prefetch buffers
(such as Apache
- * ActiveMQ), you should configure {@code prefetch=0} on your {@link
javax.jms.ConnectionFactory}
+ * ActiveMQ), you should configure {@code prefetch=0} on your {@link
jakarta.jms.ConnectionFactory}
* (e.g., via {@code ?jms.prefetchPolicy.all=0} in the broker URL or {@code
* ActiveMQPrefetchPolicy.setAll(0)}). Otherwise, unconsumed messages could be
held inside old
* consumers in low throughput scenario and could lead to message backlog.
@@ -138,8 +138,8 @@ import org.slf4j.LoggerFactory;
* <h3>Writing to a JMS destination</h3>
*
* <p>JmsIO sink supports writing text messages to a JMS destination on a
broker. To configure a JMS
- * sink, you must specify a {@link javax.jms.ConnectionFactory} and a {@link
javax.jms.Destination}
- * name. For instance:
+ * sink, you must specify a {@link jakarta.jms.ConnectionFactory} and a {@link
+ * jakarta.jms.Destination} name. For instance:
*
* <pre>{@code
* pipeline
@@ -1223,7 +1223,7 @@ public class JmsIO {
}
/**
- * Map the {@code EventT} object to a {@link javax.jms.Message}.
+ * Map the {@code EventT} object to a {@link jakarta.jms.Message}.
*
* <p>For instance:
*
@@ -1245,7 +1245,7 @@ public class JmsIO {
* .apply(JmsIO.write().withValueMapper(valueNapper)
* }</pre>
*
- * @param valueMapper The function returning the {@link javax.jms.Message}
+ * @param valueMapper The function returning the {@link
jakarta.jms.Message}
* @return The corresponding {@link JmsIO.Write}.
*/
public Write<EventT> withValueMapper(
diff --git
a/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/JmsRecord.java
b/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/JmsRecord.java
index 9f547a801d0..137ec39fe71 100644
--- a/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/JmsRecord.java
+++ b/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/JmsRecord.java
@@ -18,10 +18,10 @@
package org.apache.beam.sdk.io.jms;
import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
+import jakarta.jms.Destination;
import java.io.Serializable;
import java.util.Map;
import java.util.Objects;
-import javax.jms.Destination;
import org.checkerframework.checker.nullness.qual.Nullable;
/**
diff --git
a/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/TextMessageMapper.java
b/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/TextMessageMapper.java
index d5d85c46794..6a93c5ea639 100644
---
a/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/TextMessageMapper.java
+++
b/sdks/java/io/jms/src/main/java/org/apache/beam/sdk/io/jms/TextMessageMapper.java
@@ -17,15 +17,15 @@
*/
package org.apache.beam.sdk.io.jms;
-import javax.jms.JMSException;
-import javax.jms.Message;
-import javax.jms.Session;
-import javax.jms.TextMessage;
+import jakarta.jms.JMSException;
+import jakarta.jms.Message;
+import jakarta.jms.Session;
+import jakarta.jms.TextMessage;
import org.apache.beam.sdk.transforms.SerializableBiFunction;
/**
- * The TextMessageMapper takes a {@link String} value, a {@link
javax.jms.Session} and returns a
- * {@link javax.jms.TextMessage}.
+ * The TextMessageMapper takes a {@link String} value, a {@link
jakarta.jms.Session} and returns a
+ * {@link jakarta.jms.TextMessage}.
*/
public class TextMessageMapper implements SerializableBiFunction<String,
Session, Message> {
diff --git
a/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/CommonJms.java
b/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/CommonJms.java
index 749b215b205..106d21ead61 100644
--- a/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/CommonJms.java
+++ b/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/CommonJms.java
@@ -17,15 +17,15 @@
*/
package org.apache.beam.sdk.io.jms;
+import jakarta.jms.BytesMessage;
+import jakarta.jms.ConnectionFactory;
+import jakarta.jms.Message;
import java.io.Serializable;
import java.lang.reflect.InvocationTargetException;
import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.List;
import java.util.function.Supplier;
-import javax.jms.BytesMessage;
-import javax.jms.ConnectionFactory;
-import javax.jms.Message;
import org.apache.activemq.broker.BrokerPlugin;
import org.apache.activemq.broker.BrokerService;
import org.apache.activemq.security.AuthenticationUser;
@@ -160,7 +160,7 @@ public class CommonJms implements Serializable {
return this.connectionFactoryClass;
}
- /** A test class that maps a {@link javax.jms.BytesMessage} into a {@link
String}. */
+ /** A test class that maps a {@link jakarta.jms.BytesMessage} into a {@link
String}. */
public static class BytesMessageToStringMessageMapper implements
JmsIO.MessageMapper<String> {
@Override
diff --git
a/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/ConnectionConfigurationTest.java
b/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/ConnectionConfigurationTest.java
index 7aabc683855..11db2e9f716 100644
---
a/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/ConnectionConfigurationTest.java
+++
b/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/ConnectionConfigurationTest.java
@@ -22,9 +22,9 @@ import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertThrows;
import static org.junit.Assert.assertTrue;
-import javax.jms.Connection;
-import javax.jms.ConnectionFactory;
-import javax.jms.JMSException;
+import jakarta.jms.Connection;
+import jakarta.jms.ConnectionFactory;
+import jakarta.jms.JMSException;
import org.apache.activemq.ActiveMQConnectionFactory;
import org.apache.qpid.jms.JmsConnectionFactory;
import org.checkerframework.checker.nullness.qual.Nullable;
@@ -79,22 +79,22 @@ public class ConnectionConfigurationTest {
}
@Override
- public javax.jms.JMSContext createContext() {
+ public jakarta.jms.JMSContext createContext() {
return null;
}
@Override
- public javax.jms.JMSContext createContext(int sessionMode) {
+ public jakarta.jms.JMSContext createContext(int sessionMode) {
return null;
}
@Override
- public javax.jms.JMSContext createContext(String username, String
password) {
+ public jakarta.jms.JMSContext createContext(String username, String
password) {
return null;
}
@Override
- public javax.jms.JMSContext createContext(String username, String
password, int sessionMode) {
+ public jakarta.jms.JMSContext createContext(String username, String
password, int sessionMode) {
return null;
}
}
@@ -137,11 +137,12 @@ public class ConnectionConfigurationTest {
public void testClientsNotSlippedIntoRuntimeDependencies() {
// Verify IBM MQ client is not present on runtime classpath
assertThrows(
- ClassNotFoundException.class, () ->
Class.forName("com.ibm.mq.jms.MQConnectionFactory"));
+ ClassNotFoundException.class,
+ () -> Class.forName("com.ibm.mq.jakarta.jms.MQConnectionFactory"));
ConnectionConfiguration config =
ConnectionConfiguration.create("tcp://localhost:1414")
-
.withConnectionFactoryClassName("com.ibm.mq.jms.MQConnectionFactory");
+
.withConnectionFactoryClassName("com.ibm.mq.jakarta.jms.MQConnectionFactory");
IllegalArgumentException exception =
assertThrows(IllegalArgumentException.class,
config::createConnectionFactory);
assertTrue(exception.getCause() instanceof ClassNotFoundException);
diff --git
a/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/JmsIOIT.java
b/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/JmsIOIT.java
index 3dbb20775f7..edbb41a7ac7 100644
--- a/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/JmsIOIT.java
+++ b/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/JmsIOIT.java
@@ -23,6 +23,13 @@ import static org.apache.beam.sdk.io.jms.CommonJms.USERNAME;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotEquals;
+import jakarta.jms.Connection;
+import jakarta.jms.ConnectionFactory;
+import jakarta.jms.JMSException;
+import jakarta.jms.Message;
+import jakarta.jms.QueueBrowser;
+import jakarta.jms.Session;
+import jakarta.jms.TextMessage;
import java.io.IOException;
import java.io.Serializable;
import java.time.Instant;
@@ -34,13 +41,6 @@ import java.util.Set;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.function.Function;
-import javax.jms.Connection;
-import javax.jms.ConnectionFactory;
-import javax.jms.JMSException;
-import javax.jms.Message;
-import javax.jms.QueueBrowser;
-import javax.jms.Session;
-import javax.jms.TextMessage;
import org.apache.activemq.ActiveMQConnectionFactory;
import org.apache.activemq.command.ActiveMQTextMessage;
import org.apache.beam.sdk.PipelineResult;
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 d6cb134249a..c92794c21e5 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
@@ -45,6 +45,16 @@ import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
+import jakarta.jms.BytesMessage;
+import jakarta.jms.Connection;
+import jakarta.jms.ConnectionFactory;
+import jakarta.jms.JMSException;
+import jakarta.jms.Message;
+import jakarta.jms.MessageConsumer;
+import jakarta.jms.MessageProducer;
+import jakarta.jms.QueueBrowser;
+import jakarta.jms.Session;
+import jakarta.jms.TextMessage;
import java.io.IOException;
import java.io.Serializable;
import java.lang.reflect.Proxy;
@@ -61,16 +71,6 @@ import java.util.List;
import java.util.Set;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.function.Function;
-import javax.jms.BytesMessage;
-import javax.jms.Connection;
-import javax.jms.ConnectionFactory;
-import javax.jms.JMSException;
-import javax.jms.Message;
-import javax.jms.MessageConsumer;
-import javax.jms.MessageProducer;
-import javax.jms.QueueBrowser;
-import javax.jms.Session;
-import javax.jms.TextMessage;
import org.apache.activemq.ActiveMQConnectionFactory;
import org.apache.activemq.command.ActiveMQMessage;
import org.apache.activemq.util.Callback;
@@ -568,8 +568,8 @@ public class JmsIOTest {
assertEquals(1, jmsMark.getMessages().size());
// consume two more messages after checkpoint made
- reader.advance();
- reader.advance();
+ assertTrue(advanceWithRetry(reader));
+ assertTrue(advanceWithRetry(reader));
// the messages are still pending in the queue (no ACK yet)
assertEquals(10, count(QUEUE));
@@ -596,8 +596,8 @@ public class JmsIOTest {
assertNotNull(jmsMark.getMessages());
assertEquals(3, jmsMark.getMessages().size());
- reader.advance();
- reader.advance();
+ assertTrue(advanceWithRetry(reader));
+ assertTrue(advanceWithRetry(reader));
assertEquals(10, count(QUEUE));
mark.finalizeCheckpoint();
@@ -621,8 +621,8 @@ public class JmsIOTest {
assertNotNull(jmsMark.getMessages());
assertEquals(1, jmsMark.getMessages().size());
- reader.advance();
- reader.advance();
+ assertTrue(advanceWithRetry(reader));
+ assertTrue(advanceWithRetry(reader));
assertEquals(10, count(QUEUE));
mark.finalizeCheckpoint();
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
index 4ea8f6d317a..277c8415170 100644
---
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
@@ -26,19 +26,19 @@ import static org.junit.Assert.assertNull;
import static org.junit.Assert.assertTrue;
import static org.junit.Assert.fail;
+import jakarta.jms.Connection;
+import jakarta.jms.ConnectionFactory;
+import jakarta.jms.JMSException;
+import jakarta.jms.Message;
+import jakarta.jms.MessageConsumer;
+import jakarta.jms.Queue;
+import jakarta.jms.Session;
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;
diff --git
a/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/MockNonSerializableConnectionFactory.java
b/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/MockNonSerializableConnectionFactory.java
index 752123327e9..b9deef4991e 100644
---
a/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/MockNonSerializableConnectionFactory.java
+++
b/sdks/java/io/jms/src/test/java/org/apache/beam/sdk/io/jms/MockNonSerializableConnectionFactory.java
@@ -17,10 +17,10 @@
*/
package org.apache.beam.sdk.io.jms;
-import javax.jms.Connection;
-import javax.jms.ConnectionFactory;
-import javax.jms.JMSContext;
-import javax.jms.JMSException;
+import jakarta.jms.Connection;
+import jakarta.jms.ConnectionFactory;
+import jakarta.jms.JMSContext;
+import jakarta.jms.JMSException;
public class MockNonSerializableConnectionFactory implements ConnectionFactory
{
@Override
diff --git a/sdks/java/io/messaging-expansion-service/build.gradle
b/sdks/java/io/messaging-expansion-service/build.gradle
index 8693f4597bc..5d30a9270c0 100644
--- a/sdks/java/io/messaging-expansion-service/build.gradle
+++ b/sdks/java/io/messaging-expansion-service/build.gradle
@@ -22,6 +22,7 @@ mainClassName =
"org.apache.beam.sdk.expansion.service.ExpansionService"
applyJavaNature(
automaticModuleName: 'org.apache.beam.sdk.io.messaging.expansion.service',
+ requireJavaVersion: JavaVersion.VERSION_17,
exportJavadoc: false,
validateShadowJar: false,
shadowClosure: {},
diff --git a/sdks/java/io/mqtt/build.gradle b/sdks/java/io/mqtt/build.gradle
index 1f2d808ed7e..b828632165e 100644
--- a/sdks/java/io/mqtt/build.gradle
+++ b/sdks/java/io/mqtt/build.gradle
@@ -19,6 +19,21 @@
plugins { id 'org.apache.beam.module' }
applyJavaNature( automaticModuleName: 'org.apache.beam.sdk.io.mqtt')
+// The JMS IO module pulls ActiveMQ 6.2.5 (Java 17) into the shared dependency
+// map. Because the global resolution strategy forces every mapped version,
+// that 6.2.5 would otherwise be forced onto this Java 11 module via the shared
+// org.apache.activemq coordinates. MQTT IO only needs ActiveMQ as an embedded
+// test broker, so pin it back to the Java 11-compatible 5.x line.
+def activemq5_version = '5.19.5'
+configurations.all {
+ resolutionStrategy.dependencySubstitution {
+ substitute module('org.apache.activemq:activemq-broker') using
module("org.apache.activemq:activemq-broker:${activemq5_version}")
+ substitute module('org.apache.activemq:activemq-client') using
module("org.apache.activemq:activemq-client:${activemq5_version}")
+ substitute module('org.apache.activemq:activemq-kahadb-store') using
module("org.apache.activemq:activemq-kahadb-store:${activemq5_version}")
+ substitute module('org.apache.activemq:activemq-mqtt') using
module("org.apache.activemq:activemq-mqtt:${activemq5_version}")
+ }
+}
+
description = "Apache Beam :: SDKs :: Java :: IO :: MQTT"
ext.summary = "IO to read and write to a MQTT broker."
diff --git a/sdks/python/apache_beam/io/external/xlang_jmsio_it_test.py
b/sdks/python/apache_beam/io/external/xlang_jmsio_it_test.py
index c3f26097376..3f926f1766c 100644
--- a/sdks/python/apache_beam/io/external/xlang_jmsio_it_test.py
+++ b/sdks/python/apache_beam/io/external/xlang_jmsio_it_test.py
@@ -306,7 +306,7 @@ class IbmMqJmsIOTest(_BaseJmsIOTest):
cls.expansion_service_obj = BeamJarExpansionService(
'sdks:java:io:messaging-expansion-service:shadowJar',
classpath=[
- 'com.ibm.mq:com.ibm.mq.allclient:9.3.0.25',
+ 'com.ibm.mq:com.ibm.mq.jakarta.client:9.3.0.25',
'org.json:json:20251224'
])
cls.expansion_service = cls.expansion_service_obj.__enter__()
@@ -327,7 +327,7 @@ class IbmMqJmsIOTest(_BaseJmsIOTest):
uri += '&' + connection_param
return {
'server_uri': uri,
- 'connection_factory_class_name': 'com.ibm.mq.jms.MQConnectionFactory',
+ 'connection_factory_class_name':
'com.ibm.mq.jakarta.jms.MQConnectionFactory',
'username': 'app',
'password': 'admin123'
}
diff --git a/sdks/python/apache_beam/yaml/integration_tests.py
b/sdks/python/apache_beam/yaml/integration_tests.py
index 54be607eaba..becfaa79303 100644
--- a/sdks/python/apache_beam/yaml/integration_tests.py
+++ b/sdks/python/apache_beam/yaml/integration_tests.py
@@ -1163,7 +1163,7 @@ def temp_ibm_mq_server():
yield {
'SERVER_URI':
f'tcp://{host}:{port}?channel=DEV.APP.SVRCONN&queueManager=QM1',
- 'CONNECTION_FACTORY_CLASS_NAME': 'com.ibm.mq.jms.MQConnectionFactory',
+ 'CONNECTION_FACTORY_CLASS_NAME':
'com.ibm.mq.jakarta.jms.MQConnectionFactory',
'USERNAME': 'app',
'PASSWORD': 'admin123',
'SOURCE_QUEUE': 'DEV.QUEUE.1',
diff --git a/sdks/python/apache_beam/yaml/standard_io.yaml
b/sdks/python/apache_beam/yaml/standard_io.yaml
index 15aa9186c0f..fb26bfb6d35 100644
--- a/sdks/python/apache_beam/yaml/standard_io.yaml
+++ b/sdks/python/apache_beam/yaml/standard_io.yaml
@@ -276,7 +276,7 @@
config:
gradle_target: 'sdks:java:io:messaging-expansion-service:shadowJar'
classpath:
- - 'com.ibm.mq:com.ibm.mq.allclient:9.3.0.25'
+ - 'com.ibm.mq:com.ibm.mq.jakarta.client:9.3.0.25'
- 'org.json:json:20251224'