一、为什么要理解 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
  • FAILEDRESTARTING

这点很重要。

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、排查线上故障、设计高可用方案,都会轻松很多。

Logo

腾讯云面向开发者汇聚海量精品云计算使用和开发经验,营造开放的云计算技术生态圈。

更多推荐