Flink Checkpoint与高可用
Flink Checkpoint 受 Chandy-Lamport 分布式快照启发,可以保证数据的高可用。但是有些情况下,不见得一定有效:
Flink On Yarn 模式,某个 Container 发生 OOM 异常,这种情况程序直接变成失败状态,此时 Flink 程序虽然开启 Checkpoint 也无法恢复,因为程序已经变成失败状态,所以此时可以借助外部参与启动程序,比如外部程序检测到实时任务失败时,重新对实时任务进行拉起。
1.1. 2PC
1.1.1. Exactly-once VS At-least-once
算子做快照时,如果等所有输入端的barrier都到了才开始做快照,可保证算子的exactly-once;
如果为了降低延时而跳过对齐,从而继续处理数据,那么等barrier都到齐后做快照就是at-least-once了,因为这次的快照掺杂了下一次快照的数据,当作业失败恢复的时候,这些数据会重复作用系统,就好像这些数据被消费了两遍。
注:对齐只会发生在算子的上端是join操作以及上游存在partition或者shuffle的情况,对于直连操作类似map、flatMap、filter等还是会保证exactly-once的语义。
1.1.2. 端到端的Exactly once实现
2PC分为几个阶段: 开始事务->预提交->提交(或回滚)
为了保证Exactly once, Source和Sink必须支持Flink的2PC
当状态涉及到外部系统时,需要外部系统支持事务操作来配合Flink实现2PC协议,从而保证数据的exatly-once。
这个时候,sink算子除了将自己的state写到状态后端,还必须准备好事务提交。
- 一旦所有的算子完成了它们的pre-commit,它们会要求一个commit。
- 如果存在一个算子pre-commit失败了,本次事务失败,我们回滚到上次的checkpoint。
- 一旦master做出了commit的决定,那么这个commit必须得到执行,就算宕机恢复也有继续执行。
1.1.2.1. pre-commit
pre-commit阶段起始于一次快照的开始,即master节点将checkpoint的barrier注入source端,barrier随着数据向下流动直到sink端。barrier每到一个算子,都会出发算子做本地快照。

当所有的算子都做完了本地快照并且回复master节点时,pre-commit阶段才算结束。这个时候,checkpoint已经成功,并且包含了外部系统的状态。如果作业失败,可以进行恢复。

1.1.2.2. commit
通知所有的算子这次checkpoint成功了,即2PC的commit阶段。source节点和window节点没有外部状态,所以这时它们不需要做任何操作。
而对于sink节点,需要commit这次事务,将数据写到外部系统。

1.1.2.3. rollback
一旦任何一个算子的快照保存失败,则触发回滚,同样的sink算子也需要取消写入外部的数据


