本节摘要:物理计划的最小执行单元是 Map 与 Reduce 任务,两者之间隔着整个大数据体系里最贵的一段路——Shuffle。本节从 Map 端环形缓冲区出发,走完溢写、分区、排序、归并、HTTP 拷贝、Reduce 归并的全过程,再把 GroupBy 与 Join 两种典型算子放到这套骨架上,说明它们为什么必须付这段网络税。
第 2 章结尾说过:逻辑计划里出现 ReduceSink 算子,物理上就多一段任务边界。ReduceSink 的语义是"按某组键把数据重新分布"。GROUP BY customer_id 要求同键的行落到同一个 Reduce 任务;JOIN 要求连接键相同的行相遇。"相同键必须物理相遇"是分布式 SQL 的一条铁律,而让数据物理相遇的唯一办法,就是把它们搬运到同一台机器——这段搬运就是 Shuffle。
理解 Hive 性能,一半的功力在理解 Shuffle 的每一步开销。整个过程分 Map 端与 Reduce 端两半:

逐步展开:
Map 端。每条输出记录先写进内存里的环形缓冲区(默认约 100MB,阈值 80% 触发溢写)。溢写线程把缓冲区内容按"分区号、排序键"排序后写成磁盘片段——分区号由 Shuffle 分区函数对键哈希取模得到,决定了这条记录将来归属哪个 Reduce 任务。Map 结束前,多个片段归并成一个按分区组织的文件。这一阶段全部发生在 Map 任务本地:一次内存排序、若干次磁盘写。
Reduce 端。每个 Reduce 任务起若干拷贝线程,通过 HTTP 从所有已完成 Map 任务那里拉取属于自己的那个分区。拉完的片段再做一轮归并排序,让相同键的记录物理相邻,然后按键分组迭代,喂给聚合或连接逻辑。
把代价数一遍:Map 端至少两次磁盘写(溢写与归并),Reduce 端一次全量网络拉取加一次磁盘归并。一次 Shuffle,磁盘两三轮、网络一整轮、排序两遍。这就是为什么第 2 章的优化规则(下推、裁剪、Map 端聚合)全部指向同一个目标:让更少的数据走到这段路上。
有了物理图景,回头看 GROUP BY 的翻译。SELECT customer_id, SUM(amount) GROUP BY customer_id:
配合第 2 章的 Map 端聚合规则:Map 输出前先在本地把同键的 amount 压成部分和,网络上传的是部分和而不是原始行。两个机制叠加后,Shuffle 体量从"行数"级压到"任务数乘键数"级。
EXPLAIN 里对应的指纹是 ReduceSink 的 key expressions 与 GroupBy 的 mode 标记。看到 mode hash 加 mergepartial 成对出现,就知道两段式聚合已生效;只有 complete 一种模式时,说明原始行直接上了网络,通常意味着 Map 端聚合被关或失效,值得追问原因。
Reduce 端 Join(也叫 Common Join 或 Shuffle Join)的翻译稍复杂,因为要拼两张表:
EXPLAIN 里认它的特征:两个 TableScan 分支汇入同一个 Reduce 端的 Join Operator,且 ReduceSink 出现在两侧 Map 树的末尾。代价结构也清楚了:两张表的全量数据都要过一次 Shuffle。第 6 章会讲三条翻译路径里的另外两条——小表广播的 MapJoin 与排序合并的 SMB Join,本质上都是对"两边都上网络"这个代价结构的反击。
溢写过多。Map 输出很大(例如高基数分组没做预聚合)时,溢写片段数量暴涨,归并轮数增加,任务日志里会出现 Spilled Records 远大于 Map output records 的迹象。对策回到逻辑层:让更多过滤与预聚合发生在 Map 端,或者加大缓冲区、提高溢写线程并行度——但首选永远是缩小输出本身。
单个 Reduce 拖垮全局。键分布倾斜时,某个键占据大比例记录,全部涌向一个 Reduce 任务。进度条上 99% 的任务秒完成、剩一个跑半小时,是倾斜的标准画面。物理根源同样在 Shuffle:分区函数只保证"同键同任务",不保证"任务间均衡"。治理手段(加盐打散、两阶段聚合、MapJoin 绕过)第 6 章专门展开,这里先建立认知:倾斜不是任务慢,是分布设计遭遇了数据分布。
理论说尽,给自己安排一个十分钟的观察实验。挑一张中等大小的分区表,跑一条带 GROUP BY 的查询,任务完成后打开 YARN 应用页面的计数器区(或在任务历史里看 counters),对照下面这张清单读数:
| 计数器 | 含义 | 健康形态 |
|---|---|---|
| Map output records | Map 端输出的逻辑记录数 | 接近过滤后行数为正常 |
| Map output bytes | 进入序列化的字节数 | 与扫描量同量级 |
| Spilled Records | 溢写到磁盘的记录数 | 与 Map output records 同量级为健康 |
| Combine output records | 预聚合后的输出记录数 | 远小于 Map output records 说明聚合有效 |
| Reduce shuffle bytes | 网络上实际搬运的字节数 | 与 Combine output 挂钩 越小说明两段式越赚 |
| Reduce input records | Reduce 端收到的记录数 | 等于各任务 Combine 输出之和 |
三个读数组合出的诊断结论比任何"感觉慢"都有力:Spilled Records 数倍于 Map output records,说明缓冲区不够或输出过大,先看能不能把更多过滤推进 Map 端;Combine output records 几乎等于 Map output records,说明分组键基数极高、预聚合没压到东西(第 6 章的高基数场景);Reduce shuffle bytes 与扫描量同量级,说明网络税接近全额征收,查询设计有大刀阔斧的空间。
再做一个对照实验:同一条查询把 SELECT 列表从两列改成十列,重跑一次看 Map output bytes 与 shuffle bytes 的涨幅——你会直观看到列裁剪被破坏后的网络代价,这个数字比任何文档里的"建议显式列清单"都有说服力。把两次实验的计数器读数存进个人笔记,它们是你日后判断"这条查询的 Shuffle 是否健康"的第一组基线。
顺带把 Reduce 数量的物理直觉补全:Reduce 任务数并非由 Shuffle 决定,而是由提交时的配置推算(第 7 章参数节细讲),Shuffle 只决定"每个 Reduce 该拿哪些数据"。于是并行度与均衡是两件事:任务数给足了,键分布倾斜照样长尾;键分布均匀,任务数不够照样慢。理解这层分离,第 6 章的倾斜治理(治分布)与第 7 章的参数调优(治并行度)才不会张冠李戴。