Skip to content

Commit 3b5dd20

Browse files
committed
feat: add MongoDB driver handshake metadata to MongoDBResourceManager
Signed-off-by: Alex Bevilacqua <alex@alexbevi.com>
1 parent d6ab992 commit 3b5dd20

4 files changed

Lines changed: 28 additions & 8 deletions

File tree

it/mongodb/src/main/java/org/apache/beam/it/mongodb/MongoDBResourceManager.java

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@
2020
import static org.apache.beam.it.mongodb.MongoDBResourceManagerUtils.checkValidCollectionName;
2121
import static org.apache.beam.it.mongodb.MongoDBResourceManagerUtils.generateDatabaseName;
2222

23+
import com.mongodb.MongoDriverInformation;
2324
import com.mongodb.client.FindIterable;
2425
import com.mongodb.client.MongoClient;
2526
import com.mongodb.client.MongoClients;
@@ -57,6 +58,10 @@ public class MongoDBResourceManager extends TestContainerResourceManager<MongoDB
5758

5859
private static final String DEFAULT_MONGODB_CONTAINER_NAME = "mongo";
5960

61+
@VisibleForTesting
62+
static final MongoDriverInformation DRIVER_INFO =
63+
MongoDriverInformation.builder().driverName("Apache Beam").build();
64+
6065
// A list of available MongoDB Docker image tags can be found at
6166
// https://hub.docker.com/_/mongo/tags
6267
private static final String DEFAULT_MONGODB_CONTAINER_TAG = "4.0.18";
@@ -88,7 +93,8 @@ private MongoDBResourceManager(Builder builder) {
8893
usingStaticDatabase ? builder.databaseName : generateDatabaseName(builder.testId);
8994
this.connectionString =
9095
String.format("mongodb://%s:%d", this.getHost(), this.getPort(MONGODB_INTERNAL_PORT));
91-
this.mongoClient = mongoClient == null ? MongoClients.create(connectionString) : mongoClient;
96+
this.mongoClient =
97+
mongoClient == null ? MongoClients.create(connectionString, DRIVER_INFO) : mongoClient;
9298
}
9399