1.1.3. TwoPhaseCommitSinkFunction
为了简化2PC的实现成本,flink抽象了TwoPhaseCommitSinkFunction
- beginTransaction。开始一次事务,在目的文件系统创建一个临时文件。接下来我们就可以将数据写到这个文件。
- preCommit。在这个阶段,将文件flush掉,同时重起一个文件写入,作为下一次事务的开始。
- commit。这个阶段,将文件写到真正的目的目录。值得注意的是,这会增加数据可视的延时。
- abort。如果回滚,那么删除临时文件。
如果pre-commit成功了,但是commit没有到达算子旧宕机了,flink会将算子恢复到pre-commit时的状态,然后继续commit。
我们需要做的还有就是保证commit的幂等性,这可以通过检查临时文件是否还在来实现。
1.2. checkpoint
保留策略:
- DELETE_ON_CANCELLATION 表示当程序取消时,删除 Checkpoint 存储文件。
- RETAIN_ON_CANCELLATION 表示当程序取消时,保存之前的 Checkpoint 存储文件
默认情况下,Flink不会触发一次 Checkpoint 当系统有其他 Checkpoint 在进行时,也就是说 Checkpoint 默认的并发为1。
CheckpointCoordinator :
针对 Flink DataStream 任务,程序需要经历从 StreamGraph -> JobGraph -> ExecutionGraph -> 物理执行图四个步骤,其中在 ExecutionGraph 构建时,会初始化 CheckpointCoordinator。ExecutionGraph通过ExecutionGraphBuilder.buildGraph方法构建,在构建完时,会调用 ExecutionGraph 的enableCheckpointing方法创建CheckpointCoordinator
Flink Checkpoint 参数配置及建议:
- 当 Checkpoint 时间比设置的 Checkpoint 间隔时间要长时,可以设置 Checkpoint 间最小时间间隔 。这样在上次 Checkpoint 完成时,不会立马进行下一次 Checkpoint,而是会等待一个最小时间间隔,然后在进行该次 Checkpoint。否则,每次 Checkpoint 完成时,就会立马开始下一次 Checkpoint,系统会有很多资源消耗 Checkpoint。
- 如果Flink状态很大,在进行恢复时,需要从远程存储读取状态恢复,此时可能导致任务恢复很慢,可以设置 Flink Task 本地状态恢复。任务状态本地恢复默认没有开启,可以设置参数state.backend.local-recovery值为true进行激活。
- Checkpoint保存数,Checkpoint 保存数默认是1,也就是保存最新的 Checkpoint 文件,当进行状态恢复时,如果最新的Checkpoint文件不可用时(比如HDFS文件所有副本都损坏或者其他原因),那么状态恢复就会失败,如果设置 Checkpoint 保存数2,即使最新的Checkpoint恢复失败,那么Flink 会回滚到之前那一次Checkpoint进行恢复。考虑到这种情况,用户可以增加 Checkpoint 保存数。
- 建议设置的 Checkpoint 的间隔时间最好大于 Checkpoint 的完成时间。
下图是不设置 Checkpoint 最小时间间隔示例图,可以看到,系统一致在进行 Checkpoint,可能对运行的任务产生一定影响:
1.3. savepoint
注意:
使用DataStream进行开发,建议为每个算子定义一个 uid,这样我们在修改作业时,即使导致程序拓扑图改变,由于相关算子 uid 没有变,那么这些算子还能够继续使用之前的状态,如果用户没有定义 uid , Flink 会为每个算子自动生成 uid,如果用户修改了程序,可能导致之前的状态程序不能再进行复用。
Flink 在触发Savepoint 或者 Checkpoint时,会根据这次触发的类型计算出在HDFS上面的目录:
如果类型是 Savepoint,那么 其 HDFS 上面的目录为:Savepoint 根目录+savepoint-jobid前六位+随机数字,具体如下格式:

Checkpoint 目录为 chk-checkpoint ID,具体格式如下:

- 使用 flink cancel -s 命令取消作业同时触发 Savepoint 时,会有一个问题,可能存在触发 Savepoint 失败。比如实时程序处于异常状态(比如 Checkpoint失败),而此时你停止作业,同时触发 Savepoint,这次 Savepoint 就会失败,这种情况会导致,在实时平台上面看到任务已经停止,但是实际实时作业在 Yarn 还在运行。针对这种情况,需要捕获触发 Savepoint 失败的异常,当抛出异常时,可以直接在 Yarn 上面 Kill 掉该任务。
- 使用 DataStream 程序开发时,最好为每个算子分配 uid,这样即使作业拓扑图变了,相关算子还是能够从之前的状态进行恢复,默认情况下,Flink 会为每个算子分配 uid,这种情况下,当你改变了程序的某些逻辑时,可能导致算子的 uid 发生改变,那么之前的状态数据,就不能进行复用,程序在启动的时候,就会报错。
- 由于 Savepoint 是程序的全局状态,对于某些状态很大的实时任务,当我们触发 Savepoint,可能会对运行着的实时任务产生影响,个人建议如果对于状态过大的实时任务,触发 Savepoint 的时间,不要太过频繁。根据状态的大小,适当的设置触发时间。
- 当我们从 Savepoint 进行恢复时,需要检查这次 Savepoint 目录文件是否可用。可能存在你上次触发 Savepoint 没有成功,导致 HDFS 目录上面 Savepoint 文件不可用或者缺少数据文件等,这种情况下,如果在指定损坏的 Savepoint 的状态目录进行状态恢复,任务会启动不起来。
1.4. snapshot保存到哪里? 应该需要汇总到jobManager?
1.5. state backend

