Published on

Flink 1.18 Checkpoint Essentials

Authors
  • avatar
    Name
    Charles Chen
    Twitter

做 Flink 实时任务时,State、StateBackend、Checkpoint、CheckpointStorage 这几个概念非常容易混在一起。

尤其是从老版本 Flink 升级过来的开发者,经常会有这样的疑问:

  • State 不就是 Checkpoint 吗?
  • StateBackend 不就是决定 Checkpoint 存哪里吗?
  • EmbeddedRocksDBStateBackend(true) 里的 true 又是什么意思?

本文以 Flink 1.18 为例,把这套机制从概念、选型到常用参数完整梳理一遍。

1. 先记住最重要的一句话

在 Flink 1.18 中:

State 是程序当前正在使用的数据;Checkpoint 是某一时刻 State 的持久化快照。

因此:

State ≠ Checkpoint

Checkpoint = State 在某个时间点的恢复快照

例如一个 Flink 算子维护:

ValueState<Long> count;

当前:

count = 128

此时 128 属于运行中的 State。

12:00 做了一次 Checkpoint:

Checkpoint 100
└── count = 128

程序继续运行,到 12:01:

count = 156

那么现在:

当前 State                  = 156
Checkpoint 100 中保存的 State = 128

如果程序此时发生故障,Flink 可以利用 Checkpoint 100 恢复:

Checkpoint 100
恢复 State
count = 128
从对应的数据位置继续处理

这就是 Checkpoint 最核心的作用:故障恢复。

2. StateBackend 和 CheckpointStorage 到底是什么?

理解了 State 和 Checkpoint,另外两个概念就很简单了。可以记成:

StateBackend
运行中的 State 怎么保存

CheckpointStorage
Checkpoint 最终保存在哪里

整个关系大概是:

                    Flink Job
                 Runtime State
                StateBackend
          ┌─────────────┴─────────────┐
          ▼                           ▼
 HashMapStateBackend       EmbeddedRocksDBStateBackend
      JVM Heap                 RocksDB / Local Disk
          │                           │
          └─────────────┬─────────────┘
                   Checkpoint
               CheckpointStorage
          ┌─────────────┴─────────────┐
          ▼                           ▼
      JobManager                  FileSystem
                                  HDFS/S3/...

所以:

StateBackend 管“活着的 State”,CheckpointStorage 管“备份出来的 State”。

Flink 1.18 最主要的两个 StateBackend 是:

StateBackendState 主要在哪里特点
HashMapStateBackendJVM Heap快,但受内存限制
EmbeddedRocksDBStateBackendRocksDB能支撑超大 State

3.1 HashMapStateBackend

配置:

env.setStateBackend(
    new HashMapStateBackend()
);

它把运行中的 State 放在 JVM Heap 中。例如:

TaskManager JVM

┌───────────────────────────┐
│ Heap                      │
│                           │
│ user1001 → count=100      │
│ user1002 → count=200      │
│ user1003 → count=300      │
│ ...                       │
└───────────────────────────┘

最大的优势是访问 State 快,因为 State 就在 JVM Heap 中。

比较适合:

  • 中小规模 State
  • 低延迟任务
  • 普通聚合
  • 状态生命周期较短
  • JVM 内存比较充裕

但它的问题也很明显:State 大小受到 JVM 内存限制。

如果一个 TaskManager 只有 8 GB Heap,却需要维护几十 GB State,显然不适合 HashMapStateBackend

4. EmbeddedRocksDBStateBackend

配置:

env.setStateBackend(
    new EmbeddedRocksDBStateBackend(true)
);

它使用嵌入式 RocksDB 保存 Keyed State。大致结构:

             TaskManager
              RocksDB
          ┌──────────────┐
          │ MemTable     │
          │ Block Cache  │
          │ SST Files    │
          └──────┬───────┘
          TaskManager 本地磁盘

最大的优势是 State 可以远大于 JVM Heap。例如:

TaskManager Heap
      8 GB

State
     80 GB

HashMapStateBackend 很难处理,而 RocksDB 可以利用本地 SSD / Disk 保存大量 State。

