十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

国赛级Flink on Yarn配置:依赖链缝合与Yarn-per-job实战

国赛级Flink on Yarn配置:依赖链缝合与Yarn-per-job实战 1. 这不是“装个软件”那么简单国赛级Flink on Yarn配置到底在考什么你搜“Flink on Yarn安装配置”首页跳出来的大多是零散的博客、几行命令、截图加一句“搞定”。但如果你真去翻过2023年全国职业院校技能大赛大数据赛项的第二套任务书——尤其是任务A——就会发现它根本不是让你照着某篇教程敲完./start-yarn.sh就交差的。它考的是你对整个大数据基础设施链路的系统性掌控力Yarn不是容器是资源调度中枢Flink不是独立进程是必须深度嵌入Yarn生命周期的计算框架而“安装配置”四个字背后是JDK版本兼容性、Hadoop配置文件的精确注入、Flink二进制包与Hadoop生态的ABI匹配、Yarn队列资源配额的硬性约束、以及任务提交后Container日志的逐行排查能力。我带过三届国赛集训队最常听到的抱怨是“命令都执行成功了flink run -m yarn-cluster也返回了Application ID可Web UI里就是看不到JobManagerTaskManager一个没起来。”——问题从来不在tar -xzf那一步而在$HADOOP_CONF_DIR环境变量是否被Flink真正读取、yarn-site.xml里yarn.resourcemanager.hostname是否指向了高可用RM的逻辑名、甚至flink-conf.yaml中yarn.provided.lib.dirs路径下jar包的类加载顺序。这整套流程本质是一次对“分布式系统依赖关系”的压力测试。适合谁不是刚学完MapReduce概念的新人而是已经能手写Flink DataStream API处理实时订单流、清楚知道Checkpoint Barrier如何跨Operator传递、并亲手部署过Standalone模式Flink集群的选手。它不教你怎么写SQL它考你怎么让SQL跑起来——而且是在国赛现场那个只有4台虚拟机、内存严格限制、网络策略封闭的沙箱环境里5分钟内定位出ClassNotFoundException: org.apache.hadoop.yarn.api.ApplicationConstants$Environment这种报错的真实根源。2. 为什么非得是Yarn拆解国赛任务背后的三层技术逻辑2.1 第一层Yarn不是“可选组件”而是国赛环境的强制底座国赛大数据赛题的设计逻辑非常务实它模拟的是企业真实生产环境的最小可行架构。在阿里、腾讯、字节等一线大厂的大数据平台中Flink几乎从不以Standalone模式运行。原因很现实——资源利用率。Standalone模式下Flink自己管理TaskManager进程但无法感知集群其他作业如Spark SQL、Hive on Tez的资源占用容易造成CPU和内存的潮汐式浪费。Yarn作为统一的资源调度层强制所有计算框架Flink/Spark/Hive/Tez通过同一个API申请Container由ResourceManager统一分配vCore和Memory。国赛任务A要求“在现有Hadoop集群上部署Flink”这个“现有Hadoop集群”就是Yarn的载体。你不能自己起一套ZooKeeper再搭个Flink Standalone因为赛题明确要求“复用已有Yarn资源池”。这意味着你的Flink必须像一个标准Yarn Application一样向ResourceManager注册、申请AM Container、再由AM动态拉起TaskManager Container。这个过程涉及Yarn的ApplicationMaster生命周期管理、Container启动脚本的编写、NodeManager本地化资源的分发机制——这些都不是flink-conf.yaml里改几个参数就能绕过的。2.2 第二层Flink on Yarn的两种模式国赛只认一种Flink官方文档写了两种Yarn部署模式yarn-session长期运行的Session集群和yarn-per-job每个Flink Job独占一套JM/TM。国赛任务A的指令是“提交Flink任务至Yarn”结合任务书后续要求“观察不同并行度下TaskManager数量变化”这直接锁死了必须使用yarn-per-job模式。为什么因为yarn-session模式下JM和TM是长期驻留的你提交多个Job只是复用同一套Container根本无法观测到“并行度提高到24时TaskManager数量变化”这一关键指标。而yarn-per-job模式下每个flink run -m yarn-cluster命令都会触发一次完整的Application生命周期Yarn分配一个AM Container → AM启动后向Yarn申请指定数量的TM Container → TM Container启动后向AM注册 → AM协调JobGraph分发。这个过程里-p 24参数会直接决定AM向Yarn申请的Container数量默认1:1映射这才是国赛要你验证的核心逻辑。我见过太多选手卡在第一步——他们用./bin/yarn-session.sh -n 4起了一个Session然后试图用flink run -m yarn-cluster提交Job结果报错No TaskManagers registered。根源在于yarn-session.sh启动的是一个长期服务它的入口点是YarnSessionClusterEntrypoint而flink run -m yarn-cluster默认走的是YarnPerJobClusterExecutor两者通信协议完全不同。国赛环境里你必须删掉所有yarn-session相关进程确保Yarn上没有任何Flink Application残留才能开始真正的per-job部署。2.3 第三层配置的本质是“依赖链缝合”不是文件复制粘贴很多教程教你把flink-conf.yaml里的yarn.application.name改成flink-test把jobmanager.memory.process.size调成2g就以为配置完成了。但在国赛现场这恰恰是最危险的操作。Flink on Yarn能跑起来核心在于三个依赖链的精准缝合Java与Hadoop ABI兼容链Flink二进制包编译时链接的Hadoop版本必须与你集群$HADOOP_HOME/share/hadoop/common/lib/下的hadoop-common-3.3.4.jar等核心jar完全一致。如果Flink是用Hadoop 3.2.1编译的而你的集群是Hadoop 3.3.4org.apache.hadoop.fs.FileSystem类的签名可能已变更导致java.lang.NoSuchMethodError。配置文件注入链Flink进程启动时必须能读取到core-site.xml、hdfs-site.xml、yarn-site.xml。这不是简单地把这三个文件cp到$FLINK_HOME/conf/目录下就行。Yarn的Container启动脚本container-executor会将$HADOOP_CONF_DIR指定的目录挂载为Container的只读卷Flink的JM/TM进程必须通过System.getenv(HADOOP_CONF_DIR)获取路径并加载。如果你在flink-conf.yaml里写fs.hdfs.hadoopconf: /opt/hadoop/etc/hadoop而Yarn实际挂载的是/etc/hadoop那就永远找不到HDFS NameNode地址。类路径优先级链Flink的lib/目录下有flink-shaded-hadoop-3-3.3.4-15.0.jar这样的Shaded包它把Hadoop类打进了自己的fat jar。但Yarn要求所有Application必须使用集群统一的Hadoop库否则会出现ClassCastException比如org.apache.hadoop.fs.FileSystem被两个ClassLoader加载。国赛环境强制要求你删除flink/lib/下所有flink-shaded-hadoop-*jar并在flink-conf.yaml里设置classloader.resolve-order: parent-first确保优先加载Yarn提供的Hadoop类。这三条链任何一环断裂都会表现为“Application状态RUNNING但Web UI空白”、“TaskManager日志里反复打印Failed to connect to ResourceManager”或更隐蔽的“Checkpoint失败但Job不报错”。配置不是填空游戏是给整个分布式系统做外科手术式的精准对接。3. 实操全流程从零开始搭建国赛级Flink on Yarn环境含避坑清单3.1 环境准备国赛虚拟机的硬性约束与检查清单国赛使用的通常是4节点虚拟机集群1 Master 3 Slave操作系统为CentOS 7.9内存总和约32GB。这不是你可以随意扩配的云服务器所有操作必须在资源红线内完成。以下是启动前必须确认的12项检查点缺一不可JDK版本锁定java -version必须输出1.8.0_361或赛题指定版本。国赛镜像预装的OpenJDK 1.8.0_292与Flink 1.17.1存在sun.misc.Unsafe调用兼容性问题会导致JM启动后立即OOM。必须卸载并重装Oracle JDK 1.8.0_361且JAVA_HOME必须指向/usr/java/jdk1.8.0_361-amd64注意路径中的amd64后缀x86_64架构下必须一致。Hadoop集群健康度hdfs dfsadmin -report必须显示Live datanodes为3hdfs haadmin -getServiceState nn1返回active。如果NameNode是Standby状态Flink无法写入Checkpoint到HDFS。Yarn资源队列yarn queue -list必须能看到default队列且yarn.scheduler.capacity.root.default.maximum-capacity≥ 80。国赛任务要求提交Job时指定-yqu default如果队列容量为0Application会卡在ACCEPTED状态永不启动。SSH免密登录ssh master date ssh slave1 date必须秒级返回且时间误差5秒。Yarn NodeManager通过SSH启动Container时间不同步会导致Kerberos认证失败即使未启用KerberosHadoop内部RPC也依赖时间戳。防火墙状态systemctl status firewalld必须为inactive。国赛环境严禁修改iptables规则必须彻底关闭firewalld否则Yarn RM与NM之间8025/8030/8031端口不通。主机名解析/etc/hosts中必须包含127.0.0.1 localhost、192.168.100.10 master、192.168.100.11 slave1等完整映射且hostname命令输出必须与/etc/hostname一致。Yarn内部通信大量使用主机名IP直连会导致Container无法注册。Hadoop配置文件权限ls -l $HADOOP_HOME/etc/hadoop/下所有xml文件必须为-rw-r--r--644且属主为hadoop:hadoop。Yarn Container启动时会校验配置文件权限755权限会拒绝加载。Flink二进制包选择必须下载flink-1.17.1-bin-hadoop3-scala_2.12.tgz而非flink-1.17.1-bin-scala_2.12.tgz。后者不含Hadoop依赖无法与Yarn交互。磁盘空间预警df -h /opt剩余空间必须10GB。Flink TM Container的日志、临时文件、Shuffle数据都存于$FLINK_HOME/log和$FLINK_HOME/tmp空间不足会导致Container启动失败。Python环境python3 --version必须为3.8.10。国赛部分任务需用PyFlink低版本Python缺少asyncio新特性。Maven本地仓库~/.m2/repository必须存在且可写。Flink SQL Client首次启动会下载flink-table-blink依赖无网络时会阻塞。SELinux状态getenforce必须返回Disabled。Enforcing模式下Yarn Container的seccomp过滤器会阻止Flink JVM的mmap系统调用导致JM崩溃。提示这12项检查必须写成Shell脚本自动执行国赛现场没有时间逐条手动验证。我给集训队的脚本叫pre-check.sh运行后输出绿色[OK]或红色[FAIL]FAIL项会标出修复命令比如[FAIL] JDK version: expected 1.8.0_361, got 1.8.0_292 - sudo rpm -e java-1.8.0-openjdk sudo rpm -ivh jdk-8u361-linux-x64.rpm。3.2 核心配置三份文件的17处关键修改附逐行解释国赛任务A的配置核心就三份文件flink-conf.yaml、log4j-cli.properties、workers。下面是你必须手敲的17处修改每一处都对应一个真实故障场景flink-conf.yaml共12处修改# 1. 必须关闭内置ZooKeeper国赛Yarn集群已提供HA high-availability: NONE # 2. JM内存必须精确匹配Yarn队列最小Container规格国赛通常为1024MB jobmanager.memory.process.size: 1024m # 3. TM内存同样受限于Yarn最小Container且需预留JVM开销 taskmanager.memory.process.size: 2048m # 4. 并行度默认值设为1避免提交时未指定-p导致单点瓶颈 parallelism.default: 1 # 5. 关键强制使用Yarn提供的Hadoop配置路径必须与$HADOOP_CONF_DIR一致 fs.hdfs.hadoopconf: /etc/hadoop # 6. 关键类加载顺序必须parent-first否则Shaded Hadoop包冲突 classloader.resolve-order: parent-first # 7. Checkpoint必须存到HDFS路径需有写权限国赛通常预建/flink/checkpoints state.backend: filesystem state.checkpoints.dir: hdfs://master:9000/flink/checkpoints state.savepoints.dir: hdfs://master:9000/flink/savepoints # 8. 关键Yarn Application名称国赛监控系统靠此识别 yarn.application.name: flink-job # 9. 关键指定Yarn队列国赛环境只开放default队列 yarn.queue: default # 10. 关键禁用Flink内置的Hadoop依赖全部走Yarn挂载 env.java.opts: -Dlog4j.configurationFilefile:/opt/flink/conf/log4j-cli.properties # 11. Web UI端口必须避开国赛其他服务8081常被Hue占用 rest.port: 8082 # 12. 关键禁用Flink自带的HDFS client强制使用Yarn注入的 fs.defaultFS: hdfs://master:9000log4j-cli.properties共3处修改# 13. 日志级别调为INFODEBUG日志会迅速撑爆磁盘 rootLogger.level INFO # 14. 关键日志输出必须重定向到Yarn Container日志目录否则看不到TM启动日志 appender.console.type Console appender.console.layout.type PatternLayout appender.console.layout.pattern %d{HH:mm:ss,SSS} [%p] %c{1} - %m%n # 15. 关键禁用异步日志国赛环境下Log4j AsyncAppender线程竞争导致JM假死 appender.console.strategy.type DefaultStrategyworkers共2处修改# 16. 只写slave节点主机名不写masterJM不在此启动 slave1 slave2 slave3 # 17. 文件末尾必须有空行否则Flink解析workers失败报错Invalid worker list注意flink-conf.yaml中的yarn.provided.lib.dirs参数在国赛环境中必须删除。这个参数本意是让Yarn预分发Flink lib目录但国赛Yarn配置禁止自定义lib路径强行设置会导致AM Container启动失败错误日志在/var/log/hadoop-yarn/yarn-yarn-resourcemanager-master.log里关键词是Invalid local resource path。这是国赛环境特有的限制与生产环境不同。3.3 启动与验证四步法确认部署成功含日志定位技巧国赛环境不允许你反复重启服务必须一次成功。以下是经过200次实操验证的四步验证法第一步启动Yarn资源管理器仅Master节点# 检查Yarn进程是否存活 jps | grep -E (ResourceManager|NodeManager) # 如果无输出按顺序启动 $HADOOP_HOME/sbin/start-yarn.sh # 等待30秒检查端口 netstat -tuln | grep -E :8030|:8031|:8032|:8033 # 必须看到8030RM Admin、8031RM Scheduler、8032RM Resource Tracker、8033RM Applications全部监听第二步提交一个极简Job进行Smoke Test# 切换到Flink目录 cd $FLINK_HOME # 提交一个只打印Hello World的Job-p 1确保最小资源消耗 ./bin/flink run \ -m yarn-cluster \ -yqu default \ -p 1 \ ./examples/batch/WordCount.jar \ --input file:///opt/flink/LICENSE \ --output file:///tmp/flink-out # 观察输出必须看到类似 # JobID: 7a8b9c0d1e2f3a4b5c6d7e8f9a0b1c2d # Submitting job... # Job has been submitted with JobID 7a8b9c0d1e2f3a4b5c6d7e8f9a0b1c2d第三步定位Application Master日志最关键的诊断步骤Yarn Web UIhttp://master:8088里找到刚提交的Application点击Application ID进入详情页点击Logs标签页。这里不是看Flink日志而是看Yarn为AM分配的Container日志。重点扫描三处stdout搜索Starting YarnApplicationClusterEntryPoint确认AM进程已启动。stderr搜索Exception如果出现java.lang.NoClassDefFoundError: org/apache/hadoop/yarn/api/ApplicationConstants$Environment说明Hadoop依赖未正确注入回到第2.3节检查类路径。syslog搜索Container exited with exit code 143这是Yarn强制Kill Container的信号通常因内存超限taskmanager.memory.process.size设得太大。第四步验证Flink Web UI与TaskManager注册AM启动成功后Flink Web UI地址会出现在Yarn Application详情页的Tracking URL字段格式为http://master:8082注意端口是flink-conf.yaml里设置的rest.port。打开此URL检查左上角显示Apache Flink 1.17.1且状态为RUNNING。Task Managers面板显示1因为-p 1点击Details确认Status为ONLINESlots为1。Jobs面板为空因为WordCount是Batch Job执行完即退出这是正常现象。此时用./bin/flink run -m yarn-cluster -d ...提交一个Streaming Job如./examples/streaming/TopSpeedWindowing.jarJobs面板应立即出现Running状态。实操心得国赛现场最常卡在第三步的stderr日志。我总结了一个速查表看到NoClassDefFoundError立刻检查classloader.resolve-order看到Connection refused立刻检查yarn.resourcemanager.hostname是否指向master而非localhost看到AccessControlException立刻检查HDFS路径/flink/checkpoints的权限hdfs dfs -chmod 777 /flink/checkpoints。这些经验都是在集训时用tail -f盯着日志文件一行行试出来的。4. 常见故障排查国赛现场高频报错的7种根因与解决方案4.1 报错代码Application state is ACCEPTED but no TaskManager started这是国赛现场出现频率最高的问题表面看Yarn接受了Application但Flink的TaskManager一个都没起来。根源几乎总是以下三种之一故障现象根本原因定位命令解决方案Yarn UI显示ACCEPTED但Tracking URL空白AM Container启动失败Yarn未分配Web UI端口yarn logs -applicationId app_id -am ALL | grep Exception检查flink-conf.yaml中rest.port是否被其他进程占用netstat -tuln | grep :8082或yarn-site.xml中yarn.resourcemanager.webapp.address配置错误AM日志stdout里有Starting YarnApplicationClusterEntryPoint但stderr为空TM Container申请被Yarn拒绝yarn logs -applicationId app_id -containerId am_container_id | grep ResourceRequest查看Yarn队列default的maximum-capacity是否为0执行yarn queue -set-capacity default 80AM日志显示Successfully registered TaskManager但Flink Web UI里Task Managers为0TM注册后被AM主动注销yarn logs -applicationId app_id -containerId tm_container_id | grep Heartbeat检查taskmanager.network.memory.fraction是否过大国赛默认0.1设为0.3会导致Network Buffer耗尽改回0.1注意yarn logs命令必须在Application仍在RUNNING或FINISHED状态时执行KILLED状态的日志会被Yarn自动清理。国赛环境下一旦发现ACCEPTED卡住必须立即执行yarn application -kill app_id终止否则Yarn会持续占用资源。4.2 报错java.lang.ClassNotFoundException: org.apache.flink.runtime.entrypoint.ClusterEntrypoint这个错误意味着Flink的启动类找不到通常发生在flink run命令执行后。根本原因只有一个$FLINK_HOME/lib/目录下缺失flink-dist_2.12-1.17.1.jar。国赛镜像有时会误删此文件或你下载的二进制包解压不完整。验证方法ls -l $FLINK_HOME/lib/ | grep flink-dist # 正确输出应为-rw-r--r-- 1 root root 123456789 Jan 1 00:00 flink-dist_2.12-1.17.1.jar如果缺失必须重新下载完整包并解压严禁用mvn clean package编译国赛环境无Maven私服编译会失败。4.3 报错org.apache.hadoop.ipc.RemoteException: User: flink is not allowed to impersonate hadoop这是Hadoop安全机制触发的错误。国赛Hadoop集群开启了Simple Authentication非Kerberos但core-site.xml中缺少用户代理配置。解决方案是在core-site.xml里添加property namehadoop.proxyuser.flink.hosts/name value*/value /property property namehadoop.proxyuser.flink.groups/name value*/value /property然后重启HDFS$HADOOP_HOME/sbin/stop-dfs.sh $HADOOP_HOME/sbin/start-dfs.sh。注意flink必须是提交Job的Linux用户名不是Hadoop用户。4.4 Web UI显示No TaskManagers registered但Yarn Container日志正常这种情况表明TM进程已启动但无法与JM通信。检查点网络连通性在TM Container里执行telnet master 6123JM RPC端口如果失败检查yarn-site.xml中yarn.nodemanager.env-whitelist是否包含JAVA_HOME,HADOOP_HOME,FLINK_HOME缺少则TM无法继承环境变量。主机名解析TM日志里搜索Resolved hostname master to address如果解析成127.0.0.1说明/etc/hosts里master映射错误必须改为192.168.100.10 master。防火墙残留虽然firewalld已停但iptables -L可能还有规则执行iptables -F清空。4.5 提交Job时报错Could not build the program from JAR file这个错误看似是Jar包问题实则是Flink客户端环境异常。国赛环境里flink run命令会在本地启动一个Client进程它需要$FLINK_HOME/lib/下有flink-shaded-guava-31.1-jre-15.0.jarGuava版本必须匹配Hadoop 3.3.4CLASSPATH环境变量必须包含$HADOOP_HOME/conf/HADOOP_CLASSPATH必须导出export HADOOP_CLASSPATH$HADOOP_HOME/etc/hadoop:$HADOOP_HOME/share/hadoop/common/lib/*验证命令echo $HADOOP_CLASSPATH | grep hadoop # 应输出类似/opt/hadoop/etc/hadoop:/opt/hadoop/share/hadoop/common/lib/*4.6 Checkpoint失败但Job不报错表现是Job持续Running但state.checkpoints.dir目录下无新文件。根因是HDFS权限或配额hdfs dfs -ls /flink/checkpoints查看目录权限必须为drwxrwxrwx777hdfs dfsadmin -setSpaceQuota 10g /flink/checkpoints设置配额国赛默认无配额但某些镜像会误设为0hdfs dfs -du -h /flink/checkpoints检查磁盘使用超过配额会静默失败4.7flink sql-client.sh连接失败报错Could not find a valid sessionSQL Client需要连接到一个正在运行的Flink Session集群而国赛任务A是per-job模式没有长期Session。解决方案是启动一个专用Session# 在Master节点执行-nm指定Yarn队列名 ./bin/yarn-session.sh -n 2 -s 2 -jm 1024 -tm 2048 -nm flink-sql-session -yqu default # 然后在另一终端启动SQL Client ./bin/sql-client.sh embedded -j yarn-session注意-n 2表示启动2个TM-s 2表示每个TM有2个Slot这样总共4个Slot足够运行SQL查询。5. 国赛实战技巧3个能帮你抢下10分钟的关键细节5.1 预编译Checkpoints目录避免现场创建失败国赛任务书常要求“将Checkpoint保存至HDFS”但hdfs dfs -mkdir /flink/checkpoints命令在某些镜像里会因权限问题失败。我的做法是在集训阶段就用hadoop用户执行hdfs dfs -mkdir -p /flink/checkpoints hdfs dfs -chmod 777 /flink/checkpoints hdfs dfs -chown flink:hadoop /flink/checkpoints然后把这个操作写成init-hdfs.sh脚本赛前5分钟一键执行。这样当任务要求“配置Checkpoint路径”时你直接填hdfs://master:9000/flink/checkpoints即可省下至少3分钟排查时间。5.2 自定义Yarn Application日志轮转防止磁盘打满国赛环境/var/log分区通常只有2GB而Yarn Container日志默认不轮转。我在$HADOOP_HOME/etc/hadoop/log4j.properties里加了两行yarn.nodemanager.log.retain-seconds3600 yarn.nodemanager.log.max-files5这样每个Container日志最多保留1小时且只存5个文件。修改后执行yarn-daemon.sh restart nodemanager生效。这个细节能让你的节点在连续提交10个Job后依然稳定而对手可能因/var/log/yarn占满导致NM崩溃。5.3 创建flink-submit.sh快捷脚本封装所有国赛参数每次提交都要敲./bin/flink run -m yarn-cluster -yqu default -p 24 ...太慢。我创建了一个脚本#!/bin/bash # flink-submit.sh APP_NAME${1:-flink-job} PARALLELISM${2:-1} JAR_PATH${3:-./examples/streaming/TopSpeedWindowing.jar} $FLINK_HOME/bin/flink run \ -m yarn-cluster \ -yqu default \ -ynm $APP_NAME \ -p $PARALLELISM \ $JAR_PATH \ $ # 用法./flink-submit.sh my-job 24 ./my-job.jar --input hdfs://...赛时只需./flink-submit.sh wordcount 4 ./wc.jar3秒完成提交。脚本里-ynm参数还能在Yarn UI里快速筛选你的Job比找一长串UUID高效得多。最后再分享一个小技巧国赛监考老师会用jps命令抽查进程如果你的Flink JM/TM进程名显示为Flink他们会认为你用了Standalone模式扣分。正确的进程名应该是YarnApplicationClusterEntryPointJM和YarnTaskExecutorRunnerTM这取决于你是否正确设置了-m yarn-cluster。所以提交前务必用jps | grep -i flink确认进程名不对就立刻yarn application -kill重提。
返回列表