SparkSession读取CSV文件时如何设置分区数,解决任务过多OOM问题
问题根因
你遇到的
java.lang.OutOfMemoryError: GC overhead limit exceeded报错,本质是读CSV生成的分区数过多,Driver端需要维护的分区元数据量过大,导致GC频繁超出阈值触发。Spark读取文件的初始分区数由底层文件分片逻辑决定,大量小文件或者默认分片规则过细都会导致分区数爆炸。
可用解决方案(Spark 2.3.1均支持)
调整读取分区的最大字节阈值
通过spark.sql.files.maxPartitionBytes参数控制单个分区的最大数据量,默认值为134217728(128MB),调大该值即可直接减少分区数。可在读取文件前动态配置:// 示例将单分区最大阈值调整为256MB,分区数会直接缩减为原来的1/2左右 sparkSession.conf.set("spark.sql.files.maxPartitionBytes", 268435456L) val df = sparkSession.read.csv(dfwAbsHdfsPath)小文件场景下开启自动合并
如果你的数据源是大量小文件,可以调整spark.sql.files.openCostInBytes参数,该参数代表打开一个文件的估算开销(默认4MB),调大后Spark会自动将多个小文件合并到同一个分区中,避免生成过多分区:// 示例将打开文件开销调整为8MB,更积极的合并小文件到同一分区 sparkSession.conf.set("spark.sql.files.openCostInBytes", 8388608L) val df = sparkSession.read.csv(dfwAbsHdfsPath)读取后手动合并分区
直接在读取完DataFrame后调用coalesce方法手动指定分区数,该方法不会触发shuffle,适合单纯减少分区的场景:// 示例直接将分区数缩减到100,可根据你的集群资源灵活调整数值 val df = sparkSession.read.csv(dfwAbsHdfsPath).coalesce(100)
注意事项
- 不要将
spark.sql.files.maxPartitionBytes设置过大,否则单个分区数据量过高会导致后续计算任务出现OOM,建议根据Executor单核心可处理的数据量调整,通常设置为128MB~512MB区间即可。 - 如果你的CSV是gzip等不可拆分的压缩格式,单个不可拆分文件会作为独立分区,该场景优先使用
coalesce或者调整spark.sql.files.openCostInBytes来合并分区。
内容的提问来源于stack exchange,提问作者Anoop Deshpande
相关产品推荐
相关产品推荐

