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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 22:00:17