夜雨聆风学习资料网

ARTICLE · 1127373

Flink 1.20 源码解析:一次 Checkpoint 到底是怎么被触发的?

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 机制的重要协调者。

从源码可以看到,它负责的事情很多,包括:

  • triggerCheckpoint
  • triggerSavepoint
  • startCheckpointScheduler
  • 接收 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        ↓TaskExecutor

4. Checkpoint Request 和 Checkpoint Barrier 是两个不同概念

Request:

JobManager → TaskManager

Barrier:

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 的世界。

相关学习资料