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

向活跃Spark Session添加MongoDB配置报错,求修改方案

问题:已激活Spark Session配置MongoDB后报错的解决方法

尝试为已激活的Spark Session添加MongoDB配置,代码如下:

val spark = SparkSession.getActiveSession.get
spark.conf.set("spark.mongodb.input.uri",
  "mongodb://hello_admin:hello123@localhost:27017/testdb.products?authSource=admin")
spark.conf.set("spark.mongodb.input.partitioner" ,"MongoPaginateBySizePartitioner")
import com.mongodb.spark._

val customRdd = MongoSpark.load(sc)
println(customRdd.count())
println(customRdd.first.toJson)
println(customRdd.collect().foreach(println))

运行时抛出错误:

java.lang.IllegalArgumentException: Missing database name. Set via the 'spark.mongodb.input.uri' or 'spark.mongodb.input.database' property

而通过SparkSession.builder创建并配置的代码能正常运行:

val spark = SparkSession.builder()
  .master("local")
  .appName("MongoSparkConnectorIntro")
  .config("spark.mongodb.input.uri", "mongodb://hello_admin:hello123@localhost:27017/testdb.products?authSource=admin")
  // .config("spark.mongodb.output.uri", "mongodb://hello_admin:hello123@localhost:27017/testdb.products?authSource=admin")
  .config("spark.mongodb.input.partitioner" ,"MongoPaginateBySizePartitioner")
  .getOrCreate()
val sc = spark.sparkContext
val customRdd = MongoSpark.load(sc)
println(customRdd.count())
println(customRdd.first.toJson)
println(customRdd.collect().foreach(println))

解决方案

问题核心是:MongoSpark.load(sc)读取的是SparkContext层面的配置,而通过spark.conf.set()动态设置的参数仅存在于SparkSession层面,不会自动同步到SparkContext。以下两种修改方式均可解决:

方案一:直接用SparkSession加载(推荐)

改用SparkSession作为MongoSpark.load的参数,确保读取到动态设置的配置:

val spark = SparkSession.getActiveSession.get
spark.conf.set("spark.mongodb.input.uri",
  "mongodb://hello_admin:hello123@localhost:27017/testdb.products?authSource=admin")
spark.conf.set("spark.mongodb.input.partitioner", "MongoPaginateBySizePartitioner")
import com.mongodb.spark._

// 替换sc为spark(SparkSession)
val customRdd = MongoSpark.load(spark)
println(customRdd.count())
println(customRdd.first.toJson)
customRdd.collect().foreach(println)
// 注:原代码外层println会打印Unit,这里直接调用foreach即可

方案二:将SparkSession配置同步到SparkContext

如果必须使用SparkContext加载,需手动同步配置:

val spark = SparkSession.getActiveSession.get
val sc = spark.sparkContext

// 设置SparkSession配置
spark.conf.set("spark.mongodb.input.uri",
  "mongodb://hello_admin:hello123@localhost:27017/testdb.products?authSource=admin")
spark.conf.set("spark.mongodb.input.partitioner", "MongoPaginateBySizePartitioner")

// 同步配置到SparkContext
sc.hadoopConfiguration.set("spark.mongodb.input.uri", spark.conf.get("spark.mongodb.input.uri"))
sc.hadoopConfiguration.set("spark.mongodb.input.partitioner", spark.conf.get("spark.mongodb.input.partitioner"))

import com.mongodb.spark._
val customRdd = MongoSpark.load(sc)
println(customRdd.count())
println(customRdd.first.toJson)
customRdd.collect().foreach(println)

原因补充

  • 通过SparkSession.builder().config()设置参数时,配置会自动同步到SparkContext,因此MongoSpark.load(sc)能读取到参数。
  • 通过spark.conf.set()动态添加的配置仅作用于SparkSession,不会自动同步到SparkContext,导致直接用sc加载时找不到MongoDB的URI配置,触发缺失数据库名的报错。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 07:35:22