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'
 
 

Reply via email to