FsStateBackend
构造方法:FsStateBackend(URI checkpointDataUri,boolean asynchronousSnapshots)
1 基于文件系统的状态管理器
2 如果使用,默认是异步
3 比较稳定,3个副本,比较安全。不会出现任务无法恢复等问题
4 状态大小受磁盘容量限制
存储方式:
- State: TaskManager内存
- checkpoint: 外部文件系统(本地或HDFS)
容量限制:
- 单TaskManager上State总量不超过它的内存
- 总大小不超过配置的文件系统容量
推荐使用场景:
- 常规使用状态的作业,例如分钟级窗口聚合、join、窗口比较长、kv状态大;需要开启HA的作业
- 可以用于生产场景
RocksDBStateBackend
状态数据先写入RocksDB,然后异步的将状态数据写入文件系统。正在进行计算的热数据存储在RocksDB,长时间才更新的数据写入磁盘中(文件系统)存储,体量比较小的元数据状态写入JobManager内存中(将工作state保存在RocksDB中,并且默认将checkpoint数据存在文件系统中)
目前唯一支持incremental的checkpoints的策略
构造方法:RocksDBStateBackend(URI checkpointDataUri,boolean enableIncrementalCheckpointing)
存储方式:
- State: TaskManager上的KV数据库(实际使用内存+硬盘)
- Checkpoint: 外部文件系统(本地或HDFS)
容量限制:
- 单TaskManager上的State总量不超过他的内存+磁盘
- 单key最大2G
- 总大小不超过配置的文件系统容量
推荐使用的场景:
- 超大状态的作业,例如天级别窗口聚合;需要开启HA的作业;对状态读写性能要求不高的作业
- 可以在生产环境使用
MemoryStateBackend
构造方法:MemoryStateBackend(int maxStateSize, boolean asynchronousSnapshots)
主机内存中的数据可能会丢失,任务可能无法恢复
存储方式:
- State: TaskManager内存
- Checkpoint: JobManager内存
容量限制
- 单个State maxStateSize默认5M
- maxStateSize <= akka.frameSize 默认10M
- 总大小不超过JobManager的内存
推荐使用场景:
- 本地测试;几乎无状态的作业,比如ETL;JobManager不容易挂,或挂掉影响不大的情况
- 不推荐在生产环境使用
1.6. checkpoint 与 savepoint
Checkpoint指定触发生成时间间隔后,每当需要触发Checkpoint时,会向Flink程序运行时的多个分布式的Stream Source中插入一个Barrier标记,这些Barrier会根据Stream中的数据记录一起流向下游的各个Operator。
当一个Operator接收到一个Barrier时,它会暂停处理Steam中新接收到的数据记录。
因为一个Operator可能存在多个输入的Stream,而每个Stream中都会存在对应的Barrier,该Operator要等到所有的输入Stream中的Barrier都到达。(对齐)
当所有Stream中的Barrier都已经到达该Operator,这时所有的Barrier在时间上看来是同一个时刻点(表示已经对齐),在等待所有Barrier到达的过程中,
Operator的Buffer中可能已经缓存了一些比Barrier早到达Operator的数据记录(Outgoing Records),这时该Operator会将数据记录(Outgoing Records)发射(Emit)出去,作为下游Operator的输入,
最后将Barrier对应Snapshot发射(Emit)出去作为此次Checkpoint的结果数据。
Checkpoint 是增量做的,每次的时间较短,数据量较小,只要在程序里面启用后会自动触发,用户无须感知;Checkpoint 是作业 failover 的时候自动使用,不需要用户指定。
Savepoint 是全量做的,每次的时间较长,数据量较大,需要用户主动去触发。Savepoint 一般用于程序的版本更新(详见文档),Bug 修复,A/B Test 等场景,需要用户指定。
保存的内容
- 首先,Savepoint 包含了一个目录,其中包含(通常很大的)二进制文件,这些文件表示了整个流应用在 Checkpoint/Savepoint 时的状态。
- 以及一个(相对较小的)元数据文件,包含了指向 Savapoint 各个文件的指针,并存储在所选的分布式文件系统或数据存储中。
目标
Savepoint 和 Checkpoint 的不同之处很像传统数据库中备份与恢复日志之间的区别。Checkpoint 的主要目标是充当 Flink 中的恢复机制,确保能从潜在的故障中恢复。相反,Savepoint 的主要目标是充当手动备份、恢复暂停作业的方法。
实现
Checkpoint 被设计成轻量和快速的机制。它们可能(但不一定必须)利用底层状态后端的不同功能尽可能快速地恢复数据。例如,基于 RocksDB 状态后端的增量检查点,能够加速 RocksDB 的 checkpoint 过程,这使得 checkpoint 机制变得更加轻量。相反,Savepoint 旨在更多地关注数据的可移植性,并支持对作业做任何更改而状态能保持兼容,这使得生成和恢复的成本更高
状态文件保留策略
Checkpoint默认程序删除,可以设置CheckpointConfig中的参数进行保留 。Savepoint会一直保存,除非用户删除 。
应用
- 部署流应用的一个新版本,包括新功能、BUG 修复、或者一个更好的机器学习模型
- 引入 A/B 测试,使用相同的源数据测试程序的不同版本,从同一时间点开始测试而不牺牲先前的状态
- 在需要更多资源时扩容应用程序
- 迁移流应用程序到 Flink 的新版本上,或者迁移到另一个集群
Flink数据一致性
一、综述
flink 通过内部依赖checkpoint 并且可以通过设置其参数exactly-once 实现其内部的一致性。但要实现其端到端的一致性,还必须保证
1、source 外部数据源可重设数据的读取位置
2、sink端 需要保证数据从故障恢复时,数据不会重复写入外部系统(或者可以逻辑实现写入多次,但只有一次生效的数据sink端)
二、sink 端到端实现方式
幂等操作:
一个操作,可以重复执行多次,但只导致一次结果更改,豁免重复操作执行就不起作用了,他的瑕疵 (在系统恢复的过程中,如果这段时间内多个更新或者插入导致状态不一致,但当数据追上就可以了)
(逻辑与、逻辑或等)具体理解参照自己以前写的文章。
事务写入:
事务应该具有四个属性:原子性、一致性、隔离性、持久性等。其具体的实现方式有两种
(1)、预写日志
简单易于实现,由于数据提前在状态后端中做了缓存,所以无论什么sink系统,都能用这种方式一批搞定,DataStream API提供了一个模板类:GenericWriteAheadSink,来实现这种事务性sink;
缺点:
1)、sink系统没说他支持事务。有可能出现一部分写入集群了。一部分没有写进去(如果实表,再写一次就写重复了)
2)、checkpoint做完了sink才去真正的写入(但其实得等sink都写完checkpoint才能生效,所以WAL这个机制jobmanager确定它写完还不算真正写完,还得有一个外部系统已经确认 完成的checkpoint)
(2)、两阶段提交。 flink 真正实现exactle-once
对于每个checkpoint,sink 任务会启动一个事务,并将接下来所有接收的数据添加到事务中,然后将这些数据写入外部sink系统,但不提交他们(这里是预提交)。当checkpoint完成时的通知,它才正式提交事务,实现结果的真正写入;这种方式真正实现了exactly-once,它需要一个提供事务支持的外部sink系统,Flink提供了其具体实现(TwoPhaseCommitSinkFunction接口)
三、 2pc 对外部 sink的要求
1、外部sink系统必须事务支持,或者sink任务必须能够模拟外部系统上的事务;
2、在checkpoint的间隔期间里,必须能够开启一个事务,并接受数据写入。
3、在收到checkpoint完成通知之前,事务必须是“等待提交”的状态,在故障恢复的情况线,这可能需要一些时间。如果个时候sink系统关闭事务(例如超时了),那么未提交的数据就会丢失;
4、四年任务必选能够在进程失败后恢复事务
5、提交事务必须是幂等操作;
四、综上不同Source和sink的一致性保证:

