上一篇讲了应用程序生命周期。本篇讲 ResourceManager 怎么解决单点问题和集群规模上限问题——YARN HA 和 YARN Federation。

YARN HA 常被介绍成"两个 ResourceManager 互为热备"。这个描述对应了部署形态但漏了几个关键点:RM State Store 怎么持久化作业元数据、Standby RM 怎么与 Active 同步、Active 切换时 NM 和 AM 怎么重连。YARN Federation 常被简化成"多个 YARN 集群",但漏了 Router、AMRMProxy、SubCluster 这些核心抽象。准确的说法是:YARN HA 通过 ZooKeeper 选主 + RM State Store 共享状态实现 Active/Standby 切换;YARN Federation 通过无状态 Router + AMRMProxy 让多个独立子集群对客户端表现为单一集群,每个子集群内部仍然有自己的 HA。

本篇只抓一个问题:YARN RM HA 的切换流程是怎么工作的、RM State Store 的几种实现各自适合什么场景、YARN Federation 在什么规模下值得启用。

ResourceManager HA 的状态共享机制

Hadoop 2.4 引入了 YARN RM HA。与 HDFS HA 类似,YARN RM HA 也是 Active / Standby 双节点 + ZooKeeper 选主 + 共享状态存储。但有一个关键差异——HDFS HA 用 JournalNode 共享 EditLog(准同步),YARN HA 用 RM State Store 共享作业元数据(异步同步)。

RM State Store 的几种实现:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
FileSystemRMStateStore(HDFS)
- 把 RMApp / RMAppAttempt 元数据写到 HDFS
- 切换时新 Active 从 HDFS 加载全部状态
- 持久化最严格,但每次写状态要走 HDFS RPC,吞吐有限

ZKRMStateStore(ZooKeeper)
- 把 RMApp / RMAppAttempt 元数据写到 ZooKeeper znode
- ZooKeeper 自带多数派复制,不需要额外 HA 组件
- 写入延迟低(毫秒级),但 ZooKeeper 单 znode 数据有上限(默认 1MB)

LeveldbRMStateStore(本地 LevelDB)
- 把 RMApp / RMAppAttempt 元数据写到 RM 节点本地的 LevelDB
- 切换时新 Active 必须能访问相同 LevelDB(共享存储)
- 性能最好,但要求共享存储(NFS 等),实际部署少
- 通常用于单 RM 的"恢复"场景,不用于 HA

MemoryRMStateStore(内存)
- 测试用,不持久化

生产环境通常用 ZKRMStateStore——它复用已有的 ZooKeeper 集群,不引入新组件,性能也够。FileSystemRMStateStore 在对持久化最严格的场景使用,但 HDFS 写入延迟让 RM 处理速度受限。

RM State Store 里存的具体内容:

1
2
3
4
5
6
7
8
9
10
/yarn-leader-election            (Active 锁,临时节点)
/rmstore
/ApplicationId-1
/appattempt-1 (ApplicationSubmissionContext + attempt 状态)
/appattempt-2 (重试 attempt 的状态)
/ApplicationId-2
/appattempt-1
/epoch (RM 重启计数器,防止旧 RM 切换后冲突)
/AMRMTokenSecretManager (AM-RM 认证令牌密钥)
/ReservationSystem (reservation 数据,如果开启)

每次 RM 处理 submitApplication、createApplicationAttempt、updateApplicationAttempt 状态时都同步写 State Store。Active RM 切换为 Standby 时,新 Active 从 State Store 加载全部状态,重建 RMApp / RMAppAttempt 内存对象。

Standby RM 怎么同步

YARN HA 的设计是"Standby 不实时 tail"——Standby RM 处于"被动等待"状态,不主动处理客户端 RPC,但定期从 State Store 拉取最新状态保持内存接近 Active。

这种设计与 HDFS HA 不同。HDFS Standby NameNode 实时 tail JournalNode 上的 EditLog,与 Active 内存差异在毫秒级。YARN Standby RM 是"周期性 reload State Store",内存差异可能在秒级甚至更大。

为什么 YARN 选这个设计?因为 YARN 的工作负载不需要 Standby 实时同步——RMApp 状态变化的频率(每秒几次到几十次)远低于 HDFS EditLog(每秒几千次)。Standby RM 周期性 reload 已经足够,不需要 tail 机制。

切换时的具体行为:

