Flink 有状态流处理完全指南:从 Keyed State 到 Checkpoint 容错机制
2026/9/21 0:32:06 网站建设 项目流程
  • 大数据
  • 流处理
  • 批处理
  • 数据工程

【免费下载链接】flink

项目地址:https://gitcode.com/gh_mirrors/fli/flink
点击查看免费下载

有状态流处理是 Apache Flink 实现精确一次(exactly-once)容错与弹性扩缩容的基石。本文以 Flink 官方概念文档《有状态流处理》 为核心骨架,结合本仓库flink-runtime中的真实实现源码,系统讲解状态(State)的本质、Keyed State 与 Key Groups 的划分原理、以 Barrier 为核心的分布式快照(Checkpoint)机制、非对齐 Checkpoint、State Backend、Savepoint 以及批处理模式下的容错差异。读完本文,你将理解 Flink 为何能做到故障后状态一致,并能正确配置 Checkpoint 与 State Backend 支撑生产级应用。

什么是状态(State)

数据流中的很多算子一次只处理单个事件,例如事件解析器;但另一些算子需要跨多个事件"记住"信息,例如窗口算子(window operators)需要累积窗口内的数据。这类算子被称为有状态算子(stateful operators)

有状态操作的典型例子包括:

  • 应用在数据流中搜索特定事件模式时,状态中保存了迄今为止遇到的事件序列(例如 CEP 复杂事件处理);
  • 按分钟/小时/天对事件做聚合时,状态中保存了尚未完成的聚合结果;
  • 在数据点流上训练机器学习模型时,状态保存了模型参数的当前版本;
  • 需要管理历史数据时,状态允许高效访问过去发生的事件。

之所以让 Flink 感知状态的存在,是因为 Flink 需要借助状态来实现两件关键事情:

  1. 容错(Fault Tolerance):通过 Checkpoint 与 Savepoint 机制,让作业在故障后恢复到一致的状态;
  2. 弹性扩缩容(Rescaling):Flink 了解状态的分布方式后,可以在调整并行度时自动地把状态重新分布到各个并行实例上。

此外,不同 State Backend 决定了状态"存到哪里、怎么存",你可以在不修改应用逻辑的前提下切换 State Backend。

Keyed State 与流的 Key 严格对齐,每个并行实例只处理归属于自己 Key 的状态,保证所有状态更新都是本地操作。

Keyed State 与 Key Groups

内嵌的键值存储

Keyed State 可以被看作一个内嵌的键值(key/value)存储。关键特性在于:状态的划分与分布,严格跟随读取该状态的算子所消费的流一起进行。也就是说,只有经过 keyed/分区数据交换(即keyBy之后)的 keyed 流上才能访问键值状态,而且只能访问与当前事件 Key 相关联的值。

这种"流与状态的 Key 对齐"保证了:

  • 所有状态更新都是本地操作,无需分布式事务开销即可获得一致性;
  • Flink 可以在调整并行度时透明地重新分布状态、同步调整流的划分方式。

Key Groups:状态重分布的原子单元

Keyed State 进一步被组织为所谓的Key Groups(键组)。Key Groups 是 Flink 重分布 Keyed State 的原子单元Key Groups 的总数恰好等于作业定义的最大并行度(maximum parallelism)。执行期间,keyed 算子的每个并行实例负责一个或多个 Key Group 的 Key。

从源码可以印证这一点。KeyGroupRange(KeyGroupRange.java)注释明确写道:Key Group 是状态后端处理 keyed state 时对 key 空间进行划分的粒度,其范围是闭区间[startKeyGroup, endKeyGroup],并提供containsgetIntersectiongetNumberOfKeyGroups等操作。

而 Key 到 Key Group、Key Group 到并行算子的映射关系由KeyGroupRangeAssignment(KeyGroupRangeAssignment.java)完成:

// Key 先经过 murmurHash 散列,再对 maxParallelism 取模得到 Key Group 编号 public static int computeKeyGroupForKeyHash(int keyHash, int maxParallelism) { return MathUtils.murmurHash(keyHash) % maxParallelism; } // 根据当前并行度与最大并行度,计算某个算子实例负责的 Key Group 闭区间 public static KeyGroupRange computeKeyGroupRangeForOperatorIndex( int maxParallelism, int parallelism, int operatorIndex) { int start = ((operatorIndex * maxParallelism + parallelism - 1) / parallelism); int end = ((operatorIndex + 1) * maxParallelism - 1) / parallelism; return new KeyGroupRange(start, end); }

从源码结构还可以看到两个重要的边界约束(KeyGroupRangeAssignment.java):最大并行度的默认下界为1 << 7(128),以便用户在忘记显式配置时仍有一定扩缩容空间;最大并行度不能超过Short.MAX_VALUE + 1,否则取模分配会引入取整问题。由于并行度必须小于等于最大并行度,Key Group 总数固定为最大并行度,这使得无论当前并行度如何变化,每个 Key 归属的 Key Group 不变,从而支持任意时刻对状态进行重新划分。

状态持久化:流重放 + Checkpoint

Flink 通过流重放(stream replay)Checkpoint的组合实现容错。一个 Checkpoint 标记了每条输入流中的某个具体位置,以及每个算子对应的状态。当从 Checkpoint 恢复时,Flink 恢复算子状态并从 Checkpoint 标记的位置重放记录,从而保持一致性(精确一次处理语义)。

  • Checkpoint 间隔是一种权衡:间隔越短,容错开销越大但故障恢复时需重放的记录越少(恢复越快);间隔越长则反之。
  • 容错机制会持续对分布式数据流拍快照。对于状态很小的流式应用,这些快照非常轻量,可以高频执行而对性能影响甚微。
  • 应用状态被存储到可配置的位置,生产环境通常是一个分布式文件系统。

当程序因机器、网络或软件故障而失败时,Flink 会停止分布式数据流,重启算子并将它们重置到最近一次成功的 Checkpoint,输入流则重置到状态快照对应的位置。重启后的并行数据流所处理的任何记录,都被保证不会影响此前已 Checkpoint 的状态

⚠️注意:默认情况下 Checkpoint 是禁用的。开启与配置方法参见 Checkpointing 开发文档。

💡 该机制要兑现全部保证,要求数据源(如消息队列或 Broker)能够把流回退到某个确定的历史位置。Apache Kafka 具备这一能力,Flink 的 Kafka Connector 正是利用了这一特性。各连接器提供的具体保证参见 数据源与 Sink 的容错保证。

💡 由于 Flink 的 Checkpoint 通过分布式快照实现,文档中"快照(snapshot)"与"Checkpoint"常互换使用;"snapshot"也常被用来泛指 Checkpoint 或 Savepoint。

Checkpointing 机制详解

Flink 容错机制的核心是对分布式数据流与算子状态绘制一致的快照。这些快照作为一致的 Checkpoint,在故障时供系统回退。Flink 的快照机制论文为Lightweight Asynchronous Snapshots for Distributed Dataflows,其思想源自经典的Chandy-Lamport 分布式快照算法,并针对 Flink 的执行模型做了专门定制。

需要牢记的是,Checkpoint 相关的一切都可以异步进行:Checkpoint Barrier 不必同步齐步走,算子也可以异步地快照自己的状态。

自 Flink 1.11 起,Checkpoint 可以选择**对齐(aligned)不对齐(unaligned)**两种方式执行,下面先介绍对齐 Checkpoint。

Barrier(屏障)

Barrier 是 Flink 分布式快照的核心要素。它们被注入数据流,并作为数据流的一部分随记录一起流动。Barrier 具有以下特性:

  • 永不超越记录:Barrier 严格在流中按顺序流动,它把数据流中的记录划分为"进入当前快照的记录"和"进入下一个快照的记录"两部分;
  • 携带快照 ID:每个 Barrier 携带其所属快照的 ID(即它推动到前面的那批记录所属的快照编号);
  • 轻量且不中断:Barrier 不打断流的正常流动;同一时刻流中可存在来自不同快照的多个 Barrier,这意味着多个快照可以并发进行。