因此非常适合:

  • 大规模 Keyed State
  • 超大窗口
  • 大量用户状态
  • 长 TTL State
  • 实时去重
  • 大型 Join
  • 实时用户画像
  • TB 级 State

代价则是 RocksDB 涉及序列化、JNI、内存与磁盘访问等,相比纯 Heap State 通常有更多开销。

因此:

RocksDB 不是 HashMap 的“高级版”,而是用一定性能成本换取更强的 State 扩展能力。

5. EmbeddedRocksDBStateBackend(true) 中的 true 是什么?

来看:

new EmbeddedRocksDBStateBackend(true)

这里的 true 表示开启 Incremental Checkpoint,也就是增量 Checkpoint。

例如 RocksDB 当前有 100 GB State:

RocksDB

001.sst   20 GB
002.sst   20 GB
003.sst   20 GB
004.sst   20 GB
005.sst   20 GB

总计 ≈ 100 GB

第一次 Checkpoint 需要保存这些 State。之后 RocksDB 又产生:

006.sst
007.sst

如果开启 Incremental Checkpoint,下一次 Checkpoint 可以复用之前已经持久化的 SST 文件,只上传新产生、需要更新的文件。

因此它并不是简单地检查哪些 key/value 改了,而是充分利用 RocksDB 的 SST 文件机制进行增量快照。

对于几十 GB、几百 GB 乃至 TB 级 State,这个能力非常重要。

6. CheckpointStorage 有哪些?

Flink 1.18 主要关注两个:

CheckpointStorage
├── JobManagerCheckpointStorage
└── FileSystemCheckpointStorage

6.1 JobManagerCheckpointStorage

Checkpoint 数据主要由 JobManager 内存承担:

env.getCheckpointConfig()
   .setCheckpointStorage(
       new JobManagerCheckpointStorage()
   );

适合:

  • 本地开发
  • 测试
  • Demo
  • 极小 State

大型生产任务一般不会把它作为首选。

6.2 FileSystemCheckpointStorage

生产环境更加常见:

env.getCheckpointConfig()
   .setCheckpointStorage(
       new FileSystemCheckpointStorage(
           "hdfs:///flink/checkpoints/my-job"
       )
   );

Checkpoint 可以保存到:

  • HDFS
  • S3
  • GCS
  • 其他 Flink 支持的持久化文件系统

这也是生产环境最常见的 CheckpointStorage。

7. 生产环境 StateBackend 怎么选?

可以简单理解:

                State 大吗?
            ┌───────┴───────┐
            │               │
           不大              很大
            │               │
            ▼               ▼
         HashMap          RocksDB
       JVM Heap 能否
       稳定装下?
        能 → HashMap
        悬 → RocksDB

实际应该关注的不是整个 Job 的 State,而是:

总 State 大小
+ 并行度
+ 每个 TaskManager / Subtask 承担多少 State
+ State 是否持续增长
+ 可用内存

例如:

总 State = 100 GB
并行度 = 100

平均 ≈ 1 GB / subtask

未必必须使用 RocksDB。

但如果:

总 State = 100 GB
并行度 = 5

平均 ≈ 20 GB / subtask

就应该认真考虑 RocksDB。

8. 最常见的两种生产组合

8.1 中小 State

env.setStateBackend(
    new HashMapStateBackend()
);

env.getCheckpointConfig().setCheckpointStorage(
    new FileSystemCheckpointStorage(
        "hdfs:///flink/checkpoints/my-job"
    )
);

工作过程:

运行:

State
JVM Heap


Checkpoint:
JVM State
Snapshot
HDFS

8.2 大 State

env.setStateBackend(
    new EmbeddedRocksDBStateBackend(true)
);

env.getCheckpointConfig().setCheckpointStorage(
    new FileSystemCheckpointStorage(
        "hdfs:///flink/checkpoints/my-job"
    )
);

工作过程:

运行:

State
RocksDB
TaskManager Local SSD

Checkpoint:

RocksDB Snapshot
HDFS

这也是理解 StateBackend 和 CheckpointStorage 最直观的方式。

9. Checkpoint 最重要的参数

理解完架构,再来看真正影响生产稳定性的参数。

9.1 Checkpoint Interval

