如何在Synapse Spark Pool中安装MongoDB Spark Connector?
在Azure Synapse Spark Pool中配置MongoDB Spark Connector
步骤1:确认兼容的Connector版本
根据你的Synapse Spark Pool所使用的Spark版本,选择对应的MongoDB Spark Connector版本。例如:
- Spark 3.3.x 对应 Connector 10.x 系列
- Spark 3.2.x 对应 Connector 10.x 系列(需确认具体子版本兼容性)
- Spark 2.4.x 对应 Connector 3.x 系列
步骤2:在Notebook中配置Spark会话依赖
在Synapse Notebook的代码单元格中,通过spark.conf.set指定Connector的Maven坐标,替代--packages参数的作用:
# 替换为你需要的Connector版本,确保Scala版本与Synapse Spark一致(通常为2.12) spark.conf.set("spark.jars.packages", "org.mongodb.spark:mongo-spark-connector_2.12:10.2.1")
步骤3:配置MongoDB Atlas连接参数
同样通过spark.conf.set设置读写MongoDB所需的连接URI:
# 读取配置 spark.conf.set("spark.mongodb.read.connection.uri", "mongodb+srv://<你的用户名>:<你的密码>@<集群地址>/<数据库名>.<集合名>") # 写入配置 spark.conf.set("spark.mongodb.write.connection.uri", "mongodb+srv://<你的用户名>:<你的密码>@<集群地址>/<数据库名>.<集合名>")
注意:替换上述URI中的占位符为你的MongoDB Atlas实际信息
步骤4:使用Spark读写MongoDB
配置完成后,即可通过Spark DataFrame API进行读写操作:
读取数据
df = spark.read.format("mongodb").load() df.show()
写入数据
# 示例:创建测试DataFrame并写入 test_df = spark.createDataFrame([(1, "test"), (2, "demo")], ["id", "value"]) test_df.write.format("mongodb").mode("append").save()
注意事项
- 网络权限:确保MongoDB Atlas的IP白名单已添加Synapse Spark Pool的出站IP地址,或通过VNet对等连接实现网络互通,否则会出现连接超时。
- 依赖冲突处理:若出现依赖包冲突,可通过
spark.jars.excludes排除冲突组件,例如:spark.conf.set("spark.jars.excludes", "org.mongodb:mongodb-driver-core") - 版本兼容性:务必保证Spark版本、Scala版本与MongoDB Spark Connector版本完全兼容,否则会引发运行时错误。
内容的提问来源于stack exchange,提问作者frammnm
相关产品推荐
相关产品推荐

