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

超大数据量下PySpark分组聚合与连接的高效优化方案问询

大数据量下Spark分组聚合与连接的优化方案咨询

数据结构与需求

我有两个DataFrame:

  • df_1(activity仅包含a和b两个值):
id   item   activity
1      2       a
34    14       b
1      2       b
 .     .       .
  • df_2(activity所有值均为c):
id   item   activity
1      2       c
34    14       c
1      2       c

需求是:按id和item分组,统计两个DataFrame中各activity的频次,再通过id和item连接得到最终DataFrame,期望输出如下:

id  item  a   b   c
1     2   1   1   2
34   14   0   1   1

原处理逻辑

  1. 对df_1分组聚合:
df_1_grp = df_1.groupby("id", "item").agg(
    f.count(f.when(f.col('activity') == 'a', 1)).alias('a'),
    f.count(f.when(f.col('activity_type') == 'b', 1)).alias('b')
)

得到df_1_grp:

id  item  a   b
1     2   1   1
34   14   0   1
  1. 对df_2分组聚合:
df_2_grp = df_2.groupBy("id", "item").count().select('id', 'item', f.col('count').alias('c'))

得到df_2_grp:

id  item  c
1     2   2  
34   14   1   
  1. 连接两个聚合后的DataFrame:
df = df_1_grp.join(df_2_grp, on = ['id', 'item'], how = 'inner')

遇到的问题

数据量极大(约4TB或10亿条记录),当前处理出现磁盘存储不足问题,使用的Spark配置如下:

spark_config["spark.executor.memory"] = "32G"
spark_config["spark.executor.memoryOverhead"] = "32G"
spark_config["spark.executor.cores"] = "32"
spark_config["spark.driver.memory"] = "8G"

spark_config["spark.dynamicAllocation.minExecutors"] = "200"
spark_config["spark.dynamicAllocation.maxExecutors"] = "300"

优化建议

1. 合并DataFrame后统一聚合,避免两次分组与Join操作

原逻辑需要两次分组聚合+一次Join,会产生两次Shuffle和额外的磁盘IO开销。可以将两个DataFrame合并后一次性完成所有activity的计数,只需要一次分组聚合,大幅减少Shuffle和存储压力:

from pyspark.sql import functions as f

# 确保两个DataFrame列名一致,直接合并
df_2_processed = df_2.select('id', 'item', f.col('activity'))
combined_df = df_1.unionByName(df_2_processed)

# 一次分组完成所有activity的频次统计
result_df = combined_df.groupBy('id', 'item').agg(
    f.count(f.when(f.col('activity') == 'a', 1)).alias('a'),
    f.count(f.when(f.col('activity') == 'b', 1)).alias('b'),
    f.count(f.when(f.col('activity') == 'c', 1)).alias('c')
)

2. 修正原代码中的错误

原df_1聚合代码中误用了activity_type字段(应为activity),这会导致b的统计结果错误,修正后才能保证统计准确性。

3. 优化Spark配置,减少磁盘存储压力

  • 调整Shuffle分区数:默认spark.sql.shuffle.partitions为200,对于4TB数据来说分区过大,建议设置为每分区100-200MB左右的规模(比如40000左右),避免单分区数据量过大导致磁盘溢出:
    spark_config["spark.sql.shuffle.partitions"] = "40000"
    
  • 启用Shuffle数据压缩:开启spill数据压缩,减少磁盘占用:
    spark_config["spark.shuffle.spill.compress"] = "true"
    spark_config["spark.io.compression.codec"] = "snappy"
    
  • 启用自适应查询执行:让Spark自动调整Shuffle分区数和执行计划,优化资源利用:
    spark_config["spark.sql.adaptive.enabled"] = "true"
    spark_config["spark.sql.adaptive.shuffle.targetPostShuffleInputSize"] = "134217728"  # 128MB
    
  • 检查临时存储目录:确保spark.local.dir指向有足够剩余空间的磁盘,避免临时文件占满磁盘。

4. 预分区优化

如果原始DataFrame未按id和item分区,先对两个DataFrame按id和item预分区,这样分组聚合时无需再进行Shuffle操作,减少IO开销:

df_1_repartitioned = df_1.repartition('id', 'item')
df_2_repartitioned = df_2.repartition('id', 'item')

# 后续合并与聚合基于预分区后的DataFrame

5. 过滤不必要的列

处理前只保留需要的列(id、item、activity),减少数据传输和存储的体积:

df_1 = df_1.select('id', 'item', 'activity')
df_2 = df_2.select('id', 'item', 'activity')

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 00:45:14