上一篇讲了调度器。本篇讲一个 YARN 应用程序从提交到完成经历的全部状态变化。

YARN 应用程序的生命周期常被简化成"提交 → 运行 → 完成"。这个简化漏了三个独立状态机的协调——RM 维护 Application 状态机、RM 维护 ApplicationAttempt 状态机、NM 维护 Container 状态机。三个状态机通过 RPC 同步,共同决定一个作业从 SUBMITTED 到 FINISHED 的每一步。理解这套生命周期对调试作业卡住、诊断 AM 启动失败、设计自定义 AM 都很关键。

本篇只抓一个问题:Client → RM → NM → AM → Container 之间的完整 RPC 序列和状态机切换,AM 启动失败时 RM 怎么重试,Container 完成后 NM 怎么回收。

三层状态机的分工

YARN 应用程序生命周期由三个独立状态机协调:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
RMApp(作业级状态机)
- 由 RM 维护
- 跟踪整个作业的状态(NEW / SUBMITTED / ACCEPTED / RUNNING / FINISHED / FAILED / KILLED)
- 一个 RMApp 可以有多个 RMAppAttempt(AM 重试产生新 attempt)

RMAppAttempt(AM 尝试级状态机)
- 由 RM 维护
- 跟踪每次 AM 启动的状态(SCHEDULED / ALLOCATED / LAUNCHED / RUNNING / FINISHED / FAILED)
- 每次 AM 崩溃后重启会创建新的 RMAppAttempt

Container(NM 侧状态机)
- 由 NM 维护
- 跟踪单个 container 的状态(NEW / LOCALIZING / LOCALIZED / RUNNING / EXIT_CHECK / DONE)
- container 完成后 NM 通知 RM 和 AM

三层状态机的核心分工是:RMApp 决定作业整体命运(接受 / 拒绝 / 完成),RMAppAttempt 决定单次 AM 启动尝试的命运(启动成功 / 失败 / 重试),Container 决定单个 task 进程的命运(启动 / 运行 / 完成 / 失败)。

为什么要分三层而不是两层?因为 AM 自己会崩溃——OOM、网络分区、节点故障。如果只有 RMApp 和 Container 两层,AM 崩溃时整个作业就失败了。引入 RMAppAttempt 让 RM 可以在 AM 崩溃后启动新 attempt,作业继续跑。这就是 YARN 的 AM 容错机制。

提交流程的完整 RPC 序列

把一个作业从提交到 AM 启动完成的 RPC 序列展开:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
Client                       ResourceManager           NodeManager             AM (container-0)
─────── ──────────────── ─────────── ─────────────────
1. getNewApplication() ───►
◄─── ApplicationId, max-capability info

2. upload jar + config
to HDFS

3. submitApplication(
ApplicationSubmissionContext) ───►
RMApp state: NEW → NEW_SAVING → SUBMITTED
RM persists to state store
RMApp state: SUBMITTED → ACCEPTED
scheduler accepts the app
RMApp state: ACCEPTED → RUNNING
RM creates RMAppAttempt (attempt-1)
RMAppAttempt state: SCHEDULED
scheduler picks a NM for AM container
RMAppAttempt state: ALLOCATED
────► startContainers() ───►
NM state: NEW → LOCALIZING → LOCALIZED → RUNNING
NM launches container-0
AM process boots
◄──────────── 4. AM init

5. registerApplicationMaster() ───►
RMAppAttempt state: RUNNING
RM records AM's RPC address

6. AM heartbeat → RM:
allocate() requests N containers
RMAppAttempt scheduler processes
(more containers get scheduled...)

这个序列的关键节点:

步骤 1 的 getNewApplication 让 Client 拿到一个全局唯一的 ApplicationId。这个 ID 后续用于所有 RPC。

步骤 3 的 submitApplication 是关键 RPC。ApplicationSubmissionContext 包含:AM 的启动命令、AM 所需资源(memory / vcores)、AM 所在队列、优先级、作业 jar 在 HDFS 的路径、本地资源(jar / 配置文件)、安全令牌。

