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的工作流程:

  1. 接收提交的应用(JobGraph)
  2. 将JobGraph转换为执行图(ExecutionGraph)——包含所有可以并发执行的任务
  3. 向ResourceManager申请执行所需的资源(Slot)
  4. 获取到足够资源后,将执行图分发给TaskManager执行
  5. 运行过程中负责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的工作流程:

  1. 启动后向ResourceManager注册自己的Slot
  2. 收到ResourceManager指令后,将Slot提供给JobMaster调用
  3. JobMaster分配任务后,TaskManager执行具体计算
  4. 执行过程中可以缓冲数据,也可以与其他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):数据流维护分区和元素顺序,不需要重新分区。mapfilterflatMap等算子都是一对一关系,类似Spark的窄依赖。

重分区(Redistributing):数据流分区发生改变,比如keyBywindow之后的操作,类似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

好处:

  1. 资源密集型和非密集型任务可以在同一个Slot中自行分配资源占用比例,使负载平均分配
  2. 即使某个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

几个容易混淆的点再强调一下

  1. JobManager ≠ JobMaster:JobManager是进程,JobMaster是其中负责单个Job的组件
  2. Slot隔离内存,不隔离CPU:设置Slot数量时建议对应CPU核心数
  3. 并行度不能超过Slot总数:Slot是上限,并行度是实际使用
  4. 算子链的条件:并行度相同 + 一对一传输,缺一不可

更多推荐