Barrier 随记录流动,将数据流切分为属于当前快照与下一个快照的记录集合。

从源码看,CheckpointBarrier(CheckpointBarrier.java)本质上是一个携带idtimestampCheckpointOptions的运行时事件,其类注释说明:Barrier 由 Source 在 JobManager 的指示下发出;算子从某条输入收到 Barrier 时,就知道这是 pre-checkpoint 与 post-checkpoint 数据的分界点;Barrier 的 ID 严格单调递增

Barrier 的完整流转过程如下:

  1. 注入:Barrier 在流 Source 处被注入到并行数据流中。快照n的 Barrier 注入点记为Sₙ,即快照覆盖数据的源流位置——例如对 Kafka 而言就是分区中最后一条记录的 offset。该位置Sₙ会被上报给Checkpoint 协调器(即 JobManager)
  2. 向下游传播:当一个中间算子从它的所有输入流都收到快照n的 Barrier 后,它会向所有输出流发出快照n的 Barrier。
  3. 完成确认:当 Sink 算子(流式 DAG 的末端)从它的所有输入流都收到 Barriern后,它向 Checkpoint 协调器确认快照n。当所有 Sink 都确认后,该快照即被视为完成

快照n完成后,作业不会再向 Source 索要Sₙ之前的记录,因为此时这些记录(及其衍生记录)已经完整穿过了整个数据流拓扑。

多输入算子的 Barrier 对齐(Alignment)

接收多个输入流的算子需要在快照 Barrier 上对齐输入流。下图展示了这一过程:

算子收到部分输入的 Barrier 后暂停该输入的处理,直到所有输入都收到 Barrier n 才继续。

对齐的具体步骤为:

  • 算子从某条输入流收到快照n的 Barrier 后,在该输入收到 Barriern之前,不再处理这条流上的任何记录——否则会把属于快照n的记录与属于快照n+1的记录混在一起;
  • 当最后一条输入流收到 Barriern后,算子先发出所有挂起的输出记录,然后自己发出快照n的 Barrier;
  • 算子对自己的状态拍快照,然后恢复处理所有输入流——先处理输入缓冲区中的记录,再处理流上的记录;
  • 最后,算子把状态异步写入 State Backend。

注意:对齐对于所有多输入算子以及shuffle 之后消费多个上游子任务输出流的算子都是必需的。

算子状态快照(Snapshotting Operator State)

只要算子包含任何形式的状态,这些状态就必须纳入快照。算子在其已收到所有输入的快照 Barrier、且尚未向输出发出 Barrier 之前这一时刻对状态拍快照。此时,Barrier 之前记录对状态的全部更新都已完成,而任何依赖 Barrier 之后记录的状态更新都尚未应用。

由于快照状态可能很大,它被存储在可配置的 State Backend 中。默认情况下存放在 JobManager 的内存里,但生产环境应配置分布式可靠存储(如 HDFS)。状态存储完成后,算子确认 Checkpoint、向输出流发出快照 Barrier,然后继续执行。

Checkpoint 快照包含两部分:每个并行数据源在快照开始时的流偏移/位置,以及每个算子指向快照中已存状态的指针。

最终生成的快照包含:

  • 每个并行数据源在快照开始时的流偏移/位置
  • 每个算子指向快照中所存状态的指针

恢复(Recovery)

对齐 Checkpoint 的恢复非常直接:故障发生后,Flink 选择最近完成的 Checkpointk,然后:

  1. 重新部署整个分布式数据流;
  2. 把 Checkpointk中快照的状态赋予每个算子;
  3. 让 Source 从位置Sₖ开始读取流——例如对 Kafka,就是告诉消费者从 offsetSₖ开始拉取。

如果状态是增量快照的,算子先加载最近一次全量快照的状态,再依次应用一系列增量快照更新。更多关于重启策略的内容参见 任务故障恢复。

非对齐 Checkpoint(Unaligned Checkpointing)

