我要提问
ARTICLE DETAIL

资讯详情

前沿编程新知与开发实战干货的深度解读。

Flink Checkpoint原理详解:从Chandy-Lamport算法到生产实践

Flink Checkpoint原理详解:从Chandy-Lamport算法到生产实践 流计算里最要命的问题之一就是状态的一致性。特别是当你的作业跑了几天甚至几周中间某个节点突然挂掉重启后数据对不上账那种感觉相当难受。Flink能成为生产环境的事实标准靠的不是性能参数多好看而是那套建立在Chandy-Lamport算法之上的分布式快照机制也就是大家常说的Checkpoint。这篇文章我想把这条技术线彻底捋一遍不绕弯子从算法原理讲到Flink的工程落地再到实际调参和踩坑希望给正要深入Flink或者已经在生产环境摸爬滚打的同行一些参考。1. 分布式快照到底在解决什么问题要理解Chandy-Lamport算法先要理解它要解决什么问题。单机上的快照很简单——把进程的内存状态序列化到磁盘就行。但分布式系统里事情变得复杂得多。1.1 为什么单机快照的思路在分布式下行不通假设你有一个Flink集群3个TaskManager在跑同一个作业每个节点各自维护着自己的状态。如果你想做一个全局快照最朴素的想法是先让所有节点停止处理数据然后各自保存状态最后再恢复处理。这个方案叫“全局冻结”逻辑上没问题但代价极其昂贵——你的作业在快照期间要完全停摆。更麻烦的是如果你不冻结而是让每个节点各自找个时间点存状态那这个快照就是不一致的。拿WordCount举例Source节点已经读了100条数据但下游的聚合节点只处理到90条此时存下来的快照恢复的时候就会丢掉10条。这就是分布式系统中经典的多个独立状态副本之间难以对齐的问题。1.2 一致性快照的精确含义一个真正的一致性快照要求被保存的每个节点的状态必须对应同一条逻辑记录边界。也就是说所有节点保存的状态要像同一时刻的“合影”一样彼此是匹配的。这样从某个时间点恢复后系统才能像什么都没发生一样继续运行。这个“看起来像同一时刻”实际上是个很高的要求。因为分布式系统里没有全局时钟节点之间的通信还有延迟你没法精确地让所有节点在同一物理时刻拍这张合影。Chandy-Lamport算法的高明之处就是它不需要全局时钟不需要冻结系统只需要借助“Marker”消息就能在分布式环境中捕捉到一个一致的全局状态。1.3 现实世界里的类比全家福怎么拍理解这个算法有个很形象的类比——组织一场家庭聚会拍全家福。你没办法让所有人都精确地在同一瞬间静止不动相当于没有全局时钟但是你可以约定一个“信号”比如摄影师吹一声哨子听到哨声的人保持姿势并且把“我已经听到哨声”这个信息传递给还没听到哨子的人。中间可能会有半秒钟的时间差但只要大家都遵循这个约定拍出来的照片里每个人都是处于“听到哨子”这个事件前后的稳定状态没有人闭眼没有人转头全家福就拍成了。Chandy-Lamport算法里的Marker消息就是那声哨子。2. Chandy-Lamport算法的核心逻辑精讲算法本身其实不复杂核心就三个角色发起者Initiator、收到Marker的节点Receiver、以及Marker消息本身。整个算法的触发和传播过程很像在水中投下一颗石子涟漪一圈圈扩散开来直到覆盖整个系统。2.1 算法的三个关键步骤首先由一个或多个节点作为快照的发起者通常是作业的JobManager或者Source节点。发起者会做两件事一是记录自己的本地状态二是向自己所有的下游邻居发送一条特殊的Marker消息。收到Marker消息的节点要分为两种情况讨论一种是第一次收到某个快照发起者的Marker一种是在此之后重复收到。第一次收到时节点需要立刻记录自己的本地状态保存这个时刻所有算子状态和历史数据然后把这条Marker继续广播给下游所有邻居。如果之后再收到同一次快照的Marker就不再重复记录状态了只需要记录一下“这条边的Marker到了”因为这个信号传递的是一条逻辑边界信息。最终当所有节点都完成了状态记录并且同一个快照的所有Marker都汇聚到终点Sink时一次全局快照就形成了。2.2 为什么不会丢数据数据归边的艺术这里有个容易被忽略的细节Chandy-Lamport算法里每个节点保存的不光是自己的状态还要同步记录通道中正在传输的消息。Flink实现里其实没有这么整套去做因为Flink的Source和算子之间有更细粒度的控制流。但在纯理论层面如果A节点发了一条消息给B节点但是消息还在网络传输途中时快照启动了这条消息是在本次快照还是在下次快照里必须有明确的归属。按Chandy-Lamport的规则在A节点记录状态之前发送的消息不管是否到达B节点都要包含在本次快照里A节点记录之后发的消息算下次快照。这个归属规则保证了精确一次语义不会丢数据也不会重复。2.3 两阶段思想的影子仔细观察Chandy-Lamport你会发现它其实构成了一个分布式的两阶段提交。第一阶段是Marker的传播和状态的记录第二阶段是收集最终的快照确认。不同之处在于Chandy-Lampot不需要阻塞业务流量系统在快照期间照常运行这也就是Flink能实现“异步快照”的理论基础。3. Flink对Chandy-Lamport的工程化改造理论算法讲完了但Flink在工程实现上对算法做了很多关键改造没改的话根本没法用。最核心的一个改造就是引入了Barrier这个概念。3.1 Barrier取代Marker精确一次语义的底气Marker在算法层面只是一个信号但Flink生产环境需要正好一次语义这就要求快照边界必须在数据流中强行建立。Flink里的Barrier本质上就是Chandy-Lamport的Marker但它配套了严格的对齐协议。Barrier随着数据流一起从Source向下游流动。当一个算子节点有多个输入通道时它会等待所有输入通道的Barrier都到达之后才开始快照。这叫作Barrier对齐是确保多路输入数据逻辑边界对齐的关键操作。Barrier对齐期间已经收到Barrier的通道的数据会先被缓存起来不继续参与计算直到所有通道的Barrier都到齐才会取消阻塞继续处理。这个设计解决了多流Join或Union时状态错乱的问题。3.2 状态存储State Backend如何配合快照算法说“记录本地状态”但分布式状态怎么存是工程难题。Flink支持两类State Backend一套是历史遗留的基于内存的文件系统方案另一套是基于RocksDB的增量存储方案。个人强烈建议生产环境使用RocksDB方案因为分布式快照时RocksDB支持增量Checkpoint。新版本 RocksDB 使用 RocksDB 原生的快照能力只上传上一次Checkpoint 以来变化的SST文件大幅减少快照的数据量和网络IO压力。内存方案在小状态场景下性能不错但状态一旦增长到GB级别每次Checkpoint全量序列化内存数据会给年轻代GC造成巨大压力严重时会导致作业频繁Full GC反而拖垮整个管道。做技术选型时千万不要只看基准测试不看真实内存增长速度。3.3 分布式快照的收集与确认所有节点的状态快照完成后会把确认信息回传给JobManager。分别包含每个节点的状态句柄信息以及本次快照涉及的所有数据通道的偏移量信息。JobManager把所有分片确认完才最终生成一个全局有效的CheckpointID。这里有个经验之谈Flink分布式快照的完成时长取决于最慢的那个节点。这就是木桶效应。如果某个TaskManager负载不均衡比如数据倾斜导致一个子任务积压了海量数据Barrier会迟迟走不完那么整个Checkpoint就会超时失败。4. 实操中的关键配置与参数调节光懂原理不会调参生产环境照踩坑。我整理了自己实际操作中一整套比较稳妥的配置思路。4.1 Checkpoint的整体开关与基础配置首先是基础配置一般在作业的StreamExecutionEnvironment里设置StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); // 开启Checkpoint间隔时间根据业务容忍度来定 env.enableCheckpointing(60000); // 精确一次语义 env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE); // 两个Checkpoint之间至少要间隔多少毫秒 env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30000); // Checkpoint超时时间超过这个时间还没完成本次Checkpoint直接丢弃 env.getCheckpointConfig().setCheckpointTimeout(600000); // 最大并发Checkpoint数 env.getCheckpointConfig().setMaxConcurrentCheckpoints(1); // 外部化存储保留Checkpoint文件用于作业恢复 env.getCheckpointConfig().setExternalizedCheckpointCleanup(RETAIN_ON_CANCELLATION);这里特别提醒一下setMinPauseBetweenCheckpoints和setMaxConcurrentCheckpoints的关系。在绝大多数在线业务场景下我建议同时配置并把并发数设置为1。如果允许并发多个Checkpoint虽然理论上能提升快照吞吐但在状态量大或磁盘IO慢的场景下多个Checkpoint同时落盘会互相争抢资源反而拖慢整体性能。setMinPauseBetweenCheckpoints的值一般建议是Checkpoint间隔的一半这样能保证上一个Checkpoint完成之后至少有半个周期的空窗让主流程处理数据。4.2 增量Checkpoint与本地恢复的取舍用RocksDB作为状态后端时增量Checkpoint开启是默认的。它带来的一个副产物是Checkpoint文件会变得碎片化因为每次只上传变化的SST文件文件数量会越来越多恢复时需要读取大量小文件。解决方案是开启本地恢复功能。Flink在做Restore时从本地RocksDB的持久化副本直接恢复大幅减少远程下载Checkpoint文件的时间。如果一个TaskManager上的任务需要从Checkpoint恢复本地磁盘里的数据还在就不必全部重拉远程文件。代价是每个TaskManager需要预留一份状态副本的本地磁盘空间建议至少是状态大小的两倍。4.3 状态后端容器化的线程配置大家很容易忽略的一个参数是RocksDB的并发线程数。RocksDB内部默认使用后台线程处理压缩和flush这些线程数不会自动匹配你的CPU核数。我遇到过一次线上事故Flink作业运行一个多月都正常突然某一天Kafka积压急剧飙升。排查下来发现是RocksDB的压缩线程被突增的写入拖垮导致读写延迟暴增。后来做调优时显式指定了RocksDB的并发线程数一切稳定下来。RocksDBStateBackend rocksDBStateBackend new RocksDBStateBackend(hdfs://...); // 配置4个后台线程处理压缩 rocksDBStateBackend.setNumberOfTransferThreads(4);这个值不能拍脑袋定一般参考CPU核数除以一个合理的比例。写密集型的作业可以适当调大读多写少的场景保持默认或调小即可。4.4 Checkpoint Storage选型HDFS还是S3Checkpoint文件最终要落到分布式存储中。超大规模集群建议选择HDFS吞吐高语义强配合RETAIN_ON_CANCELLATION可以长期保留大数量级的Checkpoint。云上环境如果没有自建HDFSS3也是常见选择。但是S3的Checkpoint有个坑S3的一致性模型在过去很长一段时间里是“最终一致性”如果你在同一时间并发写同一个路径可能读到旧数据。现在大部分云厂商的S3已经支持强一致但老集群或某些兼容S3的存储服务比如对象存储服务不一定支持。如果确认存储后端是强一致的可以放心用否则建议在生产环境用HDFS或者至少加一层文件系统抽象在应用层规避。5. 实战案例模拟一次节点宕机后的恢复过程理论终归要落到实战我模拟一个最典型的生产故障场景让你直观感受Checkpoint在背后是如何救命的。5.1 任务拓扑和数据流假设有一个Flink作业拓扑是Kafka Source → 窗口聚合 → Sink到MySQL并行度4状态约50GB存于RocksDBCheckpoint间隔5分钟存储于HDFS。作业正常运行几小时后突然有一个TaskManager节点宕机。此时Flink的调度器会在几秒内感知心跳丢失把该节点的子任务调度到其他健康的TaskManager上。5.2 状态恢复的完整流程首先JobManager取出最近一次成功的CheckpointID也就是最后一次全局快照的版本。然后所有受影响的任务从Checkpoint存储下载对应的状态文件。因为有本地恢复机制任务在可能的情况下优先从本机RocksDB读取状态补拉差分数据恢复完成后从Checkpoint中记录的Kafka消费位点继续消费。这里值得强调的是Flink在恢复时会从最后一次成功的Checkpoint恢复而不是从失败那一刻恢复。所以如果Checkpoint间隔是5分钟最多会丢失5分钟的数据。如果你配置的是精确一次语义消费位点和状态文件是精确匹配的不会多读也不会少读。5.3 崩溃恢复整个过程的耗时分析我实测过恢复一个50GB状态的作业配置了本地恢复整个重启过程大约耗时3分钟左右其中大部分时间花在重新建立RocksDB实例、重建内存索引、跟上积压的Kafka数据上。如果没有本地恢复纯从HDFS上下载50GB状态在百兆带宽下可能要20分钟以上。这也是为什么强推RocksDB加本地恢复的原因——架构上多花一点磁盘成本换来的是故障恢复时间的大幅缩短对SLA保障有质的提升。5.4 一个容易被忽略的坑Kafka位移与Checkpoint的配合恢复过程中一旦涉及精确一次就必须要保证Kafka消费位点与状态是同步保存的。Flink的Kafka Connector有对应的机制会在每个Checkpoint里同时记录各分区的offset。但如果作业代码里手动提交了Kafka的offset就可能撕裂这个一致性。我见过不止一个团队因为想要“监控”消费积压额外启了一个独立Consumer去读同一个Group的offset导致这个Group的offset被外部提交破坏了Flink的精确一次语义。遇到这种情况最干净的做法是Flink作业里不设置setCommitOffsetsOnCheckpoints(true)让Flink完全掌控offset生命周期监控数据从Flink的Kafka Metric里取不要让外部Consumer掺和到同一个Group里。6. 常见问题排查与性能调优速查这一节汇总我在过去项目中反复遇到的几类问题直接对标线上实操经验。6.1 频繁Checkpoint失败的排查路径现象是作业日志里频繁报Checkpoint expired before completing任务整体在线但状态持续不稳定。排查顺序建议这样展开。第一步看JobManager日志里具体是哪几个子任务迟迟没有确认Checkpoint。如果是同一批子任务反复拖后腿优先怀疑数据倾斜——某个key的数量远大于其他key导致对应的子任务积压大量数据Barrier在它那里堵住了。验证方式也很简单看反压监控如果某些子任务的接收端持续高水位基本就是倾斜了。解决方案是加局部聚合、加盐打散或者调整KeyBy的策略。第二步看Checkpoint的具体耗时分解是同步阶段慢还是异步阶段慢。同步阶段慢说明序列化或状态拷贝有压力考虑换Kryo优化或减小单条状态量异步阶段慢则大概率是磁盘IO或网络IO瓶颈优先升级存储或调整RocksDB的压缩级别。6.2 Barrier对齐导致的数据积压问题Barrier对齐本身是有成本的——慢的通道要等快的通道汇合这个等待期间快的通道数据会被缓存不能处理。如果某个通道长期落后就会导致整个算子数据积压上游背压逐渐传导到Source。Flink提供了一些优化选项。比如setAlignmentTimeout设置对齐超时时间超时后不再等待慢通道直接进行快照代价是快照的一致性降级为至少一次。但生产环境核心作业不建议轻易使用。如果必须保证精确一次又遇到对齐时间过长核心解法还是消除数据倾斜或降低单通道的吞吐压力。很多团队一上来就想关对齐都不先排查结果精确一次形同虚设出问题更难看。6.3 RocksDB状态文件不断膨胀的治理方法用过RocksDB方案的应该都见过状态目录一天比一天大Checkpoint文件也越积越多。这不是内存泄漏本质是RocksDB的compaction跟不上写入压力的节奏或者产生了大量没被及时合并的旧版本SST文件。解决思路一般是三步走。一是调优RocksDB本身的BlockBasedTable配置把block cache大小合理分配到读路径二是调整底层压缩策略减少高压缩率带来的CPU消耗适当增加ssts文件的触发合并阈值三是实在不行就调大Checkpoint间隔降低状态写入频率。注意尽量避免频繁做全量清空类的操作比如clear()一个超大State。RocksDB对超大范围的delete操作也是要整叶子SST重建的瞬间CPU会飙升。6.4 新版Flink对默认参数的看法Flink新版本对Checkpoint的默认参数做了不少优化。比如新版中Barrier对齐默认开启且效果更好S3的默认存储路径也有了改进。但不要过度依赖默认值生产级的作业一定要显式配置核心三个参数间隔时间、超时时间、外部化删除策略。这三个参数决定了系统在异常场景下的表现边界值得结合业务容忍度去量身定制。另外值得一提的一个新特性是Checkpoint的逐步清理和版本保留数量限制。默认情况下Flink只保留最新的Checkpoint和对应的State但如果想支持按时间点回溯数据就需要保留多个历史Checkpoint。可以设置setMaxRetainedCheckpoints来控制保留数量。这个值不是越大越好每多保留一个Checkpoint就多一份存储恢复时选哪个版本也需要业务方自己清晰。7. 个人体会快照设计是分布式系统中最容易低估的部分做了多年流计算碰到过各种千奇百怪的线上问题但最终绝大多数故障都归结到同一个点——状态一致性和恢复的可靠性。Chandy-Lamport算法本身是1985年提出的老算法但它在Flink里的生命力恰恰说明了一个道理基础算法想得越透彻工程实现才能越扎实。最后分享一个我自己的调优心得。如果Checkpoint总是超时先别急着调大超时时间。超时时间设置得越大掩盖的问题越多最终Checkpoint堆积只会让作业更脆弱。先做排查从上到下确认是负载、数据倾斜、磁盘还是网络问题再有针对性地调整才是正道。还有一点经验之谈千万不要在状态大的作业上频繁手动触发Savepoint。Savepoint本质也是分布式快照但它和Checkpoint的触发路径不同不能复用增量数据每次都是全量状态一上去就是毁灭性的。整个团队都该约定好运维人员只操作Checkpoint研发人员只通过停止作业并指定Savepoint路径来做版本升级。分布式快照这条技术线从算法到实现从调优到故障排查每一环都值得做流计算的同行仔细吃透。希望这篇文章能给正在啃这块硬骨头的你省下一些自己摸索的时间。
返回列表