超大数据量下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
原处理逻辑
- 对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
- 对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
- 连接两个聚合后的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
相关产品推荐
相关产品推荐