1
2
3
4
5
6
7
1. Active RM 失去 ZooKeeper session
2. /yarn-leader-election 临时节点消失
3. Standby RM 的 EmbeddedElector 检测到锁消失
4. Standby RM 重新加载 State Store(拉取最新 RMApp 状态)
5. Standby RM fence 旧 Active(如果可能)
6. Standby RM 提升为 Active,开始接受 RPC
7. NM 和 AM 通过 RPC 重试感知新 Active

整个过程通常在 10-30 秒内完成。这个延迟里 NM 和 AM 的 RPC 失败,需要重试。

NM 和 AM 怎么重连

Active RM 切换时,NM 和 AM 的 RPC 都会失败。YARN 客户端库内置了重试逻辑:

NM 重连:每个 NM 通过 ResourceTrackerService 与 RM 通信。Active RM 切换时,NM 心跳失败,NM 自动重试其他 RM 地址(配置 yarn.resourcemanager.ha.rm-ids 列出所有 RM)。重试时 NM 会重新注册自己(newActiveRM.registerNodeManager),新 Active 从 NM 的注册信息重建该节点的资源状态。

AM 重连:AM 通过 ApplicationMasterProtocol 与 RM 通信。Active RM 切换时,AM 的 allocate 心跳失败,AM 自动重试其他 RM。重试时 AM 会重新调用 registerApplicationMaster,新 Active 从 State Store 拿到该作业的进度状态,AM 继续跑。

这种"客户端重试"让 RM HA 切换对 NM 和 AM 基本透明——RPC 失败时自动切换,业务层无感。但 allocate 心跳里未确认的请求可能丢失(例如刚发出的 container request),AM 需要重新发送。

注意一个细节——AMRMToken 在切换时可能失效。AM 与 RM 的 RPC 用 AMRMToken 认证,token 由 Active RM 签发。Active 切换时新 Active 重新加载 AMRMTokenSecretManager 状态,让旧 token 仍然有效。这个机制保证 AM 不需要在切换时重新认证。

RM Work-Preserving Restart

YARN 2.4 引入了 Work-Preserving Restart(WPR)。如果 Active RM 进程崩溃但 State Store 完好,新 Active 加载 State Store 后可以恢复所有 RMApp 的状态,作业续跑而不是失败重提。

WPR 的具体行为:

1
2
3
4
5
6
1. Active RM 进程崩溃
2. ZooKeeper 锁消失,Standby 提升
3. 新 Active 从 State Store 加载所有 RMApp / RMAppAttempt
4. NM 心跳重连,新 Active 重建 BlocksMap-like 的节点资源表
5. AM 心跳重连,新 Active 告知 AM 当前已分配的 container 状态
6. 作业继续跑

WPR 之前(YARN 2.0-2.3),RM 崩溃后所有作业直接失败,必须人工重新提交。WPR 之后,RM 崩溃只是几秒钟的 RPC 故障,作业自动续跑。这是 YARN HA 的实际价值——不只是切换,更是切换后作业不丢。

WPR 要求 NM 和 AM 配合——NM 重连时要上报完整的 container 状态(哪些 container 还在跑、哪些已完成),AM 重连时要从新 Active 拿到自己的 container 历史然后决定怎么处理未完成的 task。这种协调通过 WorkPreservingRMRestart 系列配置控制。

YARN Federation 解决什么问题

YARN 单 RM 的规模上限受 JVM 堆内存约束(类似 HDFS NameNode)。一个有 100 GB 堆的 RM 大概能管理 5000 节点、10 万并发作业。超过这个规模的集群——典型场景是大型互联网公司、公有云大数据服务——单 RM 装不下。

Hadoop 3.1 引入了 YARN Federation(YARN-5598),把多个独立 YARN 子集群组合成逻辑单一集群:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
          ┌─────────────────────────┐
│ Router (stateless) │
│ - am-rm-proxy RPC │
│ - webapp aggregator │
└─────────────────────────┘

┌───────────┴───────────┐
│ │
┌───────▼──────┐ ┌───────▼──────┐
│ SubCluster-1 │ │ SubCluster-2 │
│ RM (HA) │ │ RM (HA) │
│ + NMs │ │ + NMs │
│ + AMRMProxy │ │ + AMRMProxy │
└──────────────┘ └──────────────┘

YARN Federation 的核心组件:

