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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 04:41:08