深入 Hadoop 15 - MRv2 on YARN ApplicationMaster 与 Task Attempt
上一篇讲了 Shuffle。本篇讲 MapReduce 在 YARN 上的具体实现——MRv2。这是 MapReduce 计算模型篇的最后一篇。
MRv2 常被介绍成"MapReduce 在 YARN 上的版本"。这个描述对应了部署形态但漏了内部架构。准确的说法是:MRv2 把 Hadoop 1.x 的 JobTracker 拆成 ResourceManager(资源管理)和 MRAppMaster(作业编排)两部分,MRAppMaster 是一个普通的 YARN ApplicationMaster,运行在 container 里。每个 Task 也是一个 container(YarnChild),通过 umbilical RPC 与 MRAppMaster 通信。MRv2 的 Speculative Execution、Counter、Task Recovery 等机制都是 MRAppMaster 实现的,与 YARN 基础设施解耦。
本篇只抓一个问题:MRAppMaster 内部是怎么编排 Map / Reduce Task 的、Task Attempt 怎么向 AM 报告进度、Speculative Execution 的触发条件和代价、AM 重启时怎么恢复 Task 状态。
MRAppMaster 的内部结构
MRAppMaster 是 MapReduce 框架的 ApplicationMaster 实现。它本身是一个普通 Java 进程,跑在 YARN 分配的 container 里(container-0)。内部结构可以拆成几个核心组件:
1 | |
Job Init 阶段做几件事:
读取作业配置文件(job.xml,从 HDFS 加载)。配置包含所有 MapReduce 参数——Mapper 类、Reducer 类、InputFormat、OutputFormat、Reduce 数、Combiner 等。
调用 InputFormat.getSSplit() 计算所有 InputSplit 列表。每个 InputSplit 包含数据位置(哪些 block 在哪些 DataNode)和长度。这个列表会分发给后续 Map Task,让 MRAppMaster 可以做数据本地化调度(优先在持有数据的节点上启动 Map Task)。
初始化作业状态机——作业进入 RUNNING 状态,准备启动 task。
Job Tracker(internal)是 MRAppMaster 内部的 task 调度器。它决定哪个 task 在哪个节点上跑、什么时候启动 Reduce Task、什么时候触发 Speculative Execution。注意这个 Job Tracker 与 Hadoop 1.x 的 JobTracker 完全不同——前者是作业内调度(一个 MRAppMaster 一个),后者是集群级调度(一个集群一个)。命名巧合,不是同一个东西。
Job Tracker(internal)的调度逻辑:
1 | |
这种"Map 优先 + Reduce 提前启动"的策略是 MapReduce 性能调优的关键维度之一。slowstart 设得低(例如 0.05)让 Reduce 提前启动,可以隐藏 Reduce Task 的启动开销;设得高(例如 0.5)让 Reduce 等更久但避免浪费资源。
Task Attempt 的生命周期
每个 Task(例如某个具体的 Map Task “task_…_m_000000”)在执行过程中可能有多次 Attempt。第一次 Attempt 失败时 AM 会启动新 Attempt,直到达到最大重试次数(默认 4,mapreduce.map.max.attempts / mapreduce.reduce.max.attempts)。
Task Attempt 的状态机:
1 | |
NEW:刚创建,未分配 container。
UNASSIGNED:已分配 container,AM 准备启动。
ASSIGNED:container 已分配到某个 NM,AM 已通过 RPC 让 NM 启动。
RUNNING:container 进程已启动,开始执行 MapTask 或 ReduceTask。
COMMIT_PENDING:Task 完成,请求 AM 提交输出(防止多个 Attempt 同时写 OutputFormat)。
SUCCEEDED:Task 输出已提交,Task 成功完成。
FAILED:Task 失败,AM 决定是否启动新 Attempt。
KILLED:Task 被强制终止(Speculative、用户 kill、抢占)。
每次 Attempt 启动都创建一个全局唯一的 AttemptId(例如 attempt_1721900000000_0001_m_000000_0)。这个 ID 后缀 _0、_1、_2 表示第几次 Attempt。
Attempt 失败的常见原因:
container 进程崩溃(OOM、JVM crash、被 cgroups 杀)。
container 与 AM 心跳超时(默认 mapreduce.task.timeout = 600000ms 即 10 分钟)。这种"挂起的 Task"AM 没收到心跳,标记为 FAILED。Task 实际可能还在跑,但 AM 已经启动新 Attempt。
Task 主动抛出未捕获异常(用户 Mapper/Reducer 代码 bug)。
NM 故障(节点宕机),AM 通过 RM 知道该 NM dead,把该 NM 上所有 Task Attempt 标记 FAILED。
Umbilical RPC:Task 与 AM 的通信
Task 进程(YarnChild)启动后通过 umbilical RPC 与 MRAppMaster 通信。这是一个简单的 RPC 协议,承担三类信息:
进度报告:Task 周期性向 AM 报告当前进度(0.0 - 1.0 浮点数)。AM 用这个进度判断 Task 健康度(长时间没进度更新可能是 Task 挂起)。
Counter 上报:Task 把 Counter(Map 输入记录数、Spilled records 等)实时上报给 AM。AM 聚合所有 Task 的 Counter,在 Web UI 展示。
完成通知:Task 完成后通过 umbilical RPC 调用 AM 的 done() 方法,告诉 AM 自己的输出可以提交。
umbilical RPC 的默认心跳间隔 1 秒(mapreduce.task.progress.report.interval),超时 10 分钟(mapreduce.task.timeout)。
注意 umbilical 是 Task 主动调用 AM 的 RPC,不是 AM 拉取 Task——这与 NM 与 RM 的关系(NM 主动 RPC RM)一致。Task 是 RPC client,AM 是 RPC server。
如果 Task 不报告进度超过 timeout,AM 标记该 Task 为挂起并启动新 Attempt。这个机制防止 Task 死循环或 GC 停顿导致的无限等待。代价是误判——长时间 GC 停顿的 Task 可能被错误标记 FAILED。处理这种场景的方法是调大 timeout(mapreduce.task.timeout),或者让用户代码在长时间操作前调用 context.progress() 报告进度。
Speculative Execution:对抗慢节点
集群里某些节点比平均慢——磁盘老化、网络抖动、CPU 负载高。如果一个 MapReduce 作业有 1000 个 Map Task,999 个在 10 分钟内完成,1 个慢节点上的 Map Task 跑了 1 小时,整个作业延迟被这个慢 Task 拖到 1 小时。
Speculative Execution(推测执行)是 MapReduce 对抗慢节点的机制。AM 检测哪些 Task 比平均慢,启动备份 Attempt。哪个 Attempt 先完成就用哪个,另一个被 kill。
具体检测算法(DefaultSpeculator):
1 | |
Speculative Execution 的配置:
1 | |
Speculative Execution 的代价是资源浪费——备份 Attempt 占用额外 container。如果集群资源紧张,Speculative 抢占其他作业的资源。生产环境通常在 Map 阶段开启(Map 输出中间结果,备份代价小),Reduce 阶段关闭(Reduce 输出大文件,备份代价大)。
注意 Speculative 与 Combiner、Pipeline 容错的差异。Combiner 是数据优化(减少 Shuffle 流量),Pipeline 容错是传输层重试,Speculative 是 Task 层重试。三者独立,可以同时启用。
Counter:性能观测的核心
MapReduce 提供 Counter 机制让 Task 上报指标。Counter 是 (group, name, value) 三元组,Task 通过 context.getCounter(group, name).increment(1) 上报。AM 聚合所有 Task 的 Counter 在 Web UI 展示。
Hadoop 内置的 Counter 分几类:
1 | |
用户也可以定义自定义 Counter。在 Mapper 或 Reducer 里调用 context.getCounter("MyGroup", "MyCounter").increment(1)。这些 Counter 自动上报给 AM,在 Web UI 可见。
Counter 是 MapReduce 性能调优的核心观测点。每个 Counter 反映一个具体维度,通过对比 baseline 和调优后的 Counter 数值,可以判断调优是否有效。
例如调 io.sort.mb:调优前 Spilled records = 5M,调优后 = 1M,说明 Spill 减少 80%。配合 Filesystem Counter 的本地磁盘字节变化,可以确认 I/O 优化效果。
AM 重启时的 Task Recovery
MRAppMaster 自己是一个 YARN ApplicationAttempt。如果 AM 进程崩溃,RM 会启动新 Attempt(默认 max-attempts = 2,可配置更高)。新 Attempt 必须恢复作业状态,否则已完成的 Task 全部丢失。
MRv2 的 Recovery Manager 负责这个过程:
1 | |
Recovery 的具体行为取决于 yarn.mapreduce.am.job.recovery.enable(默认 true)和 yarn.mapreduce.am.job.recovery.task-state-limiting(默认 2000,超过这个数量 Task 的作业不做 recovery,直接重跑所有 Task)。
Recovery 的代价:
OutputFormat 必须是幂等的。Map Task 的输出在本地磁盘,AM 崩溃后这些输出可能丢失(节点没崩但 AM 进程崩了,输出还在;如果节点也崩了,输出丢失)。Recovery 时已 SUCCEEDED 的 Map Task 被认为输出已持久化(OutputCommitter.commitTask 已调用),不重跑。
Reduce Task 的输出在 HDFS。COMMIT_PENDING 但未 COMMITTED 的 Reduce Task 在 Recovery 时被重新调度(输出未提交)。
Recovery 不是完美的——某些 Task 可能重跑,造成额外开销。但相比"AM 崩溃作业失败需要重新提交",Recovery 的代价可以接受。
实验:观察 MRv2 行为
提交一个 MapReduce 作业,在 ApplicationMaster Web UI(每个作业有独立 Web UI,URL 形如 http://am-host:random-port)观察:
Job 页面展示作业整体进度、Map/Reduce 完成度、作业 Counter。Tasks 标签页列出所有 Task Attempt,每个 Attempt 显示状态、所在节点、进度、运行时间、Counter。
故意 kill 某个 Task 进程(找到 Task 所在 NM,kill -9 YarnChild 进程),观察 AM 是否在几秒内检测到并启动新 Attempt。新 Attempt 的 AttemptId 后缀 _1、_2 等。
提交一个故意慢的作业(Mapper 里加 Thread.sleep),让一个 Task 显著慢于其他 Task,观察 AM 是否启动 Speculative Attempt。Speculative Attempt 在 Web UI 显示为同一 Task ID 的两个 Attempt 并行跑。
调小 mapreduce.task.timeout 到 30 秒,提交一个 Mapper 里 Thread.sleep(60) 的作业,观察 AM 是否把 Task 标记为 FAILED(心跳超时),然后启动新 Attempt。
模式提炼
MRv2 on YARN 体现的设计模式:
1 | |
这个模式不只是 MRv2。Spark on YARN 的 Driver / Executor 结构是同构——Driver 跑在 AM container,Executor 跑在 task container,通过 RPC 通信。Spark 的 Speculative 机制(spark.speculation)与 MapReduce 的 Speculative 几乎相同。Flink on YARN 的 JobManager / TaskManager 也是同构。
差异在持久化策略。MapReduce 用 HDFS 文件做 Recovery,Spark 用 Event Log,Flink 用 Checkpoint。具体机制各不相同,但"AM 重启后从持久化状态恢复"这个模式是一致的。
工程迁移表
| MRv2 概念 | Spark on YARN | Flink on YARN | Tez on YARN | Hadoop 1.x JobTracker |
|---|---|---|---|---|
| MRAppMaster | Driver | JobManager | DAG AppMaster | JobTracker(部分) |
| YarnChild | Executor | TaskManager | Task | TaskTracker + Task |
| Umbilical RPC | Driver ↔ Executor | JM ↔ TM RPC | Tez RPC | TaskUmbilicalProtocol |
| Speculative Execution | spark.speculation | restart strategy | DAG spec | 是 |
| Counter | accumulator | metric | counter | counter |
| Recovery Manager | Event Log | checkpoint | DAG recovery | job history |
| 数据本地化 | preferred location | - | input locality | data-local slot |
注意 Spark on YARN 这一列。Spark 的 Driver 跑在 AM container(cluster 模式)或 Client 进程(client 模式)。Executor 跑在 task container,与 MapReduce 的 YarnChild 同构。差异在 Spark 的 Executor 是长期存活(Driver 与 Executor 保持长连接),MapReduce 的 YarnChild 是 Task 级别(Task 完成 YarnChild 退出)。这让 Spark 的 Executor 启动开销可以分摊到多个 Task,是 Spark 比 MapReduce 快的另一个原因。
常见误解
误解一:“MRAppMaster 是 Hadoop 1.x JobTracker 的简单移植”。MRAppMaster 与 JobTracker 是不同的代码——JobTracker 是单体进程管理集群所有作业,MRAppMaster 是每作业一个进程。共享的只是部分调度逻辑(task 分配、数据本地化)。
误解二:“Speculative Execution 总是让作业变快”。Speculative 启动备份 Attempt 占用资源,如果集群资源紧张,备份抢其他作业资源。Speculative 适合"作业里个别 Task 慢"场景,不适合"所有 Task 都慢"场景(所有 Task 都慢说明作业本身设计有问题,不是慢节点问题)。
误解三:“AM 重启会自动恢复所有 Task”。Recovery 有上限(默认 2000 Task)。超过这个数量 AM 重启时所有 Task 重跑。另外 Recovery 要求 OutputFormat 幂等,自定义 OutputFormat 必须保证多次 commit 不产生不一致。
误解四:“Task 心跳超时就是 Task 死了”。Task 可能因为长时间 GC 或长时间同步 IO(HBase scan)不报告进度,被 AM 误判 FAILED。处理方法是调大 timeout 或者用户代码主动调用 context.progress() 报告进度。
误解五:“MapReduce 的 Counter 是只读的统计信息”。Counter 不仅可以读,还可以用作控制——用户代码可以读 Counter 判断"已处理多少记录",达到阈值时主动跳出 map/reduce 循环。这种"用 Counter 做控制流"在某些场景有用(例如采样作业处理 N 条记录就停)。
练习
-
提交一个 MapReduce 作业,在 ApplicationMaster Web UI 上找到 Tasks 标签页。每个 Task 显示 Attempt 数——故意 kill 某个 Task 进程,观察 AM 是否启动新 Attempt(Attempt 数变 2)。
-
修改作业配置
mapreduce.map.speculative=true,提交一个故意让某个 Mapper 慢的作业(Mapper 里 Thread.sleep),观察 Speculative Attempt 是否启动,最终作业完成时间是否被慢 Mapper 拖累。 -
在
apache/hadoop源码里找到MRAppMaster.java、YarnChild.java、TaskAttemptListener.java、DefaultSpeculator.java(hadoop-mapreduce-client-app 模块),观察 AM 内部组件的实现。 -
思考题:如果让 MapReduce 的 Speculative 不仅对抗慢节点,还对抗"用户代码 bug 导致 Task 偶发失败"(例如某些记录触发 NPE),Speculative 的设计需要怎么改造?这种"备份对抗正确性问题"的思路在分布式系统里是否有先例?
系列导航
| 序号 | 主题 | 状态 |
|---|---|---|
| 00-06 | HDFS 存储层(00 导读 / 01-06 HDFS 各主题) | 第一阶段(已完成) |
| 07-12 | YARN 资源管理层 | 第二阶段(已完成) |
| 13 | MapReduce 编程模型:分而治之的工程化表达 | |
| 14 | Shuffle 全流程:Map 端、Reduce 端与磁盘 I/O 的代价 | 上一篇 |
| 15 | MRv2 on YARN:ApplicationMaster 与 Task Attempt | 本篇(MapReduce 篇完结) |
| 16-19 | HA、安全与 Common 基础设施(RPC / 序列化 / 安全 / 监控) | 第四阶段,下一篇开始 |
| 20-22 | 演进、生态与对比(Hadoop 生态 / 对象存储 + K8s / 设计遗产) | 第五阶段,待开始 |
参考资料
- Jeffrey Dean, Sanjay Ghemawat. MapReduce: Simplified Data Processing on Large Clusters. OSDI 2004. Section 4 “Backup Tasks” 描述了 Speculative Execution 的早期设计。
- Vinod Kumar Vavilapalli et al. Apache Hadoop YARN: Yet Another Resource Negotiator. SOCC 2013. Section 6 “MapReduce on YARN” 描述了 MRv2 的具体实现。
- Apache Hadoop 官方文档:MapReduce on YARN. https://hadoop.apache.org/docs/current/hadoop-mapreduce-client/hadoop-mapreduce-client-app/MapReduceAppMaster.html
- Apache Hadoop 源码:
MRAppMaster.java、YarnChild.java、TaskAttemptListener.java、DefaultSpeculator.java. https://github.com/apache/hadoop - Tom White. Hadoop: The Definitive Guide. O’Reilly, 4th Edition 2015. Chapter 7 详细描述了 MRv2 在 YARN 上的运行机制。
- Maysam Yabandeh et al. Speculative Execution in MapReduce.(Microsoft Research 关于 Speculative Execution 算法的改进研究)
