Fl一文吃透 Flink Jobs and Scheduling从资源调度到失败恢复
一、为什么要理解 Flink 的 Jobs and Scheduling
很多人刚接触 Flink 时,会把它理解成“提交一个 Jar,然后集群帮我跑起来”。
但实际上,Flink 在运行一个作业时,内部会做很多复杂工作:
- 解析数据流图
- 计算并行度
- 划分任务执行单元
- 分配 Slot
- 构建并行执行图
- 跟踪每个子任务状态
- 处理失败、取消、重启、恢复
这些动作并不是“附属功能”,而是 Flink 能够支撑大规模实时计算的基础能力。
你可以把它理解成这样:
- JobGraph 像是“施工蓝图”
- ExecutionGraph 像是“真正施工时拆分到每一个工位的执行计划”
- Task Slot 像是“工位资源”
- Job 状态机 像是“项目整体进度”
- Task 状态机 像是“每个施工工人的当前工作状态”
理解这套机制后,你再去看 Flink Web UI、日志、失败栈、重启过程,就不会只停留在“任务挂了”这种表层认知,而是能清楚知道:
是谁在运行、运行在哪、为什么没被调度、失败后会不会自动恢复、当前到底是 Job 失败了还是某个 Execution 在重试。
二、Flink 调度的核心起点:Task Slot
Flink 的执行资源,是通过 Task Slot 来定义的。
简单来说:
- 一个 TaskManager 可以配置一个或多个 Task Slot
- 一个 Slot 可以运行一个并行任务流水线(pipeline)
- 这个流水线里,可能包含多个可以链式或共享资源执行的连续任务
也就是说,Slot 并不是只能跑单个算子,而是可以承载一条由多个子任务组成的执行链路。
1. 什么是 Pipeline
在 Flink 中,一个 pipeline 可以看作是一组可以连续执行的并行任务组合。
例如有这样一个数据处理链路:
- Source
- Map
- Reduce
如果 Source 和 Map 的并行度是 4,Reduce 的并行度是 3,那么 Flink 会根据算子之间的数据交换关系、并行度和调度约束,去决定哪些任务实例可以组成一个 pipeline,并把它们分配到不同 Slot 中执行。
这就是为什么在 Flink 中我们经常会看到:
- 某些算子能够串在一起执行
- 某些算子必须经过网络 shuffle
- 某些任务可以共享 Slot
- 某些任务必须单独占用资源
2. SlotSharingGroup 与 CoLocationGroup
Flink 内部通过两类机制控制任务的 Slot 使用策略:
SlotSharingGroup
它定义的是:
哪些任务“允许”共享同一个 Slot
这是一个比较“宽松”的约束。只要属于同一个 SlotSharingGroup,理论上这些任务就可以被调度到同一个 Slot 内,从而提升资源利用率。
这在流处理任务中非常常见,因为上下游链式执行可以显著减少网络传输和调度开销。
CoLocationGroup
它定义的是:
哪些任务“必须”严格放在同一个 Slot 中
这比 SlotSharingGroup 更强,是一种强绑定关系。通常用于一些必须保持位置对应关系的场景,比如迭代计算、严格一一对应的上下游任务。
3. 图示:任务流水线如何分配到 Slot

