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

如何监听Flink任务的失败与重启并执行自定义逻辑?

Flink任务重启前后回调的实现方案

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 02:12:29