Scala-Spark作业中partitionBy列cardinality为1导致Executor OOM问题排查
问题分析与方案验证
核心疑问解答
- 是,partitionBy列cardinality为1是直接诱因:Spark的
partitionBy基于分区列哈希值分配任务,当wk和dt的基数都为1时,所有数据哈希值完全相同,会被集中分配到同一个Task/Executor上。生产环境解压后超500GB的数据全部压在单个Executor上,直接引发内存过载(OOM);同时Spark写入Parquet时会在本地生成临时文件,单个Executor的临时存储(ephemeral storage)被瞬间占满,触发K8s存储驱逐机制,最终导致Executor丢失、阶段失败。 - 临时存储不足是该问题的连锁反应:单Executor处理全量数据时,写入产生的临时文件远超过容器请求的存储配额,直接引发Executor被杀死。
拟尝试方案可行性验证
方案1:partitionBy前增加repartition(200)
完全可行,是针对性解决问题的最优方案:
repartition(200)会强制将数据拆分为200个并行分区,后续partitionBy时,每个分区的数据会被并行写入到同一个wk=xxx/dt=xxx目录下。单个Task处理的数据量从500GB+降至2-3GB,能有效分散Executor的内存和存储负载,彻底解决OOM和临时存储不足问题。- 最终生成的Parquet数据依然保留
wk和dt的分区目录结构,完全满足后续滚动聚合周数据的需求。 - 注意:分区数可根据集群资源调整(比如按
Executor数量 * 每个Executor核数设置),避免分区过多导致小文件问题,200是合理的起步值。
方案2:不使用partitionBy,直接写入指定文件夹路径
同样可行,适合需要更灵活控制路径的场景:
- 直接写入
s3://your-bucket/path/wk=xxx/dt=xxx/格式的路径下,最终效果和partitionBy一致,Spark后续读取时依然能识别为分区表。 - 建议配合
repartition使用,提前拆分数据为多个并行分区,避免单Task处理全量数据的问题。 - 注意:需确保路径严格符合Hive分区格式(
key=value),否则后续滚动聚合时无法自动识别分区。
额外优化建议
- 调整Executor临时存储配额:在K8s配置中增加
spark.kubernetes.executor.ephemeralStorage.request(比如设置为60Gi),作为辅助优化,但核心还是要靠数据分区拆分。 - 开启Parquet压缩:写入时设置
spark.sql.parquet.compression.codec=snappy,减少临时文件大小,降低存储压力。 - 检查广播变量:确认广播数据集确实极小,避免占用过多Executor内存(虽非本次问题核心,但能进一步优化稳定性)
内容的提问来源于stack exchange,提问作者Abhishek
相关产品推荐
相关产品推荐

