Spark读取阶段溢写问题:如何避免spark.read出现内存/磁盘溢写?
Spark数据读取阶段磁盘溢写问题分析与优化方案
为什么会触发磁盘溢写?
1. 内存分配比例不合理
Spark内存默认按5:5划分存储内存(缓存RDD、广播变量)和执行内存(Shuffle、聚合计算)。你的任务单任务内存溢写达5GB,核心原因是执行内存不足:
- 单任务输入虽仅≤137MB,但汇总操作可能会展开嵌套结构、生成大量中间键值对,数据量远超执行内存的承载上限。
- 256GB节点内存需扣除操作系统、YARN/Spark Driver的预留内存,实际分给Executor的内存被压缩;若一个节点跑多个Executor,单Executor内存会进一步拆分,执行内存空间更紧张。
2. 中间数据严重膨胀
Shuffle Write达1935GB,是原始245GB数据的近8倍,说明数据处理过程中膨胀严重:
- gzip压缩的Parquet解压后本身会变大,若汇总逻辑涉及展开数组、Map等嵌套字段,数据量会剧增。
- 若汇总前存在宽依赖操作(如join、预转换),会生成大量中间数据,直接撑爆Executor内存,只能溢写到磁盘。
3. 任务并行度不匹配
集群总共有31*32=992个vCPU,但如果Spark任务并行度设置过低,单个任务的处理负载会过大:
- 单任务Shuffle Write达1.4GB,加上5GB内存溢写,远超单任务合理处理范围(通常建议100-200MB)。并行度不足会让任务压力集中,内存自然不够用。
4. Parquet读取的限制
gzip是不可拆分的压缩格式,Spark无法将单个大Parquet文件拆分为多个任务处理;即便单任务输入≤137MB,若总文件数少、任务数不足,单任务的处理压力仍会偏高。另外,若列裁剪、谓词下推未生效,会读取不必要的字段,进一步占用内存。
具体优化建议
1. 调整内存分配参数
- 提高执行内存占比:将
spark.memory.fraction设为0.75,减少存储内存预留(你的任务以计算为主,无需大量缓存):spark.memory.fraction=0.75 - 配套调整存储内存占比:设置
spark.memory.storageFraction=0.25,确保执行内存有足够空间:spark.memory.storageFraction=0.25 - 合理分配Executor资源:每个节点32vCPU,建议每个Executor分配8vCPU、64GB内存(避免内存碎片化),单节点可跑4个Executor,总Executor数为124,对应参数:
spark.executor.cores=8 spark.executor.memory=64g spark.driver.memory=48g # 可根据Driver负载调整为32-64GB
2. 控制数据膨胀
- 优化汇总逻辑:仅读取需要的字段(开启Parquet列裁剪),提前过滤冗余行(谓词下推),避免读取不必要的数据。
- 提前局部聚合:若涉及group by操作,先在
mapPartitions内做小范围聚合,减少Shuffle阶段的数据量。
3. 优化任务并行度
- 设置Shuffle并行度:将
spark.sql.shuffle.partitions设为总CPU核心数的2-3倍(约1984),让Shuffle任务更细粒度,单任务处理的数据量更小:spark.sql.shuffle.partitions=1984 - 预处理Parquet文件:将大的gzip压缩Parquet拆分为64-128MB的小文件,或改用可拆分的压缩格式(如Snappy),让Spark生成更多并行任务。
4. 提升S3读取性能
- 开启预读:设置
spark.hadoop.fs.s3a.readahead.range=134217728(128MB),减少S3请求次数:spark.hadoop.fs.s3a.readahead.range=134217728 - 启用向量化读取:开启
spark.sql.parquet.enableVectorizedReader=true,提升Parquet读取效率:spark.sql.parquet.enableVectorizedReader=true
内容的提问来源于stack exchange,提问作者Graeme Wallace
相关产品推荐
相关产品推荐

