ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

《Flink 实战与性能优化》:如何设置 Flink Job RestartStrategy(重启策略)并打造稳定作业

《Flink 实战与性能优化》:如何设置 Flink Job RestartStrategy(重启策略)并打造稳定作业 示例工程大数据【免费下载链接】flink-learningflink learning blog. http://www.54tianzhisheng.cn/ 含 Flink 入门、概念、原理、实战、性能调优、源码解析等内容。涉及 Flink Connector、Metrics、Library、DataStream API、Table API SQL 等内容的学习案例还有 Flink 落地应用的大型项目案例PVUV、日志存储、百亿数据实时去重、监控告警分享。欢迎大家支持我的专栏《大数据实时计算引擎 Flink 实战与性能优化》项目地址https://gitcode.com/gh_mirrors/fl/flink-learning点击查看免费下载本文基于开源仓库 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不设置重启策略但开启 Checkpointenv.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-learningflink learning blog. http://www.54tianzhisheng.cn/ 含 Flink 入门、概念、原理、实战、性能调优、源码解析等内容。涉及 Flink Connector、Metrics、Library、DataStream API、Table API SQL 等内容的学习案例还有 Flink 落地应用的大型项目案例PVUV、日志存储、百亿数据实时去重、监控告警分享。欢迎大家支持我的专栏《大数据实时计算引擎 Flink 实战与性能优化》项目地址https://gitcode.com/gh_mirrors/fl/flink-learning点击查看免费下载相关推荐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后端大数据流处理批处理上一篇如何在1分钟内为Windows安装苹果USB网络共享驱动完整解决方案下一篇如何5分钟实现Windows和Office永久激活KMS智能激活完整指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表