你可以在博客中配上类似说明:
图中展示了在 2 个 TaskManager、每个 TaskManager 3 个 Slot 的情况下,Source、Map、Reduce 不同并行实例如何组合成 pipeline,并最终映射到不同的 Slot 上执行。
4. 这对生产环境意味着什么
很多线上问题,本质上都和 Slot 调度有关:
- 作业提交后一直处于等待调度状态
- 某些任务明明资源够,但就是起不来
- 算子链路太长,某个 Slot 压力过大
- SlotSharing 设置不合理,导致资源竞争严重
- 某些批任务因为上下游不能共享 Slot,导致资源需求突然膨胀
所以,理解 Slot,并不是理解一个“名词”,而是理解 Flink 资源调度的第一层入口。
三、JobManager 眼中的作业:从 JobGraph 到 ExecutionGraph
当用户提交一个 Flink 作业时,JobManager 并不是直接拿着用户代码执行,而是先构建和维护一套内部数据结构,用来管理整个作业的运行过程。
这其中最重要的两个结构,就是:
- JobGraph
- ExecutionGraph
很多初学者容易把这两个概念混淆。其实它们分别对应的是两个不同层次的视角。
四、JobGraph:逻辑层的数据流表示
JobManager 接收到的,是一个 JobGraph。
你可以把 JobGraph 理解为:
一个描述“这个作业要做什么”的逻辑执行图
它主要由两类元素构成:
- JobVertex:表示算子节点
- IntermediateDataSet:表示算子之间的中间结果集
每个 JobVertex 会携带这个算子的关键信息,例如:
- 算子代码
- 并行度
- 算子配置
- 依赖关系
此外,JobGraph 还会包含作业运行所需要的附加库,也就是执行这些算子所依赖的 Jar、Class、资源等。
1. JobGraph 的特点
JobGraph 更偏向“逻辑视图”,它关注的是:
- 作业有哪些算子
- 算子之间如何连接
- 每个算子的并行度是多少
- 数据是如何在算子之间流动的
但 JobGraph 还没有真正展开到“每个并行子任务实例”。
比如:
- 一个并行度为 100 的算子
- 在 JobGraph 中仍然只是一个 JobVertex
这也是 JobGraph 和 ExecutionGraph 最大的区别所在。
五、ExecutionGraph:真正用于执行的并行化运行图
JobManager 会把 JobGraph 转换成 ExecutionGraph。
ExecutionGraph 可以理解为:
JobGraph 的并行化、运行时版本
如果说 JobGraph 是“设计图”,那么 ExecutionGraph 就是“施工展开图”。
1. ExecutionVertex:每个并行子任务都有独立状态
在 ExecutionGraph 中:
- 每个 JobVertex 会被展开成多个 ExecutionVertex
- 每个 ExecutionVertex 对应一个并行子任务实例
例如:
- 一个算子并行度为 100
- JobGraph 中只有 1 个 JobVertex
- ExecutionGraph 中会有 100 个 ExecutionVertex
这些 ExecutionVertex 才是真正被调度、运行、失败、恢复的最小执行单元。
2. ExecutionJobVertex:算子整体视图
所有属于同一个 JobVertex 的 ExecutionVertex,会被组织在一个 ExecutionJobVertex 中。
它的作用相当于:
从“整个算子”的维度,统一跟踪所有并行子任务的状态
也就是说:
- JobVertex 是逻辑定义
- ExecutionJobVertex 是运行时的算子实例集合
- ExecutionVertex 是真正的单个并行子任务
3. IntermediateResult 与 IntermediateResultPartition
除了 Vertex 以外,ExecutionGraph 中还会维护中间结果相关的运行时对象:
- IntermediateResult:跟踪一个中间结果集的整体状态
- IntermediateResultPartition:跟踪每一个分区的状态
为什么要分这么细?
因为 Flink 的上游结果往往不是一个整体文件,而是按照并行度切分成多个 partition。下游在消费数据时,也不是“拿一个完整结果”,而是会消费这些 partition。
这对于以下场景尤其重要:
- shuffle 数据交换
- 批任务阶段性调度
- 失败恢复时的分区重用或重新生成
- 上下游依赖关系判断
4. 图示:JobGraph 与 ExecutionGraph 的关系

