HDP 2.6.5中Spark写入MongoDB遇ClassNotFoundException问题求助
问题:HDP 2.6.5沙箱中Spark写入MongoDB报ClassNotFoundException错误
在HDP 2.6.5沙箱环境下,尝试将HDFS中的JSON文件数据写入MongoDB时,抛出java.lang.ClassNotFoundException: com.mongodb.spark.sql.DefaultSource.DefaultSource错误。
原代码
# HDFS-Mongo 用于将HDFS中的JSON文件写入MongoDB from pyspark.sql import SparkSession my_spark = SparkSession \ .builder \ .appName("testdb") \ .config("spark.mongodb.input.uri", "mongodb://127.0.0.1/testdb.test1") \ .config("spark.mongodb.output.uri", "mongodb://127.0.0.1/testdb.test1") \ .getOrCreate() df = my_spark.read.option("multiline", "true").json("hdfs://sandbox-hdp.hortonworks.com:8020/user/root/output2.json") df.count() df.printSchema() df.write.format("com.mongodb.spark.sql.DefaultSource").mode("append").option("database","testdb").option("collection", "test1").save()
报错信息
Traceback (most recent call last): File "hdfs_mongo.py", line 19, in <module> df.write.format("com.mongodb.spark.sql.DefaultSource").mode("append").option("database","testdb").option("collection", "test1").save() File "/usr/local/lib/python3.6/site-packages/pyspark/sql/readwriter.py", line 738, in save self._jwrite.save() File "/usr/local/lib/python3.6/site-packages/py4j/java_gateway.py", line 1322, in __call__ answer, self.gateway_client, self.target_id, self.name) File "/usr/local/lib/python3.6/site-packages/pyspark/sql/utils.py", line 111, in deco return f(*a, **kw) File "/usr/local/lib/python3.6/site-packages/py4j/protocol.py", line 328, in get_return_value format(target_id, ".", name), value) py4j.protocol.Py4JJavaError: An error occurred while calling o44.save. : java.lang.ClassNotFoundException: Failed to find data source: com.mongodb.spark.sql.DefaultSource. Please find packages at http://spark.apache.org/third-party-projects.html at org.apache.spark.sql.errors.QueryExecutionErrors$.failedToFindDataSourceError(QueryExecutionErrors.scala:443) at org.apache.spark.sql.execution.datasources.DataSource$.lookupDataSource(DataSource.scala:670) at org.apache.spark.sql.execution.datasources.DataSource$.lookupDataSourceV2(DataSource.scala:720) at org.apache.spark.sql.DataFrameWriter.lookupV2Provider(DataFrameWriter.scala:852) at org.apache.spark.sql.DataFrameWriter.saveInternal(DataFrameWriter.scala:256) at org.apache.spark.sql.DataFrameWriter.save(DataFrameWriter.scala:247) at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) at java.lang.reflect.Method.invoke(Method.java:498) at py4j.reflection.MethodInvoker.invoke(MethodInvoker.java:244) at py4j.reflection.ReflectionEngine.invoke(ReflectionEngine.java:357) at py4j.Gateway.invoke(Gateway.java:282) at py4j.commands.AbstractCommand.invokeMethod(AbstractCommand.java:132) at py4j.commands.CallCommand.execute(CallCommand.java:79) at py4j.ClientServerConnection.waitForCommands(ClientServerConnection.java:182) at py4j.ClientServerConnection.run(ClientServerConnection.java:106) at java.lang.Thread.run(Thread.java:748) Caused by: java.lang.ClassNotFoundException: com.mongodb.spark.sql.DefaultSource.DefaultSource at java.net.URLClassLoader.findClass(URLClassLoader.java:381) at java.lang.ClassLoader.loadClass(ClassLoader.java:424) at java.lang.ClassLoader.loadClass(ClassLoader.java:357) at org.apache.spark.sql.execution.datasources.DataSource$.$anonfun$lookupDataSource$5(DataSource.scala:656) at scala.util.Try$.apply(Try.scala:213) at org.apache.spark.sql.execution.datasources.DataSource$.$anonfun$lookupDataSource$4(DataSource.scala:656) at scala.util.Failure.orElse(Try.scala:224) at org.apache.spark.sql.execution.datasources.DataSource$.lookupDataSource(DataSource.scala:656) ... 16 more
已尝试的操作
- 卸载并重装Python3.6
- 添加
spark.jars.packages配置测试,代码如下:
from pyspark.sql import SparkSession spark = SparkSession.builder.appName("myApp") \ .config("spark.mongodb.input.uri", "mongodb://xxx.xxx.xxx.xxx:27017/sample1.zips") \ .config("spark.mongodb.output.uri", "mongodb://xxx.xxx.xxx.xxx:27017/sample1.zips") \ .config('spark.jars.packages', 'org.mongodb.spark:mongo-spark-connector_2.11:2.3.2') \ .getOrCreate() df = spark.read.format("com.mongodb.spark.sql.DefaultSource").load() df.printSchema()
解决方案
核心原因
Spark无法找到MongoDB Spark连接器的依赖包,HDP沙箱默认未预装该组件,且动态加载依赖(通过spark.jars.packages)可能因沙箱网络限制或版本兼容问题失败。
具体步骤
确认版本兼容性
HDP 2.6.5对应的Spark版本为2.3.x,需匹配Scala 2.11版本的MongoDB连接器,正确包版本为org.mongodb.spark:mongo-spark-connector_2.11:2.3.2。手动部署依赖包
- 获取对应版本的
mongo-spark-connector_2.11-2.3.2.jar和依赖的mongo-java-driver-3.8.2.jar - 将两个jar包上传至HDP沙箱的Spark jars目录(通常路径为
/usr/hdp/current/spark2-client/jars/)
- 获取对应版本的
修正代码配置
可以简化数据源格式的指定,使用"mongodb"简写替代完整类名,同时避免重复配置数据库和集合(已通过uri指定):from pyspark.sql import SparkSession my_spark = SparkSession \ .builder \ .appName("testdb") \ .config("spark.mongodb.output.uri", "mongodb://127.0.0.1/testdb.test1") \ .getOrCreate() df = my_spark.read.option("multiline", "true").json("hdfs://sandbox-hdp.hortonworks.com:8020/user/root/output2.json") # 使用简写格式写入MongoDB df.write.format("mongodb").mode("append").save()验证依赖加载
运行spark-shell,执行import com.mongodb.spark.sql.DefaultSource,若未报错则说明依赖已正确加载。
内容的提问来源于stack exchange,提问作者Codew1213
相关产品推荐
相关产品推荐