Router 是无状态的网关进程,可以部署多个实例做负载均衡。客户端只配置 Router 地址,Router 内部根据作业类型、队列、用户决定路由到哪个子集群。Router 同时聚合多个子集群的 Web UI(作业列表、集群状态),让客户端看到单一视图。

State Store(ZooKeeper 或 SQL)存储路由表、子集群状态、Policy 描述。Router 共享这个状态。

AMRMProxy 是每个 NM 上可选的本地代理(不在所有部署里启用)。AM 默认直接连 RM 的 ApplicationMasterProtocol,启用 AMRMProxy 后 AM 连本地 NM 的 AMRMProxy,AMRMProxy 再转发给子集群 RM。这让 AM 可以透明地跨子集群申请 container——AMRMProxy 内部处理跨子集群 RPC。

YARN Federation 的核心收益是规模扩展——把单 RM 上限从 5000 节点扩展到任意规模(每个子集群 5000 节点,N 个子集群 5000N 节点)。代价是增加了运维复杂度,跨子集群操作(例如跨子集群 Shuffle)有额外开销。

Federation 的应用场景

Federation 适合什么场景?

第一种,超大规模集群。单集群超过 5000 节点、单 RM 内存压力大、调度延迟上升。这时把集群切成多个子集群,每个子集群独立 RM,对外提供 Router 统一入口。典型场景是大型互联网公司的全公司大数据平台。

第二种,多机房部署。集群跨多个数据中心(合规、延迟、容灾),每个数据中心一个子集群,Router 做跨机房路由。AMRMProxy 让 AM 在不同机房间调度 container,实现"逻辑单一集群 + 物理多机房"。

第三种,混合云部署。私有云子集群 + 公有云子集群(按需扩展),Router 按策略路由作业到不同云。私有云满载时溢出到公有云。

Federation 不适合什么场景?

第一种,小集群(< 1000 节点)。Federation 的运维复杂度(Router、State Store、AMRMProxy)远大于单 RM 部署,小集群收益不抵开销。

第二种,需要全局调度策略的场景。Federation 内每个子集群独立调度,没有"全局最优"的调度决策。如果作业需要精确的全局公平(所有作业不论子集群都要 1/N 资源),Federation 无法满足。

第三种,跨子集群 Shuffle 密集的场景。MapReduce / Spark 的 Shuffle 走本地磁盘 + 网络传输,跨子集群 Shuffle 意味着跨机房网络,性能显著下降。这种场景应该让作业在单个子集群内完成。

实验:观察 YARN HA 状态

部署 HA 的 YARN 集群可以用 yarn rmadmin 命令管理 HA:

1
2
yarn rmadmin -getServiceState rm1
yarn rmadmin -getServiceState rm2

输出 activestandby,标识每个 RM 当前角色。

切换演练:

1
2
yarn rmadmin -transitionToStandby rm1
yarn rmadmin -transitionToActive rm2

这两条命令手动切换 Active。生产环境用 -failover 命令做完整切换(包含 fencing):

1
yarn rmadmin -failover rm1 rm2

这条命令尝试把 Active 从 rm1 切到 rm2,过程中 fence rm1。

观察 ZooKeeper 上的状态:

1
2
3
4
5
zkCli.sh -server zookeeper-host:2181
[zk: ...] ls /yarn-leader-election
[rm1, LockNode-0000001]
[zk: ...] ls /rmstore
[ApplicationId-1, ApplicationId-2, epoch, AMRMTokenSecretManager]

/yarn-leader-election 包含参与选主的 RM 列表和当前 Active 锁。/rmstore 包含所有作业元数据。

YARN RM Web UI 的顶部有 “HA State: active” 或 “HA State: standby” 标识。两个 RM 的 Web UI 都可以访问,Standby 的 UI 上提示当前不接受 RPC。

观察 Federation 状态(如果集群启用):

1
yarn federationadmin -getSubClusters

输出所有子集群的 ID、状态、能力。Router 的 Web UI(默认 8089 端口)展示聚合视图——所有子集群的作业列表合并显示。

模式提炼

YARN HA 与 Federation 体现的设计模式:

1
2
3
4
5
6
7
8
9
10
11
12
13
模式 A:选主 + 共享状态 + 客户端重试

- 多节点通过共识算法选主,Active 锁由外部服务(ZooKeeper)持有
- 业务状态写到共享存储(State Store),切换时新主加载
- Standby 不实时 tail,按需 reload
- 客户端内置 RPC 重试,对切换基本无感

