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