TT Lab
开始
学习 学习路径 课程

Apache Flink — 用真正的引擎跑流处理

停止并恢复后序号是否连续

在 TT Lab 中继续学习

目标

把序号 Source(1、2、3 ……)接到文件 Sink 上,用检查点和保存点停止再恢复,然后在磁盘上亲自确认已确定文件中的 id 是否没有缺口、没有重复地连续。同时数一数:不用保存点重新运行时哪些内容会重叠,以及即使取消,保留下来的检查点能守住什么。

为什么重要

重新部署或故障恢复时,如果 Source 的位置、算子状态和 Sink 中尚未确定的文件没有对齐到同一时刻,结果就会缺失或被输出两遍。检查点用一个 barrier 把这三者对齐,而文件 Sink 要等检查点完成才会确定文件。本实验的评分器不会查询集群。停止时刻的行数每次运行都不同,所以没有标准数字——评分器检查的是:已确定文件的 id 是否从 1 开始没有缺口、没有重复地连续,恢复的作业是否从紧接着的下一个 id 开始,以及保存点和检查点目录中是否留下了 _metadata。

步骤

  1. 用 flink-up 启动集群,提交 /root/flink/checkpoint/run1.sql:检查点间隔 1 秒,检查点目录为 file:///root/flink/checkpoint/ckpt,保存点目录为 file:///root/flink/checkpoint/sp,作业名称为 flk-seq-1,把 datagen 序号(id 1..100000,每秒 20 行)以 csv 格式写入 file:///root/flink/checkpoint/out1;并把输出保存到 /root/flink/checkpoint/run1.out。
  2. 作业运行期间,把 ls -A /root/flink/checkpoint/out1 的结果保存到 /root/flink/checkpoint/files-running.txt。其中必须同时看到已确定的文件(part-…)和正在写入的文件(.part-…inprogress…)。
  3. 把该作业的 /jobs/<jid>/checkpoints 保存到 /root/flink/checkpoint/checkpoints.json,把 /jobs/<jid>/checkpoints/config 保存到 /root/flink/checkpoint/checkpoint-config.json。
  4. 用 /root/flink/checkpoint/stop.sql 执行 STOP JOB '<jid>' WITH SAVEPOINT,并把输出保存到 /root/flink/checkpoint/stop.out。停止之后,out1 中不应留有点文件。
  5. 提交从保存点恢复的 /root/flink/checkpoint/run2.sql(作业名称 flk-seq-2,Sink 为 file:///root/flink/checkpoint/out2),把输出放到 /root/flink/checkpoint/run2.out;几秒后通过 REST 拍下保存点并停止,把已完成的状态响应保存到 /root/flink/checkpoint/stop2.json。
  6. 提交不做恢复、从头运行的 /root/flink/checkpoint/run3.sql(作业名称 flk-seq-3,Sink 为 file:///root/flink/checkpoint/out3,检查点保留策略为 RETAIN_ON_CANCELLATION),留下 /root/flink/checkpoint/run3.out,待文件被确定之后取消作业。把已取消作业的 /jobs/<jid> 保存到 /root/flink/checkpoint/job3.json,把 /jobs/<jid>/checkpoints 保存到 /root/flink/checkpoint/checkpoints3.json。
  7. 提交从保留下来的检查点恢复、接着写入同一个 out3 的 /root/flink/checkpoint/run4.sql(作业名称 flk-seq-4),留下 /root/flink/checkpoint/run4.out;几秒后通过 REST 拍下保存点并停止,把状态响应保存到 /root/flink/checkpoint/stop4.json。
  8. 在 /root/flink/checkpoint/report.json 中写入 savepoint_path、last_id_before_stop、first_id_after_resume、fresh_run_overlap、retained_checkpoint、leftover_inprogress_files。

参考

提交启用检查点的序号作业

运行 flink-up 后,在 /root/flink/checkpoint/run1.sql 中写入:检查点间隔 1 s,execution.checkpointing.dir = file:///root/flink/checkpoint/ckpt,execution.checkpointing.savepoint-dir = file:///root/flink/checkpoint/sp,作业名称 flk-seq-1,datagen 序号 Source,写入 file:///root/flink/checkpoint/out1 的 csv Sink,以及 INSERT;再用 sql-client.sh -f run1.sql > run1.out 2>&1 生成 /root/flink/checkpoint/run1.out。

SET 语句必须放在 INSERT 之前,才会对该作业生效。序号 Source 使用 fields.id.kind = sequence,起止值用 fields.id.start 和 fields.id.end 指定。每秒行数要设得小一些(20),这样生成的文件不多,便于肉眼跟踪。输出中出现 Job ID,说明作业正在后台继续运行。

查看运行中的 Sink 目录

作业运行期间,把 ls -A /root/flink/checkpoint/out1 的输出保存到 /root/flink/checkpoint/files-running.txt。其中必须同时看到已确定的 part-… 和正在写入的 .part-…inprogress…。

检查点完成一次之后,第一个 part 文件才会被确定。间隔为 1 秒的话,等几秒就行。以点开头的文件不加 -A 就看不到,在这个 Pod 上实测,它只在两次检查点之间短暂出现——可以每隔 0.1 秒获取多次目录列表,在两者同时出现时保存。评分器会核对列表中已确定的文件现在是否仍在 out1 中。

通过 REST 获取检查点记录

把 run1 作业的 /jobs/<jid>/checkpoints 保存到 /root/flink/checkpoint/checkpoints.json,把 /jobs/<jid>/checkpoints/config 保存到 /root/flink/checkpoint/checkpoint-config.json。已完成的检查点至少有一个,且最后一个的路径必须是该作业的 ckpt/<jid>/chk-<번호>(占位符为检查点编号)。

作业 ID 在 run1.out 的 Job ID 行中。checkpoints 响应的 latest.completed 里有最后一个已完成检查点的编号和 external_path,config 响应里有模式(精确一次)、间隔(毫秒)和是否保留(externalization)。

拍下保存点并停止

在 /root/flink/checkpoint/stop.sql 中写入保存点目录的 SET 和 STOP JOB '<run1 의 jid>' WITH SAVEPOINT;(占位符为 run1 的作业 ID),并把输出保存到 /root/flink/checkpoint/stop.out。停止之后,out1 中不应有点文件,已确定文件的 id 必须从 1 开始没有缺口、没有重复地连续。

STOP JOB 是在 sql-client 中运行的语句,会把保存点路径以表格的一个单元格返回。路径形如 sp/savepoint-<作业 ID 的前 6 位>-…。保存点完成的瞬间,正在写的文件也会被确定,所以停止之后不应有点文件——如果还有,说明是被取消而停止的。

从保存点恢复并接着写

在最前面写 SET 'execution.state-recovery.path' = '<stop.out 의 세이브포인트>';(占位符为 stop.out 中的保存点),并把作业名称改为 flk-seq-2、Sink 改为 file:///root/flink/checkpoint/out2,提交得到的 /root/flink/checkpoint/run2.sql,把输出留在 /root/flink/checkpoint/run2.out。几秒后通过 REST(POST /jobs/<jid>/stop)拍下保存点并停止,把状态变为 COMPLETED 的响应保存到 /root/flink/checkpoint/stop2.json。out2 的第一个 id 必须紧接在 out1 的最后一个之后。

保存点里也以状态的形式保存着 Source 下一个要输出的编号。因此即使更换 Sink 路径,编号依然连续。通过 REST 停止时只会立即返回 request-id,所以要反复获取 /jobs//savepoints/,直到变为 COMPLETED 再保存。

不用保存点重新运行,然后取消

不设恢复路径,提交写入了作业名称 flk-seq-3、Sink file:///root/flink/checkpoint/out3 和 SET 'execution.checkpointing.externalized-checkpoint-retention' = 'RETAIN_ON_CANCELLATION'; 的 /root/flink/checkpoint/run3.sql,留下 /root/flink/checkpoint/run3.out。待 out3 中出现已确定的文件后,取消作业,并把已取消作业的 /jobs/<jid> 保存到 /root/flink/checkpoint/job3.json,把 /jobs/<jid>/checkpoints 保存到 /root/flink/checkpoint/checkpoints3.json。

没有恢复的作业,Source 会从 1 重新计数——会出现与 out1 相同的编号。取消不会拍保存点。没有保留设置时,取消的同时检查点目录会被删除,所以评分器会检查 checkpoints3.json 指向的 chk 目录中是否还留有 _metadata。

从保留的检查点接着写

把 checkpoints3.json 中最后一个已完成的检查点作为恢复路径,提交作业名称为 flk-seq-4、Sink 为同一个 file:///root/flink/checkpoint/out3 的 /root/flink/checkpoint/run4.sql,留下 /root/flink/checkpoint/run4.out。几秒后通过 REST 拍下保存点并停止,把状态响应保存到 /root/flink/checkpoint/stop4.json。out3 中已确定的文件必须从 1 开始没有缺口、没有重复地连续。

检查点目录(chk-N)也可以像保存点一样用作恢复路径。新作业会从检查点时刻的编号开始重新写,取消时写到一半的点文件不属于结果。评分器只收集不以点开头的文件。必须混有文件标签(uuid)不同的两次运行的文件。

报告——从磁盘上重新数一遍

在 /root/flink/checkpoint/report.json 中写入 savepoint_path(stop.out 中的保存点)、last_id_before_stop(out1 已确定文件的最后一个 id)、first_id_after_resume(out2 的第一个 id)、fresh_run_overlap(run3——从 1 重新运行的那次——在 out3 中确定的 id 里,同样存在于 out1、out2 的个数)、retained_checkpoint(checkpoints3.json 中最后一个已完成的路径)、leftover_inprogress_files(out3 中现在剩下的点文件数)。

所有值都可以从磁盘上重新数出来。已确定的文件名称以 part- 开头,同一次运行的文件共用同一个标签(uuid)。run3 的文件与包含 1 的文件标签相同。重叠是两份 id 列表交集的大小(comm -12)。