flink sql checkpoint的设置
·
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. 生产环境推荐配置
-
启用 Checkpoint:
确保 Flink 定期提交偏移量到 Kafka 和状态后端:sql
SET 'execution.checkpointing.interval' = '30s'; SET 'execution.checkpointing.mode' = 'EXACTLY_ONCE';
-
明确指定启动模式:
根据业务需求选择:sql
-- 从消费者组偏移量开始(推荐生产使用) 'scan.startup.mode' = 'group-offsets', 'properties.auto.offset.reset' = 'earliest' -- 后备策略
4. 如果未启用 Checkpoint 的后果
-
数据源为 Kafka:
-
若
scan.startup.mode = latest-offset(默认),重启后从最新偏移量开始,丢失未处理的数据。 -
若
scan.startup.mode = earliest-offset,重启后从头消费,导致重复处理。
-
-
状态计算(如窗口聚合):
-
所有中间状态丢失,重启后重新计算,结果可能不准确。
-
更多推荐
所有评论(0)