Flink运行时架构详解
Flink运行时架构详解:JobManager、TaskManager、Slot与并行度
前言
学Flink,运行时架构是必须搞清楚的基础。只有真正理解了JobManager、TaskManager、Slot、并行度这些概念之间的关系,后面学任务提交、状态管理、Checkpoint才不会一头雾水。本文基于尚硅谷Flink1.17教程,系统梳理Flink运行时架构的核心知识点。
一、系统架构
Flink集群由两类进程组成:JobManager(主进程)和TaskManager(工作进程)。
1.1 JobManager
JobManager是集群中任务管理和调度的核心,控制应用的执行。每个应用都有一个唯一对应的JobManager。
JobManager内部包含三个组件:
(1)JobMaster
JobMaster是JobManager中最核心的组件,与具体的Job一一对应。多个Job同时运行时,每个Job都有自己的JobMaster。
JobMaster的工作流程:
- 接收提交的应用(JobGraph)
- 将JobGraph转换为执行图(ExecutionGraph)——包含所有可以并发执行的任务
- 向ResourceManager申请执行所需的资源(Slot)
- 获取到足够资源后,将执行图分发给TaskManager执行
- 运行过程中负责Checkpoint等中央协调操作
注意:早期版本的Flink中没有JobMaster的概念,当时的JobManager实际上就是现在的JobMaster。
(2)ResourceManager
ResourceManager负责集群中资源的分配和管理,整个集群中只有一个。
这里的"资源"主要指TaskManager的任务槽(Task Slots)。每一个Task都需要分配到一个Slot上执行。
注意:要区分Flink内置的ResourceManager和YARN等资源管理平台的ResourceManager,它们是不同层面的概念。
(3)Dispatcher
Dispatcher的职责:
- 提供REST接口,用于提交应用
- 为每个新提交的作业启动一个新的JobMaster组件
- 启动Web UI,展示和监控作业执行信息
Dispatcher并不是必需组件,在不同的部署模式下可能被忽略。
1.2 TaskManager
TaskManager是Flink的工作进程,负责数据流的具体计算。
- 集群中至少需要一个TaskManager
- 每个TaskManager包含一定数量的任务槽(Task Slots)
- Slot数量限制了TaskManager能并行处理的任务数量
TaskManager的工作流程:
- 启动后向ResourceManager注册自己的Slot
- 收到ResourceManager指令后,将Slot提供给JobMaster调用
- JobMaster分配任务后,TaskManager执行具体计算
- 执行过程中可以缓冲数据,也可以与其他TaskManager交换数据
二、核心概念
2.1 并行度(Parallelism)
什么是并行度
当数据量很大时,可以把一个算子"复制"多份到多个节点上并行处理。一个算子的子任务(subtask)个数就称为该算子的并行度(Parallelism)。
一个流程序的并行度 = 所有算子中最大的并行度。
举例:
source(并行度2)→ map(并行度2)→ window(并行度2)→ sink(并行度1)
程序并行度 = 2

