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

Flink搭载无界数据源的流作业如何在代码内部实现优雅停机?

Flink无界流作业内部优雅停机实现方案

下面是三种不依赖外部服务的内部触发停机方案,均可保证Checkpoint机制正常执行:


方案1:数据源侧主动终止(全版本适配)

在无界数据源的读取逻辑中新增停机判断规则,触发条件满足后执行以下流程:

  • 停止拉取新的上游数据
  • 等待所有已拉取数据全部下发到下游算子
  • 调用CheckpointListener接口触发一次手动Checkpoint
  • 待手动Checkpoint完成回调后,调用SourceContext.close()关闭数据源,同时向JobManager发送作业完成信号
  • 提前开启execution.checkpointing.externalized-checkpoint-retention配置,保证最终Checkpoint持久化留存

Flink 1.14及以上版本提供了算子内部直接触发优雅停机的API,可在任意富函数中通过运行时上下文调用stopWithSavepoint方法实现:

// 富函数内调用示例
public class GracefulShutdownTrigger extends RichMapFunction<InputType, OutputType> {
    @Override
    public OutputType map(InputType value) throws Exception {
        // 匹配到停机触发条件时执行
        if (needShutdown(value)) {
            JobMasterGateway jobMasterGateway = ((StreamingRuntimeContext) getRuntimeContext()).getJobMasterGateway();
            // 第二个参数drain设为true,会处理完所有缓存数据后再关停
            jobMasterGateway.stopWithSavepoint(true, "/your/savepoint/path", SavepointFormatType.CANONICAL);
        }
        return value;
    }
}

该方案会自动完成最后一次全量Checkpoint,关停顺序符合Flink算子生命周期规范,不会触发重启策略。


方案3:自定义停机异常(低版本适配)

如果使用的Flink版本低于1.14,可以通过自定义不可重启异常实现:

  1. 首先定义专用的停机异常类:
    public class GracefulShutdownException extends RuntimeException {
        public GracefulShutdownException(String message) {
            super(message);
        }
    }
    
  2. 配置重启策略时,将该异常设置为不触发重启的异常类型:
    env.setRestartStrategy(RestartStrategies.failureRateRestart(
        3, // 最大重启次数
        Time.minutes(5), // 统计时间窗口
        Time.seconds(10), // 重启间隔
        // 配置停机异常不触发重启
        Collections.singletonList(GracefulShutdownException.class)
    ));
    
  3. 满足停机条件时直接抛出自定义的GracefulShutdownException即可,作业会完成当前进行中的Checkpoint后正常终止,不会自动重启。

注意:所有方案都建议提前配置外部化Checkpoint留存策略,避免作业停机后Checkpoint被自动清理,保障后续可从该Checkpoint恢复作业。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.26 11:36:03