Spark加载HDFS海量JSON文件时OutOfMemory异常的解决方法
解决Spark读取大JSON数据源时的OutOfMemoryError(GC overhead limit exceeded)
问题根源
这个问题本质不是Executor内存不足,而是Spark自动推断JSON Schema的过程会触发一个隐含Job:所有Executor会扫描各自负责的文件,提取Schema元数据并序列化后传回Driver。当数据源文件数量极多(这里对应9571个任务/文件),序列化后的结果总大小超过了spark.driver.maxResultSize的默认限制(1GB),进而引发GC overhead超限错误——哪怕你只执行了读取操作,没有触发后续Action,Schema推断本身就是一个隐含的Action。
解决方法
1. 手动指定Schema(最优方案)
直接跳过自动Schema推断,提前定义好JSON的结构,从根源上避免触发Schema推断的Job:
import org.apache.spark.sql.types._ // 根据你的JSON实际结构定义Schema val customSchema = StructType(Seq( StructField("field1", StringType, nullable = true), StructField("field2", IntegerType, nullable = false), StructField("field3", TimestampType, nullable = true) // 补充其他字段 )) val df = spark.read .option("recursiveFileLookup", "true") .schema(customSchema) // 应用手动定义的Schema .json("hdfs:///data")
2. 增大Driver的结果大小限制
如果必须使用自动Schema推断,可以调整spark.driver.maxResultSize参数,允许Driver接收更大的序列化结果:
- 在SparkSession初始化时配置:
val spark = SparkSession.builder() .config("spark.driver.maxResultSize", "2g") // 调整为合适的大小,比如2GB // 其他配置 .getOrCreate() - 或者在提交Spark任务时通过参数指定:
spark-submit --conf spark.driver.maxResultSize=2g --其他参数 你的应用.jar - 若Driver内存足够,也可以设为
unlimited取消限制:.config("spark.driver.maxResultSize", "unlimited")
3. 合并小文件减少任务数
如果数据源包含大量小JSON文件,先合并成大文件,减少分区/任务数量:
- 用HDFS命令合并:
hdfs dfs -getmerge /data /data-merged/combined.json - 或者用Spark先读取(可配合临时增大maxResultSize)后重新分区保存:
// 先读取数据源(可配合临时增大maxResultSize) spark.read.option("recursiveFileLookup", "true").json("hdfs:///data") .repartition(100) // 根据文件总大小调整分区数,比如500GB设为100-200个分区 .write.mode("overwrite").json("hdfs:///data-merged")
之后读取合并后的数据源,任务数大幅减少,Schema推断的结果大小也会降低。
4. 补充调整Driver内存
如果Driver本身内存不足,即使增大maxResultSize仍可能出现GC问题,可以同时提升Driver内存:
spark-submit --driver-memory 8g --其他参数 你的应用.jar
内容的提问来源于stack exchange,提问作者Raphael Mansuy
相关产品推荐
相关产品推荐