RM 收到 submitApplication 后立刻把 RMApp 状态置为 NEW_SAVING,把作业元数据持久化到 RM State Store(ZooKeeper / HDFS / LevelDB,第十一篇展开)。持久化成功后状态变 SUBMITTED,调度器开始处理。

调度器接受作业(检查队列 ACL、资源足够、配额未满)后状态变 ACCEPTED,然后变 RUNNING。此时 RM 创建第一个 RMAppAttempt(attempt-1),状态 SCHEDULED。

调度器找到一个 NM 能容纳 AM container,把 container 分配给这个 attempt。状态变 ALLOCATED。RM 通过 RPC 调用对应 NM 的 startContainers,NM 开始准备启动 container。

NM 启动 container 的状态机:

1
2
3
4
5
6
NEW          container 已分配但未开始处理
LOCALIZING 正在从 HDFS 下载作业 jar 等本地资源
LOCALIZED 本地资源准备完毕
RUNNING container 进程已启动
EXIT_CHECK container 进程已退出,正在收集日志
DONE container 完全清理

container 进程启动后执行 ApplicationSubmissionContext 里指定的命令(通常是 yarn-application-classpath 下的 AM 主类)。AM 进程初始化完成后通过 registerApplicationMaster RPC 向 RM 注册——这一步告诉 RM 自己的 RPC 地址(host:port)和作业的 tracking URL(Web UI)。RM 收到注册后把 RMAppAttempt 状态置 RUNNING。

注册完成后 AM 进入主循环——周期性向 RM 发 allocate 心跳,每次心跳里包含资源请求列表和已完成的 container 报告。RM 调度器处理这些请求,返回新分配的 container 和已释放的 container 状态。

AM 心跳与 container 调度

AM 与 RM 的核心 RPC 是 allocate,这个 RPC 同时承担三个职责:

第一,AM 申请新 container。AM 在 allocate 请求里传 ResourceRequest 列表,每个 request 描述"需要多少资源 + 偏好在哪个节点 + 优先级"。RM 调度器处理这些 request,把能立刻满足的 container 加到 allocate 响应里。

第二,AM 报告已完成的 container。AM 在 allocate 请求里传已完成的 ContainerStatus 列表(container ID + 退出状态 + 退出诊断)。RM 用这些信息更新作业进度。

第三,AM 与 RM 保持心跳。如果 AM 长时间不调用 allocate,RM 会判定 AM 死亡(默认 yarn.am.liveness-monitor.expiry-interval-ms = 600000ms 即 10 分钟),触发 AM 重试。所以 AM 必须周期性发 allocate,即使没有新 container 申请。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
// AM 主循环伪代码
while (!job.isDone()) {
List<ResourceRequest> requests = job.computeNextRequests();
List<ContainerId> releasedContainers = job.getCompletedContainers();

AllocateResponse response = amRMClient.allocate(
new AllocateRequest(
responseId, // 心跳序号,防止乱序
progress, // 进度 0.0 - 1.0
requests,
releasedContainers,
...));

List<Container> newContainers = response.getAllocatedContainers();
for (Container c : newContainers) {
// 直连对应 NM 启动 task
nmClient.startContainer(c, containerLaunchContext);
}

Thread.sleep(1000); // 心跳间隔
}

注意 allocate RPC 是同步阻塞的——AM 发出后阻塞等待 RM 响应。这简化了 AM 实现,但也意味着 RM 调度器的处理速度直接影响 AM 心跳间隔。如果 RM 调度器慢,AM 的心跳间隔被拉长,可能超过 liveness 阈值被判定死亡。

AM 与 NM 的 Container 启动

AM 拿到新分配的 container 后,直连对应 NM 的 ContainerManagementProtocol.startContainers 接口启动 container 进程。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
// Container 启动协议
ContainerLaunchContext clc = Records.newRecord(ContainerLaunchContext.class);
clc.setCommands(Arrays.asList(
"$JAVA_HOME/bin/java " +
"-Xmx" + containerMemory + "m " +
"com.example.TaskMain " +
"1>" + logDir + "/stdout " +
"2>" + logDir + "/stderr"
));
clc.setLocalResources(localResources); // jar / 配置文件
clc.setEnvironment(env); // CLASSPATH 等
clc.setTokens(tokens); // 安全令牌

