Spark流处理写入MongoDB报错:找不到数据源com.mongodb.spark.sql.DefaultSource
解决Spark流写入MongoDB时数据源找不到的报错
问题现象
执行Spark脚本从Kafka读取流数据并写入MongoDB时,出现如下报错:
py4j.protocol.Py4JJavaError: An error occurred while calling o66.start. : java.lang.ClassNotFoundException: Failed to find data source: com.mongodb.spark.sql.DefaultSource. Please find packages at https://spark.apache.org/third-party-projects.html
核心原因
- SparkSession配置中重复设置
spark.jars.packages,后一项覆盖了前一项,导致MongoDB Spark连接器未被加载 - 执行
spark-submit时仅指定了Kafka相关依赖包,缺少MongoDB连接器依赖 - 旧版全类名格式可能存在兼容性问题
解决方案
1. 合并SparkSession的依赖包配置
将脚本中两次设置spark.jars.packages的代码合并为一行,用逗号分隔多个依赖包,避免覆盖:
spark = SparkSession \ .builder \ .appName("Spark") \ .master('local')\ .config('spark.jars.packages', 'org.mongodb.spark:mongo-spark-connector:10.1.1,org.apache.spark:spark-sql-kafka-0-10_2.12:3.3.0') \ .getOrCreate()
2. 简化WriteStream的格式参数
使用MongoDB Spark连接器新版本支持的简洁格式"mongodb"替代全类名,避免类路径错误:
dsw = ( pktDF.writeStream .format("mongodb") # 替换原有的com.mongodb.spark.sql.DefaultSource .queryName("ToMDB") .option("checkpointLocation", "/tmp/pyspark7/") .option("forceDeleteTempCheckpointLocation", "true") .option('spark.mongodb.connection.uri','mongodb+srv://<>:<>@batchdata.niq660y.mongodb.net/?retryWrites=true&w=majority') .option('spark.mongodb.database', 'BatchData') .option('spark.mongodb.collection', 'Traffic') .outputMode("append") .start().awaitTermination());
3. 正确执行Spark Submit命令
由于脚本中已经配置了完整的依赖包,直接执行即可:
spark-submit spark-consumer-pcap.py
若仍存在依赖加载问题,可在命令行显式指定所有依赖:
spark-submit --packages org.mongodb.spark:mongo-spark-connector:10.1.1,org.apache.spark:spark-sql-kafka-0-10_2.12:3.3.0 spark-consumer-pcap.py
4. 版本兼容性检查
确保MongoDB Spark Connector版本与Spark版本匹配:
- 你使用的Spark 3.3.0与MongoDB Spark Connector 10.1.1是兼容的(Connector 10.x适配Spark 3.2+版本)
内容的提问来源于stack exchange,提问作者Muhammad Hani
相关产品推荐
相关产品推荐

