- 大数据
- 流处理
- 批处理
- 数据工程
【免费下载链接】flink
Flink 的 Web 界面提供了专门监控作业 Checkpoint 的入口,且作业终止后这些统计依然可查。本文围绕官方文档docs/content/docs/ops/monitoring/checkpoint_monitoring.md展开,系统讲解 Overview、History、Summary、Configuration 四个标签页的全部指标含义与配置方法,并结合 Flink 运行时源码(如CheckpointMetrics、WebOptions)剖析每项数据在 JobManager 侧的采集与聚合方式,帮助你在排查 Checkpoint 超时、对齐耗时过长、增量快照数据量异常等问题时,能直接依据监控面板定位根因。
监控入口与总体说明
Flink Web 界面的 Checkpoint 页面由四个标签页组成:Overview、History、Summary与Configuration。这些统计数据在作业终止后仍可访问,但有两类统计(Overview 页的计数、Summary 页的汇总)在 JobManager 丢失后不会保留,JobManager Failover 后会被重置。
Overview 标签页
Overview 页列出以下统计(注意:这些统计不跨 JobManager Failover 保留,JobManager 故障切换后会被重置):
- Checkpoint Counts(检查点计数)
- Triggered(已触发):自作业启动以来触发的 Checkpoint 总数。
- In Progress(进行中):当前正在进行的 Checkpoint 数量。
- Completed(已完成):自作业启动以来成功完成的 Checkpoint 总数。
- Failed(失败):自作业启动以来失败的 Checkpoint 总数。
- Restored(恢复):自作业启动以来的恢复操作次数,该数值也告诉你作业自提交以来重启了多少次。注意:携带 Savepoint 的初始提交也算一次 Restore;如果运行期间 JobManager 丢失,该计数会被重置。
- Latest Completed Checkpoint(最近一次完成的检查点):最近一次成功完成的 Checkpoint,点击
More details可下钻到子任务(subtask)级别的详细统计。 - Latest Failed Checkpoint(最近一次失败的检查点):最近一次失败的 Checkpoint,同样可点击
More details查看子任务级别明细。 - Latest Savepoint(最近一次 Savepoint):最近触发的 Savepoint 及其外部存储路径,可点击
More details查看明细。 - Latest Restore(最近一次恢复),分两种类型:
- Restore from Checkpoint:从常规周期性 Checkpoint 恢复;
- Restore from Savepoint:从 Savepoint 恢复。
这些计数在源码中对应 CheckpointStatsCounts.java,由 JobManager 在 Checkpoint 触发、完成、失败、恢复等事件处维护。
History 标签页
Checkpoint 历史表保留了最近触发的所有 Checkpoint 的统计,包括当前仍在进行中的 Checkpoint。需要留意:对于失败的 Checkpoint,其指标采用尽力而为(best efforts)方式更新,可能不准确。
表中的每列含义如下:
- ID:触发的 Checkpoint 的 ID,从 1 开始递增。
- Status(状态):Checkpoint 的当前状态,为In Progress、Completed或Failed。如果该 Checkpoint 是 Savepoint,会显示一个软盘(floppy-disk)符号。
- Acknowledged(已确认数):已确认(acknowledged)的子任务数 / 子任务总数。
- Trigger Time(触发时间):Checkpoint 在 JobManager 被触发的时间。
- Latest Acknowledgement(最近确认时间):JobManager 收到的任意子任务的最近一次确认时间(若尚无确认则为 n/a)。
- End to End Duration(端到端时长):从触发时间戳到最近一次确认的时长(无确认时为 n/a)。一次完整 Checkpoint 的端到端时长按最后一个确认的子任务确定,因此该时间通常大于单个子任务实际完成状态快照所需的时间。
- Checkpointed Data Size(检查点数据大小):该 Checkpoint 同步与异步阶段持久化的数据量。若启用了增量 Checkpoint 或 Changelog,该值可能与 Full Checkpoint Data Size 不同。
- Full Checkpoint Data Size(完整检查点数据大小):所有已确认子任务累计的 Checkpoint 数据量。
- Processed (persisted) in-flight data(对齐期间处理/持久化的在途数据):所有已确认子任务在对齐期间(收到第一个与最后一个 Checkpoint Barrier 之间的时间)处理/持久化的字节数近似值。只有启用非对齐 Checkpoint(unaligned checkpoint)时,持久化数据量才可能大于 0。
展开某个 Checkpoint 后,每个子任务还有更细粒度的统计:
- Sync Duration(同步阶段时长):Checkpoint 同步部分的耗时,包含对算子状态做快照,期间会阻塞该子任务上的所有其他活动(处理记录、触发定时器等)。
- Async Duration(异步阶段时长):Checkpoint 异步部分的耗时,包含将 Checkpoint 写入所选文件系统的时间。对于非对齐 Checkpoint,还包含子任务等待最后一个 Checkpoint Barrier 到达的对齐时间(alignment duration)以及持久化在途数据(in-flight data)的耗时。
- Alignment Duration(对齐时长):处理第一个与最后一个 Checkpoint Barrier 之间的时间。对齐 Checkpoint 在对齐期间,已收到 Barrier 的通道会被阻塞,不再处理更多数据。
- Start Delay(启动延迟):自 Checkpoint Barrier 创建起,到第一个 Barrier 到达该子任务所花的时间。
- Unaligned Checkpoint(是否非对齐):该子任务的 Checkpoint 是否以非对齐方式完成。对齐 Checkpoint 在对齐超时后可以切换为非对齐 Checkpoint。
这些子任务级指标在运行时由 CheckpointMetrics.java 承载,其字段与页面一一对应:bytesProcessedDuringAlignment(对齐期间处理的字节数)、bytesPersistedDuringAlignment(对齐期间持久化的字节数)、alignmentDurationNanos(流对齐耗时,纳秒)、syncDurationMillis(同步快照耗时,毫秒)、asyncDurationMillis(异步快照耗时,毫秒)、checkpointStartDelayNanos(Barrier 创建到到达子任务的延迟)、unalignedCheckpoint(是否非对齐完成)、bytesPersistedOfThisCheckpoint(本次持久化字节数)与totalBytesPersisted(累计持久化字节数)。其中“未知/未设置”的取值统一用常量UNSET = -1L表示,这解释了 UI 中部分字段在 Checkpoint 尚未完成时显示为 n/a 的原因。
历史条数配置
通过以下配置键可调整 History 页保留的近期 Checkpoint 条数,默认值为10:
# Number of recent checkpoints that are remembered web.checkpoints.history: 15在源码中,该选项定义于 WebOptions.java,配置键为web.checkpoints.history,整型、默认 10,并带有已废弃的旧键jobmanager.web.checkpoints.history(旧版本配置可直接沿用,新版本会自动兼容)。该值在 JobManager 启动时生效,例如 RestHandlerConfiguration.java 在初始化 REST 处理器时读取它来构造 Checkpoint 历史缓存。
Summary 标签页
Summary 页对所有已完成的 Checkpoint 计算简单的最小值/平均值/最大值统计,覆盖四个维度:End to End Duration(端到端时长)、Incremental Checkpoint Data Size(增量检查点数据大小)、Full Checkpoint Data Size(完整检查点数据大小)与 Bytes Buffered During Alignment(对齐期间缓冲的字节数,含义见 History 章节)。
注意:这些统计不跨 JobManager Failover 保留,JobManager 故障切换后会被重置。因此 Summary 页适合作为作业稳定运行一段时间后的趋势基线,而不适合跨 Failover 做长期对比。
Configuration 标签页
Configuration 页列出当前作业的流式 Checkpoint 配置:
- Checkpointing Mode(检查点模式):Exactly Once或At least Once。
- Interval(间隔):配置的 Checkpoint 间隔,每隔该间隔触发一次 Checkpoint。
- Timeout(超时时间):超过该超时时间后,Checkpoint 会被 JobManager 取消,并触发新的 Checkpoint。
- Minimum Pause Between Checkpoints(检查点间最小间隔):两次 Checkpoint 之间的最小暂停时间。一次 Checkpoint 成功完成后,至少等待该时间才触发下一次,可能会推迟原本的正周期触发。
- Maximum Concurrent Checkpoints(最大并发检查点数):允许同时处于进行中的 Checkpoint 最大数量。
- Persist Checkpoints Externally(外部化持久化):启用或禁用。若启用,还会列出外部化 Checkpoint 的清理策略(取消作业时删除 delete 或保留 retain)。
对照上述监控面板,可以在 Configuration 页核对作业实际生效的 Checkpointing 参数,与 History 页的 End to End Duration、Timeout 相互印证:若端到端时长长期逼近 Timeout,说明间隔或超时配置偏紧,需要从同步/异步阶段时长入手调优。
Checkpoint Details:逐算子与逐子任务明细
点击某个 Checkpoint 的More details链接,可以看到该 Checkpoint 在所有算子上的 Minimum/Average/Maximum 汇总,以及每个子任务的详细数值。
按算子汇总视图:
所有子任务统计视图:
数据流:从子任务确认到 UI 展示
从源码结构看,页面数据的生成链路是:各 TaskManager 子任务完成本地快照后向 JobManager 发送确认,JobManager 将每个子任务的CheckpointMetrics聚合进 CheckpointStatsSnapshot.java(单个 Checkpoint 的完整快照),并按条数上限保存在 CheckpointStatsHistory.java(受web.checkpoints.history约束的环形历史)中;REST 层再通过 CheckpointStatsCache.java 等处理器把数据序列化后提供给 Web 界面。这条链路也解释了文档中的两个特性:历史条数有限(超出上限的最早记录被丢弃),以及 JobManager Failover 后 Overview/Summary 统计被清零(这些聚合状态保存在 JobManager 内存中,并非持久化数据)。
小结
Flink 的 Checkpoint 监控以四个标签页覆盖了“计数趋势—单条历史—统计汇总—生效配置”四个层面:
- 用Overview快速确认作业整体健康度(Failed 计数持续增长、Restored 频繁递增都是危险信号);
- 用History逐条检查 End to End Duration、Checkpointed Data Size 与对齐期间在途数据,定位耗时集中在同步阶段、异步写盘阶段还是 Barrier 对齐阶段;
- 用Checkpoint Details下钻到算子与子任务级别,找出拖慢整体确认的“最后一名”子任务;
- 用Configuration核对间隔、超时、最小暂停与外部化策略等参数,必要时结合
web.checkpoints.history调整历史保留条数以便回溯更长时间窗口。
由于 Overview 与 Summary 统计在 JobManager Failover 后重置,且失败 Checkpoint 的指标仅为尽力而为的近似值,在做容量规划或 SLA 评估时应以 History 中逐条完成记录为准。
- 大数据
- 流处理
- 批处理
- 数据工程
【免费下载链接】flink
相关推荐
Flink Checkpoint 监控指南:深入解读 Web UI 的四个选项卡与底层实现
Flink Checkpoint 监控指南:深入解读 Web UI 的四个选项卡与底层实现 导读 Checkpoint(检查点)是 Flink 流式作业容错与恢
大数据流处理批处理数据工程LiteLLM Grafana 监控看板实战:导入 gen_ai 与 litellm_* 指标 JSON,读懂每一块 Panel 的 PromQL
LiteLLM Grafana 监控看板实战:导入 gen_ai 与 litellm_ 指标 JSON,读懂每一块 Panel 的 PromQL cookboo
后端API网关LLM 网关大模型人工智能Flink 大状态与 Checkpoint 调优实战:从监控指标、RocksDB 内存配置到任务本地恢复
Flink 大状态与 Checkpoint 调优实战:从监控指标、RocksDB 内存配置到任务本地恢复 本文围绕 Flink 官方运维手册中的大状态调优指南(
大数据流处理批处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考