row_number()必须shuffle重分布数据才能编号,因其需按partition by分组、order by排序后连续编号;原始分区与窗口分组不一致,故必须将同组行拉至同一executor完成全序排序与编号。

ROW_NUMBER() 必须重分布数据才能编号
因为 ROW_NUMBER() 是窗口函数,它需要按 PARTITION BY 分组、再在每组内按 ORDER BY 排序并连续编号。Spark 无法在原始分区中完成这个操作——原始数据的分片(partition)和窗口要求的分组完全不一致。所以必须把所有属于同一分组的行拉到同一个 Executor 上,这就触发了 Shuffle。
Shuffle 不是副作用,而是强制行为
即使你只想要每组第一条(rn = 1),Spark 也得先把整组数据 shuffle 过去、完整排序、再编号,最后才过滤。它不会“提前终止”或“跳过排序”。常见误解是“我只取一行,应该很快”,但执行计划里一定能看到 Window 节点挂在 Exchange(即 Shuffle)之后。
- 没有
PARTITION BY?那就全表当一个大组,所有数据 shuffle 到单个 task —— 极易 OOM -
PARTITION BY字段基数高(比如百万级 user_id)?会生成百万个小 shuffle 任务,调度开销压垮 Driver - 分区字段含大量 NULL 或低区分度值(如 status IN ('A','B'))?导致少数几个超大分区,sort buffer 溢出,写磁盘临时文件
为什么不能 map 端合并?
因为 ROW_NUMBER() 要求严格顺序:编号必须从 1 开始、连续、且依赖全局排序结果。Map 端看到的只是局部数据片段,无法知道“这行在整组里排第几”。除非你放弃编号语义,改用近似或聚合替代方案,比如:
- 用
MAX(CONCAT(...))+GROUP BY拼接最新记录(需确保字段无分隔符冲突) - 先
repartition(col)再mapGroups手动取 top 1(Scala/Python API 层控制,绕过 SQL 窗口) - 对小表维度 join,把
PARTITION BY和 join key 对齐,复用同一轮 shuffle(如PARTITION BY id, fid与JOIN ... ON a.id = b.id共享分区)
最容易被忽略的一点
执行计划里 Window 节点上游如果紧挨着 Scan,说明没做任何前置过滤;但真正耗时的往往不是编号本身,而是 shuffle 前没裁剪字段、没下推谓词、或者没把 ORDER BY 字段建好索引(ORC/Parquet 文件中列统计信息或 Z-Order 排序)。这些都比换函数更值得优先检查。











