暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

Flink 状态管理

李孟的博客 2020-10-20
433

一.简介

函数里所有任务去维护并用来计算结果的数据都属于任务的状态。

流计算分为无状态和有状态两种情况:

无状态计算观察每个独立事件,并根据最后一个事件输出结果。

有状态计算则会基于多个事件输出结果。

状态可以分为两类:

算子状态(operator state)。

键值分区状态(keyed state)。

state 一般指一个具体task/operator的状态。

checkpoint表示一个Flink Job,在一个特定时刻的一份全局状态快照,即包含了所有task/operator的状态。

二.算子状态

算子状态的作用范围限定为算子任务。这意味着由同一并行任务所处理的所有数据都可以访问到相同的状态,状态对于同一任务而言是共享的。算子状态不能由相同或不同算子的另一个任务访问。

Flink为算子状态提供三种基本数据结构:

列表状态(List state):将状态表示为一组数据的列表。

联合列表状态(Union list state):也将状态表示为数据的列表。它与常规列表状态的区别在于,在发生故障时,或者从保存点(savepoint)启动应用程序时如何恢复。

广播状态(Broadcast state):如果一个算子有多项任务,而它的每项任务状态又都相同,那么这种特殊情况最适合应用广播状态。

三.键值分区状态

键控状态是根据输入数据流中定义的键(key)来维护和访问的。Flink为每个键值维护一个状态实例,并将具有相同键的所有数据,都分区到同一个算子任务中,这个任务会维护和处理这个key对应的状态。当任务处理一条数据时,它会自动将状态的访问范围限定为当前数据的key。因此,具有相同key的所有数据都会访问相同的状态。Keyed State很类似于一个分布式的key-value map数据结构,只能用于KeyedStream(keyBy算子处理之后)。

Flink为键值分区状态提供三种基本数据结构:

单值状态(value state):每个键对应存储一个任意类型的值,该值也可以是某个复杂数据结构。

列表状态(list state):每个键对应存储一个值的列表,列表中条目可以是任意类型。

映射状态(map state):每个键对应存储一个键值映射,该映射的键和值可以是任意类型。

四.状态后端

每传入一条数据,有状态的算子任务都会读取和更新状态。由于有效的状态访问对于处理数据的低延迟至关重要,因此每个并行任务都会在本地维护其状态,以确保快速的状态访问。状态存储访问已经维护,由一个可插入组件决定,这个组件叫做状态后端。

负责两件事情

本地状态管理。

将检查点(checkpoint)状态写入远程存储。

分类

MemoryStateBackend

内存级的状态后端,会将键控状态作为内存中对象进行管理,将它们存储TaskManager的Jvm堆上,而将checkpoint存储在JobManager的内存中。

FsStateBackend

将checkpoint存到远程的持久化文件系统(FileSystem)上。而对于本地状态,跟MemoryStateBackend一样,也会存在TaskManager的JVM堆上。

RocksDBStateBackend

将所有状态序列化后,存入本地的RocksDB中存储。

五.有状态算子的扩容

流式应用的一项基本需求是根据输入数据到达速率的变化调整算子并行度。

对于有状态算子,改变并行度变得很复杂,因为我们需要把状态重新分组,分配到与之前数量不等的并行任务上。

键值分区扩容

带有键值分区状态算子在扩容时会根据新的任务数量对键值重新分区。

为了降低状态在不同任务之间迁移的成本,状态分组处理(key group),以组为单位分配给不同任务。

算子列表状态扩容

带有算子列表状态算子在扩容时会对列表中条目重新分配。

所有并行算子任务的列表条目统一收集起来,随后均匀分配到更少/更多的任务上。

如果列表条目数量小于算子新设置的并行度,部分任务在启动时状态就可能为空。

算子联合列表状态扩容

带有算子联合列表状态的算子会在扩缩容时把状态列表全部条目广播到全部任务上,随后由任务自己决定哪些条目该保留。

算子广播状态

带有算子广播状态的算子在扩容时会把状态拷贝到全部新任务上。

这样做原因广播状态确保所有任务的状态相同。

缩容时,由于复制不会丢失,我们可以简单地停掉多出任务。

六.状态后端场景

Flink提供三种可用的状态后端:MemoryStateBackend,FsStateBackend,和RocksDBStateBackend。

场景

MemoryStateBackend:

本地开发或调试。小状态场景。

FsStateBackend:

大状态,长窗口或大键值状态。高可用场景。

RocksDBStateBackend:

大状态,长窗口或大键值状态。高可用场景。增量checkpoint,超大状态。

设置


如果没有明确指定,将使用 jobmanager 做为默认的 state backend。你能在 flink-conf.yaml 中为所有 Job 设置其他默认的 State Backend。每一个 Job 的 state backend 配置会覆盖默认的 state backend 配置。

七.MemoryStateBackend

MemoryStateBackend 是将状态维护在Java堆上的一个内部状态后端。键值状态和窗口算子使用哈希表来存储数据(values)和定时器(timers)。当应用程序 checkpoint 时,此后端会在将状态发给 JobManager 之前快照下状态,JobManager 也将状态存储在 Java 堆上。默认情况下,MemoryStateBackend 配置成支持异步快照。异步快照可以避免阻塞数据流的处理,从而避免反压的发生。

