Spark排名与窗口函数性能优化问询:十亿级数据去重任务卡顿
大数据去重任务优化建议
业务需求
- 基于
key字段去除重复记录:重复时保留insert_date最大的记录到valid_df,其余重复记录放入invalid_df - 最终将两个DataFrame写入GCS
- 当前数据
key字段无重复,但为兼容未来必须保留去重逻辑
当前环境与问题
- 数据规模:10亿行Parquet文件,共1000个分片(每个约84MB),存储于GCS
- Spark配置:driver core=4、driver memory=8g;executor core=4、executor memory=16g;minExecutor=16、maxExecutor=40;
spark.sql.shuffle.partitions=200 - 问题:任务长时间运行无结果
现有代码
window_specs=Window.partitionBy('key').orderBy('insert_date') df_rank=df.withColumn("rank", rank().over(window_specs)) df_rank.cache() valid_df=df_rank.filter("rank = 1") invalid_df=df_rank.filter("rank > 1")
优化建议
1. 重构Window函数逻辑
- 改为降序排序:直接取第一条就是
insert_date最大的记录,避免升序后取末尾的额外开销 - 用
row_number()替代rank():计算开销更低,且在无重复key场景下结果一致,未来有重复时也能稳定生成唯一序号 - 移除不必要缓存:
df_rank.cache()会缓存全量带序号的DataFrame,占用大量内存拖慢任务,Spark会自动优化计算 lineage,无需手动缓存
优化后代码:
from pyspark.sql import functions as F from pyspark.sql.window import Window window_specs = Window.partitionBy('key').orderBy(F.col('insert_date').desc()) df_rank = df.withColumn("row_num", F.row_number().over(window_specs)) valid_df = df_rank.filter("row_num = 1").drop("row_num") invalid_df = df_rank.filter("row_num > 1").drop("row_num")
2. 用聚合+关联替代Window函数(性能更优)
对于"保留最新记录"的去重场景,聚合取最大日期再关联原表的方式,shuffle数据量更小,性能远优于Window函数:
from pyspark.sql import functions as F # 聚合每个key对应的最大insert_date max_date_df = df.groupBy('key').agg(F.max('insert_date').alias('max_insert_date')) # 关联得到有效记录 valid_df = df.join( max_date_df, (df['key'] == max_date_df['key']) & (df['insert_date'] == max_date_df['max_insert_date']), how='inner' ).drop(max_date_df['key'], 'max_insert_date') # 关联得到无效记录 invalid_df = df.join( max_date_df, (df['key'] == max_date_df['key']) & (df['insert_date'] != max_date_df['max_insert_date']), how='inner' ).drop(max_date_df['key'], 'max_insert_date')
3. 调整Spark核心配置
- 增大shuffle分区数:将
spark.sql.shuffle.partitions设置为1000(当前数据总大小约84GB,每个shuffle partition控制在80-100MB左右,避免单分区过大导致shuffle阻塞) - 优化executor内存配置:添加
spark.executor.memoryOverhead=4g(约为executor memory的25%),避免OOM - 开启动态资源分配:设置
spark.dynamicAllocation.enabled=true,让Spark根据负载自动扩容到maxExecutor=40 - 开启Parquet读取优化:确保
spark.sql.parquet.enableVectorizedReader=true和spark.sql.parquet.filterPushdown=true(默认开启,若未开启手动配置)
4. 数据读取优化
- 设置
spark.sql.files.maxPartitionBytes=256MB:当前每个文件84MB,合并为更大的初始partition,减少任务调度开销 - 固定schema时添加
mergeSchema=false:避免合并schema的额外开销
内容的提问来源于stack exchange,提问作者Sourav Chatterjee
相关产品推荐
相关产品推荐