模式 B:无状态路由层 + 多个独立子系统

- 把单一控制节点拆成多个独立子系统(子集群)
- 上层加无状态路由层(Router),对客户端屏蔽内部切分
- 共享状态(State Store)用共识系统保证一致
- 子系统内部仍然有自己的 HA,可靠性更高

这两个模式与 HDFS HA / 路由联邦完全同构(参见第五、第六篇)。HDFS 和 YARN 都用 ZooKeeper 选主、都用共享 State Store、都在大规模场景引入 Router。这是 Hadoop 项目里"用同一套思路解决相似问题"的体现。

不只 Hadoop,Kubernetes 的多集群管理(Karmada、Cluster API)也是同一思路——多个 K8s 集群上面加无状态控制层。数据库的 sharding(Vitess、Sharding-Proxy)也是同一思路——多个数据库节点上面加路由层。这种"分而治之 + 无状态路由"是分布式系统扩展性的标准模式。

工程迁移表

YARN 概念 HDFS Kubernetes Mesos Borg
Active / Standby RM Active / Standby NN API Server HA Mesos Master HA Borgmaster HA
RM State Store JournalNode + EditLog etcd ZooKeeper / etcd Paxos store
ZKFC / EmbeddedElector ZKFC leader election lease ZooKeeper Paxos
切换时间 秒级 秒级 秒级 秒级
Work-Preserving Restart 是(FSImage + EditLog) 是(etcd state) 是(framework state)
Federation Router HDFS RBF Karmada / Cluster API Marathon LB GFE
AMRMProxy - - scheduler driver -

注意 Kubernetes 这一列。Kubernetes 的 API Server HA 用 etcd 作为共享状态存储,多 API Server 通过负载均衡器接收请求。这与 YARN HA 不同——YARN 是 Active/Standby,Kubernetes 是 Active/Active(多 API Server 同时接受请求,靠 etcd 保证一致)。差异源于工作负载特性——API Server 是无状态查询为主,可以并发;RM 是有状态调度,需要单点决策。

常见误解

误解一:“HA 部署后 ResourceManager 永远不会宕机”。HA 解决的是单点切换速度,不是消除宕机。RM 仍然会宕机(OOM、节点故障),HA 让 Standby 在几十秒内接管。如果切换过程中作业 RPC 失败,需要应用层重试。

误解二:“Standby RM 实时与 Active 同步”。YARN Standby RM 是周期性 reload State Store,不实时 tail。这个设计与 HDFS Standby NameNode 实时 tail 不同。差异源于工作负载特性——YARN 状态变化频率低,不需要实时同步。

误解三:“Federation 让 YARN 像单一集群一样”。Federation 内每个子集群独立调度,没有全局调度决策。客户端通过 Router 看到单一视图,但实际调度仍然在子集群层面。跨子集群操作(Shuffle、Join)有显著开销。

误解四:“RM 切换后所有作业都要重跑”。Work-Preserving Restart 让新 Active 从 State Store 恢复所有 RMApp 状态,作业续跑。只有切换过程中正在启动的 container 可能失败重试,整体影响有限。

误解五:“YARN HA 必须用 ZooKeeper”。YARN HA 的选主和 State Store 可以用其他方案。选主可以用 EmbeddedElector(基于 ZooKeeper)或者自定义 ElectorService。State Store 可以用 ZooKeeper、HDFS、LevelDB。生产部署 ZooKeeper 最常见,但不是唯一选择。

练习

  1. 在 HA 部署的 YARN 集群上运行 yarn rmadmin -getServiceState rm1yarn rmadmin -getServiceState rm2,确认当前 Active。然后用 yarn rmadmin -failover rm1 rm2 触发切换,观察客户端在切换过程中的 RPC 行为。

  2. 在 ZooKeeper 客户端执行 ls /yarn-leader-electionls /rmstore,观察选主节点和状态存储。

  3. apache/hadoop 源码里找到 RMHAProtocolService.javaZKRMStateStore.javaFederationStateStore.javaAMRMProxy.java(hadoop-yarn-server 模块),观察 HA 与 Federation 的核心实现。

  4. 思考题:如果让 YARN Standby RM 也实时 tail Active 的状态变化(像 HDFS Standby NameNode),切换延迟会怎么变化?这种实时同步会带来什么新问题?

系列导航

序号 主题 状态
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:通用的应用历史与指标 下一篇

参考资料