尧图精选

Flink基础之Flink on Yarn原理详解:三种模式与提交流程

🕒 发布时间:2026/9/8 18:47:04 📁 来源:尧图网络
摘要讲透 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 内存抽象成容器统一分配给各种计算框架。四个核心概念必须分清ResourceManagerRM全局资源调度中枢接收应用的资源申请按调度器Capacity/Fair分配 Container。整个集群只有一个。NodeManagerNM每台机器一个负责启动、监控、回收本机上的 Container向 RM 上报资源。ApplicationMasterAM每个应用一个是应用与 YARN 之间的「项目经理」——向 RM 申请资源、协调 Container、监控应用运行。Container资源分配的最小单元一个进程 CPU/内存配额。关键提醒面试高频坑Flink 内部也有一个叫 ResourceManager 的角色它是 Flink 集群的组件负责向YARN 的 RM申请容器——两个 ResourceManager 不是一回事。一个是应用内的资源协调者一个是全局的资源调度者。二、Flink 角色与 YARN 的映射Flink on YARN 的本质是把 Flink 的进程装进 YARN 的容器里AM 容器里运行 JobManagerDispatcher JobMaster Flink 的 ResourceManagerApplication 模式下用户的 main() 也在这里执行。每个 TaskManager 一个 ContainerContainer 内存 taskmanager.memory.process.size容器内部再划分 SlotFlink 自己的资源单元。资源模型是两层的YARN 管「进程级容器」Flink 管「进程内 Slot」。YARN 分配 ContainerFlink 在 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生产环境追求隔离和可靠 → ApplicationPer-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 关键配置# AMJobManager总内存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-killapplicationId# 查看应用日志排障首选yarnlogs-applicationIdapplicationId七、四个真实踩坑把 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、计算执行留在 FlinkYARN 管进程级容器Flink 管进程内 Slot两层资源模型通过 AM 衔接。理解了「两个 ResourceManager 的区分」「三种模式的本质差异是集群生命周期和 main() 位置」「提交流程里两次资源协商」这三点Flink on YARN 的架构和排障就都通了。
上一篇/下一篇内容由系统自动关联 返回资讯列表 →