注意

默认情况,每一个独立状态大小限制是5MB。在MemoryStateBackend 的构造器中可以增加其大小。状态大小受到 akka 帧大小的限制,所以无论怎么调整状态大小配置,都不能大于 akka 的帧大小。也可以通过 akka.framesize 调整 akka 帧大小(通过配置文档了解更多)。状态的总大小不能超过 JobManager 的内存。

场景

本地开发或调试时建议使用 MemoryStateBackend,因为这种场景的状态大小的是有限的。MemoryStateBackend 最适合小状态的应用场景。例如 Kafka consumer,或者一次仅记录的函数 (Map, FlatMap,或 Filter)。

八.FsStateBackend


FsStateBackend 需要配置一个文件系统的 URL(类型、地址、路径),例如:”hdfs://namenode:40010/flink/checkpoints” 或 “file:///data/flink/checkpoints”。

当选择使用FsStateBackend时,正在进行数据会被存储TaskManager内存中。在checkpoint时,此后端会将状态快照写入配置的文件系统和目录的文件中,同时会在JobManager的内存中

(在高可用场景下会存在 Zookeeper 中)存储极少的元数据。

默认情况下,FsStateBackend 配置成提供异步快照,以避免在状态 checkpoint 时阻塞数据流的处理。该特性可以实例化 FsStateBackend 时传入 false 的布尔标志来禁用掉,例如:

new FsStateBackend(path, false);

注意

当前状态仍然会存在TaskManager中,所以状态的大小不能超过TaskManager内存。

场景

FsStateBackend 适用于处理大状态,长窗口,或大键值状态的有状态处理任务。FsStateBackend 非常适合用于高可用方案。


九.RocksDBStateBackend

RocksDBStateBackend 的配置也需要一个文件系统(类型,地址,路径),如下所示:

“hdfs://namenode:40010/flink/checkpoints” 或“s3://flink/checkpoints”

RocksDB 是一种嵌入式的本地数据库。RocksDBStateBackend 将处理中的数据使用 RocksDB 存储在本地磁盘上。在 checkpoint 时,整个 RocksDB 数据库会被存储到配置的文件系统中,或者在超大状态作业时可以将增量的数据存储到配置的文件系统中。

同时 Flink 会将极少的元数据存储在 JobManager 的内存中,或者在 Zookeeper 中(对于高可用的情况)。RocksDB 默认也是配置成异步快照的模式。

注意

RocksDB 支持的单 key 和单 value 的大小最大为每个 2^31 字节。这是因为 RocksDB 的 JNI API 是基于 byte[] 的。我们需要强调的是,对于使用具有合并操作的状态的应用程序,例如 ListState,随着时间可能会累积到超过 2^31 字节大小,这将会导致在接下来的查询中失败。

场景

RocksDBStateBackend 最适合用于处理大状态,长窗口,或大键值状态的有状态处理任务。RocksDBStateBackend 是目前唯一支持增量 checkpoint 的后端。增量 checkpoint 用于超大状态的场景。RocksDBStateBackend 非常适合用于高可用方案。


当使用RocksDB 时,状态大小只受限于磁盘可用空间的大小,这也使得 RocksDBStateBackend 成为管理超大状态的最佳选择。使用 RocksDB 的权衡点在于所有的状态相关的操作都需要序列化(或反序列化)才能跨越 JNI 边界。与上面提到的堆上后端相比,这可能会影响应用程序的吞吐量。

十.设置

如果没有明确指定,将使用 jobmanager 做为默认的 state backend。你能在 flink-conf.yaml 中为所有 Job 设置其他默认的 State Backend。每一个 Job 的 state backend 配置会覆盖默认的 state backend 配置。

设置每个Job的State Backend

StreamExecutionEnvironment 可以对每个 Job 的 State Backend 进行设置

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setStateBackend(new FsStateBackend("hdfs://namenode:40010/flink/checkpoints"));

RocksDBStateBackend

<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-statebackend-rocksdb_2.11</artifactId>
<version>1.11.0</version>
<scope>provided</scope>
</dependency>

由于 RocksDB 是 Flink 默认分发包的一部分,所以如果你没在代码中使用 RocksDB,则不需要添加此依赖。而且可以在 flink-conf.yaml 文件中通过 state.backend 配置 State Backend。

设置默认的(全局)State Backend

在 flink-conf.yaml 可以通过键 state.backend 设置默认的 State Backend。

可选值包括 jobmanager (MemoryStateBackend)、filesystem (FsStateBackend)、rocksdb (RocksDBStateBackend),或使用实现了 state backend 工厂 StateBackendFactory 的类的全限定类名,例如:RocksDBStateBackend 对应为 org.apache.flink.contrib.streaming.state.RocksDBStateBackendFactory。

state.checkpoints.dir 选项指定了所有 State Backend 写 CheckPoint 数据和写元数据文件的目录。

# 用于存储 operator state 快照的 State Backend
state.backend: filesystem
# 存储快照的目录
state.checkpoints.dir: hdfs://namenode:40010/flink/checkpoints

参考

《Stream Processing with Apache Flink》


文章转载自李孟的博客,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

评论