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 c91aa1c1d89 feat: add MongoDB driver handshake metadata for Java-based 
client connections (#39504)
c91aa1c1d89 is described below

commit c91aa1c1d89808e3f5919c9aae19b3294a838f2c
Author: Alex Bevilacqua <[email protected]>
AuthorDate: Mon Aug 10 20:57:11 2026 -0400

    feat: add MongoDB driver handshake metadata for Java-based client 
connections (#39504)
---
 it/mongodb/build.gradle                                    |  1 +
 .../org/apache/beam/it/mongodb/MongoDBResourceManager.java | 11 ++++++++++-
 .../apache/beam/it/mongodb/MongoDBResourceManagerTest.java |  5 +++++
 .../org/apache/beam/sdk/io/mongodb/MongoDbGridFSIO.java    |  8 ++++++--
 .../java/org/apache/beam/sdk/io/mongodb/MongoDbIO.java     | 14 +++++++++-----
 5 files changed, 31 insertions(+), 8 deletions(-)

diff --git a/it/mongodb/build.gradle b/it/mongodb/build.gradle
index 960e15af839..0a78bdda772 100644
--- a/it/mongodb/build.gradle
+++ b/it/mongodb/build.gradle
@@ -36,6 +36,7 @@ dependencies {
     implementation library.java.google_code_gson
     implementation library.java.mongo_java_driver
     implementation library.java.mongo_bson
+    implementation library.java.mongodb_driver_core
     implementation library.java.vendored_guava_32_1_2_jre
 
     testImplementation library.java.mockito_core
diff --git 
a/it/mongodb/src/main/java/org/apache/beam/it/mongodb/MongoDBResourceManager.java
 
b/it/mongodb/src/main/java/org/apache/beam/it/mongodb/MongoDBResourceManager.java
index 8a4f116c843..0e11f4b8e29 100644
--- 
a/it/mongodb/src/main/java/org/apache/beam/it/mongodb/MongoDBResourceManager.java
+++ 
b/it/mongodb/src/main/java/org/apache/beam/it/mongodb/MongoDBResourceManager.java
@@ -20,6 +20,8 @@ package org.apache.beam.it.mongodb;
 import static 
org.apache.beam.it.mongodb.MongoDBResourceManagerUtils.checkValidCollectionName;
 import static 
org.apache.beam.it.mongodb.MongoDBResourceManagerUtils.generateDatabaseName;
 
+import com.mongodb.ConnectionString;
+import com.mongodb.MongoDriverInformation;
 import com.mongodb.client.FindIterable;
 import com.mongodb.client.MongoClient;
 import com.mongodb.client.MongoClients;
@@ -57,6 +59,10 @@ public class MongoDBResourceManager extends 
TestContainerResourceManager<MongoDB
 
   private static final String DEFAULT_MONGODB_CONTAINER_NAME = "mongo";
 
+  @VisibleForTesting
+  static final MongoDriverInformation DRIVER_INFO =
+      MongoDriverInformation.builder().driverName("Apache Beam").build();
+
   // A list of available MongoDB Docker image tags can be found at
   // https://hub.docker.com/_/mongo/tags
   private static final String DEFAULT_MONGODB_CONTAINER_TAG = "4.0.18";
@@ -88,7 +94,10 @@ public class MongoDBResourceManager extends 
TestContainerResourceManager<MongoDB
         usingStaticDatabase ? builder.databaseName : 
generateDatabaseName(builder.testId);
     this.connectionString =
         String.format("mongodb://%s:%d", this.getHost(), 
this.getPort(MONGODB_INTERNAL_PORT));
-    this.mongoClient = mongoClient == null ? 
MongoClients.create(connectionString) : mongoClient;
+    this.mongoClient =
+        mongoClient == null
+            ? MongoClients.create(new ConnectionString(connectionString), 
DRIVER_INFO)
+            : mongoClient;
   }
 
   public static Builder builder(String testId) {
diff --git 
a/it/mongodb/src/test/java/org/apache/beam/it/mongodb/MongoDBResourceManagerTest.java
 
b/it/mongodb/src/test/java/org/apache/beam/it/mongodb/MongoDBResourceManagerTest.java
index b3ad34b70ff..73f63fcde08 100644
--- 
a/it/mongodb/src/test/java/org/apache/beam/it/mongodb/MongoDBResourceManagerTest.java
+++ 
b/it/mongodb/src/test/java/org/apache/beam/it/mongodb/MongoDBResourceManagerTest.java
@@ -71,6 +71,11 @@ public class MongoDBResourceManagerTest {
         new MongoDBResourceManager(mongoClient, container, 
MongoDBResourceManager.builder(TEST_ID));
   }
 
+  @Test
+  public void testDriverInfoHasExpectedName() {
+    
assertThat(MongoDBResourceManager.DRIVER_INFO.getDriverNames()).contains("Apache
 Beam");
+  }
+
   @Test
   public void testCreateResourceManagerBuilderReturnsMongoDBResourceManager() {
     assertThat(
diff --git 
a/sdks/java/io/mongodb/src/main/java/org/apache/beam/sdk/io/mongodb/MongoDbGridFSIO.java
 
b/sdks/java/io/mongodb/src/main/java/org/apache/beam/sdk/io/mongodb/MongoDbGridFSIO.java
index 05468a28863..86e5554f274 100644
--- 
a/sdks/java/io/mongodb/src/main/java/org/apache/beam/sdk/io/mongodb/MongoDbGridFSIO.java
+++ 
b/sdks/java/io/mongodb/src/main/java/org/apache/beam/sdk/io/mongodb/MongoDbGridFSIO.java
@@ -23,6 +23,7 @@ import static 
org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Pr
 import com.google.auto.value.AutoValue;
 import com.mongodb.ConnectionString;
 import com.mongodb.MongoClientSettings;
+import com.mongodb.MongoDriverInformation;
 import com.mongodb.client.MongoClient;
 import com.mongodb.client.MongoClients;
 import com.mongodb.client.MongoCursor;
@@ -119,6 +120,9 @@ import org.joda.time.Instant;
  */
 public class MongoDbGridFSIO {
 
+  private static final MongoDriverInformation DRIVER_INFO =
+      MongoDriverInformation.builder().driverName("Apache Beam").build();
+
   /** Callback for the parser to use to submit data. */
   public interface ParserCallback<T> extends Serializable {
     /** Output the object. The default timestamp will be the GridFSFile 
creation timestamp. */
@@ -203,13 +207,13 @@ public class MongoDbGridFSIO {
 
     MongoClient setupMongo() {
       if (uri() == null) {
-        return MongoClients.create();
+        return MongoClients.create(MongoClientSettings.builder().build(), 
DRIVER_INFO);
       }
       MongoClientSettings settings =
           MongoClientSettings.builder()
               .applyConnectionString(new 
ConnectionString(Preconditions.checkStateNotNull(uri())))
               .build();
-      return MongoClients.create(settings);
+      return MongoClients.create(settings, DRIVER_INFO);
     }
 
     GridFSBucket setupGridFS(MongoClient mongo) {
diff --git 
a/sdks/java/io/mongodb/src/main/java/org/apache/beam/sdk/io/mongodb/MongoDbIO.java
 
b/sdks/java/io/mongodb/src/main/java/org/apache/beam/sdk/io/mongodb/MongoDbIO.java
index fc2a761b99a..46c3f8fcd58 100644
--- 
a/sdks/java/io/mongodb/src/main/java/org/apache/beam/sdk/io/mongodb/MongoDbIO.java
+++ 
b/sdks/java/io/mongodb/src/main/java/org/apache/beam/sdk/io/mongodb/MongoDbIO.java
@@ -27,6 +27,7 @@ import com.mongodb.MongoBulkWriteException;
 import com.mongodb.MongoClientSettings;
 import com.mongodb.MongoClientSettings.Builder;
 import com.mongodb.MongoCommandException;
+import com.mongodb.MongoDriverInformation;
 import com.mongodb.client.AggregateIterable;
 import com.mongodb.client.MongoClient;
 import com.mongodb.client.MongoClients;
@@ -145,6 +146,9 @@ public class MongoDbIO {
 
   private static final Logger LOG = LoggerFactory.getLogger(MongoDbIO.class);
 
+  private static final MongoDriverInformation DRIVER_INFO =
+      MongoDriverInformation.builder().driverName("Apache Beam").build();
+
   public static final String ERROR_MSG_QUERY_FN =
       " class is not supported. "
           + "Please provide one of the predefined classes in the MongoDbIO 
package "
@@ -427,7 +431,7 @@ public class MongoDbIO {
                   spec.ignoreSSLCertificate())
               .applyConnectionString(new ConnectionString(uri))
               .build();
-      try (MongoClient mongoClient = MongoClients.create(settings)) {
+      try (MongoClient mongoClient = MongoClients.create(settings, 
DRIVER_INFO)) {
         return getDocumentCount(mongoClient, database, collection);
       } catch (Exception e) {
         return -1;
@@ -459,7 +463,7 @@ public class MongoDbIO {
                   spec.ignoreSSLCertificate())
               .applyConnectionString(new ConnectionString(uri))
               .build();
-      try (MongoClient mongoClient = MongoClients.create(settings)) {
+      try (MongoClient mongoClient = MongoClients.create(settings, 
DRIVER_INFO)) {
         try {
           return getEstimatedSizeBytes(mongoClient, database, collection);
         } catch (MongoCommandException exception) {
@@ -496,7 +500,7 @@ public class MongoDbIO {
                   spec.ignoreSSLCertificate())
               .applyConnectionString(new ConnectionString(uri))
               .build();
-      try (MongoClient mongoClient = MongoClients.create(settings)) {
+      try (MongoClient mongoClient = MongoClients.create(settings, 
DRIVER_INFO)) {
         MongoDatabase mongoDatabase = mongoClient.getDatabase(database);
 
         List<Document> splitKeys;
@@ -812,7 +816,7 @@ public class MongoDbIO {
                   spec.ignoreSSLCertificate())
               .applyConnectionString(new ConnectionString(uri))
               .build();
-      return MongoClients.create(settings);
+      return MongoClients.create(settings, DRIVER_INFO);
     }
   }
 
@@ -1012,7 +1016,7 @@ public class MongoDbIO {
                     spec.ignoreSSLCertificate())
                 .applyConnectionString(new ConnectionString(uri))
                 .build();
-        client = MongoClients.create(settings);
+        client = MongoClients.create(settings, DRIVER_INFO);
       }
 
       @StartBundle

Reply via email to