Flink-distributed-snapshot
Flink 的快照机制主要是为了保障作业 failover 时不丢失状态. Flink 提供了一种轻量级的快照机制,不需要停止作业就可以帮助用户持久化内存中的状态数据.
上图中的 markers(与 barrier 语义相同)通过流动来触发快照的制作,每一个编号都代表了一次快照,比如编号为 n 的 markers 从最上游流动到最下游就代表了一次快照的制作过程. 简述如下:
- 系统发送编号为 n 的
markers到最上游的算子,markers随着数据往下游流动; - 当下游算子收到
marker后,就开始将自身的状态保存到共享存储中; - 当所有最下游的算子接收到
marker并完成算子快照后,本次作业的快照制作完成.
一旦作业失败,重启时就可以从快照恢复.
下面为一个简单的 demo 说明(barrier 等同于 marker).

barrier到达 Source,将状态 offset=7 存储到共享存储;barrier到达 Task,将状态 sum=21 存储到共享存储;barrier到达 Sink,commit 本次快照,标志着快照的成功制作.

这时候突然间作业也挂掉, 重启时 Flink 会通过快照恢复各个状态. Source 会将自身的 offset 置为 7,Task 会将自身的 sum 置为 21.
现在我们可以认为 1、2、3、4、5、6 这 6 个数字的加和结果并没有丢失. 这个时候,offset 从 7 开始消费,跟作业失败前完全对接了起来,确保了 exactly-once