如何配置sparklyr连接MongoDB?附PySpark可用参考配置
sparklyr 连接 MongoDB 实现方案
问题背景
在本地搭建基于 MongoDB > Apache Spark > RStudio sparklyr 链路的大数据基础设施时,未找到sparklyr连接MongoDB的可行方案,网络上存量的少量相关旧帖均未提供有效解法。MongoDB官方连接器标注支持SparkR,但该包现已从CRAN下架。
当前可正常连通MongoDB的PySpark参考配置如下:
# import SparkSession from the pyspark package from pyspark.sql import SparkSession # initiate the connection my_spark = SparkSession \ .builder \ .appName("Analysis") \ # 注意:原配置项前缀误写了http://,实际使用时需要删除,否则会触发配置无效报错 .config("spark.mongodb.read.connection.uri", "mongodb://127.0.0.1/safricadb.vacancy") \ .config("spark.mongodb.write.connection.uri", "mongodb://127.0.0.1/safricadb.vacancy") \ .config('spark.jars.packages', 'org.mongodb.spark:mongo-spark-connector:10.0.2')\ .getOrCreate() # load df = my_spark.read.format('com.mongodb.spark.sql.connector.MongoTableProvider').load()
当前使用的组件版本:
- MongoDB Community Server 5.0.9
- 集成Hadoop 2.7的Apache Spark 3.3.0
- Mongo Spark Connector 10.0.2
具体实现步骤
sparklyr 本质是通过R接口调用Spark的JVM层能力,完全可以直接复用Spark官方Mongo连接器的所有配置,不需要依赖已经下架的SparkR包,全程不需要额外安装其他R端依赖。
1. 启动Spark会话时传入匹配的配置
启动sparklyr连接时,把和PySpark逻辑完全一致的参数传入配置即可,注意修正PySpark示例里写错的配置项前缀:
library(sparklyr) library(dplyr) # 初始化Spark配置 conf <- spark_config() # 配置MongoDB读写连接地址 conf$`spark.mongodb.read.connection.uri` <- "mongodb://127.0.0.1/safricadb.vacancy" conf$`spark.mongodb.write.connection.uri` <- "mongodb://127.0.0.1/safricadb.vacancy" # 指定Mongo Spark Connector版本,首次运行会自动拉取对应jar包 conf$`spark.jars.packages` <- "org.mongodb.spark:mongo-spark-connector:10.0.2" # 建立本地Spark连接,指定和本地安装版本匹配的Spark 3.3.0 sc <- spark_connect( master = "local", app_name = "Analysis", config = conf, version = "3.3.0" )
2. 读取MongoDB集合数据
用sparklyr内置的通用数据源读取接口,直接指定和PySpark完全一致的Mongo数据源提供者类名即可加载数据,不需要额外开发扩展:
# 读取MongoDB集合为Spark DataFrame vacancy_tbl <- spark_read_source( sc = sc, name = "vacancy", source = "com.mongodb.spark.sql.connector.MongoTableProvider", memory = FALSE ) # 验证数据读取结果,可正常执行dplyr操作 vacancy_tbl %>% select(1:5) %>% head(10)
3. 数据写入参考
如果需要把Spark计算结果写回MongoDB,调用通用写入接口即可:
# 示例:把处理后的表写回MongoDB的新集合 spark_write_source( x = vacancy_tbl, source = "com.mongodb.spark.sql.connector.MongoTableProvider", options = list( "spark.mongodb.write.connection.uri" = "mongodb://127.0.0.1/safricadb.vacancy_processed" ), mode = "overwrite" )
常见问题排查
- 首次运行拉取jar包速度慢或提示类找不到:可以手动把对应版本的mongo-spark-connector jar包放到Spark的jars目录下,删除本地maven/ivy缓存的损坏包后重启R会话重试
- 提示连接拒绝:确认本地MongoDB服务正常启动,已经开启本地免密访问,检查连接URI中的库名、集合名拼写是否正确
- 版本不匹配报错:确保使用的Mongo Spark Connector版本和Spark 3.3.0版本兼容,10.0.2版本是适配Spark 3.3.x的,不需要更换版本
内容的提问来源于stack exchange,提问作者Willian Adamczyk
相关产品推荐
相关产品推荐

