You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何按自定义时间触发Flink Checkpoint?要求对齐5分钟整点节点

默认的env.enableCheckpointing(5 * 60 * 1000)是从作业启动时间开始按固定间隔触发Checkpoint,没法对齐x:00、x:05这类整点节点。要实现对齐,需要通过自定义Checkpoint触发策略来精准控制触发时机。

实现步骤:

  1. 自定义CheckpointTrigger接口实现类,计算下一个对齐5分钟整点的触发时间
  2. 将自定义触发器绑定到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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.29 08:07:53