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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 21:40:51