开启 Checkpoint:

env.enableCheckpointing(
    60_000,
    CheckpointingMode.EXACTLY_ONCE
);

这里 60_000 ms = 60 秒,表示 Checkpoint 的触发间隔。它控制的是:希望多久触发一次 Checkpoint。

注意,它不意味着每 60 秒必然产生一个成功的 Checkpoint。如果前一个 Checkpoint 很慢,还会受到并发数、最小暂停时间等参数限制。

10. CheckpointingMode

一般配置:

env.getCheckpointConfig()
   .setCheckpointingMode(
       CheckpointingMode.EXACTLY_ONCE
   );

主要有:

  • EXACTLY_ONCE
  • AT_LEAST_ONCE

绝大多数需要一致性的生产 Stateful Streaming Job,优先使用 EXACTLY_ONCE

需要注意,Checkpoint 的 Exactly Once 与端到端 Exactly Once 不是完全等价的概念。

如果是:

Kafka → Flink → MySQL

最终能否做到端到端 Exactly Once,还取决于 Source、Sink 以及外部系统是否配合相应的一致性机制。

11. Checkpoint Timeout

config.setCheckpointTimeout(
    10 * 60 * 1000
);

表示一次 Checkpoint 最长允许执行多久。

例如:

12:00
CKPT 100 START
   │ 10 min
12:10
仍未完成
TIMEOUT

所以:

Checkpoint Interval
= 多久尝试做一次

Checkpoint Timeout
= 一次最多允许做多久

两者完全不同。

12. MinPauseBetweenCheckpoints

配置:

config.setMinPauseBetweenCheckpoints(
    30_000
);

表示上一个 Checkpoint 完成之后,至少休息多久才能开始下一个。

例如:

CKPT 100
   │ 执行 40s
完成
   │ 休息至少 30s
CKPT 101

为什么需要这个参数?因为 Checkpoint 会产生:

  • CPU
  • 序列化
  • 本地磁盘 IO
  • 网络 IO
  • HDFS / S3 IO

如果任务:

CKPT
刚完成
马上 CKPT
刚完成
又 CKPT

整个系统可能长期处于 Checkpoint 压力中。因此这个参数对于大状态任务尤其有价值。

13. MaxConcurrentCheckpoints

config.setMaxConcurrentCheckpoints(1);

表示最多允许几个 Checkpoint 同时执行。

设置为 1 意味着:

CKPT 100
███████████████

               CKPT 101
               █████████████

而不会大量重叠:

CKPT 100
████████████████

     CKPT 101
     ███████████████

生产环境设置为 1 是一个非常常见的起点。

14. Interval、MinPause 和 Concurrent 要一起看

这是三个最容易混淆的参数。

假设:

env.enableCheckpointing(60_000);

config.setMinPauseBetweenCheckpoints(30_000);

config.setMaxConcurrentCheckpoints(1);

分别控制:

参数含义
Interval希望多久触发一次
MinPause上次完成后至少休息多久
MaxConcurrent最多允许几个同时执行

因此 Checkpoint 调度并不是简单的“每 60 秒无脑执行一次”,而是这些条件共同作用的结果。

15. TolerableCheckpointFailureNumber

配置:

config.setTolerableCheckpointFailureNumber(3);

表示可以容忍一定数量的 Checkpoint Failure。例如:

CKPT 100 SUCCESS

CKPT 101 FAILED
CKPT 102 FAILED
CKPT 103 FAILED
...

这个参数的意义在于:HDFS、S3、网络等偶尔出现瞬时异常时,不一定需要立即把整个 Job 干掉。

但也不能无限容忍。否则可能出现非常危险的情况:

Job Status:

RUNNING ✓


最近一次成功 Checkpoint:

8 小时前 ✗

看起来 Job 正常运行,实际上容灾能力已经严重下降。

16. Externalized Checkpoint

生产环境另一个非常重要的配置:

config.setExternalizedCheckpointCleanup(
    ExternalizedCheckpointCleanup
        .RETAIN_ON_CANCELLATION
);

意思是:Job 被 Cancel 后,保留 Externalized Checkpoint。

例如:

Job Running

