3.3 作业运行与 Shuffle 透视 本节摘要:一个作业从提交到完成经历提交、初始化、任务分配、执行、完成五个阶段;其中 Shuffle 跨越 Map 与 Reduce,由环形缓冲、溢写、合并、拉取、归并、分组六步构成。本节逐段拆解这两条时间线,给出每一步的容量参数与调优入口,是全章乃至全教程数据流动的枢纽一节。 作业的五个阶段 提交。客户端向 ResourceManager 申请一个容器跑 MR-AppMaster;AppMaster 启动后把 jar、配置、分片信息上传到 HDFS 作业暂存目录。分片信息(每个分片的数据位置)此刻已经算好——它来自 2.3 节的块位置查询,是后面本地性调度的输入。 初始化。
本节摘要:一个作业从提交到完成经历提交、初始化、任务分配、执行、完成五个阶段;其中 Shuffle 跨越 Map 与 Reduce,由环形缓冲、溢写、合并、拉取、归并、分组六步构成。本节逐段拆解这两条时间线,给出每一步的容量参数与调优入口,是全章乃至全教程数据流动的枢纽一节。
提交。客户端向 ResourceManager 申请一个容器跑 MR-AppMaster;AppMaster 启动后把 jar、配置、分片信息上传到 HDFS 作业暂存目录。分片信息(每个分片的数据位置)此刻已经算好——它来自 2.3 节的块位置查询,是后面本地性调度的输入。
初始化。AppMaster 向 RM 注册,计算需要的资源总量,从 HDFS 取回分片信息,为每个分片生成一个 Map 任务规格、按 Reduce 数生成 Reduce 任务规格。
任务分配。AppMaster 通过心跳协议向 RM 逐个讨要容器,心跳里捎带"任务偏好节点列表"(持有分片副本的节点)。RM 的调度器尽力满足偏好:先给同节点(节点本地),退而求同机架,最后才跨机架——3.1 节的本地性原理在此落地。拿到容器后,AppMaster 通知对应 NodeManager 启动任务 JVM 并执行。
执行。任务向 AppMaster 汇报进度,AppMaster 聚合成作业进度。任务失败按类型处理:进程崩溃由 NodeManager 上报、超时由 AppMaster 检测,重试默认 4 次;作业级失败(如 AppMaster 自身)由 RM 重建,已完成的 Map 不必重算(结果还在各节点本地磁盘)。
完成。最后一个 Reduce 结束,AppMaster 标记作业成功、清理暂存目录、向 RM 注销。输出目录里此刻躺着 part-r-00000 到 part-r-N 的结果文件——注意它们每 Reduce 一个文件,下游若要单个文件需显式 getmerge 或设 Reduce 数为 1。
Shuffle 泛指"Map 输出到 Reduce 输入"之间的全部搬运。它分属两端:Map 端三分之二(分区、排序、落盘),Reduce 端三分之一(拉取、归并)。逐阶段拆开:
阶段一:环形缓冲区。 map 函数每写一个中间键值对,实际写入的是一个约 100MB(默认 io.sort.mb,现名 mapreduce.task.io.sort.mb)的环形内存缓冲。选环形而非普通队列,是为了首尾相接零拷贝地复用内存。写入的同时,后台线程在缓冲里做两件事:按分区号(hash(key) mod R)划界,分区内按 key 排序。
阶段二:溢写。 缓冲占用达到阈值(默认 0.80,即 io.sort.spill.percent)时,后台线程把已排序内容**溢写(spill)**成磁盘临时文件——注意是 Map 任务的本地磁盘,不是 HDFS。剩余 20% 缓冲让 map 继续写,写入与溢写并行,不打断计算。若配置了 Combiner,会在每次溢写前执行一次预聚合。
阶段三:合并。 map 结束时,磁盘上可能有多个溢写文件(数据量大时几十个),合并(merge)成每个分区一个、分区内有序的大文件。合并用多路归并,路数由 io.sort.factor(默认 10)控制。此刻 Map 端的工作全部结束,AppMaster 已知结果文件位置,Reduce 可以来取货。
阶段四:拉取。 每个 Reduce 任务通过 HTTP 从所有 Map 节点拉取属于自己的那一份分区数据——这是全作业唯一的大规模跨网络数据移动。拉取线程默认 5 个并行(mapreduce.reduce.shuffle.parallelcopies),先到先得。
阶段五:归并。 拉来的多路有序数据在 Reduce 端做归并排序。内存不够时同样溢写本地磁盘再归并,直到剩最多 io.sort.factor 路——注意不归并到 1 路,而是留多路在归并时流式消费,节省一轮磁盘。
阶段六:分组。 归并流上按分组比较器(默认即 key 的比较器)切分组,每组触发一次 reduce 调用。至此数据完成从"Map 端分散记录"到"Reduce 端按 key 集齐"的全部旅程。

Slowstart(慢启动)。Reduce 并非等所有 Map 完成才启动——默认 Map 完成 5%(mapreduce.job.reduce.slowstart.completedmaps=0.05)时 Reduce 就开始预热并尝试拉取已完成 Map 的输出。这样网络搬运与剩余 Map 计算重叠,作业尾延迟显著缩短。调到 1.0 意味着严格的两阶段串行,通常只在 Map 端资源紧张时使用。
Map 输出保留到作业结束。Map 完成后其输出文件不删,直到整个作业成功。若某 Reduce 失败重跑,直接重新拉取即可,不必重算 Map——这是 3.1 节"容错红利"的具体机制。代价是节点磁盘被中间结果占用,大作业要预估 Map 端磁盘水位。
拿具体数字感受各阶段的影响。设一个作业有 1000 个 Map、200 个 Reduce、每个 Map 输出 100 万条中间记录、每条 50 字节,中间数据总量约 50 GB。未优化基线:环形缓冲 100 MB、无压缩、无 Combiner——每个 Map 要溢写约 5 次(每 100MB 刷一次盘),全集群的溢写临时文件数以千计;Reduce 端 200 个任务各自从约 1000 个 Map 拉取 250 MB,网络搬运总量就是那 50 GB。三个改动逐个叠加:缓冲调到 400 MB,溢写次数降到约 1 次,磁盘写放大显著缩小;开启 Snappy 压缩,中间结果约压到 20 GB,网络段直接省六成;再上 Combiner、按聚合比 10:1 计,传输量进一步降到 2 GB 量级。三步都不是代码改动,只是配置与一个复用类,端到端耗时常见从小时级压到十分钟级。Shuffle 优化的性价比之所以高,正因为它作用在全程唯一的大规模数据搬运上。
作业完成后计数器(Counters)页面藏着诊断金矿,三个数字最关键:
| 计数器 | 含义 | 健康信号 |
|---|---|---|
| Map output records | map 写出的中间记录总数 | 与业务预估对照 |
| Spilled records | 溢写到磁盘的记录总数 | 应接近 Map output records;远大于它说明缓冲太小、反复溢写 |
| Reduce shuffle bytes | 实际跨网络拉取字节数 | 未压缩时约等于中间数据量;是网络账的直接证据 |
举例:Map output records 10 亿、Spilled records 25 亿——同一批记录被反复写盘多次,典型原因是 io.sort.mb 过小或未启用 Combiner,数据倾斜时更糟。这类诊断不需要看一行日志,计数器比值即可定案。3.4 节的优化三板斧(Combiner、压缩、推测执行)全部作用在这条 Shuffle 链上。
机制看透了,下一节上优化工具箱:Combiner、压缩、推测执行。