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

PySpark聚合后.count()正常但写入Parquet失败,求原因分析

问题描述

使用以下Spark配置处理3亿行的DataFrame:

spark = SparkSession\
        .builder\
        .appName("App")\
        .config("spark.executor.memory","10g")\
        .config("spark.executor.cores","4")\   
        .config("spark.executor.instances","6")\
        .config("spark.sql.adaptive.enabled","true")\
        .config("spark.dynamicAllocation.enabled","true")\
        .config("spark.shuffle.service.enabled","false")\
        .config("spark.shuffle.io.retryWait","60s")\
        .config("spark.shuffle.io.maxRetries","10")\
        .config("spark.network.timeout","600")\
        .config("spark.sql.shuffle.partitions","1000")\
        .enableHiveSupport()\
        .getOrCreate()

完成聚合操作后执行.count()可正常运行,但写入Parquet文件时,先多次出现Pod ephemeral local storage usage exceeds the total limit of containers 50Gi错误,最终因org.apache.spark.shuffle.FetchFailedException: Error in opening FileSegmentManagedBuffer或org.apache.spark.shuffle.MetadataFetchFailedException: Missing an output location for shuffle 0 partition 908失败。尝试在聚合前执行df = df.repartition(960)、写入前执行df = df.repartition(500),问题仍存在。

原因分析与解决方案

核心差异:.count()与Parquet写入的资源逻辑不同

  1. 计算与IO的资源消耗差

    • .count()仅需汇总聚合后各分区的计数结果到Driver,不需要在Executor本地存储完整的聚合数据,磁盘IO和临时存储占用极低,不会触发存储配额限制。
    • 写入Parquet时,Executor需要先将聚合后的数据序列化、生成Parquet临时文件,这个过程会占用大量本地临时存储。当聚合后单分区数据量过大,多个分区的临时文件会快速耗尽50Gi的Pod存储配额,导致Pod被驱逐,进而引发shuffle文件丢失、无法读取的异常。
  2. 动态分配与shuffle服务的冲突
    开启spark.dynamicAllocation.enabled=true但关闭spark.shuffle.service.enabled=false时,空闲Executor会被自动销毁,而没有shuffle服务托管的话,这些Executor上的shuffle文件会直接丢失。.count()可能在Executor被销毁前就完成了计数汇总,但写入Parquet时需要读取shuffle阶段的文件,就会出现Missing an output location错误。

  3. 分区调整未解决根本问题

    • 聚合前repartition(960)只是增加了聚合阶段的分区数,但聚合后单分区数据量依然可能过大;
    • 写入前repartition(500)反而增大了单个分区的数据量,加剧了单Executor的存储压力,依然会触发配额超限。

可行解决方案

  • 提升Pod本地存储配额:在集群允许的前提下,将Pod的ephemeral local storage限制从50Gi调高(如100Gi),直接缓解存储耗尽问题。
  • 启用shuffle服务:将spark.shuffle.service.enabled设为true,让shuffle文件由节点的shuffle服务托管,即使Executor被动态销毁,文件依然可被读取,避免MetadataFetchFailedException。
  • 精细化调整分区数:根据聚合后单条数据的大小,计算合适的写入分区数,确保每个分区大小控制在1-2Gi范围内(比如总数据量500Gi则设置500-1000个分区),避免单个分区占用过多存储。
  • 临时关闭动态分配:若暂时无法启用shuffle服务,可关闭spark.dynamicAllocation.enabled=true,固定Executor数量,防止Executor被销毁导致shuffle文件丢失。
  • 调整临时存储配置:设置spark.local.dir指向更大的存储路径,或调小spark.sql.files.maxPartitionBytes参数,减小单个分区的默认大小,间接控制临时文件的存储占用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 20:35:09