Flink作业Checkpoint恢复时运行指标重置是否正常
问题背景
- 运行环境:Flink作业部署在Kubernetes集群,Checkpoint存储对接Azure ADLS服务
- 恢复表现:作业发生故障后可正常从ADLS路径加载Checkpoint恢复,但恢复阶段所有Flink运行指标被重置,其中Bytes Received、Bytes Sent、Records Sent三类指标直接归零
- 附加观察:作业Checkpoint大小呈增量增长趋势
问题答复
指标重置是否正常,是否存在数据丢失
该指标重置现象属于Flink的正常设计表现,完全不代表Checkpoint恢复过程存在数据丢失。
- Flink的Bytes Received、Bytes Sent、Records Sent这类网络传输类运行指标,属于单进程生命周期内的临时统计值,统计维度是当前运行的TaskManager实例从进程启动开始的累计传输数据,这类计数器值不会被持久化写入Checkpoint存储。
- Flink作业故障恢复的本质是拉起全新的JobManager、TaskManager进程,从Checkpoint文件加载持久化的算子状态、消费位点等核心计算状态后重新启动计算,新启动的进程内置的统计计数器初始值为0,因此会观测到这类运行指标直接归零,该表现和状态恢复是否完整、是否丢数没有任何关联。
- 判断Checkpoint恢复是否存在数据丢失的核心校验标准是状态一致性:只要作业配置了Exactly-Once语义、Source组件支持从Checkpoint记录的消费位点重放数据、算子状态(如聚合结果、Keyed State值)恢复值和Checkpoint记录一致,就不会出现数据丢失问题。
Checkpoint增量增长的说明
如果作业使用RocksDB状态后端开启增量Checkpoint,Checkpoint总大小持续增长属于常见现象:
- 增量Checkpoint机制仅会上传和上一次Checkpoint相比新增的SST文件,不会每次全量上传所有状态数据,随着作业运行,未被Compaction合并、未被过期Checkpoint清理策略删除的历史SST文件会逐步累计,Checkpoint总存储大小就会呈上涨趋势。
- 只要单轮Checkpoint的新增上传大小没有出现异常突增,就不属于配置异常。如果需要控制Checkpoint总存储占用,可以通过调整RocksDB Compaction参数、配置合理的Checkpoint历史保留个数自动清理过期Checkpoint文件实现。
内容的提问来源于stack exchange,提问作者user1112259
相关产品推荐
相关产品推荐