配图说明建议这样写:
左侧是逻辑层的 JobGraph,包含 A、B、C、D 等 JobVertex 以及中间数据集;右侧是展开后的 ExecutionGraph,每个 JobVertex 被拆分成多个 ExecutionVertex,中间结果也被细化为多个 IntermediateResultPartition。
5. 为什么这部分很关键
很多人查线上问题时,只会看“这个算子失败了”。
但 Flink 真正运行的时候,失败的并不是抽象意义上的“算子”,而是某个:
- ExecutionVertex
- 对应当前一次具体尝试的 Execution
这意味着你在分析问题时要能区分:
- 是整个 JobVertex 出问题
- 还是某个并行子任务出问题
- 是某个 partition 卡住了
- 还是整个上游 IntermediateResult 没准备好
理解了 ExecutionGraph,很多 Flink Web UI 里的执行视图和日志信息就会一下子清晰起来。
六、Flink Job 的全局状态机:作业从提交到终止经历了什么
每个 ExecutionGraph 都会对应一个 JobStatus,也就是整个作业的全局状态。
这套状态机描述的是:
一个 Flink Job 从被创建开始,到运行、失败、取消、重启、结束的完整生命周期。
1. 正常执行路径
一个 Flink Job 最正常的状态流转路径是:
- CREATED
- RUNNING
- FINISHED
含义分别是:
CREATED
作业已经创建,但还没有真正开始稳定运行。
这个阶段更多是作业刚被接收、图结构建立、准备进入执行流程的状态。
RUNNING
作业已经进入实际执行阶段,任务正在运行。
FINISHED
所有工作都已经完成,作业正常结束。
这个状态是全局终态之一。
2. 失败路径
当作业运行过程中出现故障时,Job 并不是立刻变成 FAILED,而是会先经历一个中间阶段:
- FAILING
- FAILED 或 RESTARTING
这点很重要。
FAILING
进入 FAILING 状态时,Flink 会先取消所有还在运行的任务。
也就是说,这个阶段本质上是在做“失败收敛”。
FAILED
如果当前作业不可重启,或者已经不满足重启条件,那么当所有相关任务都进入终态后,作业最终进入 FAILED。
这是一个全局终态。
RESTARTING
如果当前作业配置了可用的重启策略,并且仍然允许重启,那么作业不会直接 FAILED,而是进入 RESTARTING。
待整个作业完成重启流程后,又会重新回到 CREATED,然后再次进入 RUNNING。
这也是为什么你有时候会在 Web UI 中看到 Job 在:
- FAILING
- RESTARTING
- CREATED
- RUNNING
之间来回切换。
3. 用户取消路径
如果不是系统失败,而是用户主动取消任务,那么作业状态会走另一条路径:
- CANCELLING
- CANCELED
CANCELLING
正在取消中。
Flink 会停止并清理所有正在运行的任务。
CANCELED
所有相关任务都已经完成取消,作业进入取消完成状态。
这同样是一个全局终态。
4. 特殊状态:SUSPENDED
除了 FINISHED、CANCELED、FAILED 这些全局终态外,Flink 还有一个很特别的状态:
- SUSPENDED
这个状态不是全局彻底结束,而是局部终态(locally terminal)。
什么意思?
它表示:
当前 JobManager 上,这个作业已经终止执行,但作业并没有被完全清理掉。
因为在开启 HA(高可用)时,其他 JobManager 仍然可以从持久化的 HA 存储中恢复这个作业,并重新启动它。
所以:
- FAILED / FINISHED / CANCELED:通常意味着全局结束,触发作业清理
- SUSPENDED:只是当前 JobManager 本地结束,不代表整个作业彻底消失
这个概念对于理解 Flink HA 模式下的主备切换特别关键。
5. 图示:Flink Job 状态流转图

