资讯详情

资讯详情

Flink on Yarn部署实战:三种运行模式与生产环境问题排查

把Flink任务跑在Yarn上是大数据集群里非常常见的部署方式。很多同学一上来就敲flink run -t yarn-session任务能起能停就觉得已经会了可真到报错的时候就开始懵连是Flink的问题还是Yarn的问题都分不清楚。这篇文章我想把Flink on Yarn这层窗户纸捅破重点聊三件事为什么要把Flink挂在Yarn上、三种运行模式到底怎么选、以及一次任务的完整提交流程背后发生了什么。文中也会穿插一些我在生产环境里踩过的坑比如JDBC连接器异常、SQL窗口不触发、CDC部署翻车这些经典问题。内容更适合正在做实时计算、离线批处理或者准备从Standalone切到Yarn部署的同学。无论你是刚接触Flink还是已经写过不少SQL任务把这套原理吃透之后排查问题会顺手很多。1. 理解Flink on Yarn的设计思路为什么大家都选它1.1 Yarn的AM机制正好补齐Flink的资源管理短板Flink本身只是一个计算引擎它只管作业的调度、状态、窗口和检查点至于“我该用多少CPU、多少内存、跑在哪台机器上”它天然不管。而Yarn的核心能力恰好是资源调度它把集群里的CPU、内存统一管理起来谁要资源就通过ApplicationMaster去申请。放到Flink on Yarn场景里Yarn的ApplicationMaster容器里跑的就是Flink的JobManager。这个角色会向Yarn的ResourceManager申请Container来启动TaskManagerTaskManager启动后再反向注册到JobManager。整个过程中Flink不需要自己实现一套资源调度系统而是把资源生命周期完全委托给Yarn。我之前见过不少团队自己维护Standalone集群机器用久了就会出现“部分节点TaskManager挂了没人管”“作业占用固定槽位闲时浪费资源”这种问题。Yarn把这些问题接过去之后NodeManager会替我们监控容器健康状态作业失败了还可以由Yarn重新拉起ApplicationMaster这是Flink on Yarn最本质的优势。1.2 为什么不用Standalone资源治理才是On Yarn的初衷很多初学者会问既然是跑Flink那用官方精简模式Standalone不就行了答案是能跑但生产不推荐。Standalone模式下的TaskManager数量和资源是固定的并行度一高就得手动加节点作业高峰低峰没法弹性伸缩。更难受的是Standalone集群里所有人的作业都抢同一批固定资源一个作业的异常GC波动可能会影响同一台机器上的其他任务。而Flink on Yarn天然继承了Yarn的队列和配额体系。公司内部可以把不同业务线划分到不同Yarn队列设置容量上限和优先级生产作业和临时作业互不干扰。比如我用Capacity Scheduler的时候yarn.scheduler.capacity.root.production.capacity可以指定生产队列最高可用容量临时队列给一个小额度。这样即使有人乱提交大作业也不会挤掉线上任务。生产环境里我基本不会再用Standalone跑核心链路。Flink on Yarn的资源隔离能力包括故障恢复、优先级调度、队列权限比我们自己维护一堆裸进程要可靠得多。1.3 三种模式诞生的逻辑线Flink on Yarn的三模式不是设计者凭空想出来的而是顺着“资源共享”和“资源隔离”这对矛盾逐步演进出来的。最开始出现的是Session模式一个Flink集群常驻在Yarn上所有人往同一个集群提交作业。好处是启动快、资源复用率高坏处是作业多了之后互相抢资源JobManager变成单点瓶颈一个作业把内存搞爆会影响同集群里的所有作业。为了隔离问题Per-Job模式应运而生。每个作业拉起一个独立的Flink集群作业结束集群自动释放。隔离性变好了但代价是客户端要参与整个提交和运行过程客户端断网或重启就可能把作业跟着带崩。到了Flink 1.11官方推出Application模式把用户main函数直接挪到集群里去执行。客户端只负责上传包和发起提交提交完就可以退出任务的后续生命周期完全由Yarn管理。这基本解决了Per-Job模式下客户端抖动导致作业失败的问题也是我现在生产环境用得最多的模式。2. Flink on Yarn三种模式怎么选会话、预作业与应用模式对比2.1 Session模式适合短小SQL的“共享厨房”Session模式先启动一个常驻的Flink集群然后在该集群中提交多个作业。这个集群的JobManager在Yarn上以一个ApplicationMaster的身份运行TaskManager则按需申请也可以预先启动一批。逻辑上它就像一个共享厨房锅碗瓢盆是公共的谁来了都能用但菜做多了就排队厨房只有一个厨师长他出问题整条流水线就停了。Session模式启动命令大概是这样的flink run -t yarn-session \ -Dyarn.application.idapplication_1700000000000_0001 \ -c com.example.MyJob \ /path/your-job.jar看起来简单但提交之前你得先启动会话。可以用yarn-session.sh -d拉起一个也可以直接通过flink run -t yarn-session让Flink自动创建一个。Session模式适合三类场景并发量不高的临时查询、Flink SQL即席分析、以及几个作业共用一个集群做低延迟验证。但不适合生产核心链路因为JobManager一旦故障上面跑的作业全部受影响且作业之间没有资源隔离抢资源问题很难避免。2.2 Per-Job模式作业越重隔离越重要Per-Job模式为每个作业拉起一套专属Flink集群。可以理解为每次做饭都租一个独立后厨食材、锅灶、厨师都是你的用完就退租。提交命令长这样flink run -t yarn-per-job \ -Djobmanager.memory.process.size2048m \ -Dtaskmanager.memory.process.size4096m \ -Dyarn.application.nameperjob-example \ -c com.example.MyJob \ /path/your-job.jar在Per-Job模式下客户端先创建JobGraph然后请求Yarn启动一个只有当前作业使用的集群。JobManagerApplicationMaster启动后等待客户端把JobGraph传过去然后申请TaskManager执行。这里有个关键问题客户端本身参与了整个执行流程的前期工作。如果提交作业的机器在客户端等待期间出问题作业很可能会失败。虽然后面Flink做了很多优化但这始终是个隐患也是为什么Flink 1.15之后官方逐步弱化了Per-Job模式。2.3 Application模式把main函数搬进集群Application模式解决的就是Per-Job的客户端依赖问题。它的设计核心是用户main函数在集群侧执行客户端只负责上传Jar包和发起提交请求提交完成之后客户端可以直接退出不影响作业运行。命令是flink run-application -t yarn-application \ -Djobmanager.memory.process.size2048m \ -Dtaskmanager.memory.process.size4096m \ -Dtaskmanager.numberOfTaskSlots4 \ -Dyarn.application.nameapp-example \ -c com.example.MyJob \ /path/your-job.jar使用Application模式时用户代码中的main()是在Yarn的ApplicationMaster容器内运行的env.execute()之后的JobGraph在集群内部直接生成并提交。也就是说本地客户端不需要再保存完整的执行环境。Application模式特别适合两类场景一类是大规模生产作业比如几十个实时ETL任务定时批量提交另一类是数据集成任务尤其现在Flink CDC的Pipeline任务很多直接用flink run-application方式提交作业稳定性和资源利用率都高出不少。2.4 一张表看清三种模式的差异对比维度Session 模式Per-Job 模式Application 模式集群生命周期常驻复用可提交多个作业每个作业独立启动结束即释放每个作业独立启动结束即释放main函数位置客户端客户端集群侧ApplicationMaster内资源隔离差所有作业共享TM和JM好作业间完全隔离好作业间完全隔离启动速度快无需等待集群创建较慢需要为每个作业创建集群较慢需要为每个作业创建集群客户端依赖低高客户端崩溃可能影响作业低提交后即可断开适用场景临时查询、小作业共享早期生产作业现已弱化生产主力模式推荐使用这张表我建议你复制下来贴在工位上每次纠结怎么提交时看一眼就够了。3. 从flink run到TaskManager一次提交背后发生了什么3.1 提交命令里的角色分配在聊全过程之前先帮你梳理一下flink run-application那条命令里每一部分在干什么。前半部分是启动方式-t yarn-application明确告诉Flink用Application模式对接Yarn中间的-D参数是给flink-conf.yaml临时覆盖配置项-c指定main类入口最后的Jar包是作业本体。这里有几个容易忽略的配置项jobmanager.memory.process.size表示整个JobManager进程可用的总内存包括了堆内、堆外、JVM自身开销等。taskmanager.memory.process.size同理决定每个TaskManager进程的内存上限。taskmanager.numberOfTaskSlots则决定单个TaskManager能并行执行多少个任务。很多新手把memory.process.size当成堆大小来用导致实际堆内存设置过高容器直接被Yarn杀掉。建议Flink 1.11以上版本直接使用jobmanager.memory.process.size主导内存配置再配合jobmanager.memory.heap.size微调堆内部分。3.2 一次Yarn Application提交的完整生命周期我们以Application模式为例一步步拆解一次任务提交的过程客户端读取Hadoop配置和Flink配置向Yarn ResourceManager发起提交请求申请一个新的ApplicationId。客户端把job.jar、Flink发行包、用户依赖Jar、日志配置文件一起上传到HDFS临时目录路径通常类似/user/xxx/.flink/application_xxx。客户端向ResourceManager发送SubmitApplication请求携带AppMaster容器的资源请求和启动命令。ResourceManager选择一个满足资源条件的NodeManager命令该节点启动ApplicationMaster容器。AppMaster容器内启动Flink的DispatcherDispatcher再启动Flink的ResourceManagerFlink内部的资源管理器和JobManager主要组件。在Application模式下Dispatcher会直接加载用户jar包并执行main()方法生成StreamGraph和JobGraph再交给对应的JobMaster。JobMaster开始调度作业向Flink内部的ResourceManager申请所需的TaskManager资源。Flink的ResourceManager再向Yarn ResourceManager发起Container请求申请TaskManager的工作容器。新容器启动TaskExecutor进程启动成功后向JobManager注册。JobManager把Task部署到具体的Slot上任务开始正常消费数据、执行转换、做Checkpoint。整个流程看起来步骤多但核心就一句话Flink把集群生命周期管理交给Yarn把作业调度牢牢握在自己手里。3.3 角色分工与资源申请顺序Yarn侧的重量级角色有三个ResourceManager负责全局资源调度NodeManager负责单机容器管理ApplicationMaster是每个应用在Yarn里的代理。Flink侧的对应关系是JobManager运行在ApplicationMaster容器里负责任务调度、Checkpoint协调、作业恢复TaskManager运行在Yarn动态申请的Container里真正执行数据流逻辑。这里需要重点理解容器与Slot的关系。Yarn的Container是资源分配的单位决定了进程能占用多少CPU和内存。Flink的Slot是任务调度的单位一个Slot对应TaskManager里一个能够并行执行任务的最小资源段。默认情况下一个TaskExecutorTaskManager会按配置把内存切分成多个Slot。我在生产环境里有个习惯不需要把每个TaskManager的Slot数量设得很大taskmanager.numberOfTaskSlots一般按CPU核数来设比如8核机器就设4或者8。Slot太少并行度不足Slot太多每个任务的内存就小状态后端的压力就大。3.4 并行度、Slot与内存的关系怎么算很多读者问并行度设置成10一个TaskManager配4个Slot那需要几个TaskManager答案是至少3个。TaskManager数量等于ceil(并行度 / 单TM的Slot数)也就是ceil(10 / 4) 3。三个TaskManager中前两个各分配4个任务第三个分配2个任务。这样可以保证所有任务并行执行。内存方面同样需要联动规划。假设每个TaskManager的内存是8GB启动3个TaskManager就需要24GB的可用内存。如果Yarn队列里没有这么多资源分配额度作业就会一直卡在ACCEPTED状态不报错也不运行。这个时候去Yarn Web界面看容器申请记录最典型的特征就是Container一直处在Pending状态等资源。我建议在提交任务前先做一次“资源预演”并行度乘以每个任务预估所需内存加上一定比例的JVM开销再除以单TM可用内存得到TM数量。这样能大幅度减少因资源不足导致的任务长时间排队问题。3.5 和Spark on Yarn的对比很多团队同时用了Flink和Spark对“工程师说Spark on Yarn提交只需要一个Spark客户端”这个说法熟悉得不行。其实Flink on Yarn也差不多只要提交机器能连接Yarn ResourceManager和HDFS就可以远程提交任务本地不需要安装Flink集群。区别在于Spark提交时Driver默认跑在客户端client模式或Yarn集群cluster模式Executor资源由Yarn动态分配Flink的Session/Application模式也有类似的选择。其中Application模式更接近Spark on Yarn的cluster模式客户端只负责提交和上传真正执行和调度都在集群内部完成。如果你的Flink任务经常因为提交机器的网络波动而挂掉建议把所有生产作业迁移到Application模式这是我在踩过几次客户端断连的坑之后的实际经验。4. 生产环境常见问题与排查实录4.1 最常见的ClassNotFound发行版与依赖冲突Flink on Yarn报ClassNotFoundException是出现频率最高的错误之一而且通常是在作业提交后运行一小会儿才出现。我遇到过场景是本地代码用了某个连接器版本跟集群Flink发行版内置的版本不一致。比如集群Flink是1.17但用户Jar里打了一个老版本Flink核心包导致Serializer、TypeInformation这些类在集群侧找不到或版本冲突。排查思路三步走看运行时ClassLoader日志确认是用户Jar找不到依赖还是Flink系统库被用户Jar覆盖。检查flink run-application时有没有使用-C或者-D把多余的Flink核心依赖打进作业。重点排查flink-shaded-hadoop-*、flink-connector-*、flink-cdc-*这类容易版本敏感的包。经验是用户代码的依赖尽量用provided作用域不要打包进用户Jar里。像flink-streaming-java、flink-clients这类包应由Flink运行时提供避免重复打包。4.2 JDBC连接器异常Driver、超时和流量洪峰另一个经常在Flink SQL任务里踩的坑是JDBC连接器异常比如报Cannot connect to MySQL或Communications link failure。原因一般有三种。第一是缺少数据库驱动。Flink官方JDBC连接器并不会自动带上MySQL驱动包需要你把mysql-connector-java-xx.jar放到Flink的lib目录或者通过-C参数引入到作业运行环境。我当时有个任务就是用Application模式提交配置里忘了带驱动结果一跑就报找不到com.mysql.cj.jdbc.Driver。第二是连接地址配错或者网络不通。Flink on Yarn运行的TaskManager容器和你的客户端机器网络环境不同如果只配置了localhost:3306TaskManager去连的时候自然连不上。正确做法是在Flink SQL里使用对集群内可达的数据库地址并用yarn.nodemanager.host所在网络段去验证连通性。第三是连接池被打爆。TaskManager并行度较高时每个task都会创建JDBC连接数据库端的max_connections一下子就被耗尽。这个时候你会看到大量Connection is not available, request timed out错误。建议在Flink SQL中配置合适的连接池参数控制每条任务的连接数量。4.3 Flink SQL的Watermark问题为什么数据总不触发窗口Flink SQL里窗口任务迟迟不触发究其原因多半是Watermark没推下去。在flink sql中water相关配置里很多人以为设置了WATER FOR ts AS ts - INTERVAL 5 SECOND就万事大吉忽略了下游Source的时间戳是否单调递增。Watermark的本质是声明“这个时间之前的数据已经到齐了”。如果上游数据偶发大乱序或者某个分区长期没有数据Watermark就推不动窗口就永远不闭合。解决办法有三个方向一是用ALLOWED LAGGING不同版本叫法不一或idle timeout让空闲分区也能推进Watermark二是定期检查上游Kafka的消息时间戳分布确认业务时间确实没有长时间断层三是用Processing Time临时验证链路通不通排查是不是事件时间设置本身的问题。我在Flink SQL窗口类任务上吃过亏调了半天SQL没反应查到最后是Kafka一个分区被积压了几百万条旧数据导致Watermark被旧数据拖死。数据源侧的消费滞后有时候比Watermark本身更值得关注。4.4 Flink CDC 3.5.0的Docker部署细节决定成败现在很多团队开始用Flink CDC做数据集成比如把MySQL数据实时同步到Kafka、Paimon或Doris。搜索热词里有“flink 2.2.1 flink cdc 3.5.0 docker 部署”这个组合其实挺典型的。Flink CDC 3.x已经不只是连接器概念而是演变成了一个数据同步框架。部署CDC Pipeline任务时通常需要把flink-cdc-pipeline-connector-mysql-3.5.0.jar放到Flink的lib下然后通过flink-cdc.sh来提交YAML pipeline文件。用Docker部署时最容易漏的是两个东西一个是MySQL的binlog格式必须设置为ROW另一个是账号权限必须包含RELOAD、REPLICATION SLAVE、REPLICATION CLIENT等。权限不足时CDC任务会启动成功但读不到任何变更数据日志里只有一串看不懂的Warning。另外容器里的时间设置也会影响CDC。默认的UTC时区会导致数据同步到下游时时间差8小时建议在Docker Compose里把TZ环境变量设为Asia/Shanghai并且MySQL链接串里加上serverTimezoneAsia/Shanghai。CDC 3.5.0的YAML配置里也支持多表同步和Schema变更路由如果你要做分库分表的汇总同步这些配置项非常关键。我建议初次部署时先跑一个最小配置确认binlog位点能正常写入再逐步加多表。4.5 从SQL Client到SQL Gateway多人提交怎么落地Flink SQL Client是本地提交SQL的轻量客户端适合单人或开发阶段使用。但在生产环境当多个团队需要提交SQL时更合理的方案是部署SQL Gateway。SQL Gateway是一个独立运行的REST服务它帮你管理多个Flink SQL会话支持多用户并发提交、Hive元数据对接、Session复用等能力。在Flink 1.16之后已经比较成熟很多平台就是基于它在做SQL化开发。我的建议是临时查询用SQL Client就够了一旦要接入Web系统、统一权限、做SQL审核和血缘记录尽早迁移到SQL Gateway。它可以把每个SQL会话绑定到不同的Yarn队列避免一个任务吃光所有资源。顺带一提Flink数据血缘这块可以选择基于SQL Lineage采集器从SQL Gateway的会话信息里解析Source、Sink、字段级依赖关系。这也是热词里“flink 数据血缘”常见的一个落地方向比完全从代码级去分析要轻量得多。4.6 问题排查速查表问题现象可能原因快速排查与解决办法作业一直ACCEPTED不运行Yarn队列资源不足查看Yarn ResourceManager确认Container Pending增加队列容量或降低并行度运行中突然TaskManager被Kill内存溢出容器超过Yarn限制调大taskmanager.memory.process.size或调整taskmanager.memory.managed.fraction报ClassNotFoundException依赖冲突或驱动缺失检查用户Jar是否把Flink核心包打进去补上缺失驱动JDBC连接超时网络不通、连接池耗尽检查集群内连通性限制并行任务连接数完善连接池参数窗口不触发Watermark推进失败检查分区空闲、事件时间最大值配置Idle TimeoutCDC同步没数据权限不足、binlog配置错误确认binlog_formatROW账号有RELOAD和复制权限容器日志里乱码时间戳时区问题设置TZAsia/ShanghaiURL加serverTimezone如果你在生产环境折腾到凌晨三点还查不出原因建议回到这张表从头到尾挨个对一遍。最后说点我自己的习惯。我现在做生产环境默认就不碰Per-Job了Flink 1.15以后官方逐步弱化这个模式Application模式才是更省心的路子。Session只用来做开发验证和多用户临时SQL绝不让核心链路跑在上面。资源上一定在Yarn侧划分队列把生产作业和开发作业分开再配合Application模式整个集群的稳定性会提升一大截。
觉得有用,分享给同行:

为您的企业打造数字门面

稳重轻奢商务风格,端正雅致视觉,长效耐看不易过时。

立即咨询 →