row_number()倾斜本质是单点排序瓶颈,需先判断业务是否允许打乱原始排序;若必须全局连续编号则优化数据源或换引擎,否则可用first_value、topn两步聚合、rank降级等方案。

不能靠加盐就解决,得先看业务是否允许打乱原始排序逻辑。
row_number() 倾斜本质是单点排序瓶颈
只要 partition by 的 key 分布不均(比如 90% 数据落在同一个 department_id),所有这些数据就必须进同一个 reducer 排序——增加 reducer 数量、调大内存、改参数都无效,因为语义上不允许拆分该分区。
- 错误认知:以为
set mapred.reduce.tasks=100就能缓解row_number() over (partition by user_id ...)的倾斜 - 真实限制:同一个
user_id的所有行必须进同一个 task,否则编号会错乱 - 典型症状:大部分 reduce 已完成,只剩 1–2 个卡在 99%,日志里反复出现
GC overhead limit exceeded或OutOfMemoryError
先判断业务场景再选解法
不是所有 row_number() 都必须保留全局连续编号。很多实际需求其实只关心“相对顺序”或“TopN”,可以降级处理。
- 全量排序(不可妥协):如审计日志要求每个用户操作严格按时间编号 → 只能优化数据源(提前过滤/归档冷数据)或换引擎(Spark SQL + adaptive query execution)
-
首末记录提取:如取每个用户的最早/最晚登录记录 → 改用
first_value()/last_value()+ignore nulls,避免排序 -
TopN(最常用可优化场景):如“每个店铺访问次数 Top3 的访客” → 必须走两步聚合:
count(1) group by user_id, shop→ 加随机盐distribute by shop, cast(rand()*100 as int)→ 再开窗 -
仅需去重编号(非连续):如标记“这是该用户第几次访问”但不要求严格 1/2/3 → 可用
rank()或dense_rank()配合预聚合,减少输入行数
加盐必须配合二次聚合,且 salt 列要进 partition by
直接对 partition by 字段加随机后缀(如 concat(user_id, '_', cast(rand()*10 as int)))会破坏业务语义——同一个用户被拆到多个分区,编号不再可比。
- 正确做法:保持原
partition by user_id不变,但在 shuffle 阶段用distribute by user_id, cast(rand()*50 as int)打散数据 - 必须补第二步:先按
user_id, salt开窗取 TopN,再按user_id二次聚合(如collect_list(struct(rank, user_id))+ UDTF 展开) - 注意 salt 范围:太小(如 *10)仍可能倾斜;太大(如 *1000)会导致 reducer 过多、小文件问题;建议从 30–80 试起
- 示例片段:
with pv_cnt as ( select user_id, shop, count(1) as cnt from visit group by user_id, shop ), salted as ( select *, cast(rand() * 50 as int) as salt from pv_cnt ), ranked as ( select *, row_number() over (partition by shop, salt order by cnt desc) as rn from salted ) select shop, user_id, cnt from ranked where rn
容易被忽略的细节:order by 字段的 NULL 和重复值
即使 partition key 均匀,order by 字段大量为 NULL 或存在高频重复值(如 event_time 精度只到天),也会导致排序阶段内部比较膨胀、reduce 拖慢。
- NULL 处理:显式写成
order by event_time nulls last,避免默认行为引发不可控排序开销 - 重复值优化:如果业务允许,把
order by event_time, user_id替换为order by event_time, md5(user_id),减少字符串比较压力 - 字段类型:确保
order by字段是timestamp或bigint,别用string存时间(隐式转换+字典序极慢) - 分区裁剪:若表按天分区,务必在 where 中限定
dt >= '2026-08-01',否则全表扫描放大倾斜影响
真正难的不是写出加盐 SQL,而是确认业务方是否真的需要那个“精确到毫秒的连续编号”。多数时候,他们要的只是“前几名”或“最新一条”,而这两者都有远比 row_number() 更轻量的实现路径。