你可以在图下方补一句说明:
该图展示了 Flink Job 从 CREATED、RUNNING 到 FINISHED 的正常路径,以及 FAILING、RESTARTING、CANCELLING、CANCELED、FAILED、SUSPENDED 等异常和控制类状态之间的转换关系。
七、为什么 Job 状态机很重要
理解 Job 状态机后,你在生产环境中看任务状态时就不会只停留在“运行中”和“失败了”。
例如:
1. 为什么任务不是直接 FAILED,而是先 FAILING
因为 Flink 不是简单地“发现异常就退出”,而是要先把所有相关运行中的任务收敛到一个可控状态,避免系统处于半执行、半失败的混乱中。
2. 为什么有时候任务失败后又自己好了
因为任务可能启用了重启策略,Job 在 FAILED 之前先走了 RESTARTING,恢复成功后重新进入 RUNNING。
3. 为什么 HA 场景下任务没了但又能回来
因为作业可能进入的是 SUSPENDED,而不是 FAILED。
SUSPENDED 说明当前 JobManager 不继续跑了,但作业元信息还在,其他 JobManager 可以接手恢复。
八、Task 级别的生命周期:Execution 才是最细粒度的执行单元
如果说 JobStatus 描述的是整个作业的状态,那么 Task 级别的状态机描述的就是:
一个具体并行子任务,是如何被创建、部署、初始化、运行、取消、失败、结束的。
这部分非常重要,因为真正发生故障时,最先出问题的往往不是整个 Job,而是某个具体的 Task Execution。
1. ExecutionVertex 与 Execution 的关系
前面说过:
- ExecutionVertex:表示一个并行子任务实例
- Execution:表示这个子任务某一次具体执行尝试
为什么要引入 Execution?
因为一个子任务可能不是只执行一次。
比如:
- 第一次执行时失败了
- 触发重启
- 第二次重新执行
- 再失败
- 第三次恢复成功
此时:
- ExecutionVertex 还是同一个逻辑子任务
- 但它会关联多个不同历史时期的 Execution
所以:
Execution 是 Flink 用来跟踪“某个子任务某一次具体执行过程”的对象。
这也是 Flink 故障恢复机制能细粒度追踪尝试次数的基础。
九、Task Execution 的状态流转过程
在运行过程中,一个 Execution 通常会经历以下状态:
- CREATED
- SCHEDULED
- DEPLOYING
- INITIALIZING
- RUNNING
- FINISHED
如果发生异常或人为干预,则还可能进入:
- CANCELING
- CANCELED
- FAILED
1. CREATED
Execution 刚刚被创建,还未调度到具体资源上。
2. SCHEDULED
已经完成调度决策,准备分配资源并下发执行。
3. DEPLOYING
正在将任务部署到对应的 TaskManager 上。
4. INITIALIZING
任务已经开始初始化运行环境,比如:
- 恢复状态
- 初始化算子实例
- 建立输入输出通道
- 做启动前准备
5. RUNNING
真正进入运行状态,开始处理数据。
6. FINISHED
执行正常结束。
这是该次 Execution 的终态之一。
7. CANCELING / CANCELED
如果任务被取消,会先进入 CANCELING,再进入 CANCELED。
8. FAILED
如果任务在执行过程中发生异常,进入 FAILED。
这也是单次 Execution 的终态之一。
十、图示:Task Execution 的完整状态机
这里建议插入你提供的第四张图。

