1. 如何验证是否启用了 Checkpoint?

方法 1:检查作业提交时的 Flink UI
  • 在 Flink Web UI 的 JobManager 日志 或 作业配置 中搜索 execution.checkpointing.interval。

  • 如果该值大于 0(如 30s),则说明 Checkpoint 已启用。

方法 2:检查作业配置
  • 在平台或者服务端的作业详情页,查看 “高级配置” 或 “Flink 配置” 部分,是否有如下参数:

    ini

    execution.checkpointing.interval: 30s
    state.backend: filesystem
    state.checkpoints.dir: hdfs:///flink/checkpoints
方法 3:在 SQL 中显式打印配置
  • 在 Flink SQL 文件中添加以下语句,通过日志观察实际生效的配置:

    sql

    SET 'execution.checkpointing.interval' = '30s';  -- 显式设置,确保覆盖默认值

2. 生产环境推荐配置

  1. 启用 Checkpoint:
    确保 Flink 定期提交偏移量到 Kafka 和状态后端:

    sql

    SET 'execution.checkpointing.interval' = '30s';
    SET 'execution.checkpointing.mode' = 'EXACTLY_ONCE';
  2. 明确指定启动模式:
    根据业务需求选择:

    sql

    -- 从消费者组偏移量开始(推荐生产使用)
    'scan.startup.mode' = 'group-offsets',
    'properties.auto.offset.reset' = 'earliest'  -- 后备策略

4. 如果未启用 Checkpoint 的后果

  • 数据源为 Kafka:

    • 若 scan.startup.mode = latest-offset(默认),重启后从最新偏移量开始,丢失未处理的数据。

    • 若 scan.startup.mode = earliest-offset,重启后从头消费,导致重复处理。

  • 状态计算(如窗口聚合):

    • 所有中间状态丢失,重启后重新计算,结果可能不准确。

Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