CKPT 100
CKPT 101
CKPT 102
   Cancel
Job 停止

但是 CKPT 仍然保留

以后可以利用它恢复。这种配置对于生产任务非常实用,但需要注意:保留下来的数据也需要进行生命周期和存储空间管理。

17. FileStateSizeThreshold

FileSystemCheckpointStorage 还有一个容易忽略的参数:

new FileSystemCheckpointStorage(
    checkpointPath,
    fileStateSizeThreshold
);

它控制:小块 State 是 inline 到 Checkpoint metadata,还是单独生成文件。

如果存在大量非常小的 State:

20 Byte
50 Byte
100 Byte
...

全部生成独立文件可能导致:

checkpoint/

state-000001
state-000002
state-000003
state-000004
state-000005
...

产生大量小文件。Threshold 可以用于控制小 State 的存储方式。

一般来说,没有明确的小文件问题时,不建议为了“优化”而随意调整这个参数。

18. Unaligned Checkpoint

这是处理反压问题时非常重要的机制:

config.enableUnalignedCheckpoints();

正常 Aligned Checkpoint 需要处理 Checkpoint Barrier。如果任务出现严重 Backpressure:

Source
Operator A
   │ ████████████
   │ 严重 Backpressure
Operator B

Barrier 可能长时间无法完成对齐。

于是:

正常:

Checkpoint = 10 秒


严重反压:

Checkpoint = 1 分钟
             3 分钟
             5 分钟
             Timeout

Unaligned Checkpoint 的重要思想是:不再等待所有 in-flight 数据完成传统意义上的对齐,而是把相关网络中的 in-flight data 也纳入快照。

优点:严重反压时,Checkpoint 时间更加可控。

代价之一:Checkpoint Size 可能增大。

因此:

Unaligned Checkpoint 主要解决的是反压导致的 Barrier Alignment 问题,而不是“让所有 Checkpoint 都变快”的万能开关。

19. AlignedCheckpointTimeout

还可以采用一种折中方案:

config.setAlignedCheckpointTimeout(...);

思路:

开始 CKPT
先尝试 Aligned
   │ 等待一定时间
Alignment 还是很慢
切换 Unaligned

这特别适合平时没有严重反压、偶尔出现反压尖峰的任务。

开发环境经常看到:

new FileSystemCheckpointStorage(
    "file:///tmp/flink/"
);

本地测试没有问题,但生产集群需要谨慎。

因为 file:// 通常代表节点本地文件系统。例如:

TaskManager 1
/local/tmp/flink

TaskManager 2
/local/tmp/flink

这两个目录不是同一个共享目录。如果机器、容器或者 Pod 消失,本地文件也可能一起消失。

Checkpoint 最重要的目的恰恰是:当前执行节点出问题之后还能恢复。

所以生产环境通常应该考虑:

  • HDFS
  • S3
  • GCS
  • 其他可靠、持久、集群可访问的存储

把前面的内容组合起来,一个 RocksDB Stateful Streaming Job 可以从下面这样的配置开始:

// ==========================================
// 1. 每 60 秒尝试进行一次 Checkpoint
// ==========================================
env.enableCheckpointing(
    60_000,
    CheckpointingMode.EXACTLY_ONCE
);

CheckpointConfig ckpt =
        env.getCheckpointConfig();


// ==========================================
// 2. 单次 Checkpoint 最长 10 分钟
// ==========================================
ckpt.setCheckpointTimeout(
    10 * 60 * 1000
);


// ==========================================
// 3. 两个 Checkpoint 至少暂停 30 秒
// ==========================================
ckpt.setMinPauseBetweenCheckpoints(
    30_000
);


// ==========================================
// 4. 最多同时执行一个 Checkpoint
// ==========================================
ckpt.setMaxConcurrentCheckpoints(1);


// ==========================================
// 5. Checkpoint Failure 容忍策略
// ==========================================
ckpt.setTolerableCheckpointFailureNumber(3);


// ==========================================
// 6. Cancel Job 后保留 Externalized CKPT
// ==========================================
ckpt.setExternalizedCheckpointCleanup(
    ExternalizedCheckpointCleanup
        .RETAIN_ON_CANCELLATION
);


