ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Flink 有状态流处理完全指南:从 Keyed State 到 Checkpoint 容错机制

Flink 有状态流处理完全指南:从 Keyed State 到 Checkpoint 容错机制 大数据流处理批处理数据工程【免费下载链接】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 需要借助状态来实现两件关键事情容错Fault Tolerance通过 Checkpoint 与 Savepoint 机制让作业在故障后恢复到一致的状态弹性扩缩容RescalingFlink 了解状态的分布方式后可以在调整并行度时自动地把状态重新分布到各个并行实例上。此外不同 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。从源码可以印证这一点。KeyGroupRangeKeyGroupRange.java注释明确写道Key Group 是状态后端处理 keyed state 时对 key 空间进行划分的粒度其范围是闭区间[startKeyGroup, endKeyGroup]并提供contains、getIntersection、getNumberOfKeyGroups等操作。而 Key 到 Key Group、Key Group 到并行算子的映射关系由KeyGroupRangeAssignmentKeyGroupRangeAssignment.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 7128以便用户在忘记显式配置时仍有一定扩缩容空间最大并行度不能超过Short.MAX_VALUE 1否则取模分配会引入取整问题。由于并行度必须小于等于最大并行度Key Group 总数固定为最大并行度这使得无论当前并行度如何变化每个 Key 归属的 Key Group 不变从而支持任意时刻对状态进行重新划分。状态持久化流重放 CheckpointFlink 通过流重放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 随记录流动将数据流切分为属于当前快照与下一个快照的记录集合。从源码看CheckpointBarrierCheckpointBarrier.java本质上是一个携带id、timestamp与CheckpointOptions的运行时事件其类注释说明Barrier 由 Source 在 JobManager 的指示下发出算子从某条输入收到 Barrier 时就知道这是 pre-checkpoint 与 post-checkpoint 数据的分界点Barrier 的 ID 严格单调递增。Barrier 的完整流转过程如下注入Barrier 在流 Source 处被注入到并行数据流中。快照n的 Barrier 注入点记为Sₙ即快照覆盖数据的源流位置——例如对 Kafka 而言就是分区中最后一条记录的 offset。该位置Sₙ会被上报给Checkpoint 协调器即 JobManager。向下游传播当一个中间算子从它的所有输入流都收到快照n的 Barrier 后它会向所有输出流发出快照n的 Barrier。完成确认当 Sink 算子流式 DAG 的末端从它的所有输入流都收到 Barriern后它向 Checkpoint 协调器确认快照n。当所有 Sink 都确认后该快照即被视为完成。快照n完成后作业不会再向 Source 索要Sₙ之前的记录因为此时这些记录及其衍生记录已经完整穿过了整个数据流拓扑。多输入算子的 Barrier 对齐Alignment接收多个输入流的算子需要在快照 Barrier 上对齐输入流。下图展示了这一过程算子收到部分输入的 Barrier 后暂停该输入的处理直到所有输入都收到 Barrier n 才继续。对齐的具体步骤为算子从某条输入流收到快照n的 Barrier 后在该输入收到 Barriern之前不再处理这条流上的任何记录——否则会把属于快照n的记录与属于快照n1的记录混在一起当最后一条输入流收到 Barriern后算子先发出所有挂起的输出记录然后自己发出快照n的 Barrier算子对自己的状态拍快照然后恢复处理所有输入流——先处理输入缓冲区中的记录再处理流上的记录最后算子把状态异步写入 State Backend。注意对齐对于所有多输入算子以及shuffle 之后消费多个上游子任务输出流的算子都是必需的。算子状态快照Snapshotting Operator State只要算子包含任何形式的状态这些状态就必须纳入快照。算子在其已收到所有输入的快照 Barrier、且尚未向输出发出 Barrier 之前这一时刻对状态拍快照。此时Barrier 之前记录对状态的全部更新都已完成而任何依赖 Barrier 之后记录的状态更新都尚未应用。由于快照状态可能很大它被存储在可配置的 State Backend 中。默认情况下存放在 JobManager 的内存里但生产环境应配置分布式可靠存储如 HDFS。状态存储完成后算子确认 Checkpoint、向输出流发出快照 Barrier然后继续执行。Checkpoint 快照包含两部分每个并行数据源在快照开始时的流偏移/位置以及每个算子指向快照中已存状态的指针。最终生成的快照包含每个并行数据源在快照开始时的流偏移/位置每个算子指向快照中所存状态的指针。恢复Recovery对齐 Checkpoint 的恢复非常直接故障发生后Flink 选择最近完成的 Checkpointk然后重新部署整个分布式数据流把 Checkpointk中快照的状态赋予每个算子让 Source 从位置Sₖ开始读取流——例如对 Kafka就是告诉消费者从 offsetSₖ开始拉取。如果状态是增量快照的算子先加载最近一次全量快照的状态再依次应用一系列增量快照更新。更多关于重启策略的内容参见 任务故障恢复。非对齐 CheckpointUnaligned CheckpointingCheckpoint 也可以不对齐执行。其基本思想是只要 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 运维文档 与 背压下的 CheckpointSavepoint 始终是对齐的。非对齐恢复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 Backendstate_backends.mdHashMapStateBackend状态以 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拍摄状态快照之前就已经处理了属于 Checkpointn1的元素。恢复时这些记录会作为重复记录出现——因为它们既被包含在 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下的一种特例——此时流是有界的元素个数有限。因此前述概念同样适用于批处理程序但有两点例外任务故障恢复文档批处理容错不使用 Checkpoint恢复通过完整重放流实现。由于输入有界全量重放是可行的。这把成本更多地推向了恢复阶段但让常规处理更便宜省去了 Checkpoint 开销批处理模式下的 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点击查看免费下载相关推荐消融实验结果ControlNet Union SDXL 1.0各模块贡献度终极评估指南消融实验结果ControlNet Union SDXL 1.0各模块贡献度终极评估指南 ControlNet Union SDXL 1.0作为AI绘图领域的多计算机视觉基础模型AI 应用Apache Flink核心组件解密Runtime、State Backends与Checkpoint机制Apache Flink核心组件解密Runtime、State Backends与Checkpoint机制 引言流处理系统的稳定性基石 在实时数据处理领域大数据流处理批处理数据工程GetJobs错误处理机制从异常捕获到状态恢复的完整指南GetJobs错误处理机制从异常捕获到状态恢复的完整指南 GetJobs作为一款全平台自动投简历脚本其强大的错误处理机制确保了在Boss直聘、前程无忧、猎聘后端前端RPAAI 应用创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表