图下说明可以这样写:
该图展示了一个 Task Execution 从创建、调度、部署、初始化到运行、完成的正常路径,以及在取消、失败等异常情况下的状态流转过程。由于一个 ExecutionVertex 可以多次执行,因此同一个子任务会经历多个不同的 Execution 实例。
十一、Job 状态机和 Task 状态机有什么区别
很多人第一次看 Flink 状态图时会疑惑:
Job 不是已经有状态了吗?为什么 Task 还有一套状态?
这里一定要区分清楚:
Job 状态机
关注的是:
- 整个作业当前所处的全局阶段
比如:
- RUNNING
- FAILING
- RESTARTING
- CANCELED
Task 状态机
关注的是:
- 单个子任务当前执行到哪一步了
比如:
- DEPLOYING
- INITIALIZING
- RUNNING
- FAILED
可以这样理解:
- Job 状态 是“项目总进度”
- Task 状态 是“每个工位上的工人干到哪一步了”
一个 Job 进入 FAILING,通常意味着其中某些 Task Execution 已经 FAILED,JobManager 正在协调整个作业进入统一失败处理流程。
十二、一个完整的故障恢复过程是怎样发生的
为了帮助大家把这些概念串起来,我们来看一个典型场景。
假设某个 Flink 流任务正在运行,某个下游算子其中一个并行子任务突然抛出异常。
此时内部大致会发生这样的事情:
第一步:某个 Execution 失败
- 某个 ExecutionVertex 当前关联的 Execution 进入 FAILED
- JobManager 收到失败事件
第二步:Job 进入 FAILING
- JobManager 判断这是影响整个作业一致性的故障
- ExecutionGraph 对应的 JobStatus 切换到 FAILING
- 系统开始取消其他相关运行中任务
第三步:判断是否可重启
如果配置了合适的 Restart Strategy,且未超过阈值:
- Job 进入 RESTARTING
否则:
- Job 进入 FAILED
第四步:重新创建执行流程
如果可重启:
- Job 完成失败清理
- 重新回到 CREATED
- 重新调度
- Task 重新经历 SCHEDULED、DEPLOYING、INITIALIZING、RUNNING
第五步:恢复成功或再次失败
如果恢复成功:
- Job 回到 RUNNING
如果再次失败:
- 继续重复失败恢复链路,直到成功或不再允许重启
这就是为什么你在实际日志里可能会看到:
- 某个 task failed
- job switched from RUNNING to FAILING
- restarting job
- reset execution graph
- scheduling tasks
- deploying task
- task switched to RUNNING
这些日志如果单独看会很碎,但你一旦掌握 JobGraph、ExecutionGraph、JobStatus、Execution 状态机,就会发现整个过程其实非常清晰。
十三、这些底层机制对排障有什么帮助
理解这套调度和状态机制后,你在线上排查问题时会非常有用。
1. 作业长时间不运行
优先看:
- Job 是否一直停留在 CREATED / SCHEDULED
- 是否存在 Slot 不足
- 是否有 SlotSharingGroup / CoLocationGroup 导致调度受限
- 是否某些上游中间结果未准备完成
2. 作业频繁重启
优先看:
- Job 是否在 RUNNING → FAILING → RESTARTING 之间循环
- 是哪个 ExecutionVertex 反复失败
- 是状态恢复失败,还是业务算子异常
- 是否超过重启策略阈值
3. 某个算子看起来“偶发挂掉”
不要只看 JobVertex 名称,要继续下钻:
- 是哪个并行子任务失败
- 当前是第几次 Execution
- 是否与某个特定 partition、数据分片、节点资源有关
4. HA 场景下任务异常消失
要区分:
- 是 FAILED
- 还是 SUSPENDED
如果是 SUSPENDED,说明还有可能被其他 JobManager 接管恢复。
十四、面向生产环境的几个理解建议
1. 不要只从“算子视角”看 Flink
很多人写代码时脑子里只有:
- Source
- Map
- KeyBy
- Window
- Sink
但运维和排障时,你必须切换到:
- JobGraph
- ExecutionGraph
- ExecutionVertex
- IntermediateResultPartition
- JobStatus
- Execution State
这才是 Flink 真正运行时的世界。
2. Web UI 上看到的不是“抽象任务”,而是运行时结构的映射
你在 Flink UI 中看到的:
- Job 状态
- Vertex 视图
- 子任务视图
- 并行实例
- 失败重启次数
本质上都是 JobManager 内部这些运行时数据结构的外在展示。
3. 调度问题、失败恢复问题、性能问题,本质都和 ExecutionGraph 有关
因为真正被调度的是它,真正失败和恢复的也是它。
十五、总结:理解 Flink 执行机制,才能真正驾驭生产任务
很多人学习 Flink 时,只关注 API 怎么写、SQL 怎么调优,但一旦进入生产环境,你会发现真正决定你技术深度的,是你是否理解 Flink 的运行时机制。
这篇文章我们系统梳理了几个核心点:
- Task Slot 是 Flink 资源调度的基础单位
- Pipeline 表示可在同一执行链中运行的一组并行任务
- SlotSharingGroup 控制哪些任务可以共享 Slot
- CoLocationGroup 控制哪些任务必须严格同槽部署
- JobGraph 是逻辑层执行图
- ExecutionGraph 是并行化后的运行时执行图
- ExecutionVertex 是真正的并行子任务实例
- Execution 表示某个子任务的一次具体执行尝试
- Job 状态机 描述整个作业的生命周期
- Task 状态机 描述每个子任务执行实例的生命周期
一句话概括就是:
JobGraph 决定“你要跑什么”,ExecutionGraph 决定“具体怎么跑”,而 Job/Task 状态机则决定“跑到哪了、出问题怎么恢复”。
如果你能把这几层真正吃透,那么无论是看 Flink 源码、分析 Web UI、排查线上故障、设计高可用方案,都会轻松很多。
更多推荐
所有评论(0)