使用PySpark读取MongoDB数据时出现Py4JJavaError报错如何解决
报错原因分析
- 缺少MongoDB Spark连接器依赖:Spark运行时未加载对应
mongo格式数据源的Java实现类,是该类报错最常见的诱因 - 版本不兼容:Spark版本、Scala编译版本与Mongo Spark连接器版本不匹配,导致类加载异常
- MongoDB连通性异常:MongoDB服务未启动、27017端口被防火墙拦截、目标库/集合名拼写错误、访问权限不足都可能触发Java侧连接报错
- 集群依赖分发失败:集群模式下仅在本地配置了连接器,执行器节点无法加载对应依赖类
可行解决方案
- 优先检查依赖配置,提交任务时主动指定匹配版本的连接器包
Spark 3.x版本(Scala 2.12编译版)提交命令示例:
如果是在交互式环境(Jupyter、pyspark shell)初始化SparkSession时就配置依赖:spark-submit --packages org.mongodb.spark:mongo-spark-connector_2.12:10.1.1 your_script.pyfrom pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("MongoLoadTest") \ .config("spark.jars.packages", "org.mongodb.spark:mongo-spark-connector_2.12:10.1.1") \ .getOrCreate() # 加载数据的代码可保持原有写法 df_train = spark.read.format('mongo')\ .option('spark.mongodb.input.uri','mongodb://127.0.0.1:27017/Quake.quakes').load() df_train.show(5) - 验证MongoDB连接可用性
本地执行Mongo shell命令验证服务可正常访问:
若Mongo开启了身份认证,需要在连接URI中补充认证信息:mongo mongodb://127.0.0.1:27017 # 进入shell后执行以下命令验证集合存在且可访问 use Quake db.quakes.find().limit(5)mongodb://<用户名>:<密码>@127.0.0.1:27017/Quake.quakes?authSource=admin - 确认版本匹配关系
连接器后缀的Scala版本需和Spark编译的Scala版本完全一致,Spark大版本和连接器版本适配规则可参考Mongo官方规范:Spark 3.0~3.3版本使用10.0.x及以上版本连接器,Spark 2.4版本使用2.4.x版本连接器。 - 集群环境额外配置
集群模式下可提前将连接器jar包上传到所有执行器节点的$SPARK_HOME/jars目录,避免运行时动态下载依赖失败,同时要放开所有执行器节点到MongoDB 27017端口的访问限制。
内容的提问来源于stack exchange,提问作者YASH PATEL
相关产品推荐
相关产品推荐

