如何监听Flink任务的失败与重启并执行自定义逻辑?
Flink任务重启前后回调的实现方案
1. 使用FailureEnricher接口(Flink 1.14+)
Flink 1.14及以上版本提供了FailureEnricher接口,专门用于在任务失败重启前后注入自定义逻辑。你可以实现该接口,重写onFailure方法——这个方法会在任务失败被检测到、重启触发前被调用;如果需要在重启完成后执行逻辑,可以结合JobListener或者任务初始化方法配合实现。
示例代码:
public class CustomFailureEnricher implements FailureEnricher { @Override public void onFailure(Throwable failureContext, Context context) { // 任务失败、重启前的自定义逻辑,比如记录失败详情、发送告警 System.out.println("任务 " + context.getTaskName() + " 即将重启,失败原因:" + failureContext.getMessage()); } }
注册该实现到Flink配置:
Configuration config = new Configuration(); config.setString(ConfigConstants.FAILURE_ENRICHERS_KEY, "com.yourpackage.CustomFailureEnricher"); StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(config);
2. 自定义RestartStrategy
如果需要更精细地控制整个重启流程,可以自定义RestartStrategy。实现RestartStrategy接口后,能在重启触发前、重启完成后分别插入自定义代码,完全掌控重启的逻辑节点。
示例思路:
public class CustomRestartStrategy implements RestartStrategy { @Override public RestartStrategyConfiguration createRestartStrategyConfiguration(Configuration configuration) { return new RestartStrategyConfiguration() { @Override public RestartStrategyFactory createRestartStrategyFactory() { return new RestartStrategyFactory() { @Override public RestartStrategy createRestartStrategy(ClassLoader classLoader) { return new RestartStrategy() { @Override public boolean canRestart() { // 判断是否允许执行重启 return true; } @Override public void notifyFailure(Throwable cause) { // 重启触发前的自定义逻辑 System.out.println("准备重启任务,失败原因:" + cause.getMessage()); } @Override public void reset() { // 重启完成后的收尾逻辑 System.out.println("任务重启完成"); } }; } }; } }; } }
注册自定义重启策略:
env.setRestartStrategy(new CustomRestartStrategy());
3. 任务级别的生命周期回调(RichFunction)
针对单个算子任务,可以继承RichMapFunction、RichFlatMapFunction等RichFunction子类,重写open()和close()方法:
open()会在任务启动(包括每次重启)时执行,可用于重启后的初始化操作;close()会在任务停止(失败前)执行,可用于失败前的清理或告警。
示例:
public class CustomRichMapFunction extends RichMapFunction<String, String> { @Override public void open(Configuration parameters) throws Exception { super.open(parameters); // 任务启动/重启后的自定义逻辑 System.out.println("任务 " + getRuntimeContext().getTaskName() + " 启动/重启完成"); } @Override public void close() throws Exception { super.close(); // 任务停止(失败前)的自定义逻辑 System.out.println("任务 " + getRuntimeContext().getTaskName() + " 即将停止"); } @Override public String map(String value) throws Exception { return value.toUpperCase(); } }
这种方式适合针对特定算子做个性化处理,灵活性很高。
注意事项
FailureEnricher是全局生效的,适合统一的失败重启前处理逻辑;- 自定义
RestartStrategy适合需要完全掌控重启流程的场景; RichFunction的生命周期方法针对单个算子,需要逐个实现但能精准控制每个任务的重启前后行为。
内容的提问来源于stack exchange,提问作者Vicky
相关产品推荐
相关产品推荐

