Checkpoint API
算子支持checkpoint
CheckpointedFunction
org.apache.flink.streaming.api.checkpoint.CheckpointedFunction
CheckpointedFunction是Stateful transformation functions的核心接口,用于跨stream维护state:
public void snapshotState(FunctionSnapshotContext functionSnapshotContext) throws Exception
在checkpoint的时候会被调用,用于snapshot state,通常用于flush、commit、synchronize外部系统
public void initializeState(FunctionInitializationContext context) throws Exception
在parallel function初始化的时候(第一次初始化或者从前一次checkpoint recover的时候)被调用,通常用来初始化state,以及处理state recovery的逻辑
从checkpoint中恢复数据时,需要判断snapshot当前的情况,
FunctionSnapshotContext实现了ManagedSnapshotContext, 父类中的方法: getCheckpointId,getCheckpointTimestampFunctionInitializationContext实现了ManagedInitializationContext接口, 实现了 isRestored、getOperatorStateStore、getKeyedStateStore方法
在初始化容器之后,我们使用上下文的isRestore()方法检查失败后是否正在恢复。如果是true,即正在恢复,则应用恢复逻辑。
样例: HBase写入OutPutFormat
1 | /** |
CheckpointListener
org.apache.flink.runtime.state.CheckpointListener
一旦所有checkpoint参与者确认完成,想要接收提交通知的功能/操作来实现。
TTL
1.8 自动清理原理
Flink 1.6.0版本引入了State TTL功能。它使流处理应用程序的开发人员配置过期时间,并在定义时间超时(Time to Live)之后进行清理。
在Flink 1.8.0中,该功能得到了扩展,包括对RocksDB和堆状态后端(FSStateBackend和MemoryStateBackend)的历史数据进行持续清理,从而实现旧条目的连续清理过程(根据TTL设置)。
RocksDB后台压缩可以过滤掉过期状态
如果你的Flink应用使用RocksDB作为状态后端存储,则可以启用另一个基于Flink特定压缩过滤器的清理策略。RocksDB定期运行异步压缩以合并状态更新并减少存储。Flink压缩过滤器使用TTL检查状态条目的到期时间戳,并丢弃所有过期值。
激活此功能的第一步是通过设置以下Flink配置选项来配置RocksDB状态后端:
state.backend.rocksdb.ttl.compaction.filter.enabled
配置RocksDB状态后端后,将为状态启用压缩清理策略,如以下代码示例所示:
1 | StateTtlConfig ttlConfig = StateTtlConfig |
TTL实现原理
State Backend 状态后端
三种状态后端:内存(MemoryStateend)、文件系统(FsStateend)和RocksDB(RocksDBStateend)

| state | 保存 | snapshot与restore | 大小 |
|---|---|---|---|
| keyed state | 堆内或堆外(RocksDB) | backend自行实现,用户不关心 | 大 |
| operator state | 堆内 | 用户自行实现 | 小 |

| State backend | snapshot保存 | checkpoint保存 |
|---|---|---|
| MemoryStateend | 内存 | 内存 |
| FsStateend | 内存 | 文件系统,如hdfs |
| RocksDBStateend | rocksdb | 文件系统,如hdfs |
Flink 的 keyed state 本质上来说就是一个键值对,所以与 RocksDB 的数据模型是吻合的。下图分别是 “window state” 和 “value state” 在 RocksDB 中的存储格式,所有存储的 key,value 均被序列化成 bytes 进行存储。

在 RocksDB 中,每个 state 独享一个 Column Family,而每个 Column family 使用各自独享的 write buffer 和 block cache,上图中的 window state 和 value state实际上分属不同的 column family。
最佳实践
operator state
慎重使用长 list
下图展示的是目前 task 端 operator state 在执行完 checkpoint 返回给 job master 端的 StateMetaInfo 的代码片段。

由于 operator state 没有 key group 的概念,所以为了实现改并发恢复的功能,需要对 operator state 中的每一个序列化后的元素存储一个位置偏移 offset,也就是构成了上图红框中的 offset 数组。
那么如果你的 operator state 中的 list 长度达到一定规模时,这个 offset 数组就可能会有几十 MB 的规模,关键这个数组是会返回给 job master,当 operator 的并发数目很大时,很容易触发 job master 的内存超用问题。我们遇到过用户把 operator state 当做黑名单存储,结果这个黑名单规模很大,导致一旦开始执行 checkpoint,job master 就会因为收到 task 发来的“巨大”的 offset 数组,而内存不断增长直到超用无法正常响应。
正确使用 UnionListState
union list state 目前被广泛使用在 kafka connector 中,不过可能用户日常开发中较少遇到,他的语义是从检查点恢复之后每个并发 task 内拿到的是原先所有operator 上的 state,如下图所示:

kafka connector 使用该功能,为的是从检查点恢复时,可以拿到之前的全局信息,如果用户需要使用该功能,需要切记恢复的 task 只取其中的一部分进行处理和用于下一次 snapshot,否则有可能随着作业不断的重启而导致 state 规模不断增长。
Keyed state 使用建议
如何正确清空当前的 state
state.clear() 实际上只能清理当前 key 对应的 value 值,如果想要清空整个 state,需要借助于 applyToAllKeys 方法,具体代码片段如下:

如果你的需求中只是对 state 有过期需求,借助于 state TTL 功能来清理会是一个性能更好的方案。
RocksDB 中考虑 value 值很大的极限场景
受限于 JNI bridge API 的限制,单个 value 只支持 2^31 bytes 大小,如果存在很极限的情况,可以考虑使用 MapState 来替代 ListState 或者 ValueState,因为RocksDB 的 map state 并不是将整个 map 作为 value 进行存储,而是将 map 中的一个条目作为键值对进行存储。
如何知道当前 RocksDB 的运行情况
比较直观的方式是打开 RocksDB 的 native metrics ,在默认使用 Flink managed memory 方式的情况下,state.backend.rocksdb.metrics.block-cache-usage ,state.backend.rocksdb.metrics.mem-table-flush-pending,state.backend.rocksdb.metrics.num-running-compactions 以及 state.backend.rocksdb.metrics.num-running-flushes 是比较重要的相关 metrics。
使用 checkpoint 的使用建议
Checkpoint 间隔不要太短
虽然理论上 Flink 支持很短的 checkpoint 间隔,但是在实际生产中,过短的间隔对于底层分布式文件系统而言,会带来很大的压力。另一方面,由于检查点的语义,所以实际上 Flink 作业处理 record 与执行 checkpoint 存在互斥锁,过于频繁的 checkpoint,可能会影响整体的性能。当然,这个建议的出发点是底层分布式文件系统的压力考虑。
合理设置超时时间
默认的超时时间是 10min,如果 state 规模大,则需要合理配置。最坏情况是分布式地创建速度大于单点(job master 端)的删除速度,导致整体存储集群可用空间压力较大。建议当检查点频繁因为超时而失败时,增大超时时间。
【参考文献】
state-management
org.apache.flink.streaming.api.checkpoint.CheckpointedFunction
- CheckpointedFunction是stateful transformation functions的核心接口,用于跨stream维护state
- snapshotState 在checkpoint的时候会被调用,用于snapshot state,通常用于flush、commit、synchronize外部系统
- initializeState 在parallel function初始化的时候(第一次初始化或者从前一次checkpoint recover的时候)被调用,通常用来初始化state,以及处理state recovery的逻辑
从checkpoint中恢复数据时,需要判断snapshot当前的情况,
FunctionSnapshotContext实现了ManagedSnapshotContext, 父类中的方法: getCheckpointId,getCheckpointTimestamp
FunctionInitializationContext实现了ManagedInitializationContext接口, 实现了isRestored、getOperatorStateStore、getKeyedStateStore方法
在初始化容器之后,我们使用上下文的isrestore()方法检查失败后是否正在恢复。如果是true,即正在恢复,则应用恢复逻辑。
样例: HBase写入OutPutFormat
1 | class PortraitOutputFormat extends RichOutputFormat<EventItem> implements CheckpointedFunction { |
org.apache.flink.runtime.state.CheckpointListener
一旦所有checkpoint参与者确认完全,该接口必须由想要接收提交通知的功能/操作来实现。
TTL
1.8 自动清理原理
Apache Flink的1.6.0版本引入了State TTL功能。它使流处理应用程序的开发人员配置过期时间,并在定义时间超时(Time to Live)之后进行清理。在Flink 1.8.0中,该功能得到了扩展,包括对RocksDB和堆状态后端(FSStateBackend和MemoryStateBackend)的历史数据进行持续清理,从而实现旧条目的连续清理过程(根据TTL设置)。
RocksDB后台压缩可以过滤掉过期状态
如果你的Flink应用程序使用RocksDB作为状态后端存储,则可以启用另一个基于Flink特定压缩过滤器的清理策略。RocksDB定期运行异步压缩以合并状态更新并减少存储。Flink压缩过滤器使用TTL检查状态条目的到期时间戳,并丢弃所有过期值。
激活此功能的第一步是通过设置以下Flink配置选项来配置RocksDB状态后端:
state.backend.rocksdb.ttl.compaction.filter.enabled
配置RocksDB状态后端后,将为状态启用压缩清理策略,如以下代码示例所示:
1 | StateTtlConfig ttlConfig = StateTtlConfig |
最佳实践
- 状态存储的所有数据,均需要考虑清理事迹
【参考文献】