Spark Structured Streaming连接MongoDB报错:找不到mongodb数据源
解决Spark Structured Streaming连接MongoDB时"Failed to find data source: mongodb"错误
核心问题分析
报错本质是Spark无法识别mongodb数据源格式,根源在于MongoDB Spark连接器的依赖包配置错误,导致连接器未被正确引入。你的代码中spark.jars.packages的版本后缀格式不符合规范,这是关键问题。
具体修复方案
1. 修正连接器依赖包的版本格式
Spark连接器的包名需要包含对应Scala版本后缀,Spark 3.3.1默认使用Scala 2.12,正确的依赖配置应为:
.config('spark.jars.packages', 'org.mongodb.spark:mongo-spark-connector_2.12:3.3.1')
你之前写的_10.0.5是错误格式,连接器命名规则为mongo-spark-connector_<Scala版本>:<连接器版本>,而非MongoDB服务器版本。
2. 调整流读取配置(可选但推荐)
结构化流读取MongoDB时,建议简化URI配置,同时注意连续处理模式(continuous)可能不被MongoDB连接器支持,可替换为处理时间触发模式:
read_from_mongo = (spark .readStream .format("mongodb") .option("uri", "mongodb://admin:admin@localhost:27017/first_db.first_collection") .load() .writeStream .format("console") .trigger(processingTime="1 second") .outputMode("append") .start())
将数据库和集合直接写入URI,能减少配置项数量。
3. 手动指定jar包(若依赖自动下载失败)
如果运行时仍无法加载依赖,可手动下载对应版本(Scala 2.12、Spark 3.3.1适配)的MongoDB Spark连接器jar包,启动Spark时通过参数指定:
spark-submit --jars mongo-spark-connector_2.12-3.3.1.jar your_script.py
验证修复后的完整代码
from pyspark.sql import SparkSession if __name__ == "__main__": spark = (SparkSession .builder .appName("Streaming from mongo db") .master("local[3]") .config('spark.jars.packages', 'org.mongodb.spark:mongo-spark-connector_2.12:3.3.1') .config("spark.streaming.stopGracefullyOnShutdown", "true") .getOrCreate()) read_from_mongo = (spark .readStream .format("mongodb") .option("uri", "mongodb://admin:admin@localhost:27017/first_db.first_collection") .load() .writeStream .format("console") .trigger(processingTime="1 second") .outputMode("append") .start()) read_from_mongo.awaitTermination()
注意添加awaitTermination(),否则程序会直接退出,无法持续监听流数据。
内容的提问来源于stack exchange,提问作者MUSTAPHAAMINE DEBBIH
相关产品推荐
相关产品推荐