并行度的设置(优先级从高到低)
方式一:代码中对单个算子设置(优先级最高)
stream.map(word -> Tuple2.of(word, 1L)).setParallelism(2);
方式二:代码中全局设置
env.setParallelism(2);
不推荐硬编码全局并行度,会导致无法动态扩容。
方式三:提交作业时通过-p参数设置
bin/flink run -p 2 -c com.atguigu.wc.SocketStreamWordCount ./FlinkTutorial-1.0-SNAPSHOT.jar
方式四:配置文件中设置(优先级最低)
# flink-conf.yaml
parallelism.default: 2
对整个集群所有作业有效,初始值为1。开发环境中没有配置文件,默认并行度 = 当前机器CPU核心数。
注意:keyBy不是算子,无法设置并行度。
2.2 算子链(Operator Chain)
算子间的数据传输方式
一对一(One-to-one / Forwarding):数据流维护分区和元素顺序,不需要重新分区。map、filter、flatMap等算子都是一对一关系,类似Spark的窄依赖。
重分区(Redistributing):数据流分区发生改变,比如keyBy、window之后的操作,类似Spark的Shuffle。
什么是算子链
满足条件:并行度相同 + 一对一的数据传输关系
满足条件的相邻算子可以合并成一个Task,由一个线程执行,这就是算子链(Operator Chain)。
优势:
- 减少线程间的切换开销
- 减少基于缓冲区的数据交换
- 降低延迟,提升吞吐量
Flink默认会按算子链原则自动合并,也可以手动控制:
// 禁用当前算子的链接
.map(word -> Tuple2.of(word, 1L)).disableChaining();
// 从当前算子开始新的算子链
.map(word -> Tuple2.of(word, 1L)).startNewChain();
算子间数据传输
合并算子链
2.3 任务槽(Task Slots)
什么是任务槽
每个TaskManager是一个JVM进程,可以启动多个线程并行执行子任务。为了控制并发量,对每个任务运行所占用的资源进行明确划分,这就是任务槽(Task Slot)。
每个Slot表示TaskManager计算资源的一个固定大小的子集,用于独立执行一个子任务。
注意:Slot目前只隔离内存,不隔离CPU。
任务槽数量设置
# flink-conf.yaml
taskmanager.numberOfTaskSlots: 8
默认为1。建议配置为机器的CPU核心数,避免任务间CPU竞争。
Slot共享
默认情况下,同一个作业的不同算子的子任务可以共享同一个Slot。
好处:
- 资源密集型和非密集型任务可以在同一个Slot中自行分配资源占用比例,使负载平均分配
- 即使某个TaskManager宕机,其他节点不受影响,作业可以继续执行(保存完整的作业管道)
手动指定Slot共享组(只有同组的子任务才共享Slot):
.map(word -> Tuple2.of(word, 1L)).slotSharingGroup("group1");
使用Slot共享组后,所需总Slot数 = 各共享组最大并行度之和。
2.4 任务槽与并行度的关系
这是容易混淆的概念,一定要区分清楚:
| 概念 | 类型 | 含义 | 配置参数 |
|---|---|---|---|
| 任务槽(Slot) | 静态概念 | TaskManager具有的并发执行能力(上限) | taskmanager.numberOfTaskSlots |
| 并行度(Parallelism) | 动态概念 | 程序运行时实际使用的并发能力 | parallelism.default |
举例:3个TaskManager,每个TM有3个Slot,共9个Slot。
- 程序为:
source → flatmap → reduce → sink - source和flatmap满足算子链条件,合并为一个Task
- 最终3个Task节点,若并行度设为9,则需要9个Slot
关系总结:Slot是资源的上限,并行度是实际使用量。并行度不能超过可用Slot总数。
三、作业提交流程与图的转换
3.1 四层图的转换
一个Flink作业从代码到实际执行,要经历四次图的转换:
用户代码
↓
逻辑流图(StreamGraph)
↓ 算子链优化
作业图(JobGraph)
↓ 按并行度拆分
执行图(ExecutionGraph)
↓ 分发到TaskManager
物理图(Physical Graph)


①逻辑流图(StreamGraph)
根据用户DataStream API代码生成的初始DAG图,表示程序的拓扑结构。一般在客户端生成。
②作业图(JobGraph)
StreamGraph经优化后的结果,是提交给JobManager的数据结构。主要优化:将符合条件的算子合并成算子链,减少数据交换消耗。一般也在客户端生成,作业提交时传给JobMaster。
在Flink Web UI中点击作业可以看到对应的作业图。
③执行图(ExecutionGraph)
JobMaster收到JobGraph后生成,是调度层最核心的数据结构。与JobGraph的最大区别:按照并行度拆分了并行子任务,并明确了任务间数据传输方式。
④物理图(Physical Graph)
JobMaster将执行图分发给TaskManager后,TaskManager部署任务形成的实际执行"图"。这不是一个具体的数据结构,而是实际执行层面的概念。主要在执行图基础上进一步确定数据存放位置和收发的具体方式。
四、总结
用一张思维导图来回顾本文的核心内容:
Flink运行时架构
├── 系统架构
│ ├── JobManager
│ │ ├── JobMaster(与Job一一对应,生成ExecutionGraph,协调Checkpoint)
│ │ ├── ResourceManager(管理Slot资源)
│ │ └── Dispatcher(REST接口,启动JobMaster,Web UI)
│ └── TaskManager(执行具体计算,提供Slot)
│
├── 核心概念
│ ├── 并行度(算子子任务数量,动态概念)
│ ├── 算子链(并行度相同+一对一 → 合并为一个Task,减少开销)
│ ├── 任务槽(TaskManager资源划分单元,静态概念,只隔离内存)
│ └── Slot共享(同作业不同算子可共享Slot)
│
└── 图的转换
└── StreamGraph → JobGraph → ExecutionGraph → Physical Graph
几个容易混淆的点再强调一下:
- JobManager ≠ JobMaster:JobManager是进程,JobMaster是其中负责单个Job的组件
- Slot隔离内存,不隔离CPU:设置Slot数量时建议对应CPU核心数
- 并行度不能超过Slot总数:Slot是上限,并行度是实际使用
- 算子链的条件:并行度相同 + 一对一传输,缺一不可
更多推荐
所有评论(0)