StartContainersRequest request = StartContainersRequest.newInstance(
Arrays.asList(new StartContainerRequest(clc, container.getId())));
StartContainersResponse response = nmClient.startContainers(request);

ContainerLaunchContext 包含启动命令、本地资源、环境变量、安全令牌。NM 收到后:

1
2
3
4
5
6
7
8
1. 把 clc 里的本地资源(jar / 配置)从 HDFS 下载到本节点
Container 状态: NEW → LOCALIZING → LOCALIZED
2. 创建工作目录、日志目录、临时目录
3. 设置 cgroups 资源限制
4. 启动 container 进程
Container 状态: LOCALIZED → RUNNING
5. 进程退出后收集日志
Container 状态: RUNNING → EXIT_CHECK → DONE

container 进程是 task 的实际执行者——MapReduce 的 map / reduce task、Spark 的 executor task、Flink 的 task slot。AM 通过定期 RPC 查询 NM 拿到 container 进度(或者 task 进程通过 RPC 主动报告给 AM)。

container 进程异常退出(exit code 非 0)时,NM 把 ContainerStatus 标记为 FAILED,附上退出码和 stderr 日志的最后几 KB。AM 在下一次 allocate 心跳里收到这个 container 的失败状态,决定是否重试。

AM 启动失败的重试机制

AM 自己也是一个 container,可能崩溃——OOM、JVM crash、NM 故障、网络分区。AM 崩溃后 RM 会自动启动新 attempt,这是 YARN 的 AM 容错机制。

重试流程:

1
2
3
4
5
6
7
8
9
1. AM 进程崩溃(或与 RM 心跳超时被判定死亡)
2. RM 标记当前 RMAppAttempt 状态 FAILED
3. RMApp 检查是否还能重试
- 当前 attempt 数 < yarn.resourcemanager.am.max-attempts(默认 2)
- 作业整体未超过 max-attempts 限制
4. 如果能重试,RM 创建新 RMAppAttempt(attempt-2)
5. 新 attempt 进入 SCHEDULED → ALLOCATED → LAUNCHED → RUNNING
6. AM 进程重启,恢复作业状态(从 RM 拿到上一次 attempt 已完成的 container 列表)
7. 作业继续跑

AM 重启后的状态恢复是关键。第一次 attempt 的 AM 内存里持有的"已派发哪些 task、哪些 task 完成、哪些 task 失败"信息丢失。新 attempt 启动后必须从 RM 拿回这些信息——RM 知道每个 attempt 已分配的所有 container,新 attempt 调用 allocate 时 RM 在响应里返回之前 attempt 已完成的 container 列表。

不同计算框架对 AM 重启的恢复策略不同:

MapReduce(MRv2):从 HDFS 上的作业历史恢复 task 进度。已完成的 map task 不重跑(输出在 HDFS),未完成的 task 在新 AM 上重新调度。

Spark:从 Driver 的 Event Log 恢复。如果 Driver(即 AM)配置为 cluster 模式,整个作业的 RDD 血统丢失,必须从头重跑。这是为什么 Spark 重要作业通常配置 checkpoint 周期性把 RDD 物化到 HDFS。

Flink:从 checkpoint / savepoint 恢复。Flink 的 AM 是 JobManager,崩溃后从最近 checkpoint 恢复 operator 状态。这是 Flink 流式计算的核心容错机制。

注意 yarn.resourcemanager.am.max-attempts 的默认值是 2,意味着 AM 第一次崩溃会重试一次,第二次崩溃才让作业整体失败。生产集群通常把这个值调高(5-10)让 AM 在节点故障场景下更稳健。

作业完成的清理

作业完成有两种路径——正常完成和异常终止。

正常完成流程:

1
2
3
4
5
6
7
8
1. AM 检测到所有 task 完成
2. AM 调用 RM 的 finishApplicationMaster(FINISHED, ...)
3. RM 标记 RMAppAttempt 状态 FINISHING → FINISHED
4. RM 标记 RMApp 状态 FINAL_SAVE → FINISHED
5. RM 通知 Client 作业完成
6. Client 拿到 FINISHED 状态后断开
7. RM 通知所有 NM 清理该作业的 container 工作目录
8. NM 后台清理 container 日志(默认保留时间 yarn.nodemanager.log.retain-seconds = 10800 秒即 3 小时)

