
1. 从一次“看起来无解”的作业失败说起先讲一个我实际踩过的场景某个凌晨线上跑着几十个MapReduce任务突然一波节点磁盘报错紧接着DataNode进程被系统OOM Killer干掉随后几十个任务像多米诺骨牌一样批量失败。我当时的第一反应是“完了要手动重跑”结果第二天早上看监控发现大部分任务自己缓过来了只有极少数任务需要人工介入而且真正丢数据的作业是零。这个反差让我意识到Hadoop的容错机制尤其是重试机制和故障恢复链路才是分布式系统真正稳如老狗的关键。如果你是个刚接触Hadoop的开发者看到“容错机制”四个字可能觉得离自己很远觉得那是运维的事。但我劝你尽早弄懂这些东西因为它直接决定了三件事第一你的任务在集群的哪个环节可能挂掉第二挂掉之后系统会用什么策略自愈第三哪些坑是你调了参数也没用的。这篇博文我打算从任务执行链路出发把MapReduce/YARN层面的重试机制、ApplicationMaster的故障恢复、ResourceManager的高可用切换以及和ZooKeeper整合的细节全部串起来。我不会只贴配置而是把每个参数背后的权衡逻辑讲明白再配上真实的故障排查记录保证你能看完就上手用也能拿去应对面试里关于“Hadoop任务容错”的连环追问。先说明一下本文默认你已经有了一套可用的伪分布式或小规模集群最好是有3个节点以上的环境。没有也没关系我会在实战环节尽量兼容你在单机上的模拟操作但涉及多节点故障切换的部分还是建议你结合虚拟机或者Docker容器搭一套最小集群来实操很多东西是伪分布式环境完全暴露不出来的。2. 先搞清楚一条任务到底要过几道关卡2.1 任务从提交到结束的完整链条想要理解容错机制首先得看清一条MapReduce任务从提交到结束到底要经历哪些关卡。我把这条链路拆成五个环节这也是故障最容易发生的五个位置第一关是客户端提交。你用hadoop jar命令提交作业客户端会先和ResourceManager通信把jar包、配置、输入分片信息上传到HDFS临时目录然后向ResourceManager提交一个Application。这一步的常见故障是网络抖动、RM地址配错、HDFS临时目录空间不足。第二关是ApplicationMaster的启动。ResourceManager收到请求后会在某个NodeManager节点上分配一个Container用来启动ApplicationMaster后面简称AM。AM是每个作业的“项目经理”负责向RM申请资源、调度map和reduce task、跟踪任务进度。这一步的故障点在于节点资源不足、AM启动脚本出错、本地目录不可写。第三关是Map Task执行。AM拿到资源后为每个map task分配ContainerTask进程去HDFS上读取对应的分片数据执行用户写的map逻辑输出中间结果到本地磁盘。这一步最常见的故障是数据倾斜导致的单task慢、节点磁盘坏道、任务抛异常。第四关是Shuffle与Reduce Task。Map阶段完成后中间结果要经过分区、排序、归并然后拉到Reduce端。这一步对网络和磁盘的压力最大经常出现fetch失败、连接超时、磁盘IO打满。第五关是作业收尾与清理。所有Task执行完后RM会回收ContainerAM退出临时文件被清理。这一步看似简单但经常因为AM在最后状态提交阶段宕机导致作业虽然计算结果正确却始终无法标记为success。整个链条看下来你会发现Hadoop把故障恢复的粒度做得很细有Task级别的重试有AM级别的重新拉起有RM级别的HA切换还有HDFS层面的数据块多副本自愈。各层机制相互独立又彼此配合这也是为什么单个节点挂掉通常不会让整个作业失败的根本原因。2.2 哪些故障是“可恢复”的哪些是“不可恢复”的谈到容错最重要的事情是先给故障分类。Hadoop内部对异常有一套隐式的分类逻辑把故障分为“可重试的”和“不可重试的”两大类。可重试的故障包括网络瞬时抖动导致的RPC超时、NodeManager节点失联、任务进程被外部kill、本地目录临时不可写、数据节点短暂宕机等。这类故障的共同点是重跑一次很大概率能成功重试成本远低于人工介入成本。不可重试的故障包括用户代码里抛出无法被捕获的RuntimeException且没过容错逻辑、输入数据本身损坏比如HDFS块校验失败且没有其他副本、任务依赖的外部系统永久性故障、AM自身的逻辑出现不可修复的异常等。这类故障重试一万次也是白搭盲目调大重试次数只会让集群资源被无效任务白白占住。我之前接过一个线上问题运维同学把mapreduce.map.maxattempts调到10以为能增强容错结果任务每次都卡在一个数据解析异常上重试10次全部失败白耗费了10倍的资源。后来把异常日志翻出来一看是数据里混了一行格式错误属于不可重试故障正确做法是在代码里做脏数据过滤而不是靠重试硬扛。判断故障是否可恢复你可以套用一个简单的思路如果故障发生在“数据读写的源端”大概率可以重试如果故障发生在“用户代码的逻辑内部”大概率重试无效。前者是基础设施抖动后者是业务逻辑问题。搞混这两类一定会被线上环境狠狠教育。3. 重试机制拆解哪一层在重试凭什么重试3.1 Task级别的重试链路Task级别重试是整个容错体系里最基础、也最常用的一环。每个Map或Reduce Task都在独立的JVM进程里运行如果这个进程运行期间报错、崩溃、或者长时间没有心跳ApplicationMaster就会把它标记为失败然后在另一个节点上重新启动一个Container来跑同一个Task。这个环节的核心配置参数是mapreduce.map.maxattempts4 mapreduce.reduce.maxattempts4默认值是4意思是同一个Task最多尝试4次如果4次都失败整个作业就进入失败状态。需要强调的是这个“4次”重试不是指同一次运行内的第4次而是指总共可以启动4个Attempt。举个例子如果第一次运行在第2分钟失败第二次运行在第3分钟失败第三次运行完了那么实际生效的尝试次数就是3第四次是可用的“保险机会”。为什么默认是4而不是2或者10原因是Hadoop的设计者一直在“尽早失败”和“尽量成功”之间找平衡。重试次数太小一群临时故障就能把作业搞死重试次数太大一个必败的任务会反复消耗资源阻塞整个队列。以我经验来看绝大多数线上作业保持默认值就够用除非你的集群硬件稳定性特别差或者某些任务就是偶发性失败比如大促期间网络抖动频繁否则不建议无脑调大。Task失败后AM不会立刻重启任务而是会等待一个时间间隔。这个间隔由节点管理器上报的心跳周期决定默认情况下Task的失败判断依赖于心跳超时。我调优过一些集群为了加速失败判断会把yarn.nodemanager.heartbeat.interval-ms从默认的1000毫秒调低到500毫秒左右代价是增加了一些心跳网络开销但对超大规模集群来说这点开销可以忽略。3.2 任务失败重试时数据是怎么保住的很多新手会问一个问题同一个Task重试那它已经算出来的中间结果还在吗答案是不在了要重新算。这里面有一个关键机制叫“TaskAttempt的确定性假设”。Hadoop默认假设同一个Task无论在哪个节点上运行只要输入相同、代码相同输出就应该相同。因此在Container跑Task之前AM会先做两件事第一把Task对应的输入分片路径告诉新的Container第二把Task需要的依赖文件在本地准备好。Task的执行进度和中间输出全部写在本地的临时目录里Task失败后这个临时目录会被清理新的Attempt从零开始。这也就是为什么MapReduce任务天然适合“一种宁可多算不愿多存”的场景——它靠重算来保证语义正确靠Shuffle阶段的分区机制来保证数据不丢失。相比之下如果换成Spark它会倾向于通过RDD血统关系和检查点来恢复数据这是两种不同的容错哲学没有好坏只有场景匹配与否。Task级别的失败还会带来一个连锁反应当一个节点上连续失败的任务达到一定数量时AM会把该节点加入“黑名单”后续的新Attempt不再调度到这个节点。这个机制叫NodeManager黑名单机制默认情况下yarn.resourcemanager.aml-黑名单相关参数是基于失败比例的。你在yarn-site.xml里可以配置yarn.resourcemanager.am.max-attempts来控制整个Application的重试次数但如果某个节点总是失败最好还是去查一下节点本身的IO、磁盘、CPU资源别指望靠重试绕过物理故障。3.3 ApplicationMaster失败后的恢复机制Task挂掉还能重试那如果“项目经理”ApplicationMaster自己挂了呢这就进入了第二层容错机制。在每个作业提交时ResourceManager会尝试启动AM并给它配置一个重试上限对应的参数是yarn.resourcemanager.am.max-attempts2默认值也是2但注意这个参数在yarn-site.xml中是全局的在mapred-site.xml中也可以通过mapreduce.am.max-attempts对单个作业覆盖。AM失败后ResourceManager不会重启整个作业而是重新在另一个NodeManager上启动一个新的AM实例并从之前写入HDFS的作业状态文件中恢复进度信息。AM的进度和状态存储在HDFS上具体路径通常是你在提交作业时指定的临时目录下的staging目录。这也是为什么HDFS的可用性直接影响作业恢复能力——如果HDFS本身出问题AM的状态恢复就成了无源之水。这里有一个值得注意的细节AM重启后已经完成的Map Task结果是直接复用的不需要重新计算但正在运行的Task如果恰好在AM重启那个瞬间被kill就需要重新调度。换句话说AM重试是“尽量保留已完成成果”的机制但没法做到完全的断点续传。它对慢任务和长任务的价值很大对那种几十秒就结束的小作业意义有限因为AM重启的时间可能比任务本身还长。4. 故障恢复实战模拟一次真实的节点宕机演练4.1 模拟前的集群准备与参数检查理论讲再多不如亲手搞一次故障演练。我建议你准备一个至少3个节点的集群每个节点装好Hadoop 3.x配好HDFS和YARN关闭防火墙或放行对应端口。如果你不想真的准备三台物理机用Docker起三个容器也是不错的选择但记得给每个容器分配至少2GB内存否则YARN很可能因为资源不足而拒绝启动Container。确认好集群状态后先用一条命令把关键配置都检查一遍hdfs dfsadmin -report yarn node -list -all确保三个DataNode都是Live状态三个NodeManager都是RUNNING状态。然后提交一个测试作业我用的是一个简单的WordCount数据量故意稍微大一点比如生成一个2GB的文本文件这样任务运行时间足够覆盖后面的故障操作窗口hadoop jar /opt/hadoop/share/hadoop/mapreduce/hadoop-mapreduce-examples-3.3.6.jar wordcount /input /output作业开始跑起来之后我用yarn application -list拿到Application ID确认它的状态是RUNNING。接下来进入正题模拟Task级别的失败。方法很简单找到某个NodeManager节点的进程直接kill掉。具体来说先登录某台节点通过jps找到NodeManager的PID然后kill -9模拟一个最极端的进程崩溃场景。4.2 观察YARN如何自动拉起失败任务kill掉NodeManager后整个集群不会立刻感知。ResourceManager的节点状态默认是通过心跳超时来感知的这里的核心参数是yarn.resourcemanager.nm.liveness-monitor.expiry-interval-ms600000这个参数默认是10分钟意味着如果某个NodeManager连续10分钟没有向ResourceManager上报心跳RM才会判定它失联。测试的时候如果不想等10分钟可以把这个值临时调小比如调到60000毫秒然后滚动重启ResourceManager让配置生效。NodeManager被判定失联后正在它上面运行的Container会被标记为失败AM会收到事件通知然后把对应的Task重新调度到其他健康的节点上。你可以在ResourceManager的Web UI界面里看到Application Attempt ID会发现有失败的Container计数同时新的Attempt会在其他节点启动。如果Task失败次数没有超过mapreduce.map.maxattempts作业会继续运行最后正常结束如果失败次数超过上限作业就会挂在等待重试或直接失败。第一次做这个实验的时候我手里正好有个作业正在跑然后我手滑把两个节点都kill了结果所有Task都集中到最后一个节点上运行时间翻了三倍但没有一个Task失败——因为Task级别的重试机制把三次失败都扛过去了。这个我在后面的“常见问题”部分还会展开讲。有一点要提醒你kill掉NodeManager不仅仅是让当前Task失败HDFS上该节点的DataNode如果也被killHDFS会开始数据块复制把副本数低于阈值的块从其他节点补齐。这两个恢复过程是并行的一个面向计算一个面向存储都很重要。4.3 模拟ResourceManager宕机ZooKeeper的救场时刻NodeManager挂掉对集群来说只是“小伤”真正的大伤是ResourceManager本身挂掉。在Hadoop 2.0之前RM是一个巨大的单点故障整个集群都得停摆现在有了YARN HA配合ZooKeeper可以实现秒级自动切换。我自己的演练过程是这样的在两台机器上分别配置了yarn-site.xml其中一台作为另一个RM的standbyproperty nameyarn.resourcemanager.ha.enabled/name valuetrue/value /property property nameyarn.resourcemanager.ha.rm-ids/name valuerm1,rm2/value /property property nameyarn.resourcemanager.hostname.rm1/name valuenode1/value /property property nameyarn.resourcemanager.hostname.rm2/name valuenode2/value /property property nameyarn.resourcemanager.zk-address/name valuenode1:2181,node2:2181,node3:2181/value /property property nameyarn.resourcemanager.zk-state-store.parent-path/name value/rmstore/value /property这里有一个容易被忽略的点RM的状态默认不写入ZooKeeper而是写在本地的LevelDB里需要在yarn-site.xml里显式指定property nameyarn.resourcemanager.store.class/name valueorg.apache.hadoop.yarn.server.resourcemanager.recovery.ZKRMStateStore/value /property上面的配置非常关键。如果不配这个RM切换后后新的active RM根本没有之前的应用、队列、容器状态之前的所有作业瞬间变成“孤儿任务”恢复工作无从谈起。配置完成后先手动启动两个RM用命令查看谁是activeyarn rmadmin -getServiceState rm1 yarn rmadmin -getServiceState rm2然后我手动把active的RM进程kill掉观察ZooKeeper的会话超时机制。在ZK里RM宕机后ZooKeeper会等待一段时间才判定会话失效这个超时时间由ZK的tickTime和sessionTimeout相关参数决定通常不会超过几十秒。之后standby RM自动升级为active开始接管集群。整个过程我不需要任何人工干预作业也能自动恢复调度唯一的影响是作业可能短暂停摆但不会整体失败。集成ZooKeeper这个环节很多教程都是一笔带过但我想单独强调一句ZooKeeper在这里起着“分布式锁元数据持久化”的双重作用。没有它两个RM可能同时认为自己是active出现“脑裂”那才是灾难。YARN HA的自动切换依赖ZK的选举机制脑裂问题是绝对要避免的红线所以生产环境请一定严格配置ZK的过半原则节点数至少3台别图省事只配1台ZK。5. 常见问题与排查技巧实录5.1 报错NoClassDefFoundError: org/apache/hadoop/crypto这个报错在跑Hive on Tez或者某些依赖Hadoop生态的组件时特别常见。代码层面它是说JVM在运行过程中找不到org.apache.hadoop.crypto这个类这不是用户代码逻辑错误的锅而是依赖冲突或类加载器隔离问题。排查思路其实很简单先确认你使用的Hadoop相关lib目录里是否有hadoop-common对应的jar确认版本号是否和集群一致。最容易出问题的场景是你本地的Hive lib目录下放了一个低版本的Hadoop common jar而Tez运行时把它后面的类加载优先级搞乱了导致在特定代码路径上引用了新版的crypto类却加载不到。我处理这类问题的标准做法是两步走第一步在Hive的hive-env.sh里手动指定HADOOP_CLASSPATH把集群真正的Hadoop classpath放到最前面第二步检查tez.lib.uris指向的tar包内容确认其中Hadoop相关jar的版本与集群一致。如果是在Docker镜像里跑还要注意镜像里的hadoop版本registry要跟随宿主机集群避免镜像内编译版本和运行时版本不一致。注意这种报错和重试机制没有直接关系但它会导致任务反复在同一位置失败如果你只盯着重试次数去调参不会有一丝效果。这种时候第一件事永远是看完整堆栈而不是盲目重试。5.2 任务一直重试但永远成功不了如何用日志定位根因很多人在Web UI上看到Task failed的红色标记第一反应是“是不是要调大maxattempts”。我劝你先别急花5分钟把日志捞出来看看再决定方向。YARN的聚合日志默认是开启的如果没有开启建议开启并配置日志保留时间。日志位置和命令如下yarn logs -applicationId application_xxx_001 -log_files stdout捞日志的时候优先看两类内容第一类是容器启动时的JVM参数确认资源设置是否符合预期第二类是Task真正跑数据时的异常堆栈特别是Caused by那一行。我遇到过一个案例Task反复失败Web UI上只显示Exit code: 1日志翻到最底部才发现是本地磁盘空间不足导致Container初始化失败。这种情况调大重试次数完全没有意义而且每一次重试都会继续把磁盘写满形成恶性循环。定位根因的一般顺序是先看聚合日志再检查节点磁盘和内存最后再确认网络。如果这三个层面都没有问题而Task总是偶发性失败再考虑是不是资源竞争导致的Container被提前kill这时就要看yarn.nodemanager.pmem-check-enabled和yarn.scheduler.maximum-allocation-mb相关的设置别让容器物理内存超限被强制清理。5.3 关于黑名单和“僵尸”任务的避坑心得最后分享一个我踩过的比较大的坑NodeManager黑名单机制虽然好用但它也有副作用。当一个节点被拉进黑名单之后AM不会再把新的Task调度给它但已经在该节点上运行且未失败的Task并不会被强制kill它们会继续把整个“运行中”记录挂着。如果节点真的是硬故障这些Task最终都会失败收场但如果节点只是短暂抖动这个“半死”状态会让作业在尾部拖很久。解决办法有两个方向一是把yarn.resourcemanager.am-scheduler.max-discipline之类参数调紧让节点在故障后更快被处理二是结合socket超时与心跳配置让NodeManager失联的判定时间缩短。但我建议不要激进地把超时时间调得太短否则网络正常波动也可能触发大规模的重排“频繁切换”本身就是一种稳定性隐患。关于“僵尸”任务我还想多说一句当AM重启后旧的AM的Container可能仍然在节点上占用资源如果没有被及时清理就会形成资源泄露。Hadoop的默认设计是让旧的AM在放弃心跳后自动退出但如果你遇到旧的AM进程一直赖着不走可以手动通过yarn application -kill强制清理或者设置yarn.resourcemanager.am.container-retry.max-retry来约束AM的重试次数避免无限重启。检查一个集群是否健康我习惯先跑一句yarn node -list -all看节点健康状态和容器数再跑一句hdfs dfsadmin -report看数据块副本数是否正常。两个都绿容错机制再复杂也不慌如果某个指标发红先处理物理故障再谈参数调优。6. 延续到生产我在容错调优上最想保留的一条经验我这些年用Hadoop踩过最大的一个认知误区是把“容错”等同于“无限重试”。真正健康的容错体系是分层配合Task失败靠重试节点故障靠黑名单AM挂掉靠重启RM挂掉靠ZooKeeper切换。每层机制解决的是一个特定范围内的故障而不是用一把“重试”的锤子去敲所有钉子。如果你现在正在搭自己的集群或者准备优化线上任务稳定性我建议你把每个重试参数都当成一种“有成本的安全网”。网络抖动频繁就适当调大Task Attempts但千万别调到两位数硬件稳定的集群保持默认值即可ZooKeeper环境一定要单独维护好最好配上独立监控。最后每次故障演练之后请把yarn logs和hdfs dfsadmin -report的输出留存几天出问题的时候这些记录就是你最有力的排查工具。容错机制本身不复杂复杂的是在真实环境里判断「该信哪个机制该等多久」这份判断力只能靠亲手制造故障喂出来。