不加partition by会触发全局shuffle至单分区,导致oom卡死;必须用高基数均匀字段(如user_id)分区,避免null或低基数字段引发倾斜。

为什么OVER不加PARTITION BY会卡死?
Spark SQL里只要窗口函数用了OVER但没写PARTITION BY,就会触发“全局排序”——所有数据被强行 shuffle 到单个分区,日志里会出现WARN WindowExec: No Partition Defined for Window operation! Moving all data to a single partition。这不是警告,是红牌。10GB数据可能直接OOM,任务卡在Stage 0不动,executor内存打满。
- ORDER BY 是必须的(否则报
AnalysisException: Window function ... requires an ORDER BY clause),但光有 ORDER BY 不够 - PARTITION BY 才是分流关键:它把大问题拆成一堆小问题,每个分区内部独立排序
- 业务上真正需要“全表排第几”的场景极少;绝大多数排名、移动平均、累计求和,都有天然分组维度(如用户ID、设备号、日期)
怎么选partitionBy字段才不翻车?
选错分区键比不选还危险。比如用PARTITION BY user_id算每个用户的点击序号,没问题;但若用PARTITION BY country算全国销量排名,而中国数据占95%,那中国分区还是单点瓶颈。
- 优先选高基数、分布均匀的字段:如
user_id、order_id、event_time::date(注意去重后基数) - 避免低基数字段:如
status(只有'active'/'inactive')、is_paid(布尔值) - 如果必须按低基数字段分组,先用
WHERE过滤掉占比过大的值,或对热点值单独处理(如country = 'CN'走盐值方案) - 验证分布:运行
SELECT country, COUNT(*) FROM t GROUP BY country ORDER BY 2 DESC LIMIT 5,看Top 5是否超过总量50%
rank()这类函数批量应用时怎么防坑?
想给多列同时加rank(),别写10个rank() OVER (ORDER BY col1)再JOIN回来——每多一个OVER就多一次shuffle。Catalyst能优化单次扫描里的多个窗口表达式,但前提是它们共享同一套PARTITION BY + ORDER BY逻辑。
- 用
selectExpr一次性生成全部表达式:df.selectExpr("rank() OVER (PARTITION BY user_id ORDER BY ts) AS rank_ts", "rank() OVER (PARTITION BY user_id ORDER BY amount DESC) AS rank_amt") - 如果各列排序逻辑不同(比如一列升序一列降序),无法合并,那就接受多次shuffle,但务必确保每次都有
PARTITION BY - 别用
withColumn链式调用多个rank():每次都会触发新job,物理计划不可合并
移动平均窗口为什么ROWS BETWEEN比RANGE安全?
写AVG(amount) OVER (PARTITION BY user_id ORDER BY event_time ROWS BETWEEN 6 PRECEDING AND CURRENT ROW)是标准解法;但若写成RANGE BETWEEN INTERVAL '6 days' PRECEDING AND CURRENT ROW,一旦某用户在7天内有1000次点击(时间戳重复或极近),窗口就可能吞入上千行,排序开销爆炸。
-
ROWS按行数截断,稳定可控;RANGE按值范围匹配,受数据分布影响极大 - 时间字段必须先保证唯一性:如有重复
event_time,先用row_number() OVER (PARTITION BY user_id, event_time ORDER BY log_id)打辅助序号,再基于该序号做ROWS窗口 - 永远不要用
ORDER BY RAND()配ROWS——结果不可复现,且失去业务意义
COUNT(*) FILTER (WHERE partition_col IS NULL)占比,超1%就得提前COALESCE(partition_col, uuid())或过滤。











