大数据流处理批处理数据工程【免费下载链接】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),仅供参考
企业数字化 ERP 产品动态
相关推荐
MXNet mxnet.image 模块实战指南:图像读取、解码、数据增强与迭代器 API 全解析 深度学习机器学习人工智能 【免费下载链接】mxnet Lightweight, Portable, Flexible Distributed/Mobile Deep Learning with Dynamic, Mutation-aware Dataflow Dep Scheduler; for Python, R, Julia, Scala, Go, Javascript and more 项目地址: https://gitcode.c… · 2026/9/21 0:32:25
marked 解析边界行为探秘:波浪线围栏代码块如何在文件末尾打断段落 marked 解析边界行为探秘:波浪线围栏代码块如何在文件末尾打断段落 【免费下载链接】marked A markdown parser and compiler. Built for speed. 项目地址: https://gitcode.com/gh_mirrors/ma/marked
导读
本篇文章聚焦 marked(A markdown pars… · 2026/9/21 0:32:25
文华财经赢顺DK多空指标优化实战:从原理到参数调试全指南 先说一个我自己的经历。早几年做螺纹钢,我在文华财经赢顺里加载了DK多空指标,蓝点买入、红点卖出,信号一目了然,当时觉得自己找到了圣杯。结果做了一整个季度,账户曲线跟过山车一样,最后算下来还是亏的。问… · 2026/9/21 0:32:25
别被网页制作模板中文坑了,懂建站报价才不亏 别被网页制作模板中文坑了,懂建站报价才不亏 网站做好了没人访问,这钱白花得冤不冤?很多老板找外包,问完建站报价,对方甩给你一个“网页制作模板中文”链接,说这是高端定制。你一看,哦,是套壳的。更坑的是,有些模板连基础的SEO结构都没做好,上线三个月,百度搜不到你公司名字。… · 2026/9/21 8:31:34
2026最新微信小程序连接wordpress:解决域名服务器搞不懂的实战指南 2026最新微信小程序连接wordpress:解决域名服务器搞不懂的实战指南 域名解析指向不对,服务器端口没开放,SSL证书配置报错——这三座大山,劝退了一半想用微信小程序展示WordPress内容的开发者。别急,2026最新的连接方案早已绕开了传统Web服务器配置的深坑,核心逻辑是:… · 2026/9/21 8:17:36
企业网站做电脑营销多少钱?揭秘防黑挂马的底层逻辑 企业网站做电脑营销多少钱?揭秘防黑挂马的底层逻辑 网站突然被黑,首页挂满赌博广告,后台密码怎么改都没用,这种绝望感做过站的都懂。很多老板第一反应是问:“清理一次病毒多少钱?”或者“换个服务器多少钱?”但真相往往扎心:单纯清理病毒的费用可能只要几百块,但重建信任、修复SEO权重、补全安全漏洞的成本,往… · 2026/9/21 8:03:27
3步搞定做品管圈网站从零搭建到上线避坑指南 3步搞定做品管圈网站从零搭建到上线避坑指南 不会写代码,但想给团队搭个品管圈展示平台?别慌。 很多河南的创业老板都卡在这一步:手里有现成的QCC成果,想做个官网放上去,结果一搜全是“前端开发教程”,看得头大。 做品管圈网站 这事儿,真没你想的那么玄乎。只要路子对,零基础也能 从零搭建… · 2026/9/21 7:45:56
Voyager 資料夾管理指南:為 Gemini 與 AI Studio 的 AI 對話打造真正的「檔案系統」 AI 应用前端 【免费下载链接】voyager Enhancement suite for Gemini, AI Studio, Claude & ChatGPT — plus a prompt manager for any websites, DeepSeek Harness included. / 面向 Gemini、AI Studio、Claude 与 ChatGPT 的增强套件;其中的提示词管理器可用… · 2026/9/21 7:41:58
Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化 直接铺开项目本身吧。这几个月我一直在折腾一件事:用Flutter给OpenHarmony做一款游戏集合类的App,说白了就是把若干小游戏塞进一个壳里,用统一入口分发。这个方向本身不算新鲜,真正让我花了不少心思的,是首页那堆游戏卡… · 2026/9/21 0:02:39
Word表格编号全攻略:从列表编号到题注交叉引用 写Word文档,最让人头疼的往往是那些“看起来不起眼”的小问题。比如表格编号这事:今天在表后面多加了两个空白行,明天给客户交稿前发现整个章节的编号全部错位,光是挨个改序号就能耗掉大半个下午。我前阵子帮人整理一份上百页的技… · 2026/9/21 0:02:39
从第一个站到第二个站:独立开发者的静态网站选型与落地实践 1. 项目概述1.1 核心需求解析做独立开发者这几年,说实话,第一个网站上线的那天晚上我兴奋得没睡着。但等它跑了半年,流量惨淡、功能臃肿、代码自己都懒得看第二遍之后,我才慢慢琢磨明白一个道理:第一个网站是练手&… · 2026/9/20 0:00:41
agents-generator 决策矩阵全解析:从项目检测到 AGENTS.md 规则生成的 16 步判定流程 agents-generator 决策矩阵全解析:从项目检测到 AGENTS.md 规则生成的 16 步判定流程 【免费下载链接】agentic-awesome-skills AAS Core is the local, agent-first control plane for complete catalog discovery, agent-owned selection, stack validation, and … · 2026/9/21 0:00:18
gin-vue-admin 前端工具函数全景指南:src/utils 复用规范与源码级解析 gin-vue-admin 前端工具函数全景指南:src/utils 复用规范与源码级解析 【免费下载链接】gin-vue-admin 🚀ViteVue3Gin拥有AI辅助的基础开发平台,企业级业务AI开发解决方案,内置mcp辅助服务,内置skills管理,… · 2026/9/21 0:00:18