检查点与保存点 — 停止后恢复为何不会重复
一句话总结
检查点是一张把 Source 读到了哪里 · 算子持有什么 · Sink 还有什么尚未确定对齐到同一时刻拍下的快照。文件 Sink 只有在收到“这张快照已经完成”的通知之后才会确定文件,所以即使作业挂掉又重新启动,已确定的结果里既没有缺口也没有重复。保存点则是用同一套机制、由人来拍下的快照。
为什么需要它
流处理作业一跑就是好几周。期间 TaskManager 会挂掉,代码要修改后重新部署,集群也要迁移。重新启动的那一刻,需要同时决定三件事:从 Source 的哪里开始重新读取,累计到目前的合计、窗口这类状态如何恢复,以及已经输出到外部的结果怎么处理。
这三件事如果分开决定,一定会错位。从头重新读,已经写出的结果会再输出一遍;从上次读到的位置接着读,状态是空的,合计就不对。即使单独把状态保存下来,只要保存的时刻与 Source 的位置相差哪怕几条,就会多算或漏算几条。需要的是“三者处于同一时刻”这样的保证。Flink 的解法,是在数据流中放入标记一起流过去。
工作原理
barrier。 JobManager 的检查点协调器会向 Source 插入 barrier n(官方文档中的 Stateful Stream Processing)。barrier 与记录走同一条路向下游流动,算子一收到 barrier,就取出自己的状态并保存,然后把 barrier 继续向下传递。Source 的“状态”就是读取位置——对序号 Source 来说,就是下一个要输出的编号。所有任务都报告拍好快照后,检查点 n 就完成了,协调器再把完成的消息通知出去。
文件 Sink 会等待这个通知。 Sink 写的文件要经过三个阶段。写入过程中,名称以点开头,形如 .part-…inprogress…;收到 barrier 后关闭并等待确定;完成通知到达后,重命名为 part-… 而被确定。按惯例,读取方会跳过以点开头的文件,它们属于隐藏文件。因此下游看到的结果始终是“截至某个检查点”的结果,检查点间隔就等于结果可见之前的延迟。在这个 Pod 上以 1 秒间隔运行,每秒会多出一个 part 文件,而点文件只在两次检查点之间短暂出现,随后消失。
存储位置。 指定 execution.checkpointing.dir 后,检查点会保存到文件系统中。文档给出的结构是 <dir>/<job-id>/chk-<n>/,其中的 _metadata 是这张快照的目录。默认情况下检查点不会保留——新的检查点完成后,旧的会被删除,取消作业时则全部删除。它用于故障恢复,并不是留给人使用的。想保留的话,把 execution.checkpointing.externalized-checkpoint-retention 设为 RETAIN_ON_CANCELLATION。
保存点用同样的机制拍下,但拥有者是人。它生成在 <savepoint-dir>/savepoint-<잡 id 앞 6자리>-<무작위>/(占位符依次为保存点目录、作业 ID 的前 6 位、随机串)下,Flink 不会自动删除。文档强调了一个陷阱:从 Flink 1.15 起,不停止作业而拍下的中间保存点不会提交副作用。负责确定文件的是检查点和 STOP ... WITH SAVEPOINT。
SET 'execution.checkpointing.savepoint-dir' = 'file:///root/flink/checkpoint/sp';
STOP JOB '<jid>' WITH SAVEPOINT; -- 찍고 멈춘다. 쓰던 파일까지 확정
SET 'execution.state-recovery.path' = 'file:/.../savepoint-xxxxxx-yyyy';
INSERT INTO sink SELECT id FROM seq; -- 같은 질의를 되살린다 — 다음 순번부터
恢复出来的作业会同时拿到保存点中的 Source 位置和 Sink 状态。因此即使更换 Sink 路径,序号依然连续,把两个目录合在一起,既没有缺口也没有重复。从保留下来的检查点也可以同样恢复(文档:通过检查点的元数据文件,像保存点一样恢复)。
在现场相遇的样子
最常见的事故是没有保存点就重新部署。稍微修改 SQL 后重新 INSERT,新作业会以空状态从 Source 的起点(或配置的起始位置)开始读取。如果是序号 Source,就会从 1 重新写起。结果目录相同的话,已有的行会再进去一遍。本实验会亲手数出这个重叠。
第二种是因取消而停止的作业。取消不会拍保存点,所以最后一次检查点之后正在写的点文件会原样留在目录里。从保留的检查点恢复时,新作业会从检查点时刻开始重新写,结果(不以点开头的文件)依然没有缺口和重复。但如果在下游挂上一个连点文件也读的工具,那一刻就会出现重复。“精确一次”是与只读取已确定的内容这一约定成对的。
第三种是恢复失败。保存点按算子 ID 保存状态。文档警告说,自动生成的 ID 对程序结构很敏感——大幅改变 SQL 的形态,就找不到可以放回状态的算子。修改有状态的作业时,要先测试“这次变更之后,保存点还能不能恢复”。
关于重启策略,如第 1 个模块所见,打开检查点后默认会改为重启一侧。本模块讲的是这次重启回到哪里。
下一项实验要做什么
把每秒拍一次检查点的序号作业接到文件 Sink 上,在运行期间留下已确定文件与点文件同时存在的目录列表。从 /checkpoints 的响应中读出最后一个检查点的位置,用 STOP JOB ... WITH SAVEPOINT 停止后,从保存点恢复,看看序号是否连续。接着不用保存点重新运行以制造重叠,再确认即使取消,也能从保留的检查点接着写;最后把所有数字从磁盘上重新数一遍,整理成报告。