窗口函数去重必须显式指定partition by字段,否则全表单分区导致oom或超时;dropduplicates适用于关键列去重但不保证顺序;分区倾斜需加盐或组合分区缓解;缓存输入dataframe可大幅减少重复shuffle。

窗口函数去重必须显式指定 PARTITION BY 字段
不写 PARTITION BY 就等于全表一个分区,所有数据被 shuffle 到单个 task,必然 OOM 或超时。哪怕你只想要最新一条记录,也得先按去重键分组,否则 row_number() 失去意义。
常见错误是误以为 ORDER BY timestamp DESC 就够了,结果执行计划里出现 Exchange SinglePartition —— 这就是危险信号。
- 正确写法:
Window.partitionBy("user_id").orderBy(col("ts").desc()) - 错误写法:
Window.orderBy(col("ts").desc())(没 partition,全局排序) - 如果去重键有多个(如
user_id+device_id),必须全部列在partitionBy中,漏一个就会漏去重
dropDuplicates 比 DISTINCT 更适合带条件的去重
DISTINCT 只能全字段比对,而 dropDuplicates(Seq("a", "b")) 允许你指定关键列,跳过无关字段(比如日志里的 trace_id、随机生成的 uuid)。这直接减少 shuffle 数据量和内存压力。
但要注意:它默认保留**第一个遇到的行**,不保证时间顺序。如果你需要“每个 user_id 的最新记录”,dropDuplicates 无法满足 —— 它不支持排序语义,必须换窗口函数。
- 适用场景:
dropDuplicates适合清洗宽表中由 ETL 错误导致的完全重复行 - 不适用场景:需要按时间/版本/状态选择保留哪条时,必须用
row_number()+where rn = 1 - 性能提示:对超大表,先
repartition("user_id")再dropDuplicates,可避免 shuffle 阶段的二次重分布
窗口函数性能瓶颈几乎都来自分区键倾斜
当 PARTITION BY user_id 遇到头部用户(比如某 uid 出现 500 万次),这个分区会卡住整个 stage。Spark 不会自动拆分热点分区,task 运行时间可能比其他 task 长 10 倍以上。
宝塔面板11.3.0是一款针对Linux服务器设计的可视化管理工具,通过重构核心模块实现资源占用显著降低,尤其适合低配置服务器环境。它将复杂的命令行操作转化为直观的图形界面,帮助开发者快速完成网站部署、环境配置及日常运维工作,无需专业技术背景即可高效管理服务器。
缓解方法不是加资源,而是改分区逻辑:
- 加盐(salting):
concat(col("user_id"), lit("_"), (rand() * 10).cast("int")),把大分区打散 - 组合分区:用
partitionBy("user_id", "date_trunc('day', ts)")把时间维度引入,天然限流 - 预过滤:先
filter掉明显异常的高频 uid(比如出现次数 > 10 万),再走窗口逻辑
缓存窗口前的 DataFrame 能省掉 60% 以上重复计算
如果你在一个 job 里多次调用不同窗口(比如既要取最新记录,又要算 7 日活跃数),每次 .over(windowSpec) 都会触发独立 shuffle。中间 DataFrame 不缓存,等于反复读磁盘、反复 shuffle。
实操上,只要窗口输入源不变,就该立刻 cache():
val baseDF = spark.read.parquet("...").filter(...).cache()- 后续所有
withColumn("rn", row_number().over(w1))和withColumn("sum_7d", sum("amt").over(w2))都基于缓存副本 - 注意:缓存后记得
unpersist(),尤其在长会话中,避免内存泄漏
窗口函数本身不难写,难的是让 Spark 真正按你设想的方式切分和调度 —— 分区键选错、没缓存、忽略倾斜,三者任一都会让优化变成负优化。