// ==========================================
// 7. RocksDB StateBackend
//    true = Incremental Checkpoint
// ==========================================
env.setStateBackend(
    new EmbeddedRocksDBStateBackend(true)
);


// ==========================================
// 8. Checkpoint 持久化到 HDFS
// ==========================================
ckpt.setCheckpointStorage(
    new FileSystemCheckpointStorage(
        "hdfs:///flink/checkpoints/my-job"
    )
);

注意:这是一套便于理解的起始模板,而不是所有生产任务都应该照抄的“最佳参数”。

真正的参数应该根据 State 大小、Checkpoint Duration、反压、磁盘 IO、网络以及 HDFS/S3 性能进行调整。

22. 最值得记住的 Checkpoint 参数

最后浓缩成一张表:

参数作用常见思路
enableCheckpointing(interval)多久尝试触发一次 CKPT根据恢复目标和开销设置
CheckpointingMode一致性语义通常 EXACTLY_ONCE
checkpointTimeoutCKPT 最长允许多久应明显高于正常 CKPT 时间
minPauseBetweenCheckpoints两次 CKPT 之间最小休息时间防止持续 CKPT
maxConcurrentCheckpoints最多几个 CKPT 并发常从 1 开始
tolerableCheckpointFailureNumber容忍多少 CKPT Failure根据 SLA 设置
RETAIN_ON_CANCELLATIONCancel 后是否保留生产常考虑
RocksDB(true)Incremental CKPT大 State 常用
CheckpointStorageCKPT 保存位置生产用可靠持久存储
fileStateSizeThreshold小 State 是否 inline通常保持默认
UnalignedCheckpoint缓解反压导致的 Alignment 问题有对应问题再使用
AlignedCheckpointTimeoutAligned 多久后转 Unaligned适合间歇性反压

23. 出现 Checkpoint 慢,应该看什么?

这是生产排障中特别重要的一点:Checkpoint 慢,不要第一反应就是把 checkpointTimeout 从 10 分钟改成 30 分钟。

Timeout 只是允许它慢得更久,并没有解决为什么慢。

应该进入 Flink Web UI 观察 Checkpoint 指标,重点关注:

  • Checkpoint Duration
  • Checkpoint Size
  • Start Delay
  • Alignment Duration
  • Sync Duration
  • Async Duration

然后大致按照这个思路排查:

Checkpoint 慢
    ├── Alignment 很慢?
    │       │
    │       └── 看 Backpressure
    │           看 Unaligned CKPT
    ├── Async Snapshot 很慢?
    │       │
    │       ├── RocksDB?
    │       ├── 本地磁盘?
    │       └── HDFS/S3?
    ├── Checkpoint Size 特别大?
    │       │
    │       ├── State 是否暴涨?
    │       ├── TTL 是否合理?
    │       └── Incremental CKPT?
    └── Start Delay 很大?
            └── Barrier 传播 / 上游反压

这比简单增加 Timeout 更有意义。

总结

Flink Checkpoint 体系看起来参数很多,但核心其实可以浓缩成四个概念:

State
    │ 当前运行的数据
StateBackend
    │ 决定 State 怎么存
    ├── HashMap
    └── RocksDB
            │ Snapshot
        Checkpoint
            │ 持久化
    CheckpointStorage
            ├── JobManager
            └── FileSystem
                 ├── HDFS
                 ├── S3
                 └── ...

选型上则可以先记住:

中小 State + 追求访问性能
HashMapStateBackend
        +
FileSystemCheckpointStorage


大 State / State 可能超过内存
EmbeddedRocksDBStateBackend(true)
        +
FileSystemCheckpointStorage

最后再记住六个最核心的调优问题:

Interval
→ 多久做一次?

Timeout
→ 一次最多做多久?

MinPause
→ 做完至少休息多久?

MaxConcurrent
→ 最多同时做几个?

Incremental
→ RocksDB 是否增量保存?

Unaligned
→ 严重反压时 Barrier 怎么处理?

掌握这套关系之后,Flink Checkpoint 的配置就不再是“背参数”,而是在控制三件事情:多久备份一次、备份如何执行、备份最终放在哪里。

官方文档