Flink搭载无界数据源的流作业如何在代码内部实现优雅停机?
Flink无界流作业内部优雅停机实现方案
下面是三种不依赖外部服务的内部触发停机方案,均可保证Checkpoint机制正常执行:
方案1:数据源侧主动终止(全版本适配)
在无界数据源的读取逻辑中新增停机判断规则,触发条件满足后执行以下流程:
- 停止拉取新的上游数据
- 等待所有已拉取数据全部下发到下游算子
- 调用
CheckpointListener接口触发一次手动Checkpoint - 待手动Checkpoint完成回调后,调用
SourceContext.close()关闭数据源,同时向JobManager发送作业完成信号 - 提前开启
execution.checkpointing.externalized-checkpoint-retention配置,保证最终Checkpoint持久化留存
方案2:官方内部API调用(Flink 1.14+ 适用)
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,可以通过自定义不可重启异常实现:
- 首先定义专用的停机异常类:
public class GracefulShutdownException extends RuntimeException { public GracefulShutdownException(String message) { super(message); } } - 配置重启策略时,将该异常设置为不触发重启的异常类型:
env.setRestartStrategy(RestartStrategies.failureRateRestart( 3, // 最大重启次数 Time.minutes(5), // 统计时间窗口 Time.seconds(10), // 重启间隔 // 配置停机异常不触发重启 Collections.singletonList(GracefulShutdownException.class) )); - 满足停机条件时直接抛出自定义的
GracefulShutdownException即可,作业会完成当前进行中的Checkpoint后正常终止,不会自动重启。
注意:所有方案都建议提前配置外部化Checkpoint留存策略,避免作业停机后Checkpoint被自动清理,保障后续可从该Checkpoint恢复作业。
内容的提问来源于stack exchange,提问作者js.
相关产品推荐
相关产品推荐

