一、前言
在掌握了Flink的部署方式之后,我们需要进一步了解Flink作业提交后内部是如何运转的。Flink作为分布式流处理引擎,其运行时架构设计直接决定了作业的执行效率、容错能力和资源利用率。
理解运行时架构,对于以下场景至关重要:
- 性能调优:知道Slot如何分配,才能合理设置并行度
- 故障排查:明白Checkpoint机制,才能快速定位恢复问题
- 面试准备:JobManager/TaskManager职责、提交流程是高频考点
本文将围绕三大核心主题展开:
二、系统架构总览
Flink运行时架构采用经典的Master-Worker模式,主要由两类进程构成:

2.1 JobManager(作业管理器)
JobManager是整个Flink集群的控制中枢,每个作业对应一个JobManager实例(严格来说是JobMaster)。它包含三个核心组件:
2.1.1 JobMaster
JobMaster是JobManager中最核心的组件,负责处理单个作业的全部生命周期:
| 作业调度 | 将JobGraph转换为ExecutionGraph,调度任务执行 |
| 资源申请 | 向ResourceManager申请Task Slot资源 |
| Checkpoint协调 | 触发、协调分布式快照的保存与恢复 |
| 故障恢复 | 作业失败时,从最近的Checkpoint重新启动 |
历史演变:早期Flink版本中,JobManager的概念范围较小,实际指的就是现在的JobMaster。从1.x版本开始,为了支持多作业同时运行,才明确区分了JobManager(进程)和JobMaster(作业实例)的概念。
2.1.2 ResourceManager
ResourceManager负责集群资源的统一管理和分配:
- Slot资源管理:维护集群中所有可用Slot的注册信息
- 资源分配:根据JobMaster的请求,分配空闲Slot给作业使用
- 资源回收:作业完成后,回收Slot供其他作业使用
- 外部对接:在YARN/K8S模式下,与外部资源管理器交互
⚠️ 注意区分:Flink内置的ResourceManager与YARN的ResourceManager是不同的组件,前者管理Flink的Slot资源,后者管理整个集群的容器资源。
2.1.3 Dispatcher
Dispatcher提供了Flink对外的REST API接口和Web UI:
- 接收作业提交请求,为每个新作业启动一个JobMaster
- 提供Web UI展示作业执行状态、Metrics指标
- 在部分部署模式下(如YARN应用模式)可能被省略
2.2 TaskManager(任务管理器)
TaskManager是Flink的工作进程,负责实际的数据计算:
核心职责:
关键特性:
- 每个TaskManager是一个独立的JVM进程
- 包含多个Task Slot,Slot数量由配置决定
- 可以缓冲数据,支持反压(Backpressure)机制
三、核心概念详解
3.1 并行度(Parallelism)
3.1.1 什么是并行度?
当数据量巨大时,Flink可以将一个算子"复制"多份到不同节点并行执行。一个算子被拆分成的并行子任务数量,就是该算子的并行度。
// 设置全局并行度
env.setParallelism(2);
// 为单个算子设置并行度
stream.map(new MyMapper()).setParallelism(4);
3.1.2 并行度设置优先级
Flink提供了四级并行度配置,优先级从高到低:

| 1 | 算子.setParallelism() | 单个算子 | 针对热点算子单独调优 |
| 2 | env.setParallelism() | 全局默认 | 开发测试环境 |
| 3 | flink run -p 2 | 提交时全局 | 生产环境推荐 |
| 4 | parallelism.default: 2 | 集群默认 | 集群基础配置 |
💡 最佳实践:生产环境不要在代码中硬编码并行度,应通过提交参数动态指定,便于根据数据量灵活调整。
3.1.3 并行度与Slot的关系

核心公式:
所需Slot数量 = 所有算子中最大的并行度(考虑Slot共享后)
以WordCount为例,假设:
- Source并行度 = 2
- FlatMap并行度 = 2(与Source合并算子链)
- KeyBy/Sum并行度 = 2
- Sink并行度 = 1
则整个作业需要 2个Slot(取最大并行度2),而非5个。
3.2 算子链(Operator Chain)
3.2.1 数据传输模式
算子之间的数据传输有两种模式:
| One-to-One | 分区不变,顺序保持 | map/filter/flatMap | 窄依赖 |
| Redistributing | 重新分区,可能乱序 | keyBy/rebalance | Shuffle |
3.2.2 算子链合并条件
Flink会自动将满足以下条件的算子合并为一个Task:
合并的好处:
- 减少线程切换开销
- 减少序列化/反序列化开销
- 基于缓冲区批量传输数据,提升吞吐量
// 禁用算子链(调试用)
stream.map(...).disableChaining();
// 从当前算子开始新链
stream.map(...).startNewChain();
3.3 任务槽(Task Slot)
3.3.1 Slot的本质
Slot是TaskManager上资源的最小分配单元,表示TaskManager拥有的计算资源子集(主要是内存)。
# flink-conf.yaml
# 每个TaskManager的Slot数量,默认1
taskmanager.numberOfTaskSlots: 4
⚠️ 注意:Slot目前只隔离内存,不涉及CPU隔离。建议将Slot数配置为CPU核心数,避免CPU竞争。
3.3.2 Slot共享机制

