公式の MongoDB Spark Connector を使用して、EMR Serverless Spark を MongoDB データベースに接続できます。このトピックでは、MongoDB からの読み取りと書き込みを行うための EMR Serverless Spark の設定方法について説明します。
前提条件
-
EMR Serverless Spark ワークスペースを作成済みであること。詳細については、「ワークスペースの作成」をご参照ください。
-
セルフマネージドの MongoDB データベースを保有していること。
詳細については、「Getting Started with MongoDB」をご参照ください。
制限事項
これらの操作は、以下の EMR Serverless Spark エンジンバージョンでのみサポートされています。
-
esr-4.x:esr-4.1.0 以降
-
esr-3.x:esr-3.1.0 以降
-
esr-2.x:esr-2.5.0 以降
操作手順
ステップ 1:MongoDB コネクタ JAR パッケージの取得と OSS へのアップロード
-
お使いの Spark と MongoDB のバージョンに基づいて、Maven から必要な依存関係をダウンロードします。詳細については、「MongoDB Spark Connector 公式ドキュメント」をご参照ください。この例では MongoDB 5.0.1 を使用しており、以下の JAR パッケージが必要です:
-
mongo-spark-connector_2.12-10.4.1.jar -
mongodb-driver-core-5.0.1.jar -
mongodb-driver-sync-5.0.1.jar -
bson-5.0.1.jar
-
-
ダウンロードした JAR パッケージを Alibaba Cloud Object Storage Service (OSS) にアップロードします。詳細については、「シンプルアップロード」をご参照ください。
ステップ 2:ネットワーク接続の作成
EMR Serverless Spark から MongoDB データベースに接続するには、ネットワーク接続が必要です。ネットワーク接続の詳細については、「EMR Serverless Spark と他の VPC の接続」をご参照ください。
セキュリティグループルールを設定する際は、送信先ポート範囲 に必要なポートのみを開くように設定してください。有効なポート範囲は 1~65535 です。
ステップ 3:MongoDB コレクションからの読み取り
-
ノートブックセッションを作成します。詳細については、「ノートブックセッションの管理」をご参照ください。
セッションの作成時には、サポートされている Engine Version とステップ 2 で作成した Normal Network Connection を選択します。次に、Spark Configuration セクションで、MongoDB Spark Connector を読み込むために以下のパラメーターを追加します。
spark.mongodb.write.connection.uri mongodb://<mongodb_host>:27017 spark.mongodb.read.connection.uri mongodb://<mongodb_host>:27017 spark.emr.serverless.user.defined.jars oss://<bucketname>/path/to/mongo-spark-connector_2.12-10.4.1.jar,oss://<bucketname>/path/to/mongodb-driver-core-5.0.1.jar,oss://<bucketname>/path/to/mongodb-driver-sync-5.0.1.jar,oss://<bucketname>/path/to/bson-5.0.1.jarプレースホルダーの値を、お使いの環境に合わせて変更してください。
パラメーター
説明
例
spark.mongodb.write.connection.uriMongoDB の読み取りおよび書き込み用の接続 URI です。
-
<mongodb_host>:MongoDB サービスの IP アドレスまたはホスト名です。 -
27017:MongoDB のデフォルトポートです。
mongodb://192.168.x.x:27017
spark.mongodb.read.connection.urispark.emr.serverless.user.defined.jarsSpark の外部依存関係をカンマ区切りで指定したリストです。
OSS にアップロードした JAR パッケージです。例:
oss://<yourBucketname>/spark/mongodb/mongo-spark-connector_2.12-10.4.1.jar -
-
Development ページで、[Python] > [notebook] タスクを作成し、作成したノートブックセッションを選択します。
詳細については、「ノートブックセッションの管理」をご参照ください。
-
以下のコードをノートブックのセルにコピーし、パラメーターを変更してから、実行 をクリックします。
df = spark.read \ .format("mongodb") \ .option("database", "<yourDatabase>") \ .option("collection", "<yourCollection>") \ .load() df.printSchema() df.show()パラメーターの説明は以下の表のとおりです。
パラメーター
説明
<yourDatabase>MongoDB データベースの名前です。例:mongo_table
<yourCollection>MongoDB コレクションの名前です。例:MongoCollection
データが正常に返された場合、設定は正しいです。出力例は以下のとおりです:
root |-- _id: string (nullable = true) |-- age: integer (nullable = true) |-- city: string (nullable = true) |-- name: string (nullable = true) +--------------------+---+-----------+-------+ | _id|age| city| name| +--------------------+---+-----------+-------+ |67e65cb60ccf7f94a...| 30|Los Angeles| Bob| |67e65cb60ccf7f94a...| 35| Chicago|Charlie| +--------------------+---+-----------+-------+
ステップ 4:MongoDB コレクションへの書き込み
以下のコードをノートブックの新しいセルにコピーし、パラメーターを変更してから、実行 をクリックします。
from pyspark.sql import Row
data = [
Row(name="Sam", age=25, city="New York"),
Row(name="Charlie", age=35, city="Chicago")
]
df = spark.createDataFrame(data)
df.show()
df.write \
.format("mongodb") \
.option("database", "<yourDatabase>") \
.option("collection", "<yourCollection>") \
.mode("append") \
.save()
次の df.show() の出力は、データが MongoDB に書き込まれる前の状態を示しています。正常に実行されれば、設定が正しく、データがコレクションに追加されたことが確認できます。
+-------+---+--------+
| name|age| city|
+-------+---+--------+
| Sam| 25|New York|
|Charlie| 35| Chicago|
+-------+---+--------+