ARTICLE DETAIL

资讯详情

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

Flink作业调度与失败恢复:从提交到运行的完整链路排查指南

Flink作业调度与失败恢复:从提交到运行的完整链路排查指南 1. 别被“调度”两个字唬住先搞清楚 Flink 作业的一生做 Flink 开发的几乎都遇到过这种场面集群资源看着还有一大堆结果提交作业就是起不来或者作业跑得好好的某天夜里突然挂了重启之后状态怼不上数据对不齐领导在群里连环问。这种时候SR 一般第一反应是查日志、看监控但真正追根溯源的时候往往都会落到一个绕不开的底层话题——Flink Jobs and Scheduling。说白了就是作业从提交到运行再到失败恢复资源是怎么分配的、任务是怎么被调度起来的、挂了以后又是怎么被捞回来的。这篇文章不打算给你搬教科书上的概念也不堆术语。我想从实际排查问题的角度把 Flink 资源调度和失败恢复这条线完整捋一遍说清楚作业从提交到挂掉再到恢复的完整链路。不管是刚开始学 Flink 的菜鸟还是已经写了几年 SQL 但没仔细抠过底层的老手这篇文章的思路都适用。看完之后你能知道为什么作业卡在 SUBMITTED 起不来为什么 Slot 明明有剩余但作业就是跑不满为什么 Checkpoint 频繁失败以及作业挂了以后 Flink 是怎么决定要不要重启、从哪儿恢复的。我尽量用大白话拆解但涉及的机制会比较深。建议你按顺序读别跳因为调度和恢复这两件事是强关联的——调度层面决定了任务怎么摆恢复层面决定了任务摆错了之后怎么办。2. 一张完整的调度路径图从提交到运行的每一步2.1 作业提交后到底发生了什么先直接给结论一个 Flink 作业从你点击提交或者执行 flink run 开始背后会经历这么一条链路用户代码 - StreamGraph - JobGraph - ExecutionGraph - 物理执行。这条链路你只要记住后面所有调度相关的问题都从这里面展开。详细拆一下。用户写的 DataStream 或者 Table/SQL 代码经过编译后会先形成一个 StreamGraph这玩意儿本质上是逻辑层面的算子拓扑里面的节点叫 StreamNode边叫 StreamEdge。它描述的是“计算逻辑长什么样”跟并行度、资源还没完全挂钩。接着 Flink 会做一个关键操作把 StreamGraph 转成 JobGraph。这个转换过程做的事情很多最核心的是算子链优化Operator Chaining。两个相邻的算子如果能满足条件——比如分区方式相同、并行度相同、没有 keyBy 之类的重分区操作——就会被合并到一个 OperatorChain 里形成 JobVertex。这一步直接影响后续的资源分配合并得越多需要的 Slot 就越少网络传输开销也越低。实际生产中经常遇到“并行度调高了反而变慢”的怪现象很多时候就是算子链被拆散导致任务被调度到不同 TaskManager 后走了一遍序列化和反序列化。JobGraph 会被提交给 JobManager。到了这里真正的调度才开始。JobGraph 里的节点会继续被展开成 ExecutionGraph这是调度层面的核心数据结构。一个 JobVertex 展开成多个 ExecutionVertex每个 ExecutionVertex 对应一个并行子任务ExecutionVertex 之间通过 ExecutionEdge 相连代表数据流转关系。2.2 调度器的工作方式两种模式的取舍逻辑ExecutionGraph 生成之后接下来问题就变成了这些 ExecutionVertex 怎么分配到 TaskManager 的 Slot 上以及什么时候分配。在 Flink 1.5 之前调度逻辑写死在 JobManager 里调整起来非常费劲。从 1.5 开始引入了 Scheduler 接口并且默认实现了两种调度模式Eager Scheduling 和 Lazy From Sources Scheduling。Eager 模式下作业一提交到 JobManager调度器就会把整个 ExecutionGraph 里的所有 ExecutionVertex 一次性申请资源全部拿到 Slot 之后才开始部署任务。好处是作业启动之后不容易因为资源不足而中途失败适合对实时性要求高、作业规模不大的场景。坏处也很明显如果集群资源不够整个作业就完全起不来哪怕你的作业逻辑上可以先跑一部分任务。Lazy 模式正好相反它从 Source 节点开始先申请部分资源部署上游任务然后随着数据往下游推进逐步申请更多资源调度下游任务。这种模式适合大规模作业资源不够时至少源端能先跑起来但代价是作业启动慢且某个下游任务调度失败时恢复逻辑会更复杂。选哪种模式其实不用我们改配置但理解这一点对排查问题非常有帮助如果作业一直处于 CREATED 状态不执行先判断它卡在哪个环节——是资源申请阶段还是任务部署阶段——再对症下药。另外补充一句生产环境大多数 Streaming 作业默认用 Eager 模式因为流作业要求所有任务都在线。3. 资源调度的核心机制Slot 到底在调度什么3.1 先搞懂 Slot、TaskManager、资源粒度之间的关系很多人对 Flink 资源调度最大的误解是——以为 Slot 就是 CPU 核数。真不是。一个 TaskManager 是一个 JVM 进程它内部被划分成若干个 Slot每个 Slot 本质上是 TaskManager 内存资源的一个固定分片。两个 Slot 共用进程里的堆内存和 CPU只是通过线程来隔离任务执行这种隔离非常弱。默认情况下一个 TaskManager 的 Slot 数量等于它的 CPU 核数可以通过 taskmanager.numberOfTaskSlots 配置。这里的逻辑是每个 Slot 能跑一个线程一个 CPU 核在同一时刻大致能跑一个线程所以核数定 Slot 数。但这不代表 Slot 跟 CPU 做了绑定实际上一个 Slot 里的任务在运行时会用到整个 TaskManager 的所有 CPU 资源没有做硬隔离。所以如果一个 TaskManager 上有 8 个 Slot而其中 6 个 Slot 的任务都是 CPU 密集型的另外 2 个 Slot 的任务就会被拖慢。内存方面是分 Slot 计算的。每个 Slot 分到的内存主要是托管内存Managed Memory的一部分这部分用于 RocksDB 状态后端、排序缓冲等操作。调整 taskmanager.memory.process.size 或 taskmanager.memory.task.off-heap.size 会影响每个 Slot 可用的内存配额进而影响作业稳定性。明确一点Slot 是 Flink 做资源调度的最小单元但不是资源隔离的最小单元。容器化部署的时候真正的隔离是靠底层 YARN 或者 Kubernetes 的 Pod 来实现的。这个区别非常重要很多分布式系统里“调度单元”和“隔离单元”的职业病如果混在一起排查问题时会走很多弯路。3.2 从资源申请到 Slot 分配的内幕一个作业请求资源的过程是这样的JobManager 中的调度器发现某个 ExecutionVertex 需要执行会根据它所属的 JobVertex 所需要的资源默认每个 Slot 的资源规格由 TaskManager 的 Slot 数、内存大小推算而来向 ResourceManager 发起 Slot 请求。ResourceManager 是 Flink 跟底层资源管理系统YARN、K8s、Standalone打交道的中介。它拿到请求以后先看已注册的 TaskManager 列表里有没有可用的 Slot。有就直接分配没有的话它会向底层资源平台申请启动新的 TaskManager。这个“申请新 TaskManager”的过程是异步的、需要时间的。YARN 模式下要经过 ResourceManager - NodeManager - 启动容器 - 注册到 JobManager 一系列流程K8s 模式要经历创建 Pod、拉镜像、启动进程。我见过很多次生产故障就是作业在高峰期追加并行度资源不足然后 TaskManager 扩容需要几分钟期间作业一直处于等待资源的状态。所以如果你发现作业运行得慢不要只盯 CPU先看监控里的 Slot 申请耗时。Slot 分配成功之后TaskManager 会向 JobManager 确认然后 JobManager 把具体的 Task也就是 ExecutionVertex 的执行体部署上去。这里有一个细节Task 的部署是以 Task 为单位但资源占用是以 Slot 为单位。一个 Slot 里可以先后运行多个 Task因为一个 Task 执行完了Slot 就会被释放然后重新分配给新的 Task。Flink 通过 SlotSharingGroup槽位共享组机制允许来自不同 JobVertex 的 Task 共享同一个 Slot。默认情况下所有节点都属于同一个默认组这样整个作业最小只需要一个 Slot 就能跑完——当然前提是并行度也要适配。3.3 资源调优时最容易踩的坑围绕 Slot 这块踩坑经验堪称丰富。先说最常见的Slot 数量设置等于并行度总和。很多新手直接把所有算子并行度加起来设 Slot 数这是错的。因为默认所有算子都在同一个 SlotSharingGroup 里一个 Slot 里可以串行执行多个不同算子的 Task。正确的做法是——Slot 数只需要大于所有 Task 中并行度最大的那个算子的并行度就行通常是大于等于 max(parallelism)。第二个坑是内存配置。Slot 数调大以后如果没同步调大 TaskManager 内存每个 Slot 分到的内存反而会变小RocksDB 状态后端很快就给你抛内存溢出。节奏应该是先确定每个 Slot 需要多少内存再确定 TaskManager 总内存最后反推 Slot 数。顺序反了作业就等着 OOM。第三个坑是 Standalone 模式的资源隔离。Standalone 集群里的 TaskManager 是提前启动好的JobManager 只是在已有的 Slot 里做分配。这种模式下就算作业资源需求很高集群资源不够作业也只会卡在等待资源状态不会自动扩容。很多公司在测试环境用 Standalone 模式跑作业遇到资源不足时第一反应是“加并行度”结果越加越起不来就是因为忘了 Standalone 不自动扩容。4. 从调度到执行ExecutionGraph 展开背后那些事4.1 ExecutionGraph 到底长什么样前面提到了 JobGraph 转 ExecutionGraph这里展开讲透。JobGraph 中的 JobVertex 是逻辑节点它包含了并行度信息和算子链信息但并不知道自己会被拆成几个并行实例。到了 ExecutionGraph 这一层每个 JobVertex 会根据并行度展开为若干个 ExecutionVertex。举个例子一个 JobVertex 并行度是 4展开后得到 4 个 ExecutionVertex编号从 0 到 3。每个 ExecutionVertex 在执行时会创建一个 Execution 对象这个对象记录了任务的当前状态RUNNING、FINISHED、FAILED 等和尝试次数。ExecutionVertex 之间通过 ExecutionEdge 相连。这些边是逻辑连接指向数据应该从哪里来、到哪里去。真正决定数据怎么传输的是中间结果IntermediateResult和分区ResultPartition的概念每个 ExecutionVertex 产生的输出数据会写入到 ResultPartition 中下游的 ExecutionVertex 从上游的 ResultPartition 消费数据。调度器在分配 ExecutionVertex 时会考虑数据本地性。如果你的作业上游和下游都分配到同一个 TaskManager 上那么数据可以直接走内存管道Pipelined传输不需要落盘和网络序列化。如果被分配到不同的 TaskManager那数据就要通过网络传输。所以你会发现并行度调整之后作业变慢往往不是 CPU 不够而是数据本地性变差了、网络开销增加了。4.2 调度策略的源码级逻辑与选择标准刚才提到 Eager 和 Lazy 两种调度器这里再往深挖一层。在 Flink 源码里调度器实现的接口是 SchedulerNG真正干活的有几个核心类SchedulerBase 负责状态管理DefaultScheduler 是 Eager 模式的实现LazyFromSourcesScheduler 是 Lazy 模式的实现而 AdaptiveScheduler 是 1.15 之后主推的适应型调度器。AdaptiveScheduler 的逻辑值得多说一句它允许作业声明一个并行度范围比如 1 到 10调度器根据当前集群可用资源量自动决定并行度。这个东西的爽点是作业不会因为并行度固定而导致资源不足失败资源多时自动加并行度资源紧张时自动降并行度。代价是作业运行途中的并行度可能变化导致需要重新分发 Key 或者状态恢复。实际选型建议流式作业稳定性优先生产环境不怕资源浪费选 Eager 模式最可靠批式作业用 Flink Batch或者源端吞吐波动大的作业Lazy 模式更合适如果你用的是 Flink 1.15自适应调度器也可以试试但做好监控和告警因为并行度动态调整带来的状态迁移问题在低版本里有可能触发 bug。4.3 数据本地性优化它的实际意义数据本地性Data Locality在调度里是个容易被忽略但影响巨大的因素。调度器在给 ExecutionVertex 分配 Slot 时会尝试把它放到已经有上游数据缓存的 TaskManager 上。比如上游任务在 TM1 上产生了数据下游任务如果能分到 TM1那就能直接从本机内存读取避免网络传输。这在计算和存储分离的架构里尤其重要。但要注意Flink 的数据本地性不是绝对的“计算跟数据放一起”而是“计算跟上一条 Task 放一起”。对于流式作业上游算子的输出数据已经分不到哪儿去了下游算子分配在同一个 TM 能大大减少网络 I/O。所以并行度调整、Slot 分配策略都会直接反映到“本地执行占比”这个监控指标上。我排查作业性能问题时最先看的就是这个指标一旦本地执行占比掉到 60% 以下优先怀疑 Slot 分配和数据倾斜。5. 失败恢复作业挂了以后 Flink 都做了哪些事5.1 失败类型和重启策略的完整决策流程作业挂掉的原因千奇百怪网络抖动、第三方系统超时、OOM、JDBC 连接池被耗尽、K8s 把 Pod 杀了等等。Flink 面对失败不是一刀切重启它有一套完整的决策流程。首先要区分失败发生在哪个层面。如果发生在 JobManager 层面整个作业都玩完了需要 JobManager 重启这个过程依赖外部系统YARN/K8s的故障转移能力。如果发生在 Task 层面也就是某个 Execution 挂了JobManager 会尝试对这个 Task 进行重启。如果同一时间多个 Task 失败了或者一个 Task 连续失败次数超过阈值那就触发整个作业级别的重启。重启策略有三个选项固定延迟重启fixed-delay、失败率重启failure-rate、无重启none。固定延迟就是失败后等 N 秒再重启最多重启 M 次。失败率策略是在一个时间窗口内如果失败次数超过阈值就停止重启。默认情况下 Flink 用的是固定延迟重启Integer.MAX_VALUE 次延迟 1 秒——这意味着如果代码有 bug作业会一直重启循环。实际操作里我最常用的是失败率重启原因很简单如果一份代码有静态 bug那固定延迟重启多少次也没用反而会让 Namenode 一类的下游系统频繁闪断。失败率策略在窗口内超过阈值就自动放弃让告警发出来人工介入处理而不是集群里像个傻子一样反复重启。这个细节建议每个 Flink 项目的默认配置都改掉。5.2 Checkpoint 与 Exactly-Once 恢复链路Task 挂了之后重启只是第一步恢复状态才是关键。Flink 靠 Checkpoint 机制实现状态恢复。Checkpoint 会把算子状态 Source 位点打包存入远端存储HDFS、S3、RocksDB 增量一旦失败就从最近一次成功的 Checkpoint 恢复。整个恢复流程是这样的作业进入 FAILED 状态后调度器重新申请 Slot重新调度 ExecutionGraph然后让每个算子从 Checkpoint 里加载状态。Source 任务从这个 Checkpoint 里记录的位点重新消费数据内部任务恢复自己的状态整个作业继续跑。这期间有几个非常重要但容易被忽略的细节。第一个是恢复粒度默认是“全量恢复”也就是说只要有一个 Task 挂了所有 Task 全部重启并从同一个 Checkpoint 恢复。这在作业规模大的时候非常慢。Flink 1.14 开始支持局部恢复Local Recovery只有受影响的任务和它的上游任务会重启其他任务不受影响前提是你用了增量 Checkpoint 和可重缩放的状态分发策略。第二个是 Checkpoint 对齐问题。如果你的作业里有一个算子处理速度特别慢会导致整个 Checkpoint 的 barrier 流转变慢Checkpoint 超时然后频繁失败最后作业挂掉。这也是生产环境最常见的失败恢复触发原因。排查思路就是看 Checkpoint 监控里的 Alignment 耗时和 Skew 指标。第三个是状态 TTL。如果 Checkpoint 里的状态数据很大恢复时间会非常长极端情况下作业重启后长时间处于 INITIALIZING 状态。建议对非核心状态设置 TTL能省去大量恢复时间。5.3 重启策略、Slot 分配与恢复失败的联动问题有一个场景值得单独拿出来讲因为它综合了调度和恢复两个主题。假设某个 TaskManager 所在节点宕机了上面有 4 个 Slot各跑着不同作业的 Task。这时候 ResourceManager 会把坏掉的 TaskManager 从列表中移除JobManager 感知到 Slot 失去后会让受影响的任务进入重启流程并重新申请 Slot——新的 Slot 大概率落在另一台 TaskManager 上。但这里有个坑如果新的 Slot 资源规格跟旧 Slot 不一致比如内存变小了根据 Task 的内存需求重新计算后有可能会导致 Slot 数量不够作业就一直处于 WAITING_FOR_RESOURCES 状态看起来像“卡死”其实是在等资源。这种情况下如果你看日志会看到类似 “Not enough free slots available” 的警告。对应的解法就是保证 TaskManager 规格统一或者开启自适应调度。还有一个坑是 Zookeeper 或者 HA 模式下JobManager 本身挂了之后新的 JobManager 恢复作业时需要重新获取所有 TaskManager 的注册信息。如果 TaskManager 数量多、注册慢作业恢复时间会很长。有时候你看到作业已经从 SUCCEEDED 变成 FAILED 了其实不是代码问题而是 JobManager 在恢复过程中给了个保守的失败判定。6. 实战排查指南从现象定位到根因6.1 作业一直 PENDING/SUBMITTED 不动怎么查这是群里问的最多的问题之一。作业提交后一直停在 SUBMITTED 或者 PENDING不进入 RUNNING。按我排查的顺序来。第一步打开 Flink Web UI看 Job 的“Task Manager”数量和“Slots”指标。如果 Slot 数为 0说明 ResourceManager 还没有可用的 TaskManager要么是资源平台层面没分配下来要么是 TaskManager 挂了。第二步看 JobManager 日志重点找 “Requesting new TaskManager” 和 “Slot request” 相关关键字。如果一直循环请求但没结果大概率是底层资源不够或者资源平台拒绝分配YARN 队列满了K8s 的 ResourceQuota 超限。第三步如果 Slot 有剩余但作业就是调度不上去那就要检查 ExecutionGraph 的调度进度。看哪些 ExecutionVertex 处在 SCHEDULED 状态但没变 RUNNING可能是它等待的上游中间结果没就绪批作业常见或者是 SlotSharingGroup 配置有问题导致调度器没法匹配。我在 K8s 环境里踩过一次因为某个 Pod 的资源请求requests设置太高超过了节点可分配量导致 Slot 永远申请不下来。6.2 任务反复失败重启如何看日志快速定段如果任务总是失败 - 重启 - 再失败 - 再重启无外乎几种原因代码里有不可恢复的异常、上下游系统连接问题、资源被打爆。先说代码异常看 TaskManager 日志里 Exception 类型如果是 NullPointerException、IllegalArgumentException 这类基本是业务逻辑写错了重启多少次都没用。这时候你的重启策略如果还是固定延迟无限次那就是灾难建议马上改成 failure-rate 并设置较低的阈值。再说连接问题热词里提到 flink 的 jdbc 连接器异常这是典型的例子。JDBC 连接器在连接池耗尽、数据库重启、驱动版本不匹配时都会抛异常。这类问题的特点是——重试一段时间后可能自己恢复。遇到这种情况重启策略配 fixed-delay但重试次数不要太高延迟时间比如 30 秒要足够给下游恢复时间同时要同步排查连接池配置max connections、validation timeout 都要调。最后是资源问题如果 TaskManager 日志里有 OutOfMemoryError或者 GC 日志显示频繁 Full GC那要立即看内存配置。注意区分堆内和堆外内存堆外内存不足会直接抛 “Direct buffer memory”堆内不足会抛 “Java heap space”。两者的解法不同前者调 taskmanager.memory.task.off-heap.size后者调 taskmanager.memory.process.size 和 JVM 堆大小。6.3 从热词“openmetadata 获取 flink 血缘关系”说开去最近 openmetadata 获取 flink 血缘关系这个话题挺火这里顺便说一句。数据血缘本质上就是从 Flink 作业解析出数据来源、处理链路、输出目标之间的关系。OpenMetadata 这类工具通常有两种接入方式一种是从 Flink Catalog 和 SQL 的解析结果中抽取血缘另一种是靠 Flink 内置的 Lineage 机制开源生态里目前在演进中或者外部解析器读取 Plan。这个场景跟调度和恢复有什么关系关系在于Flink 作业恢复之后血缘关系能不能继续保持准确取决于你的作业拓扑是否稳定。如果用了自适应调度作业运行过程中的并行度会变但血缘是逻辑层面的跟并行度无关所以血缘关系本身不会断。真正会让血缘出错的情况是同一个作业代码被多个环境共用或者通过 Java/Scala API 动态拼接算子——这个时候 OpenMetadata 抓到的血缘可能对不上实际的执行 Plan。建议是不管用不用 OpenMetadata都要在作业开发规范里强调“用 SQL 或 Table API 写逻辑避免动态生成算子链”。这不止是血缘的问题更关系到调度稳定性和恢复后的状态兼容性。动态生成算子链在极端情况下会导致 JobGraph 结构不稳定一旦任务失败恢复时执行图不一致状态恢复就会出现二次故障。6.4 面试里调度与恢复最常被问的几个题既然热词里有“flink 面试题”我顺手把调度和恢复这块高频问题整理一遍。这些问题不是死记硬背能过关的全都需要你理解机制。第一问Flink 的 Slot 是如何分配的答案里必须包含Slot 是 TaskManager 内部分配的最小资源单元SlotSharingGroup 决定哪些 Task 可以共享 Slot调度器根据 ExecutionVertex 的状态分阶段申请和分配 Slot这个分配又分 Eager 和 Lazy 两种模式。第二问一次 Task 失败后Flink 如何恢复这个问题的完整回答要覆盖失败检测TaskManager 心跳/异常回调- 调度器感知失败 - 重启策略决策固定延迟/失败率- 从 Checkpoint 加载状态 - 重新调度和执行。同时要说明如果是不可恢复异常比如 JobManager 挂了那就要靠外部 Failover 机制。第三问并行度怎么调才能让作业跑得快这个问题高频考的就是“Slot 数不是越多越好”。并行度高任务之间数据传输和协调开销也高。而且并行度增加会带来状态重新分组代价很大。合理的做法是根据数据量、单并行度处理能力和资源规格来定不要拍脑袋。第四问Checkpoint 失败会导致作业停止吗答案是不一定。如果 Checkpoint 持续失败最终会触发作业失败但在失败之前作业可能还在正常处理数据。你可以配置 decline 的容忍次数或者打开 unaligned checkpoint 来避免 barrier 对齐造成反压。核心理解是Checkpoint 失败是症状不是根因排查时得找它为什么失败。第五问状态恢复时 RocksDB 和 Heap 状态后端有什么不同这个问题看似磕碜但实际踩过的人才知道关键区别——RocksDB 的恢复是先加载本地增量文件再异步补全远端 Checkpoint 数据恢复启动快但瞬时 CPU 和磁盘 I/O 很高Heap 状态直接把数据加载到堆内存启动快但容易 OOM。选型时要结合状态大小和恢复时间要求。7. 关于 Flink 调度和恢复我的最终建议坦白说Flink 的调度和恢复机制这篇文章只能算剥了一层皮真要搞到源码级理解还得自己动手去跟踪一条 JobGraph 的完整生命周期。但作为一线使用的人你不需要把每个类名都背下来你需要的是建立完整的因果链路意识作业提交流程讲的是资源申请资源申请挂在调度器上调度器执行依赖 SlotSlot 不够就等资源或扩容任务跑起来以后监控数据流向数据积压导致 Checkpoint 慢Checkpoint 慢导致恢复失败恢复失败触发重启策略重启策略选错了就反复重启拖垮集群。链条上的每一环都有对应的配置、指标和日志可以观察。我自己在实际排查中最大的体会是——大多数调度问题都不是调度器本身的锅而是配置和设计的失衡。比如资源规格不统一、并行度设置不合理、重启策略选错、状态后端没配置好。Flink 很灵活但这种灵活是把双刃剑它把很多选择权交到你手里同时也把所有出错的代价交到你手里。最后分享一个我自己的小习惯每次新建一个 Flink 作业我都会在提交之前把一张检查表过一遍。检查表包含作业并行度是多少、每个 Slot 内存多大、状态后端用的是什么、Checkpoint 间隔和超时时间是否合理、重启策略是不是失败率模式、最大重启次数是不是超过 3 次、本地恢复开没开。这套流程跑下来作业上线之后调度和恢复层面踩雷的概率会小很多也希望你拿去就能用。
返回列表