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

