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

如何在Flink中监听Checkpoint/Savepoint的创建事件?

Flink监听Checkpoint/Savepoint创建事件的实现方法

Flink中并非直接在StreamExecutionEnvironment中挂载监听器,而是通过Checkpoint配置、算子接口或Savepoint钩子来实现事件监听,以下是三种常用方案:

1. 全局Checkpoint事件监听

通过CheckpointConfig注册全局监听器,可接收所有Checkpoint的完成/中止事件,适用于需要全局感知Checkpoint状态的场景:

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

// 启用周期性Checkpoint
env.enableCheckpointing(5000);

// 获取Checkpoint配置实例
CheckpointConfig checkpointConfig = env.getCheckpointConfig();

// 注册全局Checkpoint监听器
checkpointConfig.addCheckpointListener(new CheckpointListener() {
    @Override
    public void notifyCheckpointComplete(long checkpointId) throws Exception {
        // 处理Checkpoint完成逻辑
        System.out.printf("Checkpoint %d 已成功完成%n", checkpointId);
    }

    @Override
    public void notifyCheckpointAborted(long checkpointId) throws Exception {
        // 处理Checkpoint中止逻辑
        System.out.printf("Checkpoint %d 已中止%n", checkpointId);
    }
});

2. 算子级Checkpoint事件监听

让特定算子实现CheckpointListener接口,仅该算子会收到Checkpoint事件通知,适合需要在算子内部感知Checkpoint状态的场景:

public class CheckpointAwareMapper extends RichMapFunction<String, String> implements CheckpointListener {

    @Override
    public String map(String value) throws Exception {
        // 算子核心业务逻辑
        return value.toUpperCase();
    }

    @Override
    public void notifyCheckpointComplete(long checkpointId) throws Exception {
        // 当前算子的Checkpoint完成回调
        String taskName = getRuntimeContext().getTaskName();
        System.out.printf("算子[%s]的Checkpoint %d 已完成%n", taskName, checkpointId);
    }

    @Override
    public void notifyCheckpointAborted(long checkpointId) throws Exception {
        // 当前算子的Checkpoint中止回调
        String taskName = getRuntimeContext().getTaskName();
        System.out.printf("算子[%s]的Checkpoint %d 已中止%n", taskName, checkpointId);
    }
}

// 作业中使用该算子
env.fromElements("foo", "bar", "baz")
   .map(new CheckpointAwareMapper())
   .print();

3. Savepoint创建事件监听

Savepoint多为手动触发(CLI/REST API),可通过SavepointHook监听其创建前后的事件:

// 注册Savepoint钩子
env.registerSavepointHook(new SavepointHook<String>() {
    @Override
    public String preSavepoint(String checkpointPath) throws Exception {
        // Savepoint创建前执行,可返回自定义元数据
        return "savepoint_custom_meta_v1";
    }

    @Override
    public void postSavepoint(String checkpointPath, String metaData) throws Exception {
        // Savepoint创建完成后执行
        System.out.printf("Savepoint已创建,路径:%s,自定义元数据:%s%n", checkpointPath, metaData);
    }
});

注意事项

  • addCheckpointListener API在Flink 1.11及以上版本可用,低版本需通过算子实现接口的方式替代。
  • 监听器内的逻辑需尽量轻量化,避免阻塞Checkpoint/Savepoint的正常流程。
  • Savepoint钩子仅在手动触发Savepoint时触发,周期性Checkpoint不会触发该钩子。

内容的提问来源于stack exchange,提问作者feroze

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.23 10:45:06