尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

Flink Checkpoint 原理:概念与整体流程

Flink Checkpoint 原理:概念与整体流程 前些天发现了一个巨牛的人工智能学习网站通俗易懂风趣幽默忍不住分享一下给大家。点击跳转到网站https://www.captainai.net/dongkelun前言本文是 Flink Checkpoint 系列的第一篇从概念层面介绍 Checkpoint 是什么、为什么需要、核心算法原理以及整体执行流程。后续文章会逐层深入到源码实现。前置知识Flink SQL Checkpoint 学习总结 — Checkpoint 参数详解与重启验证版本Flink 1.15.3什么是 CheckpointCheckpoint 是 Flink 提供的一种分布式快照机制周期性地保存算子状态和数据源消费位置到持久化存储。当作业失败时Flink 可以从最近一次成功的 Checkpoint 恢复实现 Exactly-Once 语义。核心特性自动触发按配置的时间间隔周期性执行无需人工干预增量保存RocksDB StateBackend 支持增量 Checkpoint只保存上次 Checkpoint 以来的变更异步执行状态快照在后台异步写入不阻塞数据处理主流程Exactly-Once结合 Barrier 机制和两阶段提交保证故障恢复后数据不丢不重为什么需要 Checkpoint流处理作业是 7x24 小时运行的运行过程中可能发生各种故障TaskManager 宕机、网络抖动、JobManager 重启等。如果没有状态持久化机制故障后所有中间状态都会丢失需要从数据源重新消费代价巨大。Checkpoint 的作用就是将分布式作业的整体状态在某个时刻凝固下来以便故障时回到那个时刻继续处理。这和数据库的 checkpoint检查点是同一个思想。核心算法Chandy-Lamport 分布式快照Flink 的 Checkpoint 基于 Chandy-Lamport 分布式快照算法。什么是快照快照Snapshot就是把算子的当前状态保存下来。对于 Source 算子来说快照就是当前消费到的 offset如 Kafka 各分区的位点对于窗口算子来说快照就是当前窗口中的数据对于 Sink 算子来说快照就是当前已写入但未提交的事务 ID。把所有这些算子在同一时刻的状态一起保存就叫分布式快照。当作业故障时从快照恢复状态就可以接着快照时的位置继续处理数据不丢不重。Barrier 机制Barrier 的作用很简单标记快照边界。Barrier(n) 之前的数据属于 Checkpoint n之后的数据属于 Checkpoint n1。所有算子以 Barrier 为基准在同一个逻辑时刻完成状态快照保证全局一致性。Barrier 是插入数据流中的特殊标记事件由 Source Task 在收到CheckpointCoordinator运行在 JobManager 上的核心协调者负责触发、收集确认和完成 Checkpoint的触发请求后自行生成并插入到输出数据流中。Barrier 在拓扑中逐级传播CheckpointBarrier (n) — 由 Source 生成并插入 │ ┌───────────────┴───────────────┐ │ Source │ → 快照状态 转发 Barrier(n) 到输出 └───────────────┬───────────────┘ │ 同一个 Barrier(n) 往下游传播 │ ┌───────────────┴───────────────┐ │ 中间算子 │ → 收到 Barrier(n) 后触发快照 │ │ 多输入时按对齐/非对齐策略 │ │ → 转发 Barrier(n) 到下游 └───────────────┬───────────────┘ │ 同一个 Barrier(n) 继续往下游传播 │ ┌───────────────┴───────────────┐ │ Sink │ → 快照 两阶段提交 └───────────────────────────────┘Barrier 的传播和算子的状态快照是异步的Source 生成 Barrier 后立即插入输出流不等自己的异步快照写完中间算子收到 Barrier 后转发给下游也不等自己的快照写完。因此 Barrier 能在 ms 级别贯穿整个拓扑而大状态快照在后台异步执行。整体流程是Barrier 先走、快照后写。Checkpoint 模式Flink 支持两种 Checkpoint 模式通过execution.checkpointing.mode配置默认EXACTLY_ONCE模式配置值行为精确一次EXACTLY_ONCE默认Barrier 到齐后快照保证每条数据只处理一次。阻塞策略详见下方三种模式至少一次AT_LEAST_ONCEBarrier 到齐后快照但期间不阻塞通道恢复时可能有重复。不使用对齐/非对齐策略EXACTLY_ONCE下根据对齐策略又分为三种模式通过额外参数配置策略配置方式行为对齐默认execution.checkpointing.unalignedfalse多输入时 Barrier 阻塞通道等所有通道到齐后快照非对齐execution.checkpointing.unalignedtrueBarrier 到齐即快照不等含 in-flight 数据超时自动切换unalignedtruealigned-checkpoint-timeout 0先尝试对齐超时转非对齐源码中CheckpointOptions.forConfig()的决策逻辑// CheckpointOptions.forConfig() 中的决策链路if(!isExactlyOnceMode)→AT_LEAST_ONCE不用以下策略elseif(isSavepoint())→Savepoint不用对齐/非对齐elseif(!isUnalignedEnabled)→ 对齐默认elseif(timeout0)→ 非对齐无超时else→ 先对齐超时转非对齐第4行即必须开启unalignedtrue 设置aligned-checkpoint-timeoutFlink 才会进入先对齐超时转非对齐的模式。如果只设aligned-checkpoint-timeout但没设unalignedtrueisUnalignedEnabledfalse直接走对齐分支超时参数不生效。AT_LEAST_ONCE不走以上策略没有对齐和非对齐的概念。对齐Alignment当一个算子有多个输入通道时如union、join、ConnectedStreams问题来了不同通道的 Barrier(n) 到达时间可能不同。如果其中一个通道的 Barrier(n) 先到了另一个通道还在消费 Barrier(n) 之前的数据这时如果直接快照快照中就会包含跨代数据——一个通道的 Checkpoint n 数据和另一个通道的 Checkpoint n-1 数据混在一起恢复时状态不一致。对齐就是为了解决这个问题第一个通道的 Barrier(n) 到达后该通道暂停消费blockChannel等所有通道的 Barrier(n) 到齐后才执行快照。这样可以保证快照中所有算子的数据都在同一个 Barrier 边界上。输入1: ...数据...Barrier(n) → 收到 Barrier(n)阻塞通道等其余通道 输入2: ...数据...数据...Barrier(n) → 尚未收到 Barrier(n)继续消费 输入3: ...数据...Barrier(n) → 收到 Barrier(n)阻塞通道 → 三个通道 Barrier(n) 均到齐 → unblockAllChannels()执行快照 → 各通道恢复消费 Barrier(n) 之后的 Buffer → 向下游转发 Barrier(n)对齐的代价阻塞的通道会积压数据产生反压。如果下游处理慢、反压本来就存在对齐会进一步加剧反压甚至导致 Checkpoint 超时失败。非对齐Unaligned Checkpoint非对齐模式下Barrier 到达一个通道后不等其他通道立即触发快照。快照中会把 Barrier 之后还未处理的 in-flight Buffer 也一并保存下来恢复时重放这些数据。输入1: ...数据...Barrier(n)...数据... → 收到 Barrier(n)立即快照含之后的数据 输入2: ...数据...Barrier(n)...数据... → 收到 Barrier(n)立即快照 输入3: ...数据...Barrier(n)...数据... → 收到 Barrier(n)立即快照 → Barrier 到达即触发快照转发 Barrier不等其他通道非对齐的代价快照体积更大包含 in-flight 数据恢复时重放更多数据。对齐 vs 非对齐Flink两种都支持通过配置execution.checkpointing.unaligned: true开启非对齐模式默认关闭。对齐非对齐触发时机所有通道 Barrier 到齐后Barrier 到达即触发转发时机对齐完成后触发快照时Barrier 到达立即转发数据一致性精准一致不包含跨 Barrier 数据精准一致快照包含 in-flight Buffer反压影响对齐阻塞通道 → 反压加剧不等 → 不加剧反压快照大小小只包含状态大含 in-flight Buffer适用场景正常处理、低反压反压严重、大状态作业Flink 1.13 之后还引入了超时自动切换配置execution.checkpointing.aligned-checkpoint-timeout先尝试对齐超过时间未到齐则自动切换为非对齐。在源码中对应AlternatingCollectingBarriers的alignmentTimeout()方法。整体执行流程Checkpoint 的完整生命周期从触发到完成分以下几个阶段。注意各阶段是并行的——Barrier 从 Source 注入后在数据流中向下传播不同算子的快照可以同时进行。CheckpointCoordinatorJobManager 上 │ │ ① 定时触发 ▼ ┌──────────────┐ │ 异步准备 │ ② 确定参与 Task、生成 ID、创建存储目录 └──────┬───────┘ │ │ ③ 向 Source 发送触发请求RPC ▼ ┌──────────────────┐ │ Source 生成 │ ④ Barrier 插入输出流快照自己状态如 offset │ Barrier 并向下游 │ 转发 Barrier 到下游给 Coordinator 发确认 └──────┬───────────┘ │ Barrier 随数据流传播 ▼ ┌──────────────────┐ │ 中间算子 │ ⑤ 收到 Barrier → 对齐多输入时或非对齐 │ │ 快照状态 → 转发 Barrier → 发确认 └──────┬───────────┘ │ ▼ ┌──────────────────┐ │ Sink │ ⑥ 收到 Barrier → 快照状态 → 发确认 │ │ 如实现了两阶段提交在完成后提交事务 └──────┬───────────┘ │ │ ⑦ Coordinator 收齐所有 Task 确认 ▼ ┌──────────────────┐ │ 完成 Checkpoint │ 标记为 Completed存入持久化存储 │ │ 通知所有 Task清理旧快照 └──────────────────┘整体流程的关键点Barrier 是沿着数据流走的快照令牌——Source 生成它每个算子收到后触发自己的快照然后转发给下游。全部确认收齐后这个 Checkpoint 才算成功。StateBackend 与 Checkpoint StorageCheckpoint 涉及两个容易混淆的概念StateBackend决定运行时状态在 TaskManager 里怎么组织。比如用 HashMap 存在堆内存还是用 RocksDB 存在本地磁盘。Checkpoint Storage决定 Checkpoint 时把这些状态快照写到哪。比如写到 HDFS 文件系统还是 JobManager 内存。StateBackend 有 HashMapStateBackend 和 EmbeddedRocksDBStateBackend 两种StateBackend本地存储适用场景HashMapStateBackendJVM 堆内存状态量小 1GB、追求低延迟EmbeddedRocksDBStateBackendRocksDB本地磁盘状态量大GB ~ TB 级别两种 StateBackend 的 Checkpoint 写入方式也不同HashMapStateBackendCheckpoint 时序列化全部状态到 DFS全量写入EmbeddedRocksDBStateBackend默认使用增量 Checkpoint只上传上次 Checkpoint 以来变更的 SST 文件效率远高于全量Checkpoint Storage 则只有JobManagerCheckpointStorage存 JM 内存测试用和FileSystemCheckpointStorage存 HDFS/本地文件系统生产用。生产中一般配state.checkpoints.dirFlink 会自动帮你选 FileSystemCheckpointStorage。Checkpoint vs Savepoint对比维度CheckpointSavepoint触发方式自动周期性触发用户手动触发存储格式StateBackend 原生格式标准 CANONICAL 格式也可指定 NATIVE 原生格式生命周期自动清理根据保留策略用户手动管理长期保留用途故障恢复作业升级、迁移、分支开发性能快RocksDB 增量更快默认慢全量 CANONICAL指定 NATIVE 后同 Checkpoint配置state.checkpoints.dirstate.savepoints.dir跨版本兼容不保证CANONICAL 保证NATIVE 不保证Savepoint 底层复用了 Checkpoint 的完整机制两者在代码层面走的是同一套流程CheckpointCoordinator.triggerCheckpoint区别在于CheckpointType的不同CheckpointType: CheckpointType{nameCheckpoint, ...} SavepointType: SavepointType{nameSavepoint, postCheckpointActionNONE, formatTypeCANONICAL} Suspend Savepoint: SavepointType{nameSuspend Savepoint, postCheckpointActionSUSPEND, ...}更多 Savepoint 的实际操作见 Flink Savepoints 总结。关键概念汇总概念说明Checkpoint分布式快照周期性自动执行用于故障恢复Barrier插入数据流中的特殊标记标识快照边界对齐多输入算子等待所有通道 Barrier 到齐后再快照非对齐Barrier 到达即快照包含 in-flight 数据适用于反压场景PendingCheckpoint正在执行中的 CheckpointCompletedCheckpoint已完成的 Checkpoint可用于恢复StateBackend状态存储后端决定运行时状态在内存还是 RocksDBCheckpoint Storage快照存储决定 Checkpoint 文件写到哪JM 内存或 HDFSExactly-OnceCheckpoint Barrier 两阶段提交保证的语义级别CheckpointCoordinator运行在 JobManager 上的 Checkpoint 协调者下一篇Flink Checkpoint 源码一CheckpointCoordinator 触发链路 将进入源码分析从CheckpointCoordinator的startCheckpointScheduler()和triggerCheckpoint()开始逐层剖析 Checkpoint 的触发链路。
返回列表