You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark任务写入Parquet文件时容器临时存储超限问题求助

问题分析与解决方案

错误原因

错误日志明确显示Pod临时本地存储用量超出了10Gi的容器配额,核心原因是:
Spark写入AWS S3这类对象存储时,不会直接将数据写入远端,而是先在Executor的本地磁盘生成Parquet临时文件,完成序列化后再批量上传到S3。加上Shuffle过程中产生的临时文件,当单Executor处理的数据量(含临时文件)超过K8s分配的10Gi存储上限时,Pod就会被驱逐。

你的当前配置中,spark.sql.shuffle.partitions=500结合最大20个Executor,每个Executor平均要处理25个分区,再加上Parquet写入时的临时文件膨胀,很容易触发存储配额限制。

无需修改K8s配置的可行解决方案

1. 优化Spark临时存储管理

  • 指定多个临时存储路径(如果CDP工作区节点有额外可用磁盘),分散存储压力:
    SparkSession.builder
      // ... 其他配置
      .config("spark.local.dir", "/tmp,/data/tmp") // 用逗号分隔多个路径
    
  • 开启临时文件自动清理,及时回收Shuffle和输出临时文件:
    .config("spark.cleaner.referenceTracking.cleanCheckpoints", "true")
    .config("spark.shuffle.service.enabled", "true")
    .config("spark.shuffle.service.cleanup.enabled", "true")
    

2. 调整分区数量,降低单Executor负载

当前500个分区分配到20个Executor,每个Executor平均处理25个分区。可以适当增加分区数,让单个Executor处理的分区更少,临时文件体积更小:

.config("spark.sql.shuffle.partitions", 800) // 建议调整为800-1000

同时简化repartition逻辑,直接按业务字段+指定分区数拆分,避免salt列带来的额外复杂度:

df.repartition(800, "date_year", "date_month", "date_day")
  .write.partitionBy("date_year", "date_month")
  .mode("overwrite")
  .parquet(SOME__PATH)

3. 优化S3写入的本地缓存策略

启用S3快速上传优化,减少本地临时文件的占用:

.config("spark.hadoop.fs.s3a.fast.upload", "true")
.config("spark.hadoop.fs.s3a.multipart.size", "104857600") // 设置100MB分片上传,减少本地缓存
.config("spark.hadoop.fs.s3a.buffer.dir", "/tmp") // 指定可用的大临时目录

如果数据字典基数不大,可以关闭Parquet字典编码,减少临时文件体积:

.write.option("parquet.enable.dictionary", "false")

4. 调整Executor并发数,降低单节点存储压力

当前每个Executor配置4核,默认同时运行4个任务,会产生多份临时文件。可以减少核数,同时提升单Executor内存(保持总资源不变),降低并发任务数:

.config("spark.executor.cores", 2)
.config("spark.executor.memory", "12G")

这样每个Executor同时处理2个任务,本地临时文件的总占用量会显著降低。

内容的提问来源于stack exchange,提问作者Ryan

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.04 05:30:50