如何在Scala中使用Spark Context连接MongoDB多个集合
Spark 读取 MongoDB 多集合溢出问题解决方案
你遇到的溢出错误,绝大多数场景是因为你将多个集合的读取结果一次性存入List,所有DataFrame的元数据、缓存数据直接占满了Driver端堆内存。按以下步骤调整即可:
- 不要一次性批量加载所有集合,改为逐个读取、逐个处理
你已经跑通了单集合的读取逻辑,只需要把集合名称存到List里循环遍历即可,不要在List里存读取返回的DataFrame:
# PySpark 示例 from pyspark.sql import SparkSession # SparkSession 全局只初始化一次 spark = SparkSession.builder \ .appName("ReadMultiMongoColl") \ .config("spark.mongodb.input.uri", "mongodb://<你的Mongo地址>/<你的库名>") \ # 本地测试可先调大Driver内存,避免小数据量也报溢出 .config("spark.driver.memory", "4g") \ .getOrCreate() # 这里只存集合名称,不要做读取操作 coll_name_list = ["coll1", "coll2", "coll3"] for coll_name in coll_name_list: # 每次循环只读取当前一个集合 current_df = spark.read.format("mongo") \ .option("collection", coll_name) \ .load() # 此处写当前集合的处理逻辑:过滤、统计、写入存储等 current_df.show() # 若当前集合处理完后不需要再使用,主动释放缓存 # current_df.unpersist()
- 若确实需要合并多集合数据,不要在内存中一次性Union所有集合的DataFrame,逐个读取后追加写入目标存储即可。
内容的提问来源于stack exchange,提问作者Jevon Coleman
相关产品推荐
相关产品推荐

