30亿行PySpark聚合任务Pod本地存储超限优化求助
问题:30亿行Spark分组聚合引发Executor本地存储超限及数据倾斜问题
我有一个包含30亿行数据的大型DataFrame,运行以下PySpark代码及配置:
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","false")\ .enableHiveSupport()\ .getOrCreate() df = spark.read.parquet("/data") df = df.filter(col("colA").isNotNull() & col("colB").isNotNull()) df = df.withColumn("colK_udf",udf_function("colK")) df_1 = df.withColumn("newCol", when((col("colA.field") == 1) & (col("colB.field1") == 2), col("colA.field1")).otherwise(col("colB.field1")))\ # 省略其他处理逻辑 df_1 = df_1.select(...) # 仅保留聚合所需字段 df_agg = df_1.groupby("colA","colB","colC","colD").agg( count(*).alias("numRecords"), sort_array(collected_set("colE")).alias("colE"), sum("colF").alias("colF"), sum("colG").alias("colG"), sum("colH").alias("colH"), sum("colL").alias("colL"), min("colI").alias("colI"), max("colJ").alias("colJ"), countDistinct("colE").alias("colE_distinct"), sort_array(collected_set("colP")).alias("colP"), sort_array(collected_set("colQ")).alias("colQ"), max("colR").alias("colR"), max("colS").alias("colS") ) df_agg.count()
该代码在1亿行的小数据集上可正常运行,但在30亿行数据上执行df_agg.count()时出现以下错误:
ERROR org.apache.spark.scheduler.TaskSchedulerImpl - Lost executor 1 on 2.2.2.2: ... The API gave the following message: Pod ephemeral local storage usage exceeds the total limit of containers 50Gi.
已将Pod本地存储上限从30GB提升至50GB,但问题仍持续,失败任务数不断增长。通过Spark UI查看:Input最高达350GB,Shuffle Write最高达480GB(运行数小时后终止任务)。此外数据存在严重倾斜,groupby后的部分分组规模远大于其他分组。尝试在select后缓存df_1,但问题未解决,请问还能尝试哪些优化方案?
优化方案
1. 优先解决数据倾斜(核心根源)
- 加盐打散倾斜Key:对倾斜的分组Key添加随机前缀(如0-9的随机数),先做局部聚合,再去掉前缀合并结果,避免单个Task处理超大规模分组:
from pyspark.sql.functions import rand, IntegerType, flatten, collect_list # 对分组字段加盐(假设colA是主要倾斜字段,可根据实际情况调整) df_1 = df_1.withColumn("salt", (rand() * 10).cast(IntegerType())) # 局部聚合:按加盐后的完整Key分组 temp_agg = df_1.groupBy("colA", "colB", "colC", "colD", "salt").agg( count(*).alias("numRecords"), collect_set("colE").alias("colE"), sum("colF").alias("colF"), sum("colG").alias("colG"), sum("colH").alias("colH"), sum("colL").alias("colL"), min("colI").alias("colI"), max("colJ").alias("colJ"), size(collect_set("colE")).alias("colE_distinct"), collect_set("colP").alias("colP"), collect_set("colQ").alias("colQ"), max("colR").alias("colR"), max("colS").alias("colS") ) # 全局聚合:去掉salt,合并局部结果 df_agg = temp_agg.groupBy("colA", "colB", "colC", "colD").agg( sum("numRecords").alias("numRecords"), sort_array(flatten(collect_list("colE"))).alias("colE"), sum("colF").alias("colF"), sum("colG").alias("colG"), sum("colH").alias("colH"), sum("colL").alias("colL"), min("colI").alias("colI"), max("colJ").alias("colJ"), size(flatten(collect_list("colE"))).alias("colE_distinct"), sort_array(flatten(collect_list("colP"))).alias("colP"), sort_array(flatten(collect_list("colQ"))).alias("colQ"), max("colR").alias("colR"), max("colS").alias("colS") ) - 拆分倾斜Key单独处理:先找出具体的倾斜Key(通过
df_1.groupBy("colA","colB","colC","colD").count().orderBy(desc("count"))排查),将这些Key单独过滤出来做聚合,再与正常Key的聚合结果合并,避免倾斜Key拖慢整体任务。
2. 减少Shuffle数据量与磁盘占用
- 严格裁剪字段:确保
df_1.select(...)只保留聚合必需的字段,删除所有无关列,从根源减少Shuffle时的数据传输和存储量。 - 替换
countDistinct:countDistinct会触发额外Shuffle,若业务允许近似值,直接用approx_count_distinct("colE", 0.01);若需精确值,改用size(collect_set("colE"))替代,避免单独的Shuffle阶段。 - 推迟排序操作:将
sort_array推迟到全局聚合之后执行,减少局部聚合阶段的内存和磁盘存储压力。
3. 调整Spark配置优化资源利用
- 增加Shuffle分区数:设置
spark.sql.shuffle.partitions=2000(根据Executor总核数调整,建议为总核数的2-3倍),让数据更均匀分布到各个Task,避免单个Task处理过多数据。 - 开启Shuffle压缩:添加以下配置减少磁盘占用:
.config("spark.shuffle.spill.compress", "true")\ .config("spark.shuffle.compress", "true")\ .config("spark.io.compression.codec", "snappy")\ - 调整Executor内存比例:增大堆外内存,避免内存溢出导致频繁磁盘溢写:
.config("spark.executor.memoryOverhead", "4g")\
4. 优化数据处理逻辑
- 替换Python UDF:Python UDF性能差且易引发内存问题,若
udf_function逻辑可通过Spark内置函数实现,直接替换;否则改用Scala UDF,减少序列化/反序列化开销。 - 提前过滤无效数据:将
filter(col("colA").isNotNull() & col("colB").isNotNull())放在读取Parquet之后立即执行,避免后续处理无效数据。
5. 优化缓存策略
- 若
df_1没有被重复使用,直接去掉缓存操作,避免占用Executor内存和磁盘存储;若必须缓存,使用序列化存储减少内存占用:df_1.cache("MEMORY_AND_DISK_SER")
内容的提问来源于stack exchange,提问作者Rayne
相关产品推荐
相关产品推荐

