Flink源码阅读这件事,很多新手都会问从哪开始。我个人的经验是,绕过算子链、序列化那些硬骨头,先啃JobManager的HA机制,反而能把整个运行时组件串起来。JobManager是集群的调度大脑,负责作业调度、资源申请、检查点协调,一旦它宕机,作业基本就停摆了;如果没有HA机制兜底,恢复只能靠人工重启加手动恢复,这在生产环境里是不可接受的。所以JobManager的HA机制要解决的问题很明确:在多个候选节点之间维持一个稳定的单点领导权,并在任期失效后快速交接。
下面我以Flink 1.14/1.15附近的源码为主线(类名在新版本有小幅调整,但机制稳定),从为什么要做HA、骨架由哪些服务组成、启动入口在哪、ZooKeeper底层是怎么选主的、切换那一刻状态怎么恢复、以及我实际排查过的一些坑,一层层拆开讲。这篇文章适合两类人:一是准备深入Flink源码的开发者,二是正在为线上集群设计高可用方案的运维同学。
1. 为什么JobManager一定要做HA
1.1 先说清楚JobManager到底在干什么
在Flink的经典架构里,JobManager不是“一个类”,而是一组组件的合称:接收作业提交的Dispatcher、负责资源管理的ResourceManager、以及真正为每个作业生成ExecutionGraph并调度执行JobMaster,统统跑在JobManager进程里。你可以把JobManager理解成大脑,TaskManager是手脚。手脚再多,大脑一宕机,所有TaskManager会立刻发现找不到领导,作业状态卡住,Source不读了,Sink不写了,窗口也不触发了。
更麻烦的是,Flink的检查点机制依赖JobManager协调。JobManager要定期向Source端发Barrier,等所有TaskManager报告对齐完成,再提交CompletedCheckpoint。如果JobManager挂了,光这一点就无法继续。所以生产环境里JobManager绝对不能是单点,必须在多个节点上同时拉起候选进程,用外部协调服务决定谁能当领导。
这里有一个很容易误解的点:所谓的“多个JobManager”,并不是所有节点同时在对外服务。任何时刻只有一个Active JobManager持有领导权,其他节点是Standby。它们是“一主多备”,不是“多主并行”。多主同时干活在分布式系统里叫脑裂,Flink本身不解决脑裂,而是靠外部协调服务来尽量规避,后面我会专门讲。
1.2 没有HA时,一次故障处理有多痛苦
先想象一下没有HA的集群。JobManager所在机器突然断电,五个小时没人发现;就算有人发现了,手忙脚乱重启JobManager,作业还是无法直接恢复,因为上一个状态快照再哪里、作业图长什么样,普通管理员根本不知道。你需要重新提交作业,设置好Savepoint路径,手工触发一次restore。
这个过程的痛点不在于“能不能恢复”,而在于“恢复期间业务是断的”。对于7x24小时在跑的实时数仓来说,每断一分钟都可能是几十万条数据积压。引入HA机制后,故障切换时间可以从“小时级”降到“秒级到分钟级”,而且不需要人工介入。ZK检测到会话超时,待命节点自动接管,新JobManager从外部存储拿回CompletedCheckpoint,从最近一次检查点恢复作业状态,这个流程才是生产级的。
1.3 HA开关背后的运行形态差异
Flink的HA开关主要通过high-availability配置项控制,旧版本是high-availability: zookeeper,新版趋向于high-availability.mode: zookeeper。无论怎么写,本质都是让运行时在启动阶段选择不同的HighAvailabilityServices实现。
在不同部署模式下,HA的表现形式有差异。Standalone集群通常手动指定多个JobManager地址,Session集群可以多个JobManager进程同时存活;而YARN或Kubernetes模式下,Flink会自动拉起一个新的JobManager Pod作为Standby,甚至会在Active挂掉后重新创建。源码层面的选举逻辑是一样的,只是外部协调服务和进程生命周期管理不同。所以我的建议是:先在本地搭一个两节点的Standalone HA集群,看ZooKeeper里的节点变化,再看源码,比直接啃YARN/K8s模式要直观得多。
2. HA机制的骨架:两条服务链路与三种存储方案
2.1 一写一读:LeaderElectionService与LeaderRetrievalService
我第一次读HA相关源码时,最迷惑的是类名太多:ZooKeeperLeaderElectionService、ZooKeeperLeaderRetrievalService、LeaderContender、LeaderListener、LeaderElectionEventHandler……后来发现它们实际上是两条独立的链路,用“一写一读”来记就行。
第一条是LeaderElectionService,负责“选我当领导”。谁想成为候选者,就要实现LeaderContender接口,然后把自己的服务启动起来,调用leaderElectionService.start(this)。真正当选的那一刻,回调grantLeadership;失去领导权时,回调revokeLeadership。第二条是LeaderRetrievalService,负责“找到现在的领导”。非Leader节点或者客户端想连接JobManager时,通过LeaderRetrievalService注册LeaderListener,外部服务里Leader信息一变化,回调notifyLeaderAddress,拿到当前Leader的地址和Session ID。
这两条链路的底层在ZooKeeper实现里都依赖同一个路径,比如/flink/<cluster-id>/leader。选举服务往这个路径上写数据,发现服务从这个路径上读数据。把这两条链路的边界画出后,再去读源码,思路会清晰得多,不会在回调堆栈里迷路。
2.2 HighAvailabilityServices工厂在启动时做了什么
HighAvailabilityServices是HA功能的顶层抽象,Startup阶段由ClusterEntrypoint.initializeServices()创建。它的接口里暴露了一堆行为:创建LeaderElectionService、LeaderRetrievalService、CompletedCheckpointStore、CheckpointIDCounter,还有JobGraphStore。也就是说,HA不仅仅是“选主”,还包括“把元数据放到外部存储”这一整条链路。
以ZooKeeper实现为例,ZooKeeperHaServices构造时会创建CuratorFramework客户端,并用它封装出ZooKeeperLeaderElectionService和ZooKeeperLeaderRetrievalService;同时还会创建ZooKeeperCompletedCheckpointStore和ZooKeeperCheckpointIDCounter。这些服务创建完成后,后续组件的初始化都依赖它们。
从这里能看出,Flink把“高可用”的设计抽象得很好:所有跟外部协调系统打交道的逻辑,都被收拢到Success的实现类里。你在源码里看到haServices.createLeaderElectionService()这种调用,就不要往里面跳了,它只是“生产一个选举服务”;真正要看的是ZooKeeper实现里选举服务的启动和回调。
2.3 三种HighAvailabilityMode怎么选
Flink的HighAvailabilityMode枚举定义了三种模式:NONE、ZOOKEEPER、FACTORY_CLASS,新版还增加了KUBERNETES。我在表格里整理了一下它们的差异,方便对照:
| 模式 | 协调存储 | 适用场景 | 典型实现类 |
|---|---|---|---|
| NONE | 无 | 本地测试、单JobManager | DefaultHaServices |
| ZOOKEEPER | ZooKeeper | 绝大多数生产环境、混合部署 | ZooKeeperHaServices |
| KUBERNETES | Kubernetes ConfigMap | 已深度容器化的集群 | KubernetesHaServices |
| FACTORY_CLASS | 自定义 | 企业自研协调系统 | 自定义HighAvailabilityServicesFactory |
选择上不用太纠结。如果你的集群已经跑在Kubernetes里,可以用K8s模式,省掉ZK运维;否则老老实实上ZooKeeper。K8s模式底层其实也是走选举和监听机制,只是存储介质从ZK换成了ConfigMap。至于FACTORY_CLASS,一般公司不会自己造轮子,了解一下接口就好。
3. 源码追踪:JobManager是怎么把自己“卷”进选举的
3.1 从ClusterEntrypoint到DispatcherResourceManagerComponent
我读HA源码时建议从进程入口开始跟踪。Standalone集群的入口是StandaloneSessionClusterEntrypoint,它继承了ClusterEntrypoint,在main方法里最终调用了runCluster()。这个方法是HA机制的大本营:先初始化各种服务,再创建DispatcherResourceManagerComponent,然后把组件启动。
ClusterEntrypoint里的关键调用顺序大致是:先initializeServices(),生成HighAvailabilityServices;再createDispatcherResourceManagerComponent(),把haServices传给组件;最后在组件内部启动DispatcherRunner与ResourceManagerDriver。这里有个容易忽略的细节:initializeServices()执行得很早,如果ZK连不上,进程会直接报错退出。Flink在启动阶段依赖外部协调服务,不会“先启动再等ZK”,所以生产环境里ZK的可用性至关重要。
3.2 从DispatcherResourceManagerComponent到LeaderContender
DispatcherResourceManagerComponent这个类名很长,但它做的事情可以简化成两句话:创建Dispatcher和ResourceManager的Runner;把Runner包装成LeaderContender,交给选举服务。
在1.14版本里,真正的候选者实际上是DispatcherRunner。它内部持有一个LeaderElectionService,并实现了LeaderContender接口。组件启动时会调用dispatcherRunner.start(),也就是leaderElectionService.start(dispatcherRunner)。这一步意味着:候选者把自己的命运交给了选举服务,之后能不能成为Active,完全由外部协调服务决定。
读这个类时,我建议重点看grantLeadership和revokeLeadership这两个方法。它们负责创建/停止Dispatcher进程。在源码里你会看到相似的逻辑:当选后startDispatcher,失选后stopDispatcher。本质上就是把调度组件的生命周期跟领导权绑定。很多同学会误以为选完主就完了,其实选主只是一个触发器,真正的状态初始化在后面。
3.3 grantLeadership回调后的两件大事
grantLeadership回调触发后,不只是把Dispatcher拉起来,还连带做了另外两件重要的事。第一,把当前Leader信息注册到外部存储,这样其他节点通过LeaderRetrievalService能找到它;第二,从JobGraphStore和CompletedCheckpointStore里恢复之前持久化的作业。
具体到源码,DispatcherRunner当选后会调用dispatcherLeaderProcess,进程里会去创建Dispatcher,而Dispatcher实例化过程中会读取haServices.getJobGraphStore()。这一步读取完成后,才能根据作业图决定是重启还是恢复。整个过程在源码里写得很长,但拆开看就是:“选主成功—拉组件—读状态—对外服务”。理解了这一层,你再回看ZooKeeper里的节点数据,就会明白那些路径上存的是什么。
还有一个细节值得注意:revokeLeadership不一定发生在宕机场景。当ZK会话抖动、Leader节点因网络分区失联时,旧Leader的临时节点可能被ZK清理,此时旧Leader可能还没死。它收到revoke回调后,会主动关闭Dispatcher等组件,尽量避免两个JobManager同时在调度。但由于网络分区不确定性,这种“尽量避免”不能100%保证,生产上仍需要关注脑裂风险。
4. 读ZooKeeper实现:LeaderLatch背后的临时节点与监听器
4.1 ZooKeeperLeaderElectionService的启动流程
ZooKeeper实现里,选举服务的核心类是ZooKeeperLeaderElectionService。它的构造参数看起来有点吓人,实际核心只有CuratorFramework和LeaderLatch。start()方法做的事有两步:先把LeaderLatchListener注册到LeaderLatch上;然后调用leaderLatch.start()。
LeaderLatch.start()是Curator客户端的行为,它会与ZooKeeper服务端建立连接,在指定路径下创建一个临时顺序节点。节点创建后,Curator会立刻计算当前所有候选节点中谁的序号最小。序号最小的节点就会通过监听器回调notifyLeader(),从而进入Flink自己的事件处理逻辑。
这里有个并发细节容易让新手懵:notifyLeader()回调可能和start()调用发生在不同线程。源码里为了处理这个竞态,用了isLeader和synchronized关键字做同步。曾有一段时间Flink在Leader切换时偶发状态不一致,就是回调执行的时序问题。所以读者如果看到ZooKeeperLeaderElectionService里大量synchronized块,不用觉得奇怪,这是在跟Curator异步回调硬碰硬。
4.2 Curator的LeaderLatch内部原理
很多人问,为什么用LeaderLatch而不是LeaderSelector?Flink源码选型可以考虑,一个重要的原因是LeaderLatch的语义更简单:“我是小组里的1号,我就是Leader;号码变了,我就不是Leader”。
LeaderLatch的底层实现其实不复杂。每个候选者会在/flink/<cluster-id>/leader类似的路径下创建一个临时顺序节点,比如_c_xxxx000000001。ZooKeeper会保证子节点序号严格递增、唯一。LeaderLatch在启动后会读取全部子节点,如果发现自己是最小序号,就认为自己当选。同时它会注册一个监听器,监视小序号节点的删除事件。一旦序号比自己小的节点消失,它会重新检查自己的序号,如果自己变成最小,就触发Leader回调。
这个设计和现实生活里银行叫号排队几乎一模一样。每个人取一个号,最小的号先办业务;当前一个号被叫走后,下一个号自动顶上。临时顺序节点的巧妙之处在于:如果持号人突然离开(连接断开),ZooKeeper会自动作废他的号,后面的人可以往前补位。整个过程不需要人工干预,所以Flink能把故障切换做到自动。
4.3 Leader信息如何写入ZooKeeper并对外暴露
当选为Leader后,Flink还要把Leader的地址信息写到一个固定路径下,通常是/flink/<cluster-id>/leader这个节点。写入的内容是序列化后的LeaderInformation,里面包含JobManager的RPC地址和LeaderSessionID。
其他节点或客户端通过ZooKeeperLeaderRetrievalService去监听这个固定路径。它注册了一个Cache监听器,只要路径数据变化,就会触发handleLeaderChange,然后将新的地址广播给所有LeaderListener。这个过程非常关键:TaskManager启动时并也不知道谁是Leader,它必须先通过LeaderRetrievalService拿到地址,再建立RPC连接。如果直接写死JobManager地址,HA切换后TaskManager就再也不出来了。
我在实际调试中验证过这个机制:主节点切换后,去看ZK里/leader节点的数据,会发现内容跟着新节点变化;TaskManager日志里也会出现重新连接新地址的记录。这一块理解后,再去排查“TaskManager连不上JobManager”类问题,很快就能定位到是Leader信息没写对,还是RetrievalService没监听到。
5. 切换那一刻:恢复作业状态的核心链路
5.1 切换时状态从哪里来
很多人以为JobManager的HA就是“换一台机器继续跑”,但作业的运行状态怎么办?executionGraph还在老节点的内存里,如果内存状态没了,作业就只能从零开始。所以Flink的HA机制里另有一条重要链路:把关键状态持久化到外部存储,切换后重新读回来。
这个持久化重点是两个:CompletedCheckpointStore和JobGraphStore。ZooKeeper实现里,CompletedCheckpoint的元数据会写到ZK路径下,比如/flink/<cluster-id>/checkpoints/<job-id>,里面保存的是最新几个检查点的指针信息;而检查点的实际数据可能存储在文件系统、S3或HDFS。JobGraphStore则在ZK里保存作业图,这样新Leader才能知道“集群里有哪些作业,每个作业长什么样”。
5.2 新Leader从检查点重建ExecutionGraph的关键路径
新Leader选出来后,Dispatcher会启动,并逐个恢复JobGraph。恢复每条作业时,会从JobGraphStore读回作业图,从CompletedCheckpointStore找到最近的、有效的CompletedCheckpoint,然后用它来重建ExecutionGraph。
这个过程中,检查点是按“所有算子状态都对齐成功的检查点”来判断的。Flink持有一个检查点列表,恢复时从最新的往前找,找到能用的为止。如果ZK里保存的检查点元数据已经损坏,或者对应的数据文件已经丢失,恢复就会失败。源码里ZooKeeperCompletedCheckpointStore的recover()方法会做一次整理,剔除不完整或损坏的检查点,再返回可用的列表。
读这段源码时,不妨在本地做一个小实验:跑一个带状态的作业,触发几个checkpoint后强行杀掉Active JobManager,观察新Leader日志里“Restoring job X from checkpoint Y”这类信息。这会让你对“HA恢复”产生肌肉记忆,比单纯看代码有用得多。
5.3 抖动与脑裂风险
HA机制不是银弹,它最怕的就是“抖动”。为什么?假设A是Active,B是Standby。A因为网络抖动,与ZK会话超时,A的临时节点被删除;B立刻当选,写入新Leader信息。但这时A其实还活着,只是和ZK失联了片刻。A可能会继续跑几分钟,这就造成了短时间的双主。
Flink应对这个问题的方法是在RPC层使用LeaderSessionID。A会话期间签发的LeaderSessionID与B的新SessionID不同,TaskManager和客户端收到过期的LeaderSessionID后,会拒绝旧Leader的调度命令。但这并不保证绝对安全,因为如果A的网络分区只影响了部分组件,A仍然可能向部分TaskManager下发任务。所以真实生产里,ZK的sessionTimeout要结合机器负载合理设置,不宜过短,也不要指望靠代码解决一切脑裂问题。我自己的经验是:把zoo.cfg里的sessionTimeout和Flink配置里的HA超时时间拉开档位,既保证故障切换快,又减少无谓抖动。
6. 常见问题排查与源码阅读建议
6.1 典型问题速查表
我把自己和社区里遇到过的JobManager HA问题整理成一张表,排查时可以直接对着看:
| 现象 | 可能原因 | 排查手段 |
|---|---|---|
| Standby节点一直抢锁但不起作用 | 没正确配置JobManager地址列表,Standby进程没加入选举 | 检查masters文件、ZK中LeaderLatch节点数量 |
| 主节点频繁切换 | 网络抖动、ZK会话超时太短、GC停顿导致心跳延迟 | 看ZK日志、Flink日志中的SessionTimeout相关告警 |
| 作业恢复后状态不对 | 检查点恢复失败,可能从空状态启用了 | 看新Leader日志中Restoring from的检查点ID |
| TaskManager连不上新Leader | LeaderRetrievalService没有及时监听到路径变化 | 手动get ZK中的leader节点,对比新旧地址 |
| 进程启动后一直报ZK连接超时 | ZK集群不可用、网络不通 | 先单独用zkCli.sh验证客户端连接 |
这些问题的共同点,是它们都不在应用业务代码里,而在于HA组件和外部协调服务的配合上。所以排查时,第一步永远是“看ZK里现在的节点状态”,再结合Flink日志判断是哪条链路断了。
6.2 我踩过的几个坑
第一个坑是部署时没给ZK配持久化存储。ZK本身是无状态的,但如果把元数据放在本地磁盘,一旦ZK节点重建,Flink的作业恢复信息就全没了。很多生产事故并不是Flink挂了,而是ZK重建后所有HA元数据丢失。所以保证ZK数据持久化到可靠存储,尤其重要。
第二个坑是sessionTimeout设置得太短。有一次客户环境主节点每十分钟切换一次,查到最后是Flink作业GC停顿超过了ZK会话超时,导致主节点被误判为宕机。后来调整了ZK的sessionTimeout,同时优化了JobManager的堆内存和GC参数,才算稳住。
第三个坑比较隐蔽:新旧版本配置项不兼容。Flink 1.14之后,high-availability.zookeeper.quorum这类配置被整理到了high-availability.zookeeper.quorum,但很多老运维还拿着旧配置往新版本上套,导致HA一直没生效。判读的简单方法:启动日志里是否出现了ZooKeeperHaServices相关字样,如果没有,说明HA根本没被激活。
6.3 给新手的源码阅读路径
最后聊聊怎么读Flink HA源码效率最高。我的路径是:
- 先搭一个两节点Standalone HA集群,ZK里手动观察节点变化,形成感性认识;
- 从
ClusterEntrypoint入口打断电,一路走到DispatcherResourceManagerComponent; - 看
ZooKeeperLeaderElectionService,重点理解start和LeaderLatchListener的回调; - 顺着
grantLeadership进入Dispatcher,再看它怎样从JobGraphStore和CompletedCheckpointStore恢复作业; - 最后回过来看
LeaderRetrievalService,理解客户端和TaskManager是怎么发现新Leader的。
按这个顺序读,你只需要Flink runtime模块的部分源码,再加上Curator的基础知识,就能把整个HA链路串起来。我强烈建议不要一上来就看ZooKeeperHaServices的所有方法,这个类太庞大了,很容易看吐。先把选主和发现这两条主线走通,再去看存储和恢复链路,会轻松很多。
还有一个小习惯:在IDE里把断点打在LeaderContender.grantLeadership和revokeLeadership上,然后手动杀掉Active JobManager,观察断点如何被Standby节点触发。这比读一千行代码都能让你理解“选主”这件事的本质。我自己当初就是靠这个实验,彻底弄明白了整套HA机制的运转逻辑。