异常终止(用户 kill / RM 拒绝 / AM 超过 max-attempts)走类似流程,但 RMApp 状态是 KILLED 或 FAILED。

注意一个细节——container 完成后日志不会立即删除,默认保留 3 小时。如果开启 YARN Log Aggregation(第十二篇展开),日志会被聚合到 HDFS 长期保存。如果没开启,container 日志在 NM 节点本地保留 3 小时后被 NM 清理。

yarn application -kill <appId> 是用户主动 kill 作业的命令。RM 收到后通知 AM 进程退出(通过 NM 杀 container),然后清理作业元数据。这个操作不能撤销——作业状态变 KILLED 后所有 task 都停止。

实验:观察生命周期状态

提交一个作业后用 yarn application -status <appId> 周期性观察状态变化:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
Application-Id:            application_1721900000000_0001
Application-Name: Spark Pi
Application-Type: SPARK
User: magicliang
Queue: root.dev
Start-Time: 1721900001234
Finish-Time: 0 ← 还在跑
Progress: 30.0%
State: RUNNING
Final-State: UNDEFINED
Tracking-URL: http://am-host:8088/proxy/application_...
RPC Server: am-host:37123
Queue: root.dev
Aggregate Resource Allocation: 1234567 MB-seconds, 234 vcore-seconds

Progress 字段从 0.0 到 100.0 反映 AM 上报的进度。State 字段从 RUNNING → FINALING → FINISHED。Final-State 从 UNDEFINED → SUCCEEDED / FAILED / KILLED。

yarn applicationattempt -list <appId> 列出该作业的所有 attempt:

1
2
3
4
ApplicationAttempt-Id:       appattempt_1721900000000_0001_000001
State: FINISHED
AMContainer: container_1721900000000_0001_01_000001
Tracking-URL: ...

如果 AM 重启过,这里会列出多个 attempt。

yarn container -list <appAttemptId> 列出某次 attempt 的所有 container:

1
2
3
4
5
Container-Id:               container_1721900000000_0001_01_000001
NodeId: node-1:8041
State: RUNNING
Log-URL: http://node-1:8042/node/containerlogs/...
Exit-Status: -1 ← 还在跑

这三级(application → attempt → container)的状态机是 YARN 调试的核心。

模式提炼

YARN 应用程序生命周期体现的设计模式:

1
2
3
4
5
6
7
8
模式:三级状态机 + 编排器重启 + 资源生命周期解耦

- 作业级状态机(RMApp)跟踪整体命运
- 尝试级状态机(RMAppAttempt)跟踪编排器(AM)每次启动
- 容器级状态机(Container)跟踪单个执行单元
- 编排器崩溃后从中央状态机恢复,作业续跑
- 客户端不参与运行时,提交完即可断开
- 资源生命周期与编排器生命周期解耦(container 不因 AM 重启而终止)

这个模式不只是 YARN。Kubernetes 的 Job / ReplicaSet / Pod 三级与 YARN 三级有相似性——Job 是用户意图、ReplicaSet 是某次部署尝试、Pod 是执行单元。Job controller 重启后从 etcd 拿到 Pod 状态继续编排。差异是 Kubernetes 的 controller 通常是常驻进程(多个作业共享一个 controller),YARN 的 AM 是每作业独占进程。

Borg / Omega 的"作业 → task → alloc"三级也是同构。Google 在 2015 Borg 论文里描述的 BNS(Borg Naming Service)与 YARN 的 ApplicationId 类似——为每个作业分配全局唯一标识,支持跨服务发现。

Spark Standalone 模式的 Driver / Executor / Task 三级是 YARN 三级的对照。Spark on YARN 把 Driver 包装成 AM,Executor 包装成 task container,复用 YARN 的生命周期。

工程迁移表

