ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Flink批作业JobMaster故障后进度恢复实战:告别全量重跑

Flink批作业JobMaster故障后进度恢复实战:告别全量重跑 上周值班一个跑了16小时的Flink批作业在下游数据告警最密集的时候JobMaster进程被一次OOM打崩。集群配了ZK HA新的JobMaster半分钟内就接管了但作业却从源头重新读文件、重新建链、重新跑计算所有队列瞬间被顶满前十几个小时的白工一秒清零。这种事遇到一次就够记住一辈子批作业如果不主动做进度保护它在容错层面比流作业脆弱得多典型的“JM一挂全盘重跑”并不是什么罕见事故而是默认配置下的必然结果。这篇文章不绕弯子直接讲清楚批作业在JobMaster Failover后怎么保住进度先从机制上说明为什么默认会重跑再讲我用来破解这个问题的三套核心手段然后是可直接抄走的生产配置以及我自己手动杀JobMaster验证恢复链路的完整过程。最后一部分是从线上踩坑记录里挑出来的深坑清单。适合正在维护批处理集群的Flink开发和运维同学也适合那些准备把批作业接入生产、但还没认真考虑容错方案的团队。1. 先找准病根JobMaster 挂掉时批作业到底丢了什么1.1 批作业默认是裸奔状态很多批作业根本没有开启checkpoint。这一点听起来难以置信但我在大量团队的生产代码里见过太多次JobGraph运行前没有enableCheckpointingflink-conf.yaml里也没有execution.checkpointing.interval整个作业从头到尾处于“不存档”模式。为什么大家会这样干因为流作业从一开始就要考虑故障恢复生产规范一般强制开启checkpoint而批作业给人的刻板印象是“跑一遍就结束、有界数据可以随时重来”。抱着这种心态很多人从开发到上线自始至终没给批作业留任何恢复依据。结果JM一挂Flink失忆只能全部重算。Flink本身是容错的但批作业默认的容错能力恰恰是最低一档——不存状态、不存进度、不预留任何可以接续执行的“存档点”。这就像游戏打到快通关了才想起来没存档中途崩溃只能从标题画面重新开始。1.2 一个运行中作业的“进度”分别保存在哪里要理解恢复先得知道JobMaster挂掉之后到底丢了什么、哪些东西还可能留下来。一个正在运行的批作业进度信息大致分布在四层JobManager内存中的ExecutionGraph整个作业有哪些task、每个task在哪个TaskManager上、当前是RUNNING还是FINISHED、上下游通过哪种边连接。这是“施工图纸进度表”JM一崩内存里这份图纸直接没了。SourceSplit进度批作业读文件或者读消息时每个分片读到了哪一行、哪个offset。如果没配checkpoint这个进度不落盘等于日志读到哪里完全没记录。有状态算子的中间状态聚合、双流join的缓存、去重集合之类的数据通常存在TaskManager的堆内内存或RocksDB里。JM崩了TM未必立刻跟着崩但作业失败导致整个部署被清理时这些状态也会被销毁。已完成Task的中间结果批处理过程中上游算完写给下游的数据。早期pipelined shuffle时代这些数据在内存和网络Buffer里任务一失败全部蒸发Flink 1.14之后批作业默认用Blocking Shuffle中间结果会落盘这部分反而成了恢复时可争取的“残余资产”。大多数人对恢复的误解在于以为JM重启后作业原地续跑。实际不是。JM重启后它要从某个持久化位置重建ExecutionGraph和算子状态如果什么都没有就只能把作业当作一个全新的执行计划重新提交。全盘重跑就是这么来的。1.3 为什么流作业没那么慌批作业特别疼流作业几乎都会开checkpoint的另一个原因是它的Source天然带游标。Kafka offset会随checkpoint一起存下来JM挂掉新JM读最近的checkpoint从offset恢复消费位置数据接着读状态重新加载整个过程是“无缝”的。批作业的问题在于它的Source也是有界的、通常是一个个文件分片但很多Source实现压根没有把“分片进度”写进checkpoint的习惯或者作业压根没有checkpoint。两头都缺自然只能从头来。另外一个现实因素批作业通常跑得久。流作业挂了恢复点在秒级到分钟级范围内损失有限。批作业动辄几小时到几十小时跑一半挂掉等于是把整个计算周期里最昂贵的部分全废掉。所以批作业的JM Failover恢复问题不是“锦上添花”而是值得单独拿出来解决的核心运维问题。2. 破局的三个抓手checkpoint、Source 断点、Blocking Shuffle 中间结果2.1 checkpoint 是“进度”能落地的基本前提要让新的JobMaster认识旧进度唯一靠得住的方案就是定期把进度固化到外部存储。这个动作就是checkpoint。它在Flink里像游戏存档Barrier从Source流到Sink途经的每个有状态算子把状态快照下来Source把分片读进度也一起记下全部写到一个共享存储HDFS/S3/本地盘里。checkpoint完成就出现了一个“如果此刻挂掉可以从这里继续”的恢复点。我给批作业开checkpoint时很多人第一反应是“会不会特别慢”。这里要区分两种开销一是快照本身的I/O开销二是checkpoint协调对运行节奏的影响。对于批作业我见过很多人被网上流式场景的“秒级checkpoint间隔”带偏给批作业也配了3秒、5秒的间隔结果状态频繁写入整个作业的执行效率肉眼可见地下降然后他们得出结论“批处理不能开checkpoint”——这是典型的配置错了不是机制错了。批作业的checkpoint间隔我建议在5到30分钟之间具体看作业的阶段耗时。比如一个阶段要跑40分钟那中间开一次checkpoint就有意义如果一个阶段只要30秒5分钟一存也是浪费。开源社区里批作业checkpoint的最佳实践并没有一个标准答案因为它和状态大小、阶段长度、存储性能强相关。但有一件事是确定的默认关闭等于把容错底牌全丢了。2.2 Source 分片断点批作业能“续读”而不是“重读”开完checkpoint接下来要确认Source是否真正支持逐分片续读。Flink 1.14之后官方推荐的FileSource是FLIP-27架构实现的。它在读取有界文件时会把每个文件分片处理到哪一行、哪个offset作为分片状态存进checkpoint。恢复时已经处理完的分片不再读取只把未完成的分片重新拉起来跑。这意味着checkpoint恢复后Source不会把整个输入目录从头扫一遍而是接着上次的位置继续读。这一点对批作业太关键了。批作业跑十几个小时前面十几个小时的数据不一定都要重算只要Source断点生效大部分时间都能跳过去。但注意不是所有批式读取源都具备这个能力。JDBC一类的连接器在批任务里往往是把整张表按主键分段扫描但每段读到哪一行很少能精细地写入checkpoint数据越多恢复时重复读取量越大。这也是为什么很多MySQL同步到数仓的场景大家宁可用CDC或增量同步而不是每天拿JDBC全量刷一遍。如果你用的是自定义Source或者老式InputFormat更要想清楚它到底存了什么。测试方法后面会讲这里先记住一个原则凡是支持checkpoint的Source恢复能力要看它对分片进度是否真的做了持久化只看文档说“支持checkpoint”远远不够。2.3 Blocking Shuffle已算完的中间结果有机会直接复用第三个抓手可能很多文档不会重点提但它才是“JM一挂不完蛋”的进阶素材Blocking Shuffle的中间结果文件。Flink 1.14之后批作业默认把所有边切换成Blocking Shuffle。什么叫Blocking呢上游Task把输出完整写到本地磁盘的ResultPartition文件里等文件全部写完下游才开始读。在这个模型下已经Finish的上游Task它的计算结果是以文件形式存在的并不只活在内存里。如果故障范围只是JobMaster挂掉而TaskManager还活着新JM接管后是有机会复用这些中间结果文件的。上游不用重新算下游直接读已落盘的输出恢复范围可以被压缩到很小的区域。这相当于工地上半成品仓库还开着工人换了个包工头手头已经做好的部件就不需要推倒重做了。这里要泼一盆冷水这个“复用”不是百分百保证它强依赖部署形态和TM存活状态。在Standalone、K8s这种JM和TM分开管理的模式下JM挂掉TM大概率存续中间结果文件有机会保留在YARN Application模式下AM即JM挂掉后YARN经常会把整个Application重启TM一起陪葬本地文件全部消失那就只能靠checkpoint兜底。所以中间结果复用是“锦上添花”不是“雪中送炭”。真正能保证恢复的还是checkpoint。2.4 从 checkpoint 恢复到跑完理论重算窗口怎么估算做容错方案之前最好心里有本账。一次恢复需要的时间大致等于三部分之和状态加载时间RocksDB或堆内存从快照目录读取状态。状态越大越慢。Source未完成分片的重新读取时间已经完成的分片跳过只补上次checkpoint之后剩下的那部分。从checkpoint点到当前进度之间的重算时间理想情况下作业只会丢失“上次checkpoint完成之后到故障发生”这期间的进度而这个窗口的上限就是checkpoint间隔。举例作业总耗时12小时checkpoint间隔10分钟JM挂掉时最近一个checkpoint已经完成。理论上恢复后的重算损失不会超过10分钟加上状态加载和未读分片的耗时通常可以在半小时内回到故障前的进度。对比GB的“全量重跑12小时”差别是数量级的。当然这里说的是理想情况。很多批作业的Source断点做得不好、算子链重算逻辑复杂、外部系统写入需要幂等清理实际恢复窗口可能会被拉长。但方向是对的任何一次JM故障恢复都应该让重算时间远小于全量执行时间否则你的恢复方案就是假的。3. 生产配置实操HA、重启策略与 checkpoint 参数怎么配才算数3.1 没有 HAcheckpoint 开得再好也救不了你先把一个误区讲透开checkpoint只是把进度存下来不等于故障后自动续跑。如果集群没有配置HAJobMaster进程挂掉后集群根本没有第二个Leader可以接管作业只会直接FAILED。这时候进度虽在外部存储上但Flink不会自动拉起一个作业去读它需要人工拿--fromCheckpoint恢复这在生产上是不可接受的延迟。所以做JM Failover恢复的第一步不是调checkpoint而是先把HA配置好。ZooKeeper是目前最稳的方案Kubernetes上也可以用K8s原生HA。我线上用的配置大概是这样high-availability: zookeeper high-availability.storageDir: hdfs:///flink/ha high-availability.zookeeper.quorum: zk1:2181,zk2:2181,zk3:2181 zookeeper.session.timeout: 60000核心思路是HA集群里存作业的基础元数据和Leader选举信息JM挂掉后新JM能快速接管并自动从外部存储恢复ExecutionGraph。注意high-availability.storageDir一定要放在所有JM能访问的共享文件系统上不能放本地盘——否则新JM连旧元数据都读不到HA形同虚设。3.2 重启策略和 failover 策略把恢复范围尽量切小JM挂掉属于全局故障但日常碰到的更多是单个TaskManager失联或单个task反复失败。这两种情况的恢复范围完全不同需要分开配置。重启策略方面批作业我建议用fixed-delay让Task失败后以固定间隔重启给外部依赖比如网络抖动、临时文件系统故障留出恢复时间而不是立即把整个作业宣判死刑execution.restart-strategy: fixed-delay execution.restart-strategy.fixed-delay.attempts: 5 execution.restart-strategy.fixed-delay.delay: 30sfailover策略默认是region批作业里因为Blocking Shuffle把图形切成多个相对独立的区域一个区域内的task失败理论上只要重启这个区域的前驱链条即可不至于整个作业都动。如果你发现某个task失败后整个作业从头跑先排查failover-strategy是否被人改成了fulljobmanager.execution.failover-strategy: region3.3 checkpoint 参数组合一个可以直接抄的配置下面这套配置我在线上多个批作业上验证过整体效果稳定适合大多数跑1小时以上的批任务可以直接作为起点再做微调。# checkpoint 基础设置 execution.checkpointing.interval: 10min execution.checkpointing.timeout: 40min execution.checkpointing.min-pause: 2min execution.checkpointing.max-concurrent-checkpoints: 1 execution.checkpointing.tolerable-failed-checkpoints: 3 # 状态存储与保留 state.backend: rocksdb state.backend.incremental: true state.checkpoints.dir: hdfs:///flink/checkpoints state.checkpoints.num-retained: 5几个参数逐个说为什么这么设interval: 10min批作业不建议秒级甚至分钟级短间隔10分钟兼顾“恢复损失可接受”和“快照开销不过分”。timeout: 40min批作业一个阶段可能跑很久状态大时快照响应会慢超时设长一些避免checkpoint没写完就被判失败。min-pause: 2min两次checkpoint启动之间至少隔2分钟防止快照风暴持续压榨主任务。max-concurrent-checkpoints: 1批作业根本不需要并发快照并发只会增加存储I/O负载。tolerable-failed-checkpoints: 3允许一定数量的checkpoint失败而不杀作业避免一次短暂的存储抖动把大作业搞崩。num-retained: 5多保留几份checkpoint防止“最新一份已损坏但旧的可用的被自动删掉”这种惨剧。如果你在代码里更习惯用API配置对应DataStream API是StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.setRuntimeMode(RuntimeExecutionMode.BATCH); // 开启 checkpoint建议直接传 Duration env.enableCheckpointing(Duration.ofMinutes(10)); env.getCheckpointConfig().setCheckpointStorage(hdfs:///flink/checkpoints); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(Duration.ofMinutes(2)); env.getCheckpointConfig().setCheckpointTimeout(Duration.ofMinutes(40)); env.getCheckpointConfig().setMaxConcurrentCheckpoints(1); env.getCheckpointConfig().setTolerableCheckpointFailureNumber(3);老版本Flink可能只认毫秒级的long参数新版本更推荐Duration如果你维护的是老集群按你的版本调整一下即可。增量checkpoint加上RocksDB能让恢复时的状态加载明显更快因为只上传变更的SST文件不传全量。3.4 Sink 幂等和中间产物处理恢复成功不等于数据正确配置再多最后还要过一道关卡从checkpoint恢复后作业大概率会重复写一部分数据。如果Sink不是幂等的业务上就会出现重复行、重复计数、脏分区。批作业最常见的Sink是写HDFS/Hive分区。我的标准做法是输出直接写入一个临时目录作业全部成功后通过提交阶段把临时目录rename到最终路径恢复时先清理掉临时目录里的半个文件再继续写。这套“临时目录原子rename”的策略能把重复写的影响降到零。写数据库的场景优先选带主键的Upsert Sink写Kafka的要给下游约定好按事件ID去重。记住一句话checkpoint恢复解决的是“不停跑”幂等写入解决的是“不跑错”这两件事缺一个都不能上线。4. 手动杀一次 JobMaster验证恢复是真的能续跑4.1 怎么模拟才真实验证恢复链路最有效的方法就是在测试环境里手动干掉JobMaster进程。别犹豫杀进程是成本最低、效果最接近真实的演练。如果是Standalone或K8s部署找到JM进程PID直接kill -9如果是YARN模式找ApplicationMaster所在节点杀掉对应的JM进程。我建议不要用yarn application -kill那会把整个集群作业一起销毁相当于模拟“整体崩溃”而不是“JM Failover”。杀完之后观察两点HA是否完成Leader切换以及新JM是否自动恢复作业而不是把作业标记为失败。如果等了很久作业没有从checkpoint恢复先看ZK里HA存储目录有没有内容再看日志里是否出现了“Trying to restore job”相关的关键字。4.2 怎么确认是新JM从checkpoint恢复而不是自动全量重跑这是我验证时必看的三类证据日志关键词新JM启动日志里会出现从checkpoint恢复的标记。搜Restoring job、Completed checkpoint、Recovered checkpoint之类关键词。Flink版本不同措辞不完全一样但一定有一行明确说明恢复来源。Web UI上的Last Checkpoint新JM接管后作业详情页能看到“Last Checkpoint”时间和路径。如果是从checkpoint恢复恢复时间点应该和故障前最近一次checkpoint对应上如果显示没有任何checkpoint说明故障前压根没存过作业必然在裸跑。Source侧的实际行为这是最直观的。恢复后的作业如果已经处理完的分片不再重新读那Source的读取吞吐会立刻恢复到高位而不是从零开始相邻分片全扫一遍。我习惯在Source端加一个计数器记录“本次作业实际读取的分片数量”恢复后如果计数远小于全量分片数就说明断点确实生效了。4.3 恢复后的一致性验证光看作业“跑起来”不算完数据验证才说明问题。方法很简单记录故障前的输出行数和样本数据等恢复后的作业跑完后再做一次分区分时段的对账。重点不是总数一致而是最终数据里没有重复行、没有脏数。我自己的恢复验证清单就两条一是作业跑完时间明显小于全量重跑时间二是下游表数据经过幂等约束去重后与业务预期一致。两条都过了这个作业才算真正具备“JM挂掉不重跑”的能力。5. 线上批作业恢复最容易踩的五个深坑5.1 算子没写 uid代码一变恢复直接报错Flink自动生成的算子ID依赖算子链结构和代码写法一旦你在上游加了个map、调整了算子顺序自动ID就会变恢复时Flink拿ID去匹配旧状态就失败了作业直接恢复失败。所有有状态算子从一开始就显式写uid()stream .keyBy(...) .process(new MyStatefulProcessFunction()) .uid(my-stateful-process)这个习惯在批作业里尤其重要因为批作业很多是长期运行的Daily任务代码会不断迭代。没写uid改一次代码就丢一次旧状态那你开不开checkpoint意义都有限。5.2 checkpoint 频繁让人误以为“批作业不能开”这个坑前面提过这里再多说一句。判断checkpoint间隔是否合理不要拍脑袋看两个指标一个是checkpoint完成耗时一个是间隔期间主任务是否有明显变慢。如果完成耗时接近间隔说明状态太大或者存储太慢先加间隔状态裁剪或者上RocksDB增量而不要直接把checkpoint关掉。5.3 恢复时改了并发度比你想的更容易翻车从checkpoint恢复时如果作业并行度和保存checkpoint时不一致Keyed State理论上支持按Key Group重新划分但实际生产里风险很大。尤其当算子用了Operator State非Keyed State并发变化后状态重新分配的逻辑非常脆弱一旦分配不均匀或者触发状态迁移失败死得很难看。我自己的原则是从checkpoint恢复时不改并行度想调整并发先做一次savepoint再基于savepoint调并发。这句话救过我很多次。5.4 整体集群重启本地中间结果瞬间清零很多人把希望押在Blocking Shuffle的中间结果文件上但忽略了一个问题这些文件在TaskManager本地磁盘。如果故障是“整个集群重启”而不是“JM单点挂掉”TM全陪葬本地文件全消失中间结果复用直接归零。YARN Application模式尤其如此AM挂掉后YARN倾向重启整个Application所有本地资产一并清空。所以底线永远是checkpoint必须放在HDFS/S3这类外部持久化存储上Source断点状态必须跟着checkpoint走。本地任何文件都只是加速项不是安全项。5.5 非确定性的中间算子会让恢复变成“跑错”最后一个坑最隐蔽如果你在中间算子里用了System.currentTimeMillis()、随机数、或者调用外部接口拿动态值那么从checkpoint恢复后这部分计算结果和故障前版本是有偏差的。幂等写入能防重复但防不了“错值”。办法是尽量把这类非确定因素放在Source阶段注入或者在Sink阶段再解析中间算子保持纯函数式处理。实在躲不掉的使用确定性种子或把时间字段作为记录内容而不是计算依据。批作业恢复要想结果可对账确定性是硬要求。最后分享一个小技巧在大批作业上线恢复方案前先用一个小规模的代表性作业做演练——开好checkpoint跑一段时间后手动kill JM观察恢复日志确认从“正常进度”到“恢复后进度”之间重算的时间窗口明显小于全量执行时间再把这个方案推广到所有核心批作业。我在实际项目中靠这个两步走的流程两周内把主力批作业全部接上了进度恢复之后JobMaster再挂值班群里第一条消息不再是“全量重跑告警”而是“已从checkpoint恢复无新增异常”。这套思路和配置希望能让你的批作业也早点摆脱“JM一挂全盘重跑”的宿命。
返回列表