
第一次在生产环境把JobManager所在节点直接拔电是在一次季度故障演练里。当时心里其实没底提前配了ZooKeeper做HA但真到拔电那一刻你根本不知道Flink内部会不会按文档说的那样自动把Standby节点顶上来。结果作业确实没挂几秒钟的延迟抖动之后整个集群恢复了调度作业从最近的checkpoint继续跑。也就是从那次开始我决定把JobManager的HA机制源码认真啃一遍——与其靠文档猜不如直接看代码里每个组件是怎么协作的。这个机制的核心并不复杂但涉及的东西很散Leader选举、元数据持久化、TaskManager注册信息管理、JobGraph存取甚至还有Blob存储。这些环节全部串起来才构成了JobManager的高可用。这篇博文我会从源码层面拆一遍重点说清楚三个问题HA到底把哪些状态持久化了、Leader选举是怎么做的、故障恢复时各组件按什么顺序工作。适合已经跑过Flink任务、想深入理解集群内部原理的读者也适合准备排查HA相关问题的运维同学。1. 从运行架构看JobManager的HA到底要解决什么问题1.1 为什么JobManager必须做主备先抛开源码回到Flink集群的最基本模型。一个集群里通常有一个或几个JobManagerJobManager负责接收作业、调度任务、协调checkpoint、管理TaskManager的资源。如果它挂了整个集群就失去了调度中枢所有运行中的作业都会失败且无法提交新作业。这跟TaskManager挂掉有本质区别。TaskManager挂掉之后JobManager可以把任务重新调度到其他存活节点因为作业的完整执行计划、当前检查点信息都还掌握在JobManager手里。但JobManager自身挂掉就没人知道作业应该怎么继续调度了。所以HA机制的核心命题是怎样让另一个独立的JobManager进程能够无缝接管集群且不能出现两个大脑同时发号施令的情况。这就引出了两个必须解决的能力Leader选举同一时刻只能有一个Active JobManager其他实例都必须处于Standby状态。元数据共享新的Active JobManager必须能拿到之前调度所需的所有状态否则接管之后也无法恢复作业。源码里负责这两件事的正是HighAvailabilityServices这个接口以及它的一堆实现类。搞清楚这个接口就找到了阅读HA机制的入口。1.2 HA机制到底持久化了哪些状态很多人在配置HA时只记得写high-availability: zookeeper却忽略了Flink在后台到底往ZooKeeper里塞了什么东西。我读源码后整理了一张清单重要状态分四类JobGraph作业的执行计划逻辑包括作业顶点、连接关系、并行度、算子配置、用户jar包引用。JobManager接管后要重新调度作业第一件事就是从持久化存储里读回JobGraph。CompletedCheckpoint已经完成的checkpoint的元数据。这里存的是指向外部存储的路径和状态句柄不是全量状态数据。没有这些元数据新的JobManager不知道从哪个checkpoint恢复作业。Blob用户jar包、任务所需的二进制大对象。如果新的JobManager进程在同一台机器上还容易处理跨节点恢复就需要从共享的Blob存储里拉取。RunningJobRegistry当前正在运行的作业注册表里面标记了作业运行状态、对应的JobManager地址以及一些并发版本信息。这四类状态的持久化位置会随着HA后端的不同而变化。ZooKeeper做HA时JobGraph和CompletedCheckpoint元数据都存在ZooKeeper节点上Kubernetes做HA时对应信息放在ConfigMap或者外部持久卷中。但不管存储介质怎么变HighAvailabilityServices接口提供的抽象方法始终一致这也是为什么Flink能轻易扩展多种HA后端的原因。2. 源码主干HighAvailabilityServices接口解剖2.1 接口里藏着的核心方法我推荐直接打开flink-runtime模块下的org.apache.flink.runtime.highavailability.HighAvailabilityServices接口这里定义了HA机制对外提供的所有能力。从方法命名就能看出它不是在单个服务里做一件事而是把整个JobManager所需的高可用能力做了一个统一抽象。核心方法大致有以下几组createLeaderElectionService()为JobManager或ResourceManager创建Leader选举服务。createLeaderRetrievalService()创建Leader查询服务供客户端、TaskManager等组件获取当前Leader地址。getJobGraphStore()获取JobGraph持久化存储服务的句柄。getCompletedCheckpointStore()获取已完成checkpoint元数据存储服务的句柄。getRunningJobsRegistry()获取运行作业注册表服务句柄。getBlobStore()获取Blob存储句柄。这些方法返回值几乎都是接口类型比如LeaderElectionService、LeaderRetrievalService、JobGraphStore、CompletedCheckpointStore。实际底层是ZooKeeper还是Kubernetes全部隐藏在实现类里。调用方只依赖接口做交互不需要关心底层细节。这其实是阅读Flink HA源码时最重要的心智模型HA不是某一个类而是一整套服务集合。JobManager在启动时通过HighAvailabilityServices拿到这些服务的实例然后把它们注入到各个组件中。如果想看某一块的持久化逻辑直接找对应的Store实现很快就能定位。2.2 ZooKeeperHaServices如何把组件组装起来以最经典的ZooKeeper实现为例类名是ZooKeeperHaServices。这个类继承自AbstractHaServices构造的时候接收ZooKeeper客户端相关的配置以及用于区分HA集群的ResourceID、ClusterID。创建ZooKeeper服务时关键点是它内部使用Curator Framework操作ZooKeeper。Flink的ZooKeeperHA不是自己裸写ZooKeeper协议而是基于CuratorFramework来做连接管理、重试和监听。所有ZooKeeper节点路径都有一个固定的前缀结构默认形如/flink/cluster-id/component。ZooKeeperHaServices里有一个很重要的内部方法createRecoverableStateStoreHelper它负责创建各类持久化Store。以getCompletedCheckpointStore()为例它最终会构造一个ZooKeeperCompletedCheckpointStore内部维护着一个ZooKeeperStateHandleStore后者负责把序列化后的checkpoint元数据写入ZooKeeper的指定节点。去看ZooKeeperStateHandleStore时能看到它每次写数据都会加锁还带版本校验。这是为了防止多个JobManager同时写同一个checkpoint元数据造成冲突——虽然理论上只有一个Leader会写但Flink在底层还是做了防御性设计。这些细节如果只是看配置文档根本不会注意到。从宏观来看ZooKeeperHaServices好比一个工厂把所有高可用组件都组装好交给JobManager启动流程使用。而JobManager启动时具体怎么消费这些组件就涉及下一个重点JobManagerRunner中的Leader选举。3. 深入Leader选举从接口到ZooKeeper实现3.1 LeaderElectionService与LeaderRetrievalService的分工我见过不少人在排查HA问题时分不清选举和发现这两个概念导致日志看不懂。其实源码里分得很清楚LeaderElectionService是给候选者用的它负责让当前进程参与选举当选后通知回调。JobManager启动时会调用start(LeaderContender)如果进程成为Leader服务会回调LeaderContender.grantLeadership()并把sessionID传给调用方。如果失去Leader身份则回调revokeLeadership()。LeaderRetrievalService是给依赖者用的比如TaskManager、JobClient、ResourceManager它们需要知道当前Leader的地址通过LeaderRetrievalListener接收Leader变化通知。通知里包含Leader的地址和sessionID。用一个生活化的类比选举服务好比是员工自己参加竞聘当选后拿到聘书查询服务则是外部客户查公司现在谁是总经理然后按这个地址去对接业务。为什么Flink要把这两个服务拆开因为参与者和观察者天然就是两类不同的角色。JobManager自身只需要关心自己有没有当选不需要知道另一个JobManager的地址而TaskManager不参与选举却必须知道当前Active JobManager连接到哪里。拆开后两边可以各自使用不同的ZooKeeper监听策略也方便做权限控制。3.2 ZooKeeperLeaderElectionService的源码细节ZooKeeper实现下ZooKeeperLeaderElectionService内部用的是Curator的LeaderLatch。核心代码逻辑简洁但值得细看LeaderLatch会在ZooKeeper上创建一个临时顺序节点所有参与选举的JobManager都在同一个路径下创建节点临时顺序节点编号最小者成为Leader。其他节点监听这个路径一旦Leader节点的会话失效或节点被删除就会触发重新选举。Flink里的ZooKeeperLeaderElectionService.start()会调用leaderLatch.start()并注册一个LeaderLatchListener。当选为Leader后Listener会回调leaderLatchEvent最终由LeaderElectionDriver封装成Flink自己的回调。这里有一个很多初学者容易踩坑的点Curator的LeaderLatch在失去Leader后会进入重新选举的等待状态但Flink的回调并不总是立即重置。如果JobManager实例在短时间内连续当选、失去、再当选底层回调的时序处理稍有改动就会导致重复的suspendLeadership和grantLeadership。所以Flink在ZooKeeperLeaderElectionService里专门加了suspendedLeadership状态用于标识临时失去连接而不是直接失去Leader身份。这个状态在日志里通常表现为ZooKeeper连接断开后不会立刻切换Leader。LeaderContender接口里还有一个关键方法是handleError。如果选举过程遇到不可恢复的异常比如会话超时后节点被删除且重连失败服务就会调用这个回调通知JobManager放弃Leader职位。JobManager收到回调后会停止正在执行的调度动作重新进入Standby状态。3.3 JobManagerRunner中的选举回调执行流程只看选举服务还不够要真正理解当选之后发生什么就得看JobManagerRunnerImpl这个类。它是Flink 1.15之后统一的任务运行器负责启动JobManager并管理其生命周期。当ZooKeeperLeaderElectionService回调grantLeadership后LeaderContender会执行以下步骤先通过HighAvailabilityServices获取当前自己的地址并注册到Leader检索服务。从HighAvailabilityServices.getJobGraphStore()中读取作业列表尝试恢复之前提交的作业。如果存在需要恢复的JobGraph则创建相应的JobMaster并触发JobManagerRunner接管。将自身状态变为Running这时候TaskManager通过LeaderRetrievalService查询到新Leader地址开始注册。如果你去看JobManagerRunnerImpl的grantLeadership方法会发现它会先检查当前jobManagerRunner是否处于终止状态然后调用leaderElectionService.confirmLeadership(sessionID, leaderAddress)确认身份。确认成功后才真正启动JobManager。我特别想提醒一点选举成功不代表立刻接管所有作业。中间还有一堆恢复动作任何一步失败都会导致本次接管失败。比如JobGraph恢复不出来、checkpoint元数据不可读、Blob下载失败都可能让新的JobManager启动过程卡住。排查时要看JobManagerRunnerImpl的完整日志不要只看成了Leader就开始判断。4. 恢复链路的源码视角JobGraphStore与CompletedCheckpointStore4.1 JobGraphStore把作业藏到哪了JobGraphStore在ZooKeeper实现下的名字是ZooKeeperJobGraphStore。每次提交作业时Dispatcher会调用jobGraphStore.putJobGraph(new JobGraph)把作业图数据写入持久化节点。存储的路径由HA配置的high-availability.zookeeper.path.jobgraphs决定目录结构类似/flink/cluster-id/jobgraphs/job-id。ZooKeeper是典型的CP系统节点里的数据量必须控制在一定范围。我见过有人把几百MB的大JobGraph往ZooKeeper里塞结果整个集群的ZooKeeper性能都受影响。普遍经验是JobGraph本身不应该太大如果一个作业的JobGraph序列化后超过几MB就该考虑是不是用户代码里塞了不该塞的大对象。从源码看ZooKeeper存储JobGraph之前会序列化为JobGraph的字节流存储时小于1MB的数据直接存在节点里过大的数据其实会有失败风险。新增版本的Flink对ZooKeeper节点的数据大小也做了校验超限会直接抛异常。这是很多人在作业提交失败后才发现的隐藏坑。4.2 CompletedCheckpointStore的恢复顺序作业恢复时JobMaster需要从CompletedCheckpointStore里找到最近一次成功的checkpoint。在ZooKeeper实现下ZooKeeperCompletedCheckpointStore维护了一个按顺序排列的checkpoint句柄列表。每次checkpoint完成新的元数据会写入ZooKeeper旧的超过保留数量后会被清理。从源码看recover()方法它会遍历已存储的所有checkpoint句柄逐个校验是否完整。校验包括句柄对应的路径是否存在、能否从外部存储中读取状态元数据。如果最近的一个checkpoint元数据因为各种原因读不出来Flink并不会直接失败而是会继续往前找可用的更早的checkpoint。这种回退恢复的设计极大提升了恢复的成功率。不过也正因为有这个回退机制如果你看到作业从某个较旧的checkpoint恢复别惊讶。日志里会明确记录使用了哪个checkpoint ID。排查时把checkpoint ID和最近完成时间对照基本能判断恢复是否合理。4.3 BlobStore与运行作业注册表在故障转移中的角色BlobStore是容易被忽略的一环。作业恢复时新的JobManager可能运行在不同的节点上用户代码的jar包如果只存在于旧的JobManager本地就会因为环境不同而找不到类。BlobStore解决的就是这个问题提交作业时jar包会上传到共享Blob存储中恢复作业时从共享存储拉取。ZooKeeper HA模式下共享存储路径由high-availability.storageDir配置决定可以是HDFS、S3或本地文件系统。RunningJobsRegistry则负责标记作业运行实例的归属。ZooKeeper实现里它会记录当前运行的作业ID与JobManager的对应关系。这个注册表的一个关键作用是防止并发恢复同一个作业——如果两个JobManager同时尝试恢复注册表里的锁机制会让其中一个失败。综合来看恢复链路是多个组件按顺序协作的过程先读JobGraphStore得到作业逻辑再读CompletedCheckpointStore拿到checkpoint句柄从BlobStore拉取依赖jar包最后通过RunningJobsRegistry锁定作业归属并启动JobMaster。只要其中一个环节故障作业就可能无法自动恢复。5. Kubernetes HA与ZooKeeper HA的源码差异5.1 KubernetesLeaderElectionService的实现思路Flink从1.9开始支持原生的Kubernetes HA类名叫KubernetesLeaderElectionService。它的底层不是用Curator而是直接利用Kubernetes的Lease资源对象做leader选举。每个候选JobManager进程都会尝试更新Lease对象的Spec.HolderIdentity字段持有Lease的对象即被认为是Leader。从源码看Kubernetes实现里有个关键设计只有当前持有Lease的进程才有资格更新renewTime。其他候选者定期读取Lease对象如果发现renewTime超过租约期限就可以发起抢占。这里期间涉及到Kubernetes API Server的监听机制所有候选者通过Watch实时感知Lease变化从而触发重新选举。相比ZooKeeper实现Kubernetes HA不需要单独维护一个ZooKeeper集群它依托于Kubernetes自带的高可用能力。但从源码角度看选举的语义和ZooKeeper版本基本等价有leader、有lease、有callbackgrantLeadership。转换为Flink原生的LeaderElectionService也只是把底层的监听机制换成了Kubernetes的Informer机制。5.2 从源码看两种后端的选型依据这几年在群里经常看到有人问到底选ZooKeeper HA还是Kubernetes HA。我的建议是回到团队基础设施来看不要盲目追新。如果团队已经有一个成熟的ZooKeeper集群ZooKeeper HA最稳妥。它的元数据存储逻辑比较统一排查手段也成熟出问题可以用zk命令行直接看节点内容。而且在非Kubernetes环境部署Flink比如用YARN或原生集群ZooKeeper依然是标配。如果集群整体跑在Kubernetes上而且希望尽量少运维外部依赖Kubernetes HA值得优先考虑。它不需要额外组件直接使用API Server和Lease资源。不过要注意Kubernetes HA的元数据持久化依赖ConfigMap或存储卷如果你没有给Flink集群配置持久卷JobGraph等关键元数据会存在内存态一旦Pod重建后可能丢失。从源码看KubernetesHighAvailabilityServices在配置持久化存储时会读取high-availability.storageDir这个目录必须落在持久化存储上才能做到真正高可用。这里还有一个经验在Kubernetes里跑HA模式JobManager实例数建议为2或3远超3个意义不大因为Leader选举的决策节点过多反而会增大API Server压力。配置时通过replicas控制副本数同时保证每个副本的资源配置一致避免因节点资源不足导致频繁重启。6. 实战排查JobManager HA常见问题6.1 ZooKeeper会话超时导致的主备频繁切换生产环境里最常见的HA故障就是JobManager主备频繁切换。现象是作业周期性地失败、恢复、再失败日志里能看到反复的grantLeadership和revokeLeadership。从源码角度分析LeaderLatch的会话超时由ZooKeeper服务器决定默认sessTimeout和连接参数由high-availability.zookeeper.client.session-timeout控制。如果JobManager所在节点和ZooKeeper集群之间的网络延迟波动明显或者ZooKeeper压力过大导致心跳处理不及时会话就会超时。会话超时后临时节点被删除Curator会触发重新选举。排查时要先看ZooKeeper三端的指标网络延迟、ZooKeeper节点负载、JobManager节点的GC情况。GC停顿如果超过会话超时时间也会导致会话失效。我建议把high-availability.zookeeper.client.session-timeout从默认值适当调大比如50000毫秒以上给GC停顿留缓冲。但调得太大会延长真实故障的探测时间所以要根据业务对RTO的要求做权衡。6.2 JobGraph过大拖垮ZooKeeper前面提过JobGraph过大的问题这里再强调一下排查方法。如果作业提交时日志出现类似ZooKeeper data size exceeds limit或者提交卡住先去查flink-conf.yaml里的high-availability.zookeeper.path.jobgraphs节点大小。有些作业会把较大的配置对象直接塞进JobGraph的全局配置里比如几百MB的字典数据。我们线上就遇到过因为某个作业塞了字典导致提交超时的情况。解决办法是把大对象放进外部存储作业运行时通过读取外部配置拉取而不是硬编码到作业图里。从源码角度看Flink在ZooKeeper上写入JobGraph时没有强制分片数据超限后会抛出ZooKeeperNodeLimitException。这个异常在作业提交时表现为DispatcherResourceManagerComponent启动失败。排查时结合提交端的日志和ZooKeeper端日志基本能定位。6.3 通过日志定位Leader变更要快速定位Leader变更最直接的办法是看JobManager日志里的LeaderElectionService相关输出。ZooKeeper实现里当选后日志会打印Leadership granted失去时会打印Leadership revoked。Kubernetes实现里也有类似日志。另外建议开启org.apache.flink.runtime.leaderretrieval和org.apache.flink.runtime.leaderelection这两个包的DEBUG日志这样可以看到每次Leader地址解析、sessionID变化等细节。曾经有一次我们排查TaskManager连不上JobManager打开DEBUG日志后发现原来是TaskManager端的LeaderRetrievalService还缓存着旧的JobManager地址必须等监听事件刷新。看到日志里缓存地址和实际地址对不上思路立刻清晰了。6.4 恢复期TaskManager没注册上新的JobManager接管后所有TaskManager都需要重新注册。这个阶段如果TaskManager迟迟注册不上作业恢复就会卡住。最常见原因是TaskManager端的resourcemanager.address或jobmanager.address配置指向了旧地址或者TaskManager侧的网络策略限制了对新JobManager的访问。源码里TaskManager通过TaskExecutorRegistration向ResourceManager注册注册成功后会收到确认并建立心跳。如果新JobManager和TaskManager之间网络不通日志会持续出现Could not resolve ResourceManager address或Registration timed out。排查完网络配置后重启TaskManager进程让它重新走注册流程即可。结束语把JobManager HA源码翻过一遍之后最大的感受是高可用不是单点技术而是一整套分布式协作机制的组合。写完这篇之后我再去看线上问题思路清晰了不少——Leader选举只是入口后面紧跟的JobGraph恢复、checkpoint解析、Blob拉取、TaskManager重注册每一步都可能成为瓶颈。最后再分享一个小经验读HA源码时不要一头扎进实现类先看接口把组件之间的边界搞清楚再逐个看具体实现。读HighAvailabilityServices接口比直接读ZooKeeperHaServices更重要因为它定义的是问题的边界。之后再看ZooKeeper或Kubernetes实现你会发现所有后端都在做同一件事只是换了存储介质和容错方式。如果你也在做Flink集群的稳定性治理建议从这套接口入手收益会比追着某个具体报错大得多。