From ac9bf8ae9ddfc65b584f373ab891a44ceb8cba0b Mon Sep 17 00:00:00 2001 From: Alex Bevilacqua Date: Mon, 27 Jul 2026 10:04:40 -0400 Subject: [PATCH 1/2] feat: add MongoDB driver handshake metadata to MongoDBResourceManager Signed-off-by: Alex Bevilacqua --- .../beam/it/mongodb/MongoDBResourceManager.java | 11 ++++++++++- .../it/mongodb/MongoDBResourceManagerTest.java | 5 +++++ .../beam/sdk/io/mongodb/MongoDbGridFSIO.java | 8 ++++++-- .../org/apache/beam/sdk/io/mongodb/MongoDbIO.java | 14 +++++++++----- 4 files changed, 30 insertions(+), 8 deletions(-) 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 8a4f116c8436..0e11f4b8e293 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 @@ 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 extends Serializable { /** Output the object. The default timestamp will be the GridFSFile creation timestamp. */ @@ -203,13 +207,13 @@ static ConnectionConfiguration create( 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 fc2a761b99a8..46c3f8fcd58a 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.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 @@ long getDocumentCount() { 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 long getEstimatedSizeBytes(PipelineOptions pipelineOptions) { 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 List> split( 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 splitKeys; @@ -812,7 +816,7 @@ private MongoClient createClient(Read spec) { spec.ignoreSSLCertificate()) .applyConnectionString(new ConnectionString(uri)) .build(); - return MongoClients.create(settings); + return MongoClients.create(settings, DRIVER_INFO); } } @@ -1012,7 +1016,7 @@ public void createMongoClient() { spec.ignoreSSLCertificate()) .applyConnectionString(new ConnectionString(uri)) .build(); - client = MongoClients.create(settings); + client = MongoClients.create(settings, DRIVER_INFO); } @StartBundle From aaae3b75ed39c198ad870ebecc2e22c556f9b4c8 Mon Sep 17 00:00:00 2001 From: Alex Bevilacqua Date: Mon, 27 Jul 2026 15:42:13 -0400 Subject: [PATCH 2/2] fix: add mongodb-driver-core dependency to it:mongodb module Resolves Gradle dependency analysis violation caused by using ConnectionString and MongoDriverInformation classes without explicitly declaring mongodb-driver-core as a dependency. Co-Authored-By: Claude Sonnet 4.5 --- it/mongodb/build.gradle | 1 + 1 file changed, 1 insertion(+) diff --git a/it/mongodb/build.gradle b/it/mongodb/build.gradle index 960e15af8394..0a78bdda7724 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