摘要
讲透 Flink 跑在 YARN 上的完整原理:YARN 核心概念与 Flink 角色映射、Session/Per-Job/Application 三种运行模式的差异与选型、Application 模式下从上传 JAR 到 TaskManager 启动的完整提交流程、Container 与 Slot 的两层资源模型,并给出容错机制、常用命令与四个真实踩坑点。
关键词
Flink、YARN、Application Master、Session、Application、Per-Job、Container、资源调度、run-application、Slot
前两篇分别讲了 JobManager 和 TaskManager,组合起来就是一个完整的 Flink 集群。但生产环境里,Flink 集群很少自己裸奔——绝大多数跑在 Hadoop 的 YARN 上。原因很实际:公司里已经有一套 YARN 集群在跑 Hive、Spark 等任务,Flink 再搭一套独立集群,运维成本翻倍、资源还浪费。把 Flink 装进 YARN,就能和其他计算框架共享一套资源池。
这一篇把「Flink on YARN」拆开讲透:YARN 是什么、Flink 在它上面有几种跑法、提交一个作业背后发生了什么、以及为什么有人会在「两个 ResourceManager」上栽跟头。
一、YARN 核心概念:先搞清楚它的四个角色
YARN 是 Hadoop 生态的「资源操作系统」,它把集群资源(CPU + 内存)抽象成容器,统一分配给各种计算框架。四个核心概念必须分清:
- ResourceManager(RM):全局资源调度中枢,接收应用的资源申请,按调度器(Capacity/Fair)分配 Container。整个集群只有一个。
- NodeManager(NM):每台机器一个,负责启动、监控、回收本机上的 Container,向 RM 上报资源。
- ApplicationMaster(AM):每个应用一个,是应用与 YARN 之间的「项目经理」——向 RM 申请资源、协调 Container、监控应用运行。
- Container:资源分配的最小单元,一个进程 + CPU/内存配额。
关键提醒(面试高频坑):Flink 内部也有一个叫 ResourceManager 的角色,它是 Flink 集群的组件,负责向YARN 的 RM申请容器——两个 ResourceManager 不是一回事。一个是应用内的资源协调者,一个是全局的资源调度者。
二、Flink 角色与 YARN 的映射
Flink on YARN 的本质,是把 Flink 的进程装进 YARN 的容器里:
- AM 容器里运行 JobManager(Dispatcher + JobMaster + Flink 的 ResourceManager);Application 模式下用户的 main() 也在这里执行。
- 每个 TaskManager 一个 Container,Container 内存 =
taskmanager.memory.process.size;容器内部再划分 Slot(Flink 自己的资源单元)。 - 资源模型是两层的:YARN 管「进程级容器」,Flink 管「进程内 Slot」。YARN 分配 Container,Flink 在 Container 里用 Slot 调度任务。
理解了这个映射,后面所有配置和排障都有了解释框架。
三、Flink on YARN 的三种运行模式
Flink 在 YARN 上有三种跑法,差异就两个问题:集群生命周期多长、main() 在哪里执行。
3.1 Session 模式:集群常驻,多作业共享
启动一个 YARN 应用(AM 内 JobManager 常驻 + 一批 TaskManager),之后多个作业提交到这个已存在的集群上,共享 TM 资源池。
# 先启动常驻 session 集群bin/yarn-session.sh-d# 再把作业提交到已有集群bin/flink run-tyarn-session-ccom.example.MainClass ./app.jar优点:提交快(不用每次起集群)、资源复用率高。缺点:隔离差——作业之间争抢资源,一个作业拖垮整个集群是常态;JobManager 单点,挂了所有作业一起遭殃。适合作业多且小、对隔离不敏感的团队。
3.2 Per-Job 模式:每作业一个集群(已弃用)
每次提交都向 YARN 申请一个临时集群,作业结束集群销毁。但 main() 在本地 Client执行,意味着提交机器的 classpath 必须完整,且要先「起集群再跑代码」,启动慢。Flink 1.15 起官方不再推荐,被 Application 模式取代——新项目不要再用。
3.3 Application 模式:每作业一个集群,main() 在集群内(推荐)
与 Per-Job 的区别只有一个:main() 在集群内的 AM 里执行,本地只上传 JAR,没有任何 Client 进程负担,也不依赖本地的 classpath。
bin/flink run-application-tyarn-application\-ccom.example.MainClass ./app.jar隔离性好(每作业独立集群)、资源按需申请、故障影响面小,是目前生产环境的主流选择。
选型一句话:作业多且小、要秒级提交 → Session;生产环境追求隔离和可靠 → Application;Per-Job 直接跳过。
四、Application 模式提交流程:一个作业的完整旅程
以推荐的 Application 模式为例,一个作业从敲下命令到跑起来,背后经历七步:
- Client 上传 JAR 到 HDFS,并向 YARN RM 提交应用请求。
- RM 选择一个 NM,在该节点启动 AM 容器。
- AM 内执行用户 main(),并启动 Flink 的 Dispatcher 和 JobMaster(即 JobManager 组件)。
- JobMaster 构建 ExecutionGraph 后,Flink 的 ResourceManager 汇总 Slot 需求,向 YARN RM 申请 Container。
- YARN RM 按调度器分配 Container(资源不足则排队)。
- NM 在分配的节点上启动 TaskManager 容器,容器内存 =
taskmanager.memory.process.size。 - TaskManager 注册到 JobManager,任务部署到 Slot 开始运行,Checkpoint、故障恢复等机制照常工作。
注意第 4-5 步是两层资源协商:Flink 内部先决定要多少 Slot,再换算成需要多少 Container 去找 YARN 要。TaskManager 的扩缩容(容器级别)就发生在这条链路上。
五、容错与资源管理
- AM 重启:AM(含 JobManager)挂掉后,YARN 按
yarn.application-attempts(默认 2)重启 AM,新 JobManager 从持久化存储恢复 JobGraph、从 Checkpoint 恢复作业状态。 - TaskManager 故障:容器被杀或 NM 宕机 → 该 TM 上任务失败 → JobManager 重新调度(如果还有剩余 Slot),或触发 Flink 的 ResourceManager 重新申请 Container。
- 资源弹性:YARN 上 TaskManager 可以按需增减(作业扩并行度时申请更多容器),这是 Standalone 模式做不到的。
六、配置与常用命令
# flink-conf.yaml 关键配置# AM(JobManager)总内存jobmanager.memory.process.size:2048m# 每个 TaskManager 总内存(= 一个 Container 的内存)taskmanager.memory.process.size:4096m# 每个 TM 的 Slot 数taskmanager.numberOfTaskSlots:4# AM 重启次数(YARN 层面)yarn.application-attempts:2# 应用名(YARN 页面上显示)yarn.application.name:flink-demo常用命令:
# Application 模式提交bin/flink run-application-tyarn-application-cMainClass ./app.jar# 启动 Session 集群(后台)bin/yarn-session.sh-d-nmmy-session# 查看/终止 YARN 应用yarnapplication-listyarnapplication-kill<applicationId># 查看应用日志(排障首选)yarnlogs-applicationId<applicationId>七、四个真实踩坑
- 把 Flink 的 ResourceManager 当成 YARN 的 RM。面试和排障都容易在这里卡壳:Flink 的 ResourceManager 是 JobManager 进程里的组件,负责向 YARN RM 申请容器。看到「ResourceManager 申请资源」的日志,先分清是哪个。
- 内存配置超了队列上限,作业一直 ACCEPTED。Container 内存 = TM 的
process.size,如果队列剩余内存不够分配一个 TM 容器,作业就卡在 ACCEPTED/RUNNING 但不推进。调小 TM 内存,或检查队列容量。 - Session 模式作业相互拖累。一个 Session 集群跑多个作业,某个作业状态膨胀或背压会把整个集群拖垮,其他作业一起变慢甚至失败。对隔离有要求的作业,别图省事塞进共享 Session。
- 还在用 Per-Job 的老写法。
flink run -t yarn-per-job在新版本会警告已废弃,且 main() 在本地执行容易踩 classpath 缺失的坑。统一改run-application,把 main() 交给 AM,问题自然消失。
Flink on YARN 的本质,是把 Flink 的进程装进 YARN 的容器,让资源管理交给 YARN、计算执行留在 Flink:YARN 管进程级容器,Flink 管进程内 Slot,两层资源模型通过 AM 衔接。理解了「两个 ResourceManager 的区分」「三种模式的本质差异是集群生命周期和 main() 位置」「提交流程里两次资源协商」这三点,Flink on YARN 的架构和排障就都通了。