向活跃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
相关产品推荐
相关产品推荐

