- 示例工程
- 大数据
【免费下载链接】flink-learning
flink learning blog. http://www.54tianzhisheng.cn/ 含 Flink 入门、概念、原理、实战、性能调优、源码解析等内容。涉及 Flink Connector、Metrics、Library、DataStream API、Table API & SQL 等内容的学习案例,还有 Flink 落地应用的大型项目案例(PVUV、日志存储、百亿数据实时去重、监控告警)分享。欢迎大家支持我的专栏《大数据实时计算引擎 Flink 实战与性能优化》
本文基于开源仓库 flink-learning(《大数据实时计算引擎 Flink 实战与性能优化》专栏代码库)中 第十章 10.1 节 的完整内容,系统讲解 Flink RestartStrategy 的配置方法、源码实现与实战踩坑经验。你将掌握flink-conf.yaml与应用程序两种配置方式、FixedDelay / FailureRate / None 三类重启策略的选型原则、Flink 默认 Fallback 策略的行为,以及结合 Checkpoint 与监控告警构建高稳定 Flink 作业的完整思路,并能在 flink-learning-examples 的 restartStrategy 示例代码上直接动手复现。
10.1.1 常见错误导致 Flink 作业重启:一个凌晨两点的教训
作者从使用 Flink 至今,解决过大量来自生产与微信好友的问题,其中**"整个 Job 一直在重启,并伴随各种异常报错(可在 Web UI 的 Exceptions 日志中查看)"**是最常见的一类。生产环境中最典型的三类报错场景包括:脏数据不符合规范、字段为 null 触发空指针(NPE)、数组越界、数据类型转换错误等。
作者曾因其中一个异常导致作业持续重启,在深夜线上发版时,同事发现问题后凌晨两点打电话将其叫醒修复 BUG——这正说明合理的重启策略配置是生产 Flink 作业稳定性的第一道防线。
有人可能会说:"只要过滤掉脏数据、做好 try/catch 异常捕获,Job 就不会不断重启了。" 确实如此,但需要注意:复杂的 Job 下每个算子都可能产生脏数据(包括 Source 本身也可能产生 null 或非法数据),不可能在每个算子中套一个大 try/catch。因此,一方面要尽力保证代码健壮性,另一方面必须配置好 Flink Job 的 RestartStrategy(重启策略),二者缺一不可。
10.1.2 RestartStrategy 简介
RestartStrategy(重启策略)是 Flink 的容错机制核心组件之一。在遇到机器故障、代码异常等不可预知的问题导致 Job 或 Task 挂掉时,Flink 会根据配置的重启策略,将 Job 或受影响的 Task 拉起来重新执行,使作业恢复到之前的正常执行状态。
Flink 中的重启策略决定了三件事:
- 是否要重启Job 或 Task;
- 重启次数(尝试多少次);
- 每次重启的时间间隔(相邻两次重启之间等待多久)。
在 flink-learning-common 的 ExecutionEnvUtil.prepare() 方法 中可以看到,项目公共工具类默认通过env.getConfig().setRestartStrategy(RestartStrategies.fixedDelayRestart(4, 60000))为所有作业设置了"最多重启 4 次、每次间隔 60 秒"的固定延迟策略,这是仓库在生产场景中沉淀下来的默认兜底配置,读者可以直接参考。
10.1.3 为什么需要 RestartStrategy?
重启策略的价值主要体现在三个方面:
- 状态一致性恢复:重启会让 Job 从上一次完整的 Checkpoint 处恢复状态,保证 Job 重启前后状态保持一致(前提是已开启 Checkpoint,对应源码可参考 EnableCheckpointMain);
- 避免消息堆积:重启后 Job 可以继续处理数据,不会因为 Job 挂掉导致消息在 Kafka 等消息队列中大量堆积;
- 降低运维成本:合理的重启策略可以减少 Job 不可用时间,避免人工介入处理故障的运维成本。
因此,重启策略对于 Flink Job 的稳定性有着举足轻重的作用。
10.1.4 如何配置 RestartStrategy?
配置方式遵循"Job 级配置覆盖集群级配置"的原则:
- 若 Flink Job 没有单独设置重启策略,则使用集群启动时加载的默认重启策略;
- 若 Flink Job 中单独设置了重启策略,则覆盖默认的集群重启策略。
默认重启策略在 Flink 的配置文件flink-conf.yaml中通过restart-strategy参数控制,共有三种可选值:fixed-delay(固定延时重启策略)、failure-rate(故障率重启策略)、none(不重启策略),选择不同的策略会对应不同的配套参数。下面逐一介绍。
FixedDelayRestartStrategy(固定延时重启策略)
FixedDelayRestartStrategy按照集群配置文件中或程序中额外设置的重启次数尝试重启作业,若尝试次数超过给定的最大次数后作业仍未成功启动,则停止作业;同时可配置连续两次重启之间的等待时间。
在flink-conf.yaml中配置:
restart-strategy: fixed-delay restart-strategy.fixed-delay.attempts: 3 # 表示作业重启的最大次数,启用 checkpoint 的话是 Integer.MAX_VALUE,否则是 1。 restart-strategy.fixed-delay.delay: 10 s # 如果设置分钟可以类似 1 min,该参数表示两次重启之间的时间间隔,当程序与外部系统有连接交互时延迟重启可能会有帮助,启用 checkpoint 的话,延迟重启的时间是 10 秒,否则使用 akka.ask.timeout 的值。在应用程序中设置固定延迟重启策略:
ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment(); env.setRestartStrategy(RestartStrategies.fixedDelayRestart( 3, // 尝试重启的次数 Time.of(10, TimeUnit.SECONDS) // 延时 ));仓库中的完整可运行示例在 FixedDelayRestartStrategyMain,它配置的是RestartStrategies.fixedDelayRestart(3, 5000),即"每隔 5 秒重启一次,尝试三次如果 Job 还没有起来则停止",随后通过一个持续向 map 算子发送 null 值的 SourceFunction 触发空指针异常,用来真实复现"Job 失败 → 重启 → 再失败"的完整链路。
FailureRateRestartStrategy(故障率重启策略)
FailureRateRestartStrategy在发生故障之后重启作业,但如果在固定时间间隔之内发生的故障次数超过设置的值,作业就会失败停止。该策略同样支持设置连续两次重启之间的等待时间。
在flink-conf.yaml中配置:
restart-strategy: failure-rate restart-strategy.failure-rate.max-failures-per-interval: 3 # 固定时间间隔内允许的最大重启次数,默认 1 restart-strategy.failure-rate.failure-rate-interval: 5 min # 固定时间间隔,默认 1 分钟 restart-strategy.failure-rate.delay: 10 s # 连续两次重启尝试之间的延迟时间,默认是 akka.ask.timeout在应用程序中设置故障率重启策略:
ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment(); env.setRestartStrategy(RestartStrategies.failureRateRestart( 3, // 固定时间间隔允许 Job 重启的最大次数 Time.of(5, TimeUnit.MINUTES), // 固定时间间隔 Time.of(10, TimeUnit.SECONDS) // 两次重启的延迟时间 ));仓库示例 FailureRateRestartStrategyMain 使用的是RestartStrategies.failureRateRestart(3, Time.minutes(2), Time.seconds(10)),语义为"每隔 10 秒重启一次,如果两分钟内重启过三次则停止 Job"。
NoRestartStrategy(不重启策略)
NoRestartStrategy作业不重启,直接失败停止。在flink-conf.yaml中配置:
restart-strategy: none在应用程序中设置不重启:
ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment(); env.setRestartStrategy(RestartStrategies.noRestart());仓库示例 NoRestartStrategyMain 通过RestartStrategies.noRestart()配置后,作业一旦出现空指针异常就会直接 FAILED,不会进行任何重启。
Fallback(备用重启策略)
如果程序没有启用 Checkpoint,则采用不重启策略;如果开启了 Checkpoint 且没有设置重启策略,则采用固定延时重启策略,最大重启次数为 Integer.MAX_VALUE。这就是 Flink 的 Fallback 逻辑,也是理解"为什么默认行为不同"的关键。
在应用程序中配置好固定延时重启策略后,可以测试代码异常导致 Job 失败后重启的情况,观察日志可以看到 Job 重启相关的输出:
[flink-akka.actor.default-dispatcher-5] INFO org.apache.flink.runtime.executiongraph.ExecutionGraph - Try to restart or fail the job zhisheng default RestartStrategy example (a890361aed156610b354813894d02cd0) if no longer possible. [flink-akka.actor.default-dispatcher-5] INFO org.apache.flink.runtime.executiongraph.ExecutionGraph - Job zhisheng default RestartStrategy example (a890361aed156610b354813894d02cd0) switched from state FAILING to RESTARTING. [flink-akka.actor.default-dispatcher-5] INFO org.apache.flink.runtime.executiongraph.ExecutionGraph - Restarting the job zhisheng default RestartStrategy example (a890361aed156610b354813894d02cd0).最后重启次数达到配置的最大重启次数后 Job 还没有起来,则会停止 Job 并打印日志:
[flink-akka.actor.default-dispatcher-2] INFO org.apache.flink.runtime.executiongraph.ExecutionGraph - Could not restart the job zhisheng default RestartStrategy example (a890361aed156610b354813894d02cd0) because the restart strategy prevented it.日志中的
zhisheng default RestartStrategy example正是仓库 DefaultRestartStrategyMain 中env.execute("zhisheng default RestartStrategy example")指定的作业名,说明这些日志就是该示例作业在未显式设置策略、由集群默认 Fallback 逻辑兜底时的真实输出。
如何选择合适的重启策略?以空指针异常为例:如果程序抛出 NPE 而你配置的是无限重启,会导致 Job 一直在重启,白白浪费机器资源。此时建议配置固定延时重启策略(固定重试次数 + 固定间隔),在重试一定次数后 Job 就会停止;如果对 Job 的状态做了监控告警,你会第一时间收到告警信息,从而及时发现问题并修复 Job。
仓库中监控告警的最佳实践可参考 flink-learning-monitor-alert,它提供了完整的 Flink 作业监控告警实现,可与重启策略形成"自动恢复 + 人工兜底"的完整闭环。
10.1.5 RestartStrategy 源码分析
从上面的程序配置代码可以看到,设置重启策略使用的都是RestartStrategies类,通过该类的方法即可创建不同的重启策略。在RestartStrategies类中提供了五个方法用来创建四种不同的重启策略(其中两个方法是创建 FixedDelay 重启策略的,只是参数不同)。
在每个方法内部,实际调用的是RestartStrategies中的内部静态配置类:
NoRestartStrategyConfigurationFixedDelayRestartStrategyConfigurationFailureRateRestartStrategyConfigurationFallbackRestartStrategyConfiguration
这四个配置类都继承自RestartStrategyConfiguration抽象类。
在 Flink 中,RestartStrategyResolving类的resolve方法负责解析RestartStrategies.RestartStrategyConfiguration,然后根据配置使用RestartStrategyFactory创建RestartStrategy。
RestartStrategy是一个接口,定义了canRestart和restart两个核心方法,它有四个实现类:
FixedDelayRestartStrategyFailureRateRestartStrategyThrowingRestartStrategyNoRestartStrategy
从接口设计可以推断:canRestart用于判断当前是否还允许重启(如是否超过最大次数/时间窗口内故障率是否超限),restart用于实际执行重启动作,而ThrowingRestartStrategy这类实现则对应配置解析出错等异常场景下的兜底行为。
结合仓库源码看 RestartStrategies 的真实用法
- FixedDelayRestartStrategyMain 与 AEMain:均使用
RestartStrategies.fixedDelayRestart(3, 5000),后者通过map(aLong -> aLong / 0)制造除零异常来触发重启; - FailureRateRestartStrategyMain:使用
failureRateRestart(3, Time.minutes(2), Time.seconds(10)); - NoRestartStrategyMain:使用
noRestart(); - DefaultRestartStrategyMain:不显式设置策略,用于观察集群默认(Fallback)重启策略的行为;
- EnableCheckpointMain:不设置重启策略但开启 Checkpoint(
env.enableCheckpointing(10000)+MemoryStateBackend),用来验证"开启 Checkpoint 后默认采用 FixedDelay 且最大重启次数为 Integer.MAX_VALUE"的 Fallback 规则。
这一组示例恰好覆盖了"配置策略 vs 不配置策略""开 Checkpoint vs 不开 Checkpoint"两个维度,是复现本文全部结论的最小实验集。
另外,flink-learning-project 的 FlinkJobScaffold 作为生产级作业模板,将重启策略与 Checkpoint 配置放在一起呈现:先配置 Checkpoint(间隔 60 秒、Exactly-Once、最小间隔 30 秒、超时 10 分钟、取消时保留),再配置fixedDelayRestart(3, Time.seconds(10))固定延迟重启策略,最后设置并行度与 Kafka Source——这是"重启策略 + Checkpoint + 状态恢复"协同工作的标准生产姿势,读者可以直接以此为模板落地。
10.1.6 Failover Strategies(故障恢复策略)
除 RestartStrategy 之外,Flink 还提供 Failover Strategies(故障恢复策略),用于决定Task 失败后如何恢复(RestartStrategy 解决的是"要不要重启、重启多少次",Failover Strategy 解决的是"重启哪些 Task")。主要包含两类:
重启所有的任务
默认的故障恢复策略,Task 失败后重启作业的所有 Task(当作业开启 Checkpoint 后,会从最近一次 Checkpoint 恢复所有算子状态)。该策略实现简单、语义直观,适合作业规模不大或对恢复速度要求不苛刻的场景。
基于 Region 的局部故障重启策略
Flink 1.9 之后引入的 Region 级故障恢复策略。将作业按照算子连接关系划分为多个 Region(上游与下游共享数据交换的算子处于同一 Region),当某个 Task 失败时,只重启故障所在 Region 及其依赖的上游 Region,其他 Region 的 Task 不受影响继续运行,从而显著减少故障恢复的代价、降低重启对整体作业的影响面。该策略适合作业链条长、并行度大、对可用性要求高的生产场景。
从 Flink 1.9 起,基于 Region 的故障恢复策略已作为默认值,读者可在flink-conf.yaml中通过jobmanager.execution.failover-strategy参数显式指定(可选值full或region),结合自身的重启策略一起规划作业的容错行为。
10.1.7 小结与反思
- 配置是兜底,不是替代:脏数据和异常防不胜防,尽量保证代码健壮性(过滤脏数据、合理异常捕获),但每个算子都做防御式编程不现实,RestartStrategy 是 Job 稳定性的最后防线;
- 策略选择要克制:空指针这类确定性 Bug 配无限重启只会空耗集群资源,推荐固定延时重启策略(有限次数 + 间隔),让作业"重试几次即停",再配合监控告警第一时间人工介入;
- 与 Checkpoint 深度联动:重启后会从最近一次完整 Checkpoint 恢复状态,因此"开启 Checkpoint + 合理的重启策略 + 外部化 Checkpoint 保留"是生产环境的黄金组合;
- 故障恢复粒度可选:Region 级局部故障恢复可以最小化故障影响面,作业规模大时应优先考虑;
- 动手验证:直接运行 flink-learning-examples 中 restartStrategy 包 下的六个示例(分别覆盖默认 Fallback、FixedDelay、FailureRate、None、开启 Checkpoint 等场景),结合 10.1.4 节的两段 ExecutionGraph 日志对比观察作业状态迁移(FAILING → RESTARTING → 重启成功或失败停止),是对本文结论最好的印证。
通过本节的学习,读者应能根据自身业务特征为每个 Flink 作业选择并配置合适的重启策略,并与 Checkpoint、监控告警协同,构建高可用的实时计算作业。
- 示例工程
- 大数据
【免费下载链接】flink-learning
flink learning blog. http://www.54tianzhisheng.cn/ 含 Flink 入门、概念、原理、实战、性能调优、源码解析等内容。涉及 Flink Connector、Metrics、Library、DataStream API、Table API & SQL 等内容的学习案例,还有 Flink 落地应用的大型项目案例(PVUV、日志存储、百亿数据实时去重、监控告警)分享。欢迎大家支持我的专栏《大数据实时计算引擎 Flink 实战与性能优化》
相关推荐
Flink SQL JOB 语句实战:SHOW JOBS / DESCRIBE JOB / STOP JOB 作业生命周期管理
Flink SQL JOB 语句实战:SHOW JOBS / DESCRIBE JOB / STOP JOB 作业生命周期管理 JOB 语句是 Flink SQ
后端大数据流处理批处理Flink并行度与资源分配:如何优化TaskManager与Slot配置
Flink并行度与资源分配:如何优化TaskManager与Slot配置 Apache Flink作为业界领先的流处理框架,其并行度配置与资源分配策略直接影响作
示例工程大数据Flink SQL JOB 语句详解:SHOW JOBS / DESCRIBE JOB / STOP JOB 作业生命周期管理实战指南
Flink SQL JOB 语句详解:SHOW JOBS / DESCRIBE JOB / STOP JOB 作业生命周期管理实战指南 在 Flink Tabl
后端大数据流处理批处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考