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 51ed48738b37..703cae762172 100644 --- a/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy +++ b/buildSrc/src/main/groovy/org/apache/beam/gradle/BeamModulePlugin.groovy @@ -627,6 +627,7 @@ class BeamModulePlugin implements Plugin { def log4j2_version = "2.20.0" def nemo_version = "0.1" // [bomupgrader] determined by: io.grpc:grpc-netty, consistent with: google_cloud_platform_libraries_bom + def mongodb_version = "5.4.0" def netty_version = "4.1.110.Final" def postgres_version = "42.2.16" // [bomupgrader] determined by: com.google.protobuf:protobuf-java, consistent with: google_cloud_platform_libraries_bom @@ -840,7 +841,10 @@ class BeamModulePlugin implements Plugin { log4j2_log4j12_api : "org.apache.logging.log4j:log4j-1.2-api:$log4j2_version", mockito_core : "org.mockito:mockito-core:4.11.0", mockito_inline : "org.mockito:mockito-inline:4.11.0", - mongo_java_driver : "org.mongodb:mongo-java-driver:3.12.11", + mongodb_driver_legacy : "org.mongodb:mongodb-driver-legacy:$mongodb_version", + mongodb_driver_sync : "org.mongodb:mongodb-driver-sync:$mongodb_version", + mongodb_driver_core : "org.mongodb:mongodb-driver-core:$mongodb_version", + bson : "org.mongodb:bson:$mongodb_version", nemo_compiler_frontend_beam : "org.apache.nemo:nemo-compiler-frontend-beam:$nemo_version", netty_all : "io.netty:netty-all:$netty_version", netty_handler : "io.netty:netty-handler:$netty_version", diff --git a/it/mongodb/build.gradle b/it/mongodb/build.gradle index 6be9b91f5b34..fc6584f6dffa 100644 --- a/it/mongodb/build.gradle +++ b/it/mongodb/build.gradle @@ -34,7 +34,8 @@ dependencies { implementation library.java.testcontainers_base implementation library.java.testcontainers_mongodb implementation library.java.google_code_gson - implementation library.java.mongo_java_driver + implementation library.java.mongodb_driver_sync + implementation library.java.bson implementation library.java.vendored_guava_32_1_2_jre testImplementation library.java.mockito_core diff --git a/sdks/java/extensions/sql/build.gradle b/sdks/java/extensions/sql/build.gradle index 6f34891c2d3f..cbd2a61f0e44 100644 --- a/sdks/java/extensions/sql/build.gradle +++ b/sdks/java/extensions/sql/build.gradle @@ -87,7 +87,9 @@ dependencies { implementation "org.codehaus.janino:janino:3.0.11" implementation "org.codehaus.janino:commons-compiler:3.0.11" implementation library.java.jackson_core - implementation library.java.mongo_java_driver + implementation library.java.mongodb_driver_sync + implementation library.java.mongodb_driver_core + implementation library.java.bson implementation library.java.slf4j_api implementation library.java.joda_time implementation library.java.vendored_guava_32_1_2_jre diff --git a/sdks/java/io/mongodb/build.gradle b/sdks/java/io/mongodb/build.gradle index b9e90082f0dc..028aa27e3990 100644 --- a/sdks/java/io/mongodb/build.gradle +++ b/sdks/java/io/mongodb/build.gradle @@ -27,7 +27,7 @@ ext.summary = "IO to read and write on MongoDB." dependencies { implementation project(path: ":sdks:java:core", configuration: "shadow") implementation library.java.joda_time - implementation library.java.mongo_java_driver + implementation library.java.mongodb_driver_legacy implementation library.java.slf4j_api implementation library.java.vendored_guava_32_1_2_jre testImplementation library.java.junit 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 07cc238c7e6b..5ac3dffa6c88 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 @@ -21,6 +21,7 @@ import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions.checkNotNull; import com.google.auto.value.AutoValue; +import com.mongodb.BasicDBObject; import com.mongodb.DB; import com.mongodb.DBCursor; import com.mongodb.DBObject; @@ -29,7 +30,6 @@ import com.mongodb.gridfs.GridFS; import com.mongodb.gridfs.GridFSDBFile; import com.mongodb.gridfs.GridFSInputFile; -import com.mongodb.util.JSON; import java.io.BufferedReader; import java.io.IOException; import java.io.InputStreamReader; @@ -333,10 +333,10 @@ public void teardown() { public void processElement(final ProcessContext c) throws IOException { Preconditions.checkStateNotNull(gridfs); ObjectId oid = c.element(); - GridFSDBFile file = gridfs.find(oid); + GridFSDBFile file = gridfs.find(Preconditions.checkArgumentNotNull(oid)); Parser parser = Preconditions.checkStateNotNull(parser()); parser.parse( - file, + Preconditions.checkArgumentNotNull(file), new ParserCallback() { @Override public void output(T output, Instant timestamp) { @@ -380,7 +380,7 @@ protected static class BoundedGridFSSource extends BoundedSource { private DBCursor createCursor(GridFS gridfs) { if (spec.filter() != null) { - DBObject query = (DBObject) JSON.parse(spec.filter()); + DBObject query = BasicDBObject.parse(spec.filter()); return gridfs.getFileList(query); } return gridfs.getFileList(); 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 905c7418e26c..a0736dd28aeb 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 @@ -47,7 +47,6 @@ import java.util.NoSuchElementException; import java.util.Optional; import java.util.stream.Collectors; -import javax.net.ssl.SSLContext; import org.apache.beam.sdk.coders.Coder; import org.apache.beam.sdk.coders.SerializableCoder; import org.apache.beam.sdk.io.BoundedSource; @@ -372,9 +371,7 @@ private static MongoClientOptions.Builder getOptions( if (sslEnabled) { optionsBuilder.sslEnabled(sslEnabled).sslInvalidHostNameAllowed(sslInvalidHostNameAllowed); if (ignoreSSLCertificate) { - SSLContext sslContext = SSLUtils.ignoreSSLCertificate(); - optionsBuilder.sslContext(sslContext); - optionsBuilder.socketFactory(sslContext.getSocketFactory()); + optionsBuilder.sslContext(SSLUtils.ignoreSSLCertificate()); } } return optionsBuilder;