ARTICLE DETAIL

资讯详情

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

Flink on YARN 作业无限重启:appattempt 疯涨根因排查与解决

Flink on YARN 作业无限重启:appattempt 疯涨根因排查与解决 你有没有半夜被告警短信吵醒登上 YARN 的 ResourceManager 页面看到一个 Flink application 的 appattempt 像复读机一样疯狂往上蹦眼看着从 attempt 1 跳到 attempt 30任务却始终进不了 RUNNING 状态我遇到过。而且不止一次。这种无限次 appattempt现象本质就是 Flink 作业提交到 YARN 后ApplicationMaster 进程反复拉起、反复死亡YARN 的 ResourceManager 又按照最大尝试次数不断重新调度。只要根因没找到它就能一直循环下去既占着队列资源又让排查的人头皮发麻。这篇文章我把这套问题的现场确认、底层机制、日志定位方法和高频根因完整梳理一遍看完至少能让你在下次遇到flink-yarn提交任务失败时不再靠猜而是有章法地把它摁住。适合正在部署或维护 Flink on YARN 作业的工程师也包括刚接触 Flink 提交任务、被application、appattempt这些概念绕晕的同学。1. 先把无限重启的现象定位到正确的层次1.1 看到 attempt 疯涨先搞清楚它到底是谁在重启很多人一看到 appattempt 数字不断增长第一反应就是任务在重启。这话只对了一半。YARN 里一个 application 从提交到结束真正负责跟 ResourceManager 通信、协调整个作业运行的是 ApplicationMaster也就是 AM。Flink on YARN 场景下这个 AM 进程同时承载了 JobManager 的角色或者至少承载了启动 JobManager 的入口逻辑。所以每次 appattempt 上涨都是 YARN 在尝试重新拉一个 AM 容器起来。这里有一个关键认知AM 重启不等于 Flink 作业重启。如果 AM 一直没起来任务压根谈不上重启它只是在一个接一个地失败。只有 AM 成功起来、并且跟 TaskManager 建立连接之后AM 的反复重启才会间接导致作业状态的丢失和恢复。所以排查的第一步是确认到底卡在哪个阶段——AM 起不来、AM 起来后自杀、还是 AM 起来了但 TaskManager 活不了。方向错了后面全是浪费时间。然后说回无限次。YARN 集群层面有个参数叫yarn.resourcemanager.am.max-attempts默认值通常是 2也就是说正常情况下你最多看到两个 appattempt 就 FAILED 了。你能看到无限次基本可以断定集群把该值配成了很大的数字甚至设成了 -1表示无限制。Flink 侧还有一个yarn.application-attempts配置项提交时可以通过-D传进去但它最终生效值不能超过 YARN 集群的上限两者比较后取较小者。1.2 用一条命令把 YARN 诊断信息拿全遇到这类问题别急着翻日志。先执行yarn application -status application_xxx这条命令会一次性把该 application 的 State、Diagnostics、提交时间、开始时间、结束时间、Tracking-URL 全部打出来。注意看Diagnostics字段它往往比茫茫日志更直接地告诉你上一次失败的原因。比如Container is running beyond physical memory limits. Current usage: xxx. Killing container.——内存超限被 NodeManager 斩杀Application application_xxx failed 2 times due to AM Container for appattempt_xxx exited with exitCode: 127——随后往往跟着一段 stderr 里的关键异常信息什么都没写就一句AM container has been killed——这种反而难搞通常是抢占或心跳超时得另找线索。再配合 Web UI 看进入 RM 页面http://RM_IP:8088/cluster/app/application_xxx最上方有 Diagnostics 摘要点进 AM Container 的 Logs 可以直达当前 attempt 的 stdout、stderr。要注意的是如果 AM 失败过快UI 上的日志链接可能来不及聚合特别是 YARN 默认的日志聚合有延迟。所以最稳的做法是等 application 真正 FAILED 之后再用聚合日志去捞。1.3 用 ExitCode 判断谁杀了他从yarn application -status的输出里还能看到一个有用的字段就是上一次 attempt 的 AM Container 退出码。我的习惯是退出码为 0进程被常规退出但任务却失败了。这通常是 Flink 的入口程序主动退出比如 JobManager 初始化阶段抛出致命异常后走完了错误处理流程或者 HDFS 路径不存在导致系统直接 fallback 退出。退出码非 01、127、134 等JVM 异常崩溃典型的是类加载失败、主类找不到、JVM 内存错误。退出码为负值-1000、-2000 等基本是被外部系统杀掉或抢占最常见的 -1000 表示 Container 被 NodeManager/ApplicationMaster 主动终止具体原因要去 Diagnostics 里找。这一段判断做完基本能区分Flink 自己不行还是YARN 不让它活。这两条路线的排查工具和侧重点完全不同。2. AM 从提交到稳定要闯五道关看不明白机制就永远在盲猜我在排障时习惯把 AM 的整个生命周期拆成五关。每一关的失败特征不一样能快速把问题收敛到一个很小的范围。2.1 第一关提交报文与 AM Container 拉起flink run -m yarn-cluster老写法或flink run-application -t yarn-application新写法提交时客户端会把 Flink 发行包、用户 JAR、flink-conf.yaml 等打成资源包向 RM 申请一个容器来启动 AM。这一关出问题通常在客户端就报错比如用户没有提交权限、队列不存在、客户端无法连接 RM。如果客户端显示提交成功说明这一关过了。值得留意的是AM 容器申请的内存大小Flink 是根据jobmanager.memory.process.size或jobmanager.memory.heap.size来算的。这个值如果超过 YARN 集群单容器上限会很不好看——提交时报错还算是幸运的不少情况是提交成功但调度器永远满足不了这个请求AM 容器一直分配不下来YARN 就那么干等着直到某次超时把 application 判死重启。2.2 第二关Flink 进程初始化与依赖加载AM 容器在某个 NodeManager 上真正启动之后Flink 入口类开始加载配置、初始化 JobManager 上下文。这一关的失败案例大多是类冲突。用户把 Hadoop 客户端、Guava、Log4j 等打包进了 fat jar跟 Flink 自带的 shaded 版本冲突于是NoSuchMethodError、ClassNotFoundException被一遍遍抛出来进程退出attempt 重来再退出再重来。另外一个常见的隐性坑是内存模型配置。Flink 1.11 之后引入了jobmanager.memory.process.size这类新内存参数但很多老项目还在用jobmanager.heap.mb、jobmanager.heap.size。这些旧参数在新版本里会被忽略或提示异常最终导致 AM 进程的内存跟预期完全对不上进程启动后直接被 NodeManager 判定超限杀掉。2.3 第三关外部系统连通性与状态恢复JobManager 起来之后要连接外部系统。开 HA 的要连 ZooKeeper或 Kubernetes、其他高可用服务有状态恢复需求的话还要读 HDFS 上的 Checkpoint 元数据。这一关失败的特征非常典型AM 起来了进程没崩溃但日志停在 Recovering job 或者连接外部服务的阶段然后因为处理不了异常退出或者一直卡着不动心跳发不出去被 RM 判死。特别是大作业、大状态的恢复场景。状态文件多、Checkpoint 元数据大JobManager 恢复状态非常慢而 YARN 对 AM 的心跳有超时检测yarn.am.liveness-monitor.expiry-interval-ms。一旦认为 AM 不健康RM 就杀容器重启。重启之后又是同样的大状态恢复于是变成一种不是 OOM 也不是类冲突的循环重启特别容易被误判成资源问题。2.4 第四关向 RM 申请 TaskManager 容器AM 稳定下来的第一件事是向 RM 发起容器请求用来启动 TaskManager。如果作业配置的 TM 内存超过队列剩余资源或超过yarn.scheduler.maximum-allocation-mb的上限这些容器请求会一直 PENDING。AM 等在那边等到不耐烦或者触发了某个超时就打包重来。这种局面很讨厌因为 RM 页面上你不一定看得到明显报错。很可能你看到的只是AM 起来了 → 申请资源 → 等 → 被杀 → 再起来。判断方法也很直接看队列剩余资源和作业的资源请求是否匹配。执行yarn queue -status 队列名看队列的 UsedCapacity、MaxCapacity再对照 Flink 作业申请的总内存很快就能算明白。2.5 第五关TaskManager 反向注册与心跳最后一道关是 TaskManager 起来后反连 JobManager。这一步对网络环境敏感比如集群开启了 Kerberos 认证、或节点之间有防火墙只放开了部分端口TM 的 RPC 端口连不上 AM就会一直重试最终导致作业无法 RUNNINGAM 重试。这里要提醒一下TaskManager 连接 AM 走的是 Flink 的 Akka/RPC 和 Netty 端口不是 YARN 分配的端口那么简单。遇到AM 活着、TM 容器也在跑、但作业就是起不来的诡异现象优先怀疑网络 ACL、Kerberos ticket 过期以及taskmanager.host相关配置是否被错误设置。3. 一套可复制的日志定位命令与对齐方法前面讲了一堆理论这里给实战命令。这套命令我在不同集群上用过很多次基本够用。3.1 先列应用再拉聚合日志# 列出所有 flink 相关 application含运行中和已结束的 yarn application -list -appStates ALL | grep -i flink # 查看特定 application 的状态和诊断信息 yarn application -status application_xxx # 拉取该 application 所有 container 的聚合日志 yarn logs -applicationId application_xxx full.log 21yarn logs一次性会把 stdout、stderr、prelaunch.err 全部打出来内容可能非常大所以先落到文件再针对性搜索。如果是别人的 application、当前用户无权限记得加-appOwner 用户名。还有一类情况是 YARN 开启了日志聚合但还没把日志刷到 HDFS此时yarn logs可能提示找不到日志。处理方式是等 application 彻底结束后再跑一次如果始终拉不到就直接去 NodeManager 本地目录找${yarn.nodemanager.log-dirs}/userlogs/application_xxx/。3.2 用时间线对齐判断谁先崩拿到 full.log 之后我很少从头看到尾。更高效的做法是抓时间线。先找到日志最后一行的时间戳再和yarn application -status显示的 Finish-Time 对齐。如果日志最后停在申请 TaskManager 容器的阶段说明 AM 在等资源如果最后停在初始化某个连接器的阶段说明外部依赖有问题如果日志尾部有一大堆异常堆栈直接搜Caused by。我常用的快速检索命令# 找关键异常 grep -a -E Exception|Error|FATAL|Caused by|KryoException|ClassNotFound|NoSuchMethod full.log | tail -100 # 看最后 50 行判断进程死前的最后一口气 tail -50 full.log这里要特别强调日志尾部很多时候比异常堆栈更有价值。有些失败不会抛异常只是 log 停在某个阶段然后进程被杀。比如 JobManager 在等待 TM 注册时被打断日志尾部可能就是一句 Waiting for taskmanagers to connect这种情况去查资源分配或网络而不是去搜异常。3.3 常见异常关键词与根因对照我把这些年在 Flink on YARN 排障中碰到的高频日志片段整理成了一张对应表。看到什么关键词就优先去查什么方向。日志/诊断关键词大概率根因方向Container is running beyond physical memory limits物理内存超限被 NodeManager 强制杀掉java.lang.NoSuchMethodError/ClassNotFoundException依赖冲突、类加载失败JobSubmissionException/Slot request ... is not fulfilledSlot 不足TaskManager 起不来或资源不够Could not connect to the ResourceManager网络不通、AM 配置错误、HA 元数据异常AccessControlException/Permission deniedHDFS 权限或 Kerberos 票据问题The specified filesystem ... is not accessibleHDFS 路径不存在、挂载异常或权限拒绝exitCode 为 -1000Container 被杀优先看 DiagnosticsContainer ... was preempted by the scheduler资源被抢占通常是队列配额不够或超卖java.lang.OutOfMemoryError: Java heap space堆内存配置过小java.lang.OutOfMemoryError: Direct buffer memory堆外内存/Direct Memory 超限4. 六种高频无限重启根因与我的完整处置过程这一节是我个人经验的精华部分。每一种都是真实踩过、并且确认会反复发生的场景。4.1 内存超配被 NodeManager 直接斩杀这是我遇到最多的一种情况表现非常经典AM 起来过几十秒到几分钟退出Diagnostics 里写着Container is running beyond physical memory limits。根因就是 AM 进程实际占用的物理内存超过了它申请的容器内存。Flink 1.11 的内存模型里进程总内存由堆、直接内存、Metaspace、JVM 开销共同组成。很多人只设置了jobmanager.memory.heap.size忽略了进程总内存的概念结果堆占了一大部分再加上各种 off-heap实际内存超了容器限制。我的处理方式是先把jobmanager.memory.process.size显式配置出来并且保证它小于 YARN 容器可分配内存。比如想让 AM 用 8G 容器我会把 process size 配成 6G 左右留 20%~25% 给系统开销和 k8s/yarn 本身。第二步再检查yarn.scheduler.maximum-allocation-mb如果集群限制单容器最大 8G你死活申请 16G那就是政策和配置的冲突只能找管理员或拆并行度。4.2 用户 JAR 依赖与 Flink 框架类冲突有段时间我接手一个作业AM 每次起来都抛java.lang.NoSuchMethodError: org.apache.hadoop.fs.FileSystem.create进程退出码非 0过几秒 YARN 再拉一个新 attempt再挂。排查后发现用户用 maven-shade-plugin 打 fat jar 时把 Hadoop 相关类也打了进去并且 class 优先级比 Flink 自带的 shaded Hadoop 更靠前。解决办法是打开dependency:tree分析把 hadoop-client、flink-core、flink-runtime 等以provided方式引入或者直接在 shade 配置里排除。另外一个有效手段是调整 Flink 的类加载顺序classloader.resolve-order: child-first改为parent-first或者反向调整取决于你的冲突是哪一类。这个参数不是万能的但很多时候能缓解。4.3 HDFS 权限与目录问题还有一个高频场景是Flink 作业配置了 Checkpoint 或状态恢复目录比如state.checkpoints.dir: hdfs:///flink/checkpoints。AM 启动后要访问这个目录结果权限不够或目录不存在直接抛FileNotFoundException或AccessControlException。这类问题的隐蔽之处在于如果目录不存在且 HDFS 允许自动创建Flink 会帮你建但如果用户没有写权限或者开启了 Kerberos 但 delegation token 没正确传递AM 就会一直失败。排查手段很直接先用hdfs dfs -ls验证路径是否存在、是否可写再检查提交任务的用户是否就是集群里的合法用户。如果是 Kerberos 环境看 AM 日志里有没有No token、AuthenticationException有的话就要在提交命令里加上 token 传递相关配置。4.4 队列资源配额不足AM 反复等死有一种循环重启特别容易让人误判成网络问题或作业问题。RM UI 里 AM 明明起来了但 TaskManager 的容器一直没分配出去等几分钟后 AM 被 YARN 杀掉重来。杀的时候 Diagnostics 经常只有简单一句甘地式提示没有任何报错。我处理过一个并行度 144、每个 TM 配 32G 的作业提交到默认队列。用yarn queue -status default一看队列 MaxCapacity 只有 96G也就是说只能放下两个 TM。作业申请的资源远远超出队列上限AM 永远是等不到容器的状态。解决办法要么调小并行度并同步调小 TM 内存要么把作业提交到更大配额的功能队列。在提交大作业之前先算一笔账总内存 JobManager 内存 TaskManager 内存 × TM 数必须小于队列剩余资源。这条铁律能省掉一半的深夜告警。4.5 客户端与集群版本不一致版本不一致引发的失败日志表现极其混乱。有的是NoClassDefFoundError有的是序列化异常还有的是 AM 起来后跟客户端通信协议对不上一直处于假死状态。你很难从单条报错里直接判断出是版本问题。我的建议很朴素提交任务的节点上尽量使用和 YARN 集群上 Flink 发行版一致的版本。特别要注意 Flink 1.11 前后Application 模式、入口类、内存模型都有变化客户端和集群跨大版本时AM 的启动命令甚至可能直接找不到类。另一个高频坑是 Hadoop 版本不一致客户端 Hadoop 3.x集群 Hadoop 2.xAM 起来后访问 HDFS 时经常会碰到协议或接口不兼容。4.6 大状态恢复卡死被 AM 健康检查杀掉这类案例最隐蔽。现象是作业开启 HA或者配置了从 Checkpoint/Savepoint 恢复且状态很大。AM 起来后 JobManager 花了十几分钟在恢复状态。此时 YARN 的 AM 健康检查认为这个进程心跳停滞直接把容器杀掉然后开始新的 attempt。新 attempt 再去恢复同样的状态又花了十几分钟……如此循环。这种情况日志里往往没有异常堆栈你能看到的只是 YARN 反复把 AM 从 RUNNING 变 DEAD同时 JobManager 的日志停在恢复阶段的进度。解决思路有三个方向一是调大 RM 侧 AM 过期判定时间yarn.am.liveness-monitor.expiry-interval-ms二是调大 Flink 侧 AM 心跳配置yarn.heartbeat.interval-ms确保恢复过程中 YARN 能持续收到心跳三是优化状态恢复本身比如给 JobManager 更多堆内存、减少大状态的并发恢复压力。第三种才是治本的前两种只是给 JobManager 争取时间。5. 把无限重启变成有限重启配置与提交规范最后聊一聊预防。说实话无限次 appattempt这种状态本身就是一种集群配置上的失守。5.1 合理设置 attempt 上限Flink 提交时可以通过-D指定重试次数./bin/flink run-application -t yarn-application \ -D yarn.application-attempts3 \ -p 16 \ /path/to/your-job.jar这个值不能超过 YARN 集群的yarn.resourcemanager.am.max-attempts。生产环境我建议设成 2~3 次就好。设置成无限次看起来高可用实际上一旦进入死循环每一次重启都会占用 RM 调度资源和队列资源还可能影响同队列的其他作业。把次数限制住快速失败反而更容易暴露问题。5.2 内存配比给一个可参考的经验值Flink on YARN 下内存配置是否正确直接决定了 AM 是否会被杀掉。我个人的经验配置表如下适合大多数中等规模作业配置项建议范围说明jobmanager.memory.process.size容器内存的 70%~85%留出进程自身和 YARN 容器的开销taskmanager.memory.process.size单个 TaskManager 容器内存的 80% 左右不要超过yarn.scheduler.maximum-allocation-mbtaskmanager.memory.managed.fraction0.3~0.5数据处理量大、状态多的场景调高yarn.application-attempts2~3不建议设成 0 或无穷大另外一个实测心得不要同时混用新旧两套内存参数。如果你在 flink-conf.yaml 里写了jobmanager.heap.mb又写了jobmanager.memory.heap.size新版 Flink 会忽略旧参数但你会误以为配置生效最后看到的内存使用跟预期完全不符。提交前用./flink info或直接看 Web UI 的配置核对页确认实际生效值。5.3 提交作业前的自查清单基于这些年踩过的坑我每次提交 Flink on YARN 作业前都会过一遍这张清单版本确认客户端 Flink 版本与集群发行版一致Hadoop 版本匹配内存预算JM 总内存 TM 总内存 × 并行度 ≤ 队列剩余资源队列权限确认当前用户有向目标队列提交的权限依赖检查用户 JAR 里不带 flink-core、hadoop-client 等框架类冲突依赖排除干净HDFS 路径Checkpoint/Savepoint 目录存在且有写权限类加载策略如果用过classloader.resolve-order确认和实际依赖结构匹配HA 配置如果开启 HA确认 ZooKeeper/HDFS 连通性检查恢复目录可读重试策略yarn.application-attempts明确设置不要依赖集群默认值。最后说一个我自己的习惯遇到AM 无限重启别急着改配置、删依赖先花十分钟把yarn application -status的 Diagnostics 和 exitCode 完整抄下来。很多时候真相就藏在那几行人类可读的诊断文本里。真正需要看海量日志的场景其实没有想象中那么多。这套方法论帮我在不同环境里解决了很多次同类告警也希望你下次碰见appattempt疯涨的时候能比我当年更快一步找到答案。
返回列表