ARTICLE · 1127373
Flink 1.20 源码解析:一次 Checkpoint 到底是怎么被触发的?
本文基于 Flink 1.20 源码,从 CheckpointCoordinator 出发,一路跟踪到 StreamTask,看看一次 Checkpoint 到底是如何从 JobManager 侧的一个“请求”,变成 Task 侧真正执行的 Checkpoint。这一篇暂时不深入 State Backend、RocksDB 和 Incremental Snapshot,先把 Checkpoint 的触发链路 搞清楚。
一、先抛一个问题:Checkpoint 到底是谁触发的?
刚开始看 Flink Checkpoint 源码时,很容易产生一个错觉:
“Checkpoint 不就是
CheckpointCoordinator触发的吗?”
这句话不能说错,但还不够准确。
因为一个完整的 Checkpoint 至少涉及三个不同层次:
CheckpointCoordinator ↓Task ↓StreamTask它们分别解决不同的问题:
CheckpointCoordinator:决定什么时候发起一次 Checkpoint,并负责协调整个 Checkpoint。Task:接收来自 TaskManager 侧的 Checkpoint 请求。StreamTask:真正进入流处理执行逻辑,创建并传播CheckpointBarrier,随后执行 State Snapshot。
所以这篇文章真正要回答的是:
CheckpointCoordinator 发起的 Checkpoint 请求,究竟是怎么一路到达 StreamTask 的?
二、先看 CheckpointCoordinator:它到底负责什么?
源码中的CheckpointCoordinator是整个 Checkpoint 机制的重要协调者。
从源码可以看到,它负责的事情很多,包括:
triggerCheckpointtriggerSavepointstartCheckpointScheduler接收 Task 的 Checkpoint ACK 接收 Checkpoint Decline 管理 Pending Checkpoint Restore 相关逻辑 Checkpoint Metrics 等
也就是说,CheckpointCoordinator并不是负责具体执行 State Snapshot 的地方。
它更像是一个:
Checkpoint 的总协调器。
三、周期性 Checkpoint 是怎么启动的?
先从最容易理解的地方开始:
startCheckpointScheduler()这个方法负责启动周期性 Checkpoint 调度。
其中首先拿到:
long baseInterval = chkConfig.getCheckpointInterval();也就是用户配置的 Checkpoint Interval。
然后还会结合:
minPauseBetweenCheckpoints来计算实际的调度间隔。
之后会进入:
scheduleTriggerWithDelay(...)最终由定时任务触发:
triggerCheckpoint(...)所以最简单的理解就是:
Checkpoint Interval ↓Checkpoint Scheduler ↓定时触发 ↓triggerCheckpoint()这里有一个很重要的认识:
Checkpoint Interval 只是“多久尝试触发一次”,并不意味着整个 Checkpoint 恰好每隔这个时间完成一次。
Checkpoint 本身可能因为:
Task 响应慢 State Snapshot 慢 Checkpoint 排队 Checkpoint 超时 其他协调条件
而持续更长时间,所以:Trigger和Complete是两个完全不同的时间点。
四、Checkpoint 真正开始准备:startTriggeringCheckpoint()
周期调度最终会进入 CheckpointCoordinator 的 Checkpoint 触发流程。
这里最值得关注的是:
startTriggeringCheckpoint(...)这个方法并不是简单地:
“通知所有 Task 做 Checkpoint”相反,它首先是在准备这一次 Checkpoint。
整个过程可以概括成:
startTriggeringCheckpoint() │ ├── 检查当前是否允许触发 │ ├── 生成 Checkpoint Plan │ ├── 生成 Checkpoint ID │ ├── 创建 PendingCheckpoint │ ├── 初始化 Checkpoint Storage Location │ ├── 处理 Operator Coordinator State │ ├── 处理 Master State │ └── triggerCheckpointRequest()这里出现了一个非常重要的对象:PendingCheckpoint
五、PendingCheckpoint 到底是什么?
很多人第一次看到:PendingCheckpoint
会把它理解成:“Checkpoint ID+Task 列表”
实际上它承担的角色更重要。
它代表的是:
一个已经被触发、但还没有完成的 Checkpoint。
它里面会保存很多和这个 Checkpoint 生命周期相关的信息,例如:
Job Checkpoint ID Timestamp Checkpoint Plan Operator Coordinator 信息 Master Hook 信息 Checkpoint Properties Completion Promise Checkpoint Metrics Master Trigger Completion Promise
因此可以把它理解成:
PendingCheckpoint │ ├── 这次 CP 是谁? ├── CP ID 是多少? ├── 哪些 Task 要参与? ├── CP 存储在哪里? ├── Master State 怎么处理? ├── Coordinator State 怎么处理? ├── Metrics 怎么记录? └── 最终什么时候 Complete?所以:
PendingCheckpoint 是 CheckpointCoordinator 对“这一次正在进行中的 Checkpoint”的状态管理对象。
它不是最终 State,也不是 Snapshot 本身。
六、什么时候才真正通知 Task?
这是阅读源码时非常容易迷糊的地方。
前面我们一直在:startTriggeringCheckpoint()里面看到:
Checkpoint Plan Checkpoint ID PendingCheckpoint Storage Location Master State
很容易产生一个疑问:
“到底在哪里真正让 Task 开始做 Checkpoint?”
答案就在:
triggerCheckpointRequest() ↓triggerTasks()七、triggerCheckpointRequest:从协调阶段进入 Task 触发阶段
源码首先进入:
CheckpointCoordinator#triggerCheckpointRequest(...)如果 PendingCheckpoint 已经被 dispose,就走失败逻辑。
否则:
triggerTasks(request, timestamp, checkpoint)这一步非常关键。
因为从这里开始,CheckpointCoordinator 不再只是“准备 Checkpoint”,而是:
开始向参与 Checkpoint 的 Task 发起请求。
八、triggerTasks:真正遍历需要触发 Checkpoint 的 Task
核心代码非常直观:
for (Execution execution : checkpoint.getCheckpointPlan().getTasksToTrigger()) { acks.add( execution.triggerCheckpoint( checkpointId, timestamp, checkpointOptions));}这里终于可以回答一个核心问题:
CheckpointCoordinator 是怎么通知所有 Task 的?
答案:
CheckpointPlan ↓tasksToTrigger ↓Execution ↓Execution.triggerCheckpoint()同时,这里还会根据配置生成:CheckpointOptions
其中包含:
Snapshot 类型 Checkpoint Storage Location Exactly-Once 配置 Unaligned Checkpoint 配置 Alignment Timeout
所以到这里可以画出:
CheckpointCoordinator ↓PendingCheckpoint ↓CheckpointPlan ↓tasksToTrigger ↓Execution九、Execution 并不直接执行 Checkpoint
进入:Execution#triggerCheckpoint(...)
这里并没有看到 State Snapshot,它做的事情其实非常简单:
final LogicalSlot slot = assignedResource;final TaskManagerGateway taskManagerGateway = slot.getTaskManagerGateway();return taskManagerGateway.triggerCheckpoint( attemptId, jobId, checkpointId, timestamp, checkpointOptions);所以:
Execution 更像是连接 ExecutionGraph 与实际 TaskManager 上 Task 的桥梁。
调用链变成:
CheckpointCoordinator ↓Execution ↓LogicalSlot ↓TaskManagerGateway十、TaskManagerGateway:进入 TaskManager 侧
接下来经过:
RpcTaskManagerGateway#triggerCheckpoint()然后把请求继续转给:
TaskExecutor#triggerCheckpoint()到这里,Checkpoint 请求已经从 JobManager 侧进入 TaskManager 侧。
十一、TaskExecutor 如何找到具体的 Task?
这是整个链路中非常关键的一步,TaskExecutor会根据:ExecutionAttemptID 从 taskSlotTable 找到真正对应的 Task,也就是:
final Task task = taskSlotTable.getTask(executionAttemptID);然后:
task.triggerCheckpointBarrier( checkpointId, checkpointTimestamp, checkpointOptions);于是:
TaskExecutor ↓ExecutionAttemptID ↓Task到这里,一个 JobManager 侧的 Checkpoint 请求,终于落到了 TaskManager 上实际运行的 Task。
十二、Task 也不是最终执行 Checkpoint 的地方
进入Task#triggerCheckpointBarrier() 这里有一个非常重要的类型判断:
if (executionState == ExecutionState.RUNNING) { checkState( invokable instanceof CheckpointableTask,"invokable is not checkpointable");}然后:
((CheckpointableTask) invokable) .triggerCheckpointAsync( checkpointMetaData, checkpointOptions);这里的invokable最终会进入:StreamTask
所以又多了一层:
Task ↓TaskInvokable ↓CheckpointableTask ↓StreamTask这也解释了为什么我们不能简单地说:
“Task 执行 Checkpoint。”
更准确的说法是:
Task 接收 Checkpoint 请求,并把请求交给具体的 CheckpointableTask;对于流任务来说,最终进入 StreamTask。
十三、StreamTask:真正进入 Checkpoint 执行逻辑
进入:StreamTask#triggerCheckpointAsync(...) 这里有一个非常值得注意的设计:
mainMailboxExecutor.execute( () -> { ... triggerCheckpointAsyncInMailbox(...); });也就是说,Checkpoint 请求进入:Mailbox,然后由 StreamTask 的主执行逻辑继续处理,于是:
Task ↓StreamTask ↓Mailbox ↓triggerCheckpointAsyncInMailbox()到这里,我们终于进入真正的 StreamTask Checkpoint 处理过程。
十四、triggerCheckpointAsyncInMailbox开始执行一次 Task Checkpoint
这个方法首先初始化一些 Checkpoint Metrics:
CheckpointMetricsBuilder checkpointMetrics =new CheckpointMetricsBuilder() .setAlignmentDurationNanos(0L) .setBytesProcessedDuringAlignment(0L) .setCheckpointStartDelayNanos(...);然后:
subtaskCheckpointCoordinator.initInputsCheckpoint( checkpointId, checkpointOptions);最后:
boolean success = performCheckpoint( checkpointMetaData, checkpointOptions, checkpointMetrics);所以这里又出现一个非常重要的对象:SubtaskCheckpointCoordinator,它负责的是:
单个 Subtask 内部的 Checkpoint 执行协调。
这和前面的:CheckpointCoordinator 不是一个层次,可以简单区分:
CheckpointCoordinator ↓Job 级别协调整个 CheckpointSubtaskCheckpointCoordinator ↓Task/Subtask 级别协调这个 Task 的 Checkpoint十五、真正看到 CheckpointBarrier 的地方
继续进入:SubtaskCheckpointCoordinatorImpl#checkpointState()
这里终于看到了整个 Checkpoint 机制最核心的一个东西:CheckpointBarrier
源码把流程明确分成了几个 Step。
Step 1:Barrier 发送前准备
operatorChain.prepareSnapshotPreBarrier( metadata.getCheckpointId());源码注释明确说明:
Prepare the checkpoint,allow operators to do some pre-barrier work.
也就是说,在 Barrier 发出去之前,Operator 可以进行必要的准备。
Step 2:创建 CheckpointBarrier
接下来:
CheckpointBarrier checkpointBarrier =new CheckpointBarrier( metadata.getCheckpointId(), metadata.getTimestamp(), options);然后:
operatorChain.broadcastEvent( checkpointBarrier, options.isUnalignedCheckpoint());这两行非常重要,因为它回答了两个问题:
CheckpointBarrier 在哪里创建?
SubtaskCheckpointCoordinatorImpl#checkpointState()Barrier 怎么进入 Operator Chain?
operatorChain.broadcastEvent(...)因此:
StreamTask ↓SubtaskCheckpointCoordinator ↓CheckpointBarrier ↓OperatorChain十六、一个非常容易产生的误区:RPC 请求 ≠ Checkpoint Barrier
现在我们可以明确区分两个东西:
Checkpoint Request 发生在:
CheckpointCoordinator ↓Execution ↓TaskManager ↓Task ↓StreamTask它是控制面上的请求。
Checkpoint Barrier发生在:
StreamTask ↓CheckpointBarrier ↓OperatorChain ↓数据流处理链它是数据流执行过程中的控制事件。
所以整个过程实际上是:
控制面CheckpointCoordinator │ │ triggerCheckpoint ↓ Task ↓ StreamTask │ │ 创建 Barrier ↓ CheckpointBarrier │ ↓ 数据流 OperatorChain这两个概念一定不要混。
十七、为什么 Barrier 要尽快发送?
源码在checkpointState()的注释里有一句非常重要的话:
我们通常会尽可能早地发出 Checkpoint Barrier,以减少对下游 Checkpoint Alignment 的影响。
所以源码的执行顺序非常值得注意:
Step 1prepareSnapshotPreBarrier() ↓Step 2broadcast CheckpointBarrier ↓Step 3register alignment timer ↓Step 4处理 Channel State ↓Step 5takeSnapshotSync()注意:Barrier 是在 Snapshot 之前发送的,不是:
Snapshot ↓Snapshot 完成 ↓发送 Barrier而是:
准备 ↓发送 Barrier ↓Alignment / Channel State ↓Snapshot这是理解 Flink Checkpoint 非常重要的一个转折点。
十八、那非对齐 Checkpoint 又是什么?
这里需要特别谨慎,不能简单理解成:
“非对齐 Checkpoint=把 inflight 数据写下来。”
从源码来看,Checkpoint 是否需要 Channel State,是通过:
options.needsChannelState()判断的。如果需要:
channelStateWriter.finishOutput( metadata.getCheckpointId());Channel State 的核心作用之一,就是处理 Checkpoint 时仍然存在于输入/输出通道中的数据,也就是通常所说的in-flight data。因此更准确的理解是:
Unaligned Checkpoint ↓需要考虑 Channel State ↓记录 Checkpoint 时刻通道中的 in-flight data而不是简单地把“Unaligned”和“写 inflight”完全画等号。
十九、到这里,一次 Checkpoint 的触发链终于完整了
把前面的源码全部串起来:
Checkpoint Scheduler ↓CheckpointCoordinator ↓startTriggeringCheckpoint() │ ├── Checkpoint Plan ├── Checkpoint ID ├── PendingCheckpoint ├── Storage Location └── Master / Coordinator State ↓triggerCheckpointRequest() ↓triggerTasks() ↓Execution ↓TaskManagerGateway ↓RpcTaskManagerGateway ↓TaskExecutor ↓Task ↓StreamTask ↓Mailbox ↓SubtaskCheckpointCoordinator ↓checkpointState() │ ├── prepareSnapshotPreBarrier() │ ├── 创建 CheckpointBarrier │ ├── OperatorChain.broadcastEvent() │ ├── Alignment / Channel State │ └── takeSnapshotSync()到这里,我们暂时停在:takeSnapshotSync(),也就是说:
本文我们只解决“Checkpoint 是怎么被触发并进入 StreamTask”的问题。
至于:
takeSnapshotSync()到底怎么让每个 Operator 做 Snapshot?
那就是后续文章的主题。
二十、源码阅读过程中,我认为最值得记住的 5 个结论
1. CheckpointCoordinator 不负责真正的 State Snapshot
它负责的是:
触发+协调+跟踪+完成+失败处理真正的 State Snapshot 会进入 Task/Operator 侧。
2. PendingCheckpoint 是一次“正在进行中的 Checkpoint”
它不是 Snapshot 本身,它负责保存这次 Checkpoint 在 Coordinator 侧的各种状态。
3. Execution 是连接调度侧和 TaskManager 的桥梁
调用关系:
CheckpointCoordinator ↓Execution ↓TaskManagerGateway ↓TaskExecutor4. Checkpoint Request 和 Checkpoint Barrier 是两个不同概念
Request:
JobManager → TaskManagerBarrier:
StreamTask → OperatorChain前者是控制请求,后者是数据流中的控制事件。
5. Barrier 在 Snapshot 之前发送
这是本篇最值得记住的一条源码结论:
prepare ↓Barrier ↓Alignment / Channel State ↓Snapshot而不是:
Snapshot ↓Barrier二十一、最后用一句话总结这篇文章
用一句话概括 Flink 1.20 一次 Checkpoint 的“启动过程”:
CheckpointCoordinator 决定并准备一次 Checkpoint,通过 Execution 和 TaskManagerGateway 把请求送到 Task,再由 StreamTask 进入 Mailbox,在 SubtaskCheckpointCoordinator 中创建 CheckpointBarrier 并广播到 OperatorChain,之后才进入真正的 State Snapshot 阶段。
而下一篇,我们继续沿着最后这一行:
takeSnapshotSync()往下走。
看看一个看似简单的:
“做一次 Checkpoint”最终是怎么变成:
Operator ↓State ↓State Snapshot ↓SnapshotResult ↓StateHandle以及最终我们最关心的:
HDFS 上到底生成了什么文件?这才是真正进入 Flink State Snapshot 的世界。