YARN 概念 Kubernetes Borg Spark Standalone Flink Cluster
RMApp(作业状态) Job / Deployment job Driver Job JobGraph
RMAppAttempt(编排器尝试) ReplicaSet generation job attempt Driver attempt ExecutionAttempt
Container(执行单元) Pod alloc Executor Task slot
AM(编排器) Job controller job master Driver JobManager
状态机持久化 etcd Borgmaster state ZooKeeper ZK HA
重试上限 backOffLimit max retry spark.yarn.maxAppAttempts restart strategy

注意 Kubernetes 这一列。Kubernetes 的 Job spec 有 backoffLimit 字段控制重试次数,与 YARN 的 yarn.resourcemanager.am.max-attempts 语义相同。差异是 Kubernetes 的"AM 重启"是 Job controller 自动重启 Pod,不需要每个作业独占 controller——controller 是常驻的,作业 spec 在 etcd 里持久化。

常见误解

误解一:“AM 启动后立即开始跑 task”。AM 启动后还要做几件事——初始化作业上下文、向 RM 注册、申请 container、等 container 启动。AM 启动到第一个 task 跑起来通常需要 10-30 秒。这个延迟让 YARN 上的短作业(几秒)启动开销显著,不如 Kubernetes Pod 启动快。

误解二:“AM 心跳失败立即让作业失败”。RM 等待 10 分钟(yarn.am.liveness-monitor.expiry-interval-ms)才判定 AM 死亡。这个延迟避免了网络抖动造成的误判。10 分钟内 AM 重连上 RM,作业继续跑。

误解三:“作业完成时所有 container 立即清理”。作业完成后 container 工作目录被立即清理,但 container 日志默认保留 3 小时(yarn.nodemanager.log.retain-seconds)。如果开启 YARN Log Aggregation,日志被聚合到 HDFS 后本地日志立即删除,HDFS 上长期保存。

误解四:“用户 kill 作业是异步操作”。yarn application -kill 是同步命令——RM 收到后立刻标记作业 KILLED,通知 NM 杀 container。命令返回时作业状态已经是 KILLED。但 container 进程的实际退出可能需要几秒(NM 收到指令 → kill 进程 → 进程响应 SIGTERM)。

误解五:“AM 重启后所有 task 都要重跑”。重试策略取决于计算框架。MapReduce 已完成的 map task 不重跑(输出在 HDFS)。Spark 在 AM 重启时整个作业从头开始(除非配置了 checkpoint)。Flink 从最近 checkpoint 恢复,仅重跑未 checkpoint 的部分。这是为什么流式计算框架对 checkpoint 周期敏感——checkpoint 越频繁,AM 重启代价越小。

练习

  1. 提交一个 MapReduce 或 Spark 作业,在 RM Web UI 上观察作业状态从 SUBMITTED → ACCEPTED → RUNNING → FINISHED 的变化。同时用 yarn applicationattempt -list <appId> 观察 attempt 状态。

  2. 故意在测试集群上 kill AM 进程(找到 AM container 所在 NM,kill -9 AM 进程),观察 RM 是否自动启动新 attempt。新 attempt 启动后作业是否继续跑(取决于 AM 是否支持状态恢复)。

  3. apache/hadoop 源码里找到 RMAppImpl.javaRMAppAttemptImpl.javaContainerImpl.java(hadoop-yarn-server 模块),观察三个状态机的状态转移图。

  4. 思考题:如果让 YARN 支持"应用级 checkpoint"(用户作业周期性把状态快照保存到 RM,AM 重启时从快照恢复),需要扩展哪些接口?这种机制对哪些计算框架(MapReduce / Spark / Flink)最有价值?

系列导航

序号 主题 状态
00-06 HDFS 存储层(00 导读 / 01-06 HDFS 各主题) 第一阶段(已完成)
07 YARN 架构:ResourceManager、NodeManager、ApplicationMaster 的三方契约
08 资源模型:Resource、Container 与 NodeLabel
09 调度器对比:FIFO、Capacity、Fair 的设计取舍 上一篇
10 应用程序生命周期:提交、调度、启动、运行、完成 本篇
11 YARN HA 与 Federation 下一篇
12 YARN Timeline Service v2:通用的应用历史与指标

参考资料