知识卡片
Checkpoint基于异步Barrier快照的原理
内容
Flink的Checkpoint机制是Chandy-Lamport分布式快照算法的一个变体:Checkpoint协调器让所有Source算子记录当前读取的偏移量,并往数据流里插入一种特殊标记——Checkpoint barrier,barrier像普通数据一样沿着数据流向下游传递;每个算子收到barrier时,就基于自己当前的状态生成一份快照,再把barrier转发给下游,如此层层传递直到所有算子都完成快照。这样做的巧妙之处在于,不需要停止整个作业去做一次全局同步的”世界暂停”式快照,而是靠一个标记在数据流中的传播顺序,天然地把”这个时刻之前的所有数据都已处理”这个边界一致地划在了每个算子上。一旦作业失败,只需要从最近一次成功的快照恢复所有算子状态,配合能重放的数据源(如Kafka),就能保证不丢数据地继续计算。这套自动周期性快照正是[[状态检查点是exactly-once语义的前提|精确一次语义]]的底层实现;由人工手动触发、生命周期不同的版本见[[Savepoint与Checkpoint的本质区别]]。
参考来源
《Flink入门与实战》第4章《状态管理及容错机制》