Checkpoint 也可以不对齐执行。其基本思想是:只要 in-flight(在途)数据成为算子状态的一部分,Checkpoint 就可以超越所有在途数据

值得说明的是,这种方法实际上更接近 Chandy-Lamport 算法本身,但 Flink 仍然在 Source 处插入 Barrier,以避免 Checkpoint 协调器过载。

算子遇到第一条非对齐 Barrier 时立即转发,被超越的记录被标记为异步存储。

非对齐方式下,算子处理非对齐 Checkpoint Barrier 的流程为:

  • 算子对存储在输入缓冲区中的第一条Barrier 立即作出反应;
  • 它立刻把 Barrier 转发给下游算子——通过把它追加到输出缓冲区的末尾
  • 算子把所有被超越的记录标记为异步存储,并对自己其余的状态创建快照。

因此,算子只会短暂地暂停输入处理(用于标记缓冲区)、转发 Barrier、创建其余状态的快照。

适用场景与限制:

  • 非对齐 Checkpoint 能保证 Barrier尽可能快地到达 Sink,特别适合至少存在一条慢速数据路径、对齐时间可能长达数小时的应用;
  • 但由于它增加了额外的 I/O 压力,当State Backend 的 I/O 本身就是瓶颈时,非对齐并不能带来帮助;
  • 更深入的讨论与其他限制参见 Checkpoints 运维文档 与 背压下的 Checkpoint;
  • Savepoint 始终是对齐的

非对齐恢复(Unaligned Recovery):算子先恢复 in-flight 数据,再开始处理来自上游算子的数据;除此之外,与对齐 Checkpoint 的恢复步骤相同。

开启非对齐 Checkpoint 的方式(配置execution.checkpointing.unaligned: true,或编程式开启,详见 Checkpointing 开发文档):

execution.checkpointing.unaligned: true
// 编程式开启(需配合 EXACTLY_ONCE 模式,且并发 Checkpoint 数为 1) env.getCheckpointConfig().enableUnalignedCheckpoints();

State Backends

键值索引(key/value indexes)底层采用何种数据结构,取决于所选的 State Backend:一种 State Backend 将数据存放在内存哈希表(HashMap)中,另一种使用 RocksDB 作为键值存储。除定义保存状态的数据结构外,State Backend 还实现了对键值状态进行时间点快照、并将快照作为 Checkpoint 一部分存储的逻辑。State Backend 可以在不修改应用逻辑的前提下替换。

Flink 开箱即用地提供两种 State Backend(state_backends.md):

  • HashMapStateBackend:状态以 Java 对象形式保存在堆中,读写极快,但状态大小受限于集群可用内存,且重用对象数据不安全;
  • EmbeddedRocksDBStateBackend:运行中的状态保存在内嵌 RocksDB 数据库中(默认存储在 TaskManager 数据目录),数据以序列化字节数组存储,Key 的比较按字节序进行而非 Java 的hashCode/equals();支持异步快照、状态大小仅受磁盘限制,且是唯一支持增量 Checkpoint的 State Backend;代价是每次读写都需要序列化/反序列化,最大吞吐量低于堆内存方案。

选择两者本质上是在性能与可扩展性之间权衡。若不显式配置,默认使用 HashMapStateBackend。可以通过 Flink 配置文件state.backend.type(可选值hashmap/rocksdb)做集群级默认配置,也可以在作业中编程覆盖:

Configuration config = new Configuration(); config.set(StateBackendOptions.STATE_BACKEND, "rocksdb"); env.configure(config);

自 Flink 1.13 起,所有 State Backend 生成统一的 savepoint 二进制格式,因此可以在生成 savepoint 后用另一种 State Backend 读取它(建议先升级到新版本再切换)。

状态快照被写入 State Backend,并作为 Checkpoint 的一部分持久化存储。

Savepoints

所有使用 Checkpoint 的程序都可以从Savepoint恢复执行。Savepoint 允许在完全不丢失状态的前提下更新程序或升级 Flink 集群。