94100
public static Builder builder(String testId) {

it/mongodb/src/test/java/org/apache/beam/it/mongodb/MongoDBResourceManagerTest.java

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,7 @@
2727
import static org.mockito.Mockito.verify;
2828
import static org.mockito.Mockito.when;
2929

30+
import com.mongodb.MongoDriverInformation;
3031
import com.mongodb.MongoBulkWriteException;
3132
import com.mongodb.client.MongoClient;
3233
import com.mongodb.client.MongoCollection;
@@ -71,6 +72,11 @@ public void setUp() {
7172
new MongoDBResourceManager(mongoClient, container, MongoDBResourceManager.builder(TEST_ID));
7273
}
7374

75+
@Test
76+
public void testDriverInfoHasExpectedName() {
77+
assertThat(MongoDBResourceManager.DRIVER_INFO.getDriverNames()).contains("Apache Beam");
78+
}
79+
7480
@Test
7581
public void testCreateResourceManagerBuilderReturnsMongoDBResourceManager() {
7682
assertThat(

sdks/java/io/mongodb/src/main/java/org/apache/beam/sdk/io/mongodb/MongoDbGridFSIO.java

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@
2323
import com.google.auto.value.AutoValue;
2424
import com.mongodb.ConnectionString;
2525
import com.mongodb.MongoClientSettings;
26+
import com.mongodb.MongoDriverInformation;
2627
import com.mongodb.client.MongoClient;
2728
import com.mongodb.client.MongoClients;
2829
import com.mongodb.client.MongoCursor;
@@ -119,6 +120,9 @@
119120
*/
120121
public class MongoDbGridFSIO {
121122

123+
private static final MongoDriverInformation DRIVER_INFO =
124+
MongoDriverInformation.builder().driverName("Apache Beam").build();
125+
122126
/** Callback for the parser to use to submit data. */
123127
public interface ParserCallback<T> extends Serializable {
124128
/** Output the object. The default timestamp will be the GridFSFile creation timestamp. */
@@ -203,13 +207,13 @@ static ConnectionConfiguration create(
203207

204208
MongoClient setupMongo() {
205209
if (uri() == null) {
206-
return MongoClients.create();
210+
return MongoClients.create(DRIVER_INFO);
207211
}
208212
MongoClientSettings settings =
209213
MongoClientSettings.builder()
210214
.applyConnectionString(new ConnectionString(Preconditions.checkStateNotNull(uri())))
211215
.build();
212-
return MongoClients.create(settings);
216+
return MongoClients.create(settings, DRIVER_INFO);
213217
}
214218

215219
GridFSBucket setupGridFS(MongoClient mongo) {

sdks/java/io/mongodb/src/main/java/org/apache/beam/sdk/io/mongodb/MongoDbIO.java

Lines changed: 9 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@
2626
import com.mongodb.MongoBulkWriteException;
2727
import com.mongodb.MongoClientSettings;
2828
import com.mongodb.MongoClientSettings.Builder;
29+
import com.mongodb.MongoDriverInformation;
2930
import com.mongodb.MongoCommandException;
3031
import com.mongodb.client.AggregateIterable;
3132
import com.mongodb.client.MongoClient;
@@ -145,6 +146,9 @@ public class MongoDbIO {
145146

146147
private static final Logger LOG = LoggerFactory.getLogger(MongoDbIO.class);
147148

149+
private static final MongoDriverInformation DRIVER_INFO =
150+
MongoDriverInformation.builder().driverName("Apache Beam").build();
151+
148152
public static final String ERROR_MSG_QUERY_FN =
149153
" class is not supported. "
150154
+ "Please provide one of the predefined classes in the MongoDbIO package "
@@ -427,7 +431,7 @@ long getDocumentCount() {
427431
spec.ignoreSSLCertificate())
428432
.applyConnectionString(new ConnectionString(uri))
429433
.build();
430-
try (MongoClient mongoClient = MongoClients.create(settings)) {
434+
try (MongoClient mongoClient = MongoClients.create(settings, DRIVER_INFO)) {
431435
return getDocumentCount(mongoClient, database, collection);
432436
} catch (Exception e) {
433437
return -1;
@@ -459,7 +463,7 @@ public long getEstimatedSizeBytes(PipelineOptions pipelineOptions) {
459463
spec.ignoreSSLCertificate())
460464
.applyConnectionString(new ConnectionString(uri))
461465
.build();
462-
try (MongoClient mongoClient = MongoClients.create(settings)) {
466+
try (MongoClient mongoClient = MongoClients.create(settings, DRIVER_INFO)) {
463467
try {
464468
return getEstimatedSizeBytes(mongoClient, database, collection);
465469
} catch (MongoCommandException exception) {
@@ -496,7 +500,7 @@ public List<BoundedSource<Document>> split(
496500
spec.ignoreSSLCertificate())
497501
.applyConnectionString(new ConnectionString(uri))
498502
.build();
499-
try (MongoClient mongoClient = MongoClients.create(settings)) {
503+
try (MongoClient mongoClient = MongoClients.create(settings, DRIVER_INFO)) {
500504
MongoDatabase mongoDatabase = mongoClient.getDatabase(database);
501505

502506
List<Document> splitKeys;
@@ -812,7 +816,7 @@ private MongoClient createClient(Read spec) {
812816
spec.ignoreSSLCertificate())
813817
.applyConnectionString(new ConnectionString(uri))
814818
.build();
815-
return MongoClients.create(settings);
819+
return MongoClients.create(settings, DRIVER_INFO);
816820
}
817821
}
818822

@@ -1012,7 +1016,7 @@ public void createMongoClient() {
10121016
spec.ignoreSSLCertificate())
10131017
.applyConnectionString(new ConnectionString(uri))
10141018
.build();
1015-
client = MongoClients.create(settings);
1019+
client = MongoClients.create(settings, DRIVER_INFO);
10161020
}
10171021

10181022
@StartBundle

0 commit comments

Comments
 (0)