五、应用(flink+kafka 端到端一致性保证)
flink 和kafka 端到端一致性(kafka(source+flink+kafka(sink)))
1、内部 – 利用checkpoint机制,把状态存盘,发生故障的时候可以恢复,保证内部的状态一致性
2、source – kafka consumer作为source,可以将偏移量保存下来,如果后续任务出现了故障,恢复的时候可以由连接器重置偏移量,重新消费数据,保证一致性;
1 | kafka 0.8 和kafka 0.11 之后 通过以下配置将偏移量保存,恢复时候重新消费 |
3、sink FlinkkafkaProducer作为Sink,采用两阶段提交的sink,由下图可以看出flink 0.11 已经默认继承了TwoPhaseCommitSinkFunction
但我们需要在参数种传入指定语义,它默认时还是at-least-once
此外我们还需要进行一些producer的容错配置:
(1)除了启用Flink的检查点之外,还可以通过将适当的semantic参数传递给FlinkKafkaProducer011(FlinkKafkaProducer对于Kafka> = 1.0.0版本)
(2)来选择三种不同的操作模式
1)、Semantic.NONE 代表at-mostly-once语义
2)、Semantic.AT_LEAST_ONCE(Flink默认设置
3)、Semantic.EXACTLY_ONCE 使用Kafka事务提供一次精确的语义,每当您使用事务写入Kafka时
(3)、请不要忘记消费kafka记录任何应用程序设置所需的设置isolation.leva(read_committed 或者read_uncommitted-后者是默认)
read_committed,只是读取已经提交的数据。
应用;
Semantic.EXACTLY_ONCE依赖与下游系统能支持事务操作.以0.11kafka为例.
transaction.max.timeout.ms 最大超市时长,默认15分钟,如果需要用exactly语义,需要增加这个值。(因为它小于transaction.timeout.ms )
isolation.level 如果需要用到exactly语义,需要在下级consumerConfig中设置read-commited [read-uncommited(默认值)]
transaction.timeout.ms 默认为1hour
其参数对应关系为 和一些报错问题
checkpoint间隔<transaction.timeout.ms<transaction.max.timeout.ms
参考:https://www.cnblogs.com/createweb/p/11971846.html
注意:
1、Semantic.EXACTLY_ONCE 模式每个FlinkKafkaProducer011实例使用一个固定大小的KafkaProducers池。每个检查点使用这些生产者中的每一个。如果并发检查点的数量超过池大小,FlinkKafkaProducer011 将引发异常,并使整个应用程序失败。请相应地配置最大池大小和最大并发检查点数。
2、Semantic.EXACTLY_ONCE采取所有可能的措施,不要留下任何挥之不去的数据,否则这将有碍于消费者更多地阅读Kafka主题。但是,如果flink应用程序在第一个检查点之前失败,则在重新启动此类应用程序后,系统种将没有有关先前池大小信息,因此,在第一个检查点完成前按比例缩小Flink应用程序的FlinkKafkaProducer011.SAFE_SCALE_DOWN_FACTOR
1 | //1。设置最大允许的并行checkpoint数,防止超过producer池的个数发生异常 |
完整应用代码:
1 | package com.shufang.flink.connectors |