如何按自定义时间触发Flink Checkpoint?要求对齐5分钟整点节点
如何让Flink Checkpoint严格对齐5分钟整点触发?
默认的env.enableCheckpointing(5 * 60 * 1000)是从作业启动时间开始按固定间隔触发Checkpoint,没法对齐x:00、x:05这类整点节点。要实现对齐,需要通过自定义Checkpoint触发策略来精准控制触发时机。
实现步骤:
- 自定义
CheckpointTrigger接口实现类,计算下一个对齐5分钟整点的触发时间 - 将自定义触发器绑定到Flink执行环境中
代码实现示例:
首先实现自定义触发器:
import org.apache.flink.api.common.JobID; import org.apache.flink.runtime.checkpoint.CheckpointTrigger; import org.apache.flink.runtime.executiongraph.ExecutionGraph; import org.apache.flink.runtime.jobgraph.JobStatus; import java.time.LocalDateTime; import java.time.temporal.ChronoUnit; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.ScheduledThreadPoolExecutor; import java.util.concurrent.TimeUnit; public class AlignedCheckpointTrigger implements CheckpointTrigger { private static final long CHECKPOINT_INTERVAL = 5 * 60 * 1000; // 5分钟间隔 private final ScheduledThreadPoolExecutor scheduler = new ScheduledThreadPoolExecutor(1); private ScheduledFuture<?> scheduledTrigger; @Override public void initialize(ExecutionGraph executionGraph, JobID jobID) { // 计算第一个对齐整点的触发延迟 long delay = calculateInitialDelay(); scheduledTrigger = scheduler.scheduleAtFixedRate( () -> triggerCheckpoint(executionGraph), delay, CHECKPOINT_INTERVAL, TimeUnit.MILLISECONDS ); } private long calculateInitialDelay() { LocalDateTime now = LocalDateTime.now(); // 计算下一个5分钟整点:分钟数取整到最近的5的倍数,当前刚好是整点则延后5分钟 int currentMinute = now.getMinute(); int nextMinute = ((currentMinute + 4) / 5) * 5; LocalDateTime nextTriggerTime = now.withMinute(nextMinute).withSecond(0).withNano(0); if (nextTriggerTime.isBefore(now) || nextTriggerTime.isEqual(now)) { nextTriggerTime = nextTriggerTime.plusMinutes(5); } return ChronoUnit.MILLIS.between(now, nextTriggerTime); } private void triggerCheckpoint(ExecutionGraph executionGraph) { if (executionGraph.getState() == JobStatus.RUNNING) { executionGraph.triggerCheckpoint(false, false); } } @Override public void stop() { if (scheduledTrigger != null) { scheduledTrigger.cancel(false); } scheduler.shutdown(); } }
然后在作业中配置这个触发器:
import org.apache.flink.configuration.Configuration; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; public class CheckpointAlignmentJob { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // 配置Checkpoint基础参数(超时、并发、模式等) env.getCheckpointConfig().setCheckpointTimeout(10 * 60 * 1000); env.getCheckpointConfig().setMaxConcurrentCheckpoints(1); // 按需添加其他Checkpoint配置(如持久化策略、恢复模式等) // 绑定自定义对齐触发器 Configuration config = new Configuration(); config.setString("execution.checkpoint.trigger.class", AlignedCheckpointTrigger.class.getName()); env.configure(config); // 编写你的作业逻辑... env.execute("Aligned Checkpoint Job"); } }
注意事项:
- 时区适配:如果需要基于特定时区(如UTC)计算整点,将
LocalDateTime.now()替换为ZonedDateTime.now(ZoneId.of("UTC"))即可 - 重启兼容性:作业重启时,触发器会重新计算下一个触发时间,依然能保持整点对齐
- 版本要求:确保使用Flink 1.12及以上版本,该版本开始支持自定义CheckpointTrigger扩展点
内容的提问来源于stack exchange,提问作者gfytd
相关产品推荐
相关产品推荐

