You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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)可能因沙箱网络限制或版本兼容问题失败。

具体步骤

  1. 确认版本兼容性
    HDP 2.6.5对应的Spark版本为2.3.x,需匹配Scala 2.11版本的MongoDB连接器,正确包版本为org.mongodb.spark:mongo-spark-connector_2.11:2.3.2。

  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/)
  3. 修正代码配置
    可以简化数据源格式的指定,使用"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()
    
  4. 验证依赖加载
    运行spark-shell,执行import com.mongodb.spark.sql.DefaultSource,若未报错则说明依赖已正确加载。


内容的提问来源于stack exchange,提问作者Codew1213

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.09 01:10:27