Savepoint 本质上是手动触发的 Checkpoint:它使用常规 Checkpoint 机制对程序拍快照,并写入 State Backend。它与 Checkpoint 的相似之处在于都依赖同一套快照机制;区别在于两点(Savepoint 运维文档):

  • 由用户触发,而非周期性自动执行;
  • 不会自动过期:即使更新的 Checkpoint 完成,Savepoint 也不会被删除。

为了正确使用 Savepoint,理解 Checkpoint 与 Savepoint 的区别非常重要,详见 Checkpoints 与 Savepoints 对比。

Exactly Once vs. At Least Once

对齐步骤可能给流处理程序增加延迟。通常额外延迟只有几毫秒,但也出现过部分异常记录延迟明显增大的情况。对于要求**所有记录都保持超低延迟(几毫秒级)**的应用,Flink 提供了一个开关:在 Checkpoint 期间跳过流对齐。此时,只要算子从每条输入都看到 Checkpoint Barrier,就会立即绘制快照。

跳过对齐时,即使 Checkpointn的部分 Barrier 已到达,算子也会继续处理所有输入。这样算子在为 Checkpointn拍摄状态快照之前,就已经处理了属于 Checkpointn+1的元素。恢复时,这些记录会作为重复记录出现——因为它们既被包含在 Checkpointn的状态快照中,又会在 Checkpointn之后作为数据被重放。这就是 at-least-once 语义的来源。

💡重要提示:对齐只发生在有多个前驱算子(如 join)以及有多个发送方(如流重分区/shuffle 之后)的算子上。因此,仅包含可并行度极高的简单流式操作(map()flatMap()filter()等)的数据流,即使在 at-least-once 模式下,实际上也提供 exactly-once 保证。

在实际工程中,可以通过enableCheckpointing(interval, mode)显式选择语义模式,完整的配置示例参见 Checkpointing 开发文档:

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // 每 1000ms 开始一次 checkpoint env.enableCheckpointing(1000); // 设置模式为精确一次(默认值) env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE); // 确认 checkpoints 之间的最小间隔为 500ms env.getCheckpointConfig().setMinPauseBetweenCheckpoints(500); // Checkpoint 必须在一分钟内完成,否则被抛弃 env.getCheckpointConfig().setCheckpointTimeout(60000); // 允许两个连续的 checkpoint 错误 env.getCheckpointConfig().setTolerableCheckpointFailureNumber(2); // 同一时间只允许一个 checkpoint 进行 env.getCheckpointConfig().setMaxConcurrentCheckpoints(1);

批处理程序中的状态与容错

Flink 将批处理程序视为流处理程序在BATCH ExecutionMode下的一种特例——此时流是有界的(元素个数有限)。因此前述概念同样适用于批处理程序,但有两点例外(任务故障恢复文档):

  1. 批处理容错不使用 Checkpoint:恢复通过完整重放流实现。由于输入有界,全量重放是可行的。这把成本更多地推向了恢复阶段,但让常规处理更便宜(省去了 Checkpoint 开销);
  2. 批处理模式下的 State Backend 使用简化的内存/外置(in-memory/out-of-core)数据结构,而非键值索引结构。

小结

状态是 Flink 一切高级语义的根基:Keyed State 通过 Key Groups 实现可重分布的状态分区,Checkpoint 借助流 Barrier 与分布式快照把"流位置 + 算子状态"固化下来,配合 State Backend 与 Savepoint 共同构成了完整的容错体系。理解这套机制,不仅有助于正确开启与调优 Checkpoint,也能在遇到背压、超低延迟需求或批量升级场景时做出合理的架构决策。

延伸阅读(仓库内文档与源码)

  • 开发视角:Checkpointing 开发文档、Working with State
  • 运维视角:Checkpoints 运维文档、State Backends、Savepoints、背压下的 Checkpoint
  • 源码实现:KeyGroupRange.java、KeyGroupRangeAssignment.java、CheckpointBarrier.java
  • 大数据
  • 流处理
  • 批处理
  • 数据工程

【免费下载链接】flink

项目地址:https://gitcode.com/gh_mirrors/fli/flink
点击查看免费下载

相关推荐

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询