Skip to content

Flink checkpoint 间隔怎么定:写 Iceberg 为什么是 5~15 分钟,秒级为什么不划算

整体建议:

  • 当下游 Sink 为 iceberg 时, checkpoint 间隔设置为 5-15 min
  • 当下游 Sink 未使用事务写时(如 Kafka, rocketmq ),可以适当延长 checkpoint (与数据新鲜度无关)
  • 当状态量比较大时 (包含去重、聚合、join 算子),checkpoint 设置为至少 3min 以上
  • 所有作业都不建议设置低于分钟级别的 Checkpoint 间隔
  • 提升单并发的处理速度来减少并发度, 用于减少 Checkpoint 文件数量

Checkpoint 作用和概念介绍

事务 sink 与非事务 sink 的数据可见时刻不同

Checkpoint 作用

Checkpoint 是流处理系统中保证故障恢复和数据一致性的重要机制。通常情况下, Flink 的每个算子都会维护自身的状态 (source operator 会保存消费的位点, window operator 保存中间的计算结果, sink operator 保存写入的事务 id) , 因此如果作业出现重启, Flink 作业直接从上一个位点重新消费即可。

客观上: Good:

  1. Checkpoint 间隔越小,作业重启后需要重新消费的数据量就越少,作业恢复的速度越快
  2. 在事务提交时,Checkpoint 间隔越短,链路的数据新鲜度越高

Bad:

  1. Checkpoint 间隔越短,Checkpoint 本身产生的状态文件数量会越多
  2. 在下游是写文件系统时 (HDFS/Iceberg)Sink 产生的小文件越多

对于 checkpoint 的作用, 可以参考文章 https://flink.apache.org/2018/02/28/an-overview-of-end-to-end-exactly-once-processing-in-apache-flink-with-apache-kafka-too/

Checkpoint 间隔和数据新鲜度

Checkpoint 间隔直接影响下游数据的可见性,尤其在流处理系统中数据传递的实时性至关重要。根据下游 sink 是否支持事务写入,数据的可见性有不同的表现。

2.1 事务写入 Exactly Once(Iceberg、Kafka 事务 等)

下游数据的可见性和 checkpoint 间隔相关,** checkpoint 间隔越小, 数据越早可见**。

  • (3) Sink 接受到上游传递过来的 Checkpoint barrier 开始做 Checkpoint

  • (4) Sink 算子会进行 precommit 结束本次事务的写入, 并将事务 ID 记录到 state 中

  • (1) Checkpoint 完成之后, JobManager 通知所有的 operator Checkpoint 已经完成

  • (2) Sink Operator 接受到通知, 进行 commit 操作, 此时数据可见

2.2 非事务写入 At-Least Once ( Kafka, jdbc, rocketmq)

下游的数据的可见性和单次写的 batch 数或者刷新间隔相关。

  • (3) Sink 接收到 checkpoint barrier 进入 checkpoint 状态
  • (4) Sink 算子执行 flush 操作, 将 buffer 中的数据全部写入 kafka ,Sink operator 不保存任何状态

Checkpoint 间隔过短产生的问题

3.1 状态文件过多,对 HDFS 的 namenode 压力较大

Flink 状态文件的特点:

  • Flink 会频繁的写入和删除状态文件, 只有在状态恢复时才会读取, 对 name node 的压力较大。
  • Checkpoint 越频繁, 单位时间内产生的状态文件越多 (线性增长)
  • 每个算子都保存独立的 checkpoint 文件,算子越多, 并发越多, checkpoint 文件数越多 (线性增长)

状态比较大的算子

  • 主要: join, 去重 (缓存历史所有数据)
  • 其次: group (缓存当前的状态文件)

3.2 下游的文件系统可能产生较多的小文件

当下游写入 HDFS/Iceberg 这类文件系统时:

  • 每次事务操作, 每个 subtask 至少产生一个文件 (一个并发写多个 bucket 则会更多)
  • 并发越大,产生的小文件越多 (线性增长)
  • Checkpoint 间隔越短, 产生的小文件越多 (线性增长)

优化思路

减少 checkpoint 文件数量:

  • 尽量减少并发数,提高单并发处理速度,减少整体的文件数量
    • 当作业瓶颈在状态访问上时,可以考虑增大内存,来提升状态访问速度
  • 对于状态量不大的作业,可以考虑使用 FileSystem Statebackend 减少小文件数量
  • 特定场景 join 和 去重优化
    • Regular join 转换为使用 paimon partial update
    • 去重节点下放至 paimon (使用 paimon changelog-mode 产生完整 CDC 下游作业无需再次去重)

减少湖存储的小文件数量:

  • 对于新鲜度要求比较高的作业,中间数据使用消息队列 (Kafka/rocketmq)
  • 对实时性比要求比较高的带主键表,更换为 paimon
  • Kafka2Iceberg 小文件优化
    • 增大 checkpoint 间隔,减少文件数量
    • 在二级分区能保证均匀 hash 场景, 将 distribute mode 改成 hash ,将具有相同 bucket 的数据 hash 至同一个 subtask 减少小文件

问题1: 为什么不建议设置秒级的 checkpoint 间隔 ?

答: 秒级的 checkpoint 无论是状态文件还是, 还是下游的系统,都会产生较大压力, 而对作业本身, 恢复的耗时依旧是分钟级, 这样作业整体失败恢复的流程依旧是分钟级, 设置 checkpoint interval 太小意义不大。

问题2: 数据新鲜度和 checkpoint 的选择 ?

答: 当我们对数据新鲜度要求比较高时,建议使用支持 at-least-once 的系统 (Kafka/doris) 而不是 exactly-once 系统 (iceberg/paimon) 会受到 checkpoint interval 限制。


本文是「我的数据空间」实时计算实践笔记的一篇。平台把这类链路做成了托管作业(Flink SQL 流作业、实时集成、Iceberg 表管理与维护),支持私有化部署与 OEM 合作 —— 见核心能力总览技术博客,交流合作 QQ:1559851993。