Flink默认允许不同算子的子任务共享Slot,只要它们属于同一作业。这种设计的优势:
Slot共享组:
// 手动指定Slot共享组,实现资源隔离
stream.map(...).slotSharingGroup("group1");
stream.filter(...).slotSharingGroup("group2");
不同共享组的任务必须分配到不同Slot,总Slot需求 = 各组最大并行度之和。
四、作业提交流程
4.1 四张图的转换
Flink作业从代码到物理执行,会经历四次图转换:

4.1.1 StreamGraph(逻辑流图)
- 生成位置:客户端
- 特点:根据DataStream API代码直接生成,节点=算子
- 作用:表示程序的拓扑结构
4.1.2 JobGraph(作业图)
- 生成位置:客户端
- 优化:合并符合条件的算子链(Operator Chain)
- 特点:节点=Task,边=数据交换方式
- 提交方式:通过REST API或命令行提交给JobManager
4.1.3 ExecutionGraph(执行图)
- 生成位置:JobMaster
- 转换:按并行度拆分为并行子任务
- 特点:节点=ExecutionVertex(子任务),是调度核心数据结构
- 作用:明确任务间数据传输方式,指导任务调度
4.1.4 Physical Graph(物理流图)
- 生成位置:TaskManager
- 特点:具体执行层面的图,非正式数据结构
- 作用:确定数据存放位置和收发方式
4.2 Standalone会话模式提交流程
1. 用户提交作业 → Client生成StreamGraph → 优化为JobGraph
2. Client将JobGraph提交给Dispatcher
3. Dispatcher启动JobMaster,将JobGraph传递给JobMaster
4. JobMaster生成ExecutionGraph,向ResourceManager申请Slot
5. ResourceManager分配Slot,JobMaster将任务分发到TaskManager
6. TaskManager执行任务,定期汇报心跳和状态
4.3 YARN应用模式提交流程
YARN应用模式与Standalone的主要区别在于:
- main()方法在JobManager执行:客户端只需发送作业jar包
- 动态资源申请:TaskManager按需启动,非预先分配
- YARN负责容器管理:Flink的ResourceManager与YARN RM交互
五、代码实战:观察并行度与Slot
5.1 实验代码
以下代码演示如何观察算子链和并行度的实际效果:
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
public class ParallelismDemo {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env =
StreamExecutionEnvironment.getExecutionEnvironment();
// 设置全局并行度为2
env.setParallelism(2);
// 读取数据源
DataStreamSource<String> source = env.socketTextStream("hadoop102", 7777);
// Map算子:与Source形成算子链(One-to-One + 相同并行度)
SingleOutputStreamOperator<Tuple2<String, Integer>> mapped = source
.map(new MapFunction<String, Tuple2<String, Integer>>() {
@Override
public Tuple2<String, Integer> map(String value) {
return Tuple2.of(value, 1);
}
});
// keyBy会改变分区,触发重分区,算子链在此处断开
// sum聚合算子
SingleOutputStreamOperator<Tuple2<String, Integer>> result = mapped
.keyBy(value -> value.f0)
.sum(1);
// Sink单独设置并行度为1
result.print().setParallelism(1);
env.execute("Parallelism Demo");
}
}
5.2 Web UI观察要点
提交作业后,打开Flink Web UI(http://jobmanager:8081),重点关注:
六、常见问题与注意事项
6.1 Slot不足导致作业无法启动
现象:作业提交后一直处于CREATED状态,无法进入RUNNING。
原因:申请的Slot总数超过了集群可用Slot数。
解决:
# 查看可用Slot数量
# 在Web UI Overview页面查看"Task Managers"和"Available Slots"
# 方案1:增加TaskManager数量或每个TM的Slot数
# 方案2:降低作业并行度
flink run -p 2 your-job.jar # 将全局并行度设为2
6.2 算子链断开导致性能下降
现象:作业吞吐量低,CPU利用率不高。
原因:算子之间频繁序列化/反序列化,网络传输开销大。
排查:在Web UI的JobGraph中查看算子链是否被意外断开。
常见断开原因:
- 并行度不一致
- 显式调用了startNewChain()或disableChaining()
- 算子之间是Redistributing模式(如keyBy)
6.3 状态过大导致OOM
现象:TaskManager频繁崩溃,日志显示内存溢出。
原因:状态数据持续增长,超过TaskManager可用内存。
解决:
// 1. 使用RocksDB状态后端,将状态存储到磁盘
env.setStateBackend(new EmbeddedRocksDBStateBackend());
// 2. 配置状态TTL,自动清理过期状态
StateTtlConfig ttlConfig = StateTtlConfig
.newBuilder(Time.hours(24))
.setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
.setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
.build();
总结
本文深入解析了Flink运行时架构的核心内容:
理解这些底层原理,将帮助你在实际工作中更好地进行性能调优和故障排查。下一篇我们将深入Flink的DataStream API,从执行环境到Source算子,开启Flink编程的核心篇章。




