尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

深入解析MapReduce容错机制:任务重试、推测执行与高可用设计

深入解析MapReduce容错机制:任务重试、推测执行与高可用设计 1. 集群故障是常态先承认这一点才能理解MapReduce的容错设计我最早接触MapReduce时心里一直有个疑问为什么这套计算框架要把容错设计得这么重Map任务跑到一半节点宕机了重新跑一遍不就行了吗后来真的在几百台节点的集群上跑过作业才发现事情没那么简单。一个处理TB级数据的作业Map阶段可能有上万个任务在同时跑任何一台机器出问题都会牵一发而动全身。如果容错机制设计得不够精细一个作业可能跑十几个小时最后因为某个节点上的临时故障整体失败那用户体验是非常糟糕的。MapReduce容错机制的核心目标就是在大规模集群中保证作业可靠运行。所谓大规模集群意味着参与计算的机器数量动辄几十台、几百台甚至更多。在这个规模下故障不是会不会发生的问题而是什么时候发生、一次发生几台的问题。硬盘坏道、内存ECC纠错失败、网卡松动、电源模块老化、内核panic、JVM卡死……任何一层的小概率事件乘以几百台机器和十几个小时的运行时间都会变成一个大概率事件。所以MapReduce的设计哲学从第一步就和单体计算不一样它默认所有参与者都不可靠然后在这个前提下通过一系列机制把作业的整体完成时间变得可控。这套容错体系可以分为几个层次来理解底层是分布式文件系统如HDFS提供的数据冗余解决的是数据不丢的问题中间层是任务级别的失败重试与推测执行解决的是计算不卡的问题上层是作业调度器的高可用设计解决的是大脑不死的问题。很多人聊MapReduce容错只盯着任务重试忽略了数据冗余和调度器HA这是不完整的。本文就按这个层次把MapReduce在大规模集群中保证作业可靠运行的全套机制拆开讲一遍。内容偏原理但我会结合实际排障经验尽量让每个机制都能落地到具体的参数、日志和运维动作上。2. 任务粒度与中间结果落盘容错的第一道地基是让重跑足够便宜2.1 为什么一个输入分片就是一个最小恢复单元MapReduce把所有计算抽象成Map和Reduce两个阶段。Map阶段的输入是HDFS上的文件每个文件会被切分成若干块默认128MB一个Block框架为每个Block启动一个Map任务。这意味着什么意味着整个作业的计算被切分成了大量独立的小任务每个任务只负责一小块数据。这个粒度选择是容错设计的第一块基石。假设一份输入数据有1TB被切成8192个Block就会有8192个Map任务。某个Map任务挂了框架只需要重新调度这一个任务——它对应的输入数据还在HDFS上而且因为是分布式存储同一份数据通常有3个副本只要集群里还有任意一个节点持有该Block的副本新的任务就可以在那里重新运行。相比之下如果整个作业只有一个巨大的任务那它一失败全部数据都要重新处理。这就是为什么MapReduce要选择小任务数据分片的模型它把失败的影响范围从整个作业缩小到了一个Block恢复成本降低了三个数量级。Map任务的失败重试还有一个细节值得注意理想情况下框架会优先把重试任务调度到持有输入数据副本的节点上这样可以省掉跨网络拉取数据的时间。Hadoop的TaskScheduler在调度时会参考数据本地性Data Locality在重试场景下同样适用。如果你发现集群里某个节点挂掉后大量任务被重新调度到了少数几个节点上那就是数据副本只剩两份、且刚好在一个机架上的情况。这也是为什么一般建议副本数保持3个的原因——副本太少重试时的数据本地性会变得很差。2.2 Map端的Spill与Merge中间结果必须先落到本地磁盘Map任务执行的时候并不是直接把结果写给Reducer而是先写到一个内存缓冲区。Hadoop MapReduce中每个Map任务有一个环形缓冲区默认大小是100MB通过mapreduce.task.io.sort.mb配置。当缓冲区的使用率达到一定阈值mapreduce.map.sort.spill.percent默认0.8后台线程就会把缓冲区里的数据Spill溢写到本地磁盘。这个Spill动作是MapReduce容错的重要组成部分。为什么因为Map任务执行过程中输入数据还在源源不断地被处理如果中间结果只存在内存里一旦Map的JVM进程崩溃这些结果就全部丢失了。落盘之后即使Map任务所在节点宕机框架也只是重新执行整个任务——但如果中间结果在内存里那就连重新执行的机会都没有。反过来说正是因为Map端会持续把中间结果spill到本地磁盘整个Map阶段才算是有检查点的不然Map任务跑到90%崩溃前90%的工作就白费了重跑代价高出不少。Reduce端拉取Map输出时也是先写内存写不下再落盘多个Map输出的数据在Reduce端会进行merge最终才送入reduce函数。整个shuffle过程涉及大量磁盘I/O而磁盘恰恰是最容易出故障的组件之一。所以你会发现MapReduce作业的容错设计很大程度上是在和磁盘这个最不稳定的硬件做斗争数据放HDFS上是多副本Map端中间结果尽量落本地盘Reduce端拉取的数据也是一边消费一边落盘就是为了在任何一个环节崩溃时都只需要从最近的已持久化状态重新开始而不是从头再来。2.3 Combine与Reduce的幂等性要求在Map和Reduce之间经常还会有一个Combine操作本地聚合类似Mini Reducer。Combine函数的调用次数是不确定的可能在Map端对部分数据调用也可能在Reduce端对已合并的数据再次调用。这就要求Combine函数必须是幂等的——对同一组数据执行多少次结果都一样。如果用户写的Combine函数不满足这个要求比如用了全局计数器、写了外部状态那么在任务失败重试后结果可能就是错的。Reduce函数同理。框架不会保存Reduce函数的中间状态一个Reduce任务失败后框架会重新执行整个Reduce任务重新拉取所有Map输出从头计算。所以Reduce函数在设计上必须是给定相同输入产出相同输出的纯函数。这其实是一个约束但也是一个恩赐正因为有这个约束框架才敢放心大胆地重试任何任务而不必担心状态不一致的问题。我在实际中见过不少因为Combine非幂等导致的诡异问题某次作业输出结果和上一次对不上排查到最后发现是用户在Combine里写了一个累加外部计数器的逻辑。任务在正常执行时结果是对的但只要有一次任务失败并重跑计数器就被累加了两次最终输出就变了。这种问题非常隐蔽因为它不是必现的只有容错机制真正触发时才会暴露。所以做MapReduce开发时有一条铁律Map、Combine、Reduce函数的执行次数对最终结果必须无影响所有的副作用都必须限定在函数内部。3. 心跳检测与黑名单机制故障发现的速度决定恢复的速度3.1 心跳是MapReduce集群的生命线任务重试的前提是发现任务失败了。在大规模集群里任务和节点数量太多任何集中式的状态收集都做不到实时。MapReduce采用心跳机制来解决这个问题TaskTracker或YARN中的NodeManager每隔固定时间向JobTracker或ResourceManager发送一次心跳报告自己的存活状态、资源使用情况和正在运行的任务进度。Hadoop 1.x时代TaskTracker默认每3秒发送一次心跳到了YARN时代NodeManager的心跳间隔可以通过yarn.nodemanager.heartbeat.interval-ms配置默认是1000毫秒。心跳间隔决定了故障发现时间的下限一个节点宕机最坏情况下要等到下一次心跳该来却没来时调度器才会意识到它失联了。如果节点是直接宕机心跳会直接中断如果只是网络分区心跳可能会延迟。还有一种情况是节点活着但某个任务卡死了比如死循环、长时间GC这时候心跳里会带上任务进度信息调度器可以判断这个任务是否长时间没有进展从而触发超时终止。这里有一个需要注意的细节MapReduce的超时判断不只看心跳还看任务进度更新。Hadoop里有一个参数叫mapreduce.task.timeout默认值是600000毫秒10分钟。如果一个任务在10分钟内没有任何进度更新既没读写数据也没报告新状态框架会判定它卡住了主动kill掉并重新调度。我在调优时经常把这个参数调小比如对于纯计算型任务调到3~5分钟因为这种任务长时间无进度基本就是死循环了。但对于某些数据倾斜严重的场景Reduce任务可能在等最后一个Map任务的输出这时候进度更新本来就会暂停超时时间设太短反而会把正常的等待误判成故障。3.2 任务失败后发生了什么一个任务被判定失败后框架的处理流程很固定但很多细节值得展开。首先是重试次数。Hadoop MapReduce里Map任务默认最多重试4次mapreduce.map.maxattemptsReduce任务也是4次mapreduce.reduce.maxattempts。为什么要设次数限制而不是无限重试因为失败可能不是偶然的而是由于数据损坏、代码bug、或者是某个任务一直打不赢的脏数据。无限重试只会无限期占用集群资源。所以框架的策略是少量重试是给偶发故障留余地超过次数就把这个任务标记为失败整个作业也标记为失败。其次是失败归因。当TaskTracker报告一个任务执行失败时JobTracker会记录失败日志并把这个节点标为不可信。如果同一个节点上失败的次数达到一定阈值该节点会被放入黑名单blacklist短期内不再向它调度新任务。这个机制非常关键在大规模集群中经常存在带病工作的节点——磁盘快满了、内存条有问题、风扇故障导致CPU降频这类机器在硬件彻底坏掉之前会反复导致任务失败。如果不把它们隔离出去作业会陷入任务失败-重试-再失败的循环。黑名单机制的阈值可以通过mapreduce.job.maxtaskfailures.per.tracker参数配置默认是3次。也就是说一个节点上只要有3个不同的任务失败这个节点就会被拉黑。实际运维时黑名单机制偶尔会误伤。比如某次作业里出现了一个对所有节点都不友好的毒任务某个输入数据触发了用户代码的bug导致JVM崩溃这时候每个节点上跑这个任务都会失败最终这个任务耗尽了所有重试次数作业失败并且很多无辜的节点也被拉黑了。这种毒任务其实是数据或代码问题和节点健康状况无关。遇到这种情况不要急着从黑名单里移除节点先去查那个任务对应的输入分片是什么大概率能发现脏数据的影子。3.3 日志与诊断信息容错机制的事后诸葛容错不只是失败后重跑还要能回答为什么会失败。每个Map/Reduce任务的失败都会在JobHistory里留下记录包括退出码、stderr输出、JVM崩溃信息等。Hadoop会把任务运行日志以日志聚合log-aggregation的方式集中存放到HDFS上方便事后统一查看。我处理过的一次典型故障是这样的一个作业的Map任务频繁失败任务重试4次后作业整体失败。从JobHistory看失败原因清一色是Container killed by ApplicationMaster for exceeding memory limits。这个信息表面上是说容器超内存被杀了但实际上是Map任务处理某个分片时加载了异常多的数据到内存。最后定位到是某个输入文件里混入了一个超大字段导致解析时内存暴涨。如果当时没有日志聚合这种问题几乎无法排查——几十台节点上几百个容器你不可能一台一台登录去看日志。所以日志聚合本身也是容错体系的重要一环它让分析失败这件事实实在在变得可操作。4. 推测执行用多余的算力换取作业的确定性4.1 为什么作业的正确性不依赖单台机器的速度集群里几百台机器跑同一个作业理论上所有任务应该差不多同时完成。但现实是总会有那么一两个任务拖后腿。原因很多——某台机器磁盘I/O变慢、CPU被其他作业抢占、内存不足导致频繁GC、甚至就是机器本身性能老化。这些慢任务本身没有失败也没有报错可它们会拖慢整个作业的完成时间。MapReduce给出方案是推测执行Speculative Execution当一个任务运行时间明显超过同类任务的中位时间时框架会在另一台机器上启动一个相同的任务副本两个副本谁先完成就用谁的结果另一个会被杀掉。这个方案的本质是用多余的算力来对冲不确定性。在大规模集群上算力通常是有富余的如果没富余慢任务也不会拖慢整体进度——因为资源都被占满了。既然富余算力闲着也是闲着不如让慢任务跑一个竞争副本反正最坏情况也只是多消耗一份资源而已。4.2 推测执行的触发阈值和判定逻辑Hadoop早期版本的推测执行逻辑比较粗糙主要是看任务平均进度和当前任务进度的差距。以Map任务为例当某个Map任务的已完成比例显著低于所有Map任务的平均完成比例且运行时间超过阈值时JobTracker会为该任务启动备份任务。YARN版本中逻辑有所演进但还是基于任务进度和运行时间的启发式判断。相关的核心参数包括mapreduce.map.speculative是否对Map任务启用推测执行默认truemapreduce.reduce.speculative是否对Reduce任务启用推测执行默认truemapreduce.job.speculative.slowtaskthreshold判定为慢任务的平均进度阈值默认0.3意思是低于平均进度30%会触发mapreduce.job.speculative.slownodethreshold判定为慢节点的阈值默认0.3意思是节点上任务平均进度低于整体30%该节点会被标记为慢节点实际调优中有一个反直觉的经验推测执行在Reduce阶段往往弊大于利。原因很简单Reduce任务要拉取所有Map的输出数据任务之间高度耦合启动一个并发副本意味着两个Reduce任务同时拉取同一批数据shuffle阶段的网络I/O会翻倍。如果集群的网络带宽本来就紧张推测执行不但不能加快作业反而可能让整个集群的shuffle变慢。所以很多生产环境会把mapreduce.reduce.speculative设为false只保留Map阶段的推测执行。4.3 推测执行不是万能的这些场景要主动关闭推测执行有一个隐含前提任务必须是确定性的即两个副本的执行结果一致。如果用户代码里有外部副作用写数据库、调用外部接口、操作共享文件那么两个副本同时执行副作用就会被执行两次可能导致数据重复或状态错乱。比如在Map端直接调用外部API获取数据推测执行一旦触发同一个Map任务的两个副本会同时调用API两次。如果API不是幂等的结果就很麻烦。所以当MapReduce作业的用户代码包含外部副作用时一定要主动关闭推测执行或者修改代码把外部调用挪到Map/Reduce函数之外。另一个场景是数据倾斜。如果Reduce端的数据分布严重不均某个Reduce任务处理的key特别多它的运行时间天然会比别的Reduce任务长很多。这时候推测执行会误判这是一个慢任务启动副本却发现两个副本一样慢白白消耗资源。这种情况的正确做法是处理数据倾斜比如加盐、二次聚合而不是靠推测执行解决。我记得有个数据清洗的实训项目用MapReduce处理网约车数据时把订单数据按城市分组结果某几个超大城市的key严重倾斜。按默认配置开了Reduce推测执行作业反而比关掉推测执行时慢了20%——这就是典型的推测执行在数据倾斜场景下帮倒忙的案例。5. 从JobTracker到YARN作业级高可用的演进5.1 单点故障MapReduce容错的最后一块短板前面讲的任务重试、黑名单、推测执行都解决了Worker挂了怎么恢复的问题。还有一个更严重的隐患整个作业的调度中枢——JobTracker——如果挂了怎么办在Hadoop 1.x时代JobTracker是整个MapReduce集群的单点。TaskTracker只和JobTracker通信所有任务调度、资源分配、状态管理都集中在JobTracker内存里。JobTracker一旦宕机所有正在运行的作业全部失败而且没有自动恢复机制需要人工重启并重新提交作业。对于运行十几个小时的大作业来说这个风险是不可接受的。JobTracker单点问题的本质是它既要管资源哪些节点有多少可用的Map/Reduce槽位又要管作业作业划分、任务调度、进度收集两个职责耦合在一起状态量极大做高可用就非常困难。状态太大意味着备份的成本很高恢复时需要重建所有作业和任务的状态时间长且容易出错。5.2 YARN如何拆解问题资源管理和作业调度分离YARNHadoop 2.x起重新设计了架构把原来JobTracker的两个核心职责拆开ResourceManagerRM负责集群资源和应用调度ApplicationMasterAM负责单个作业的任务划分和调度。这样一来单个作业的失败只会影响这个作业本身——AM挂了RM可以重新启动一个新的AM从之前保存的状态恢复作业进度。这就是作业级容错的本质提升。在YARN中AM本身也受容错机制保护。MRAppMaster是MapReduce作业的ApplicationMaster它的状态会定期写入到HDFS上的状态存储中通过yarn.app.mapreduce.am.job.store.class配置默认是org.apache.hadoop.mapreduce.v2.app.job.impl.FileSystemJobStoreImpl。AM如果挂掉RM会在另一个节点上重新启动MRAppMaster实例从HDFS读取之前保存的状态恢复所有任务执行情况然后继续调度剩余任务。这个恢复过程中已经完成的任务不需要重跑只需记录它们的完成结果即可。这就解决了一个很实际的问题以前JobTracker挂了作业只能从头跑现在AM挂了作业可以从断点继续跑。对十几个小时的大作业来说省下的时间可能是一整个晚上的窗口期。5.3 ResourceManager的HA设计RM本身仍然是集群级别的单点但它可以做成Active-Standby模式。两个RM节点一个Active一个Standby它们通过ZooKeeper进行选主。Active RM把内部状态包括所有应用的列表、资源分配信息写入到一个共享存储ZooKeeper或HDFSStandby RM持续监控并同步这个状态。当Active RM故障时Standby RM通过ZooKeeper的选举机制自动接管成为新的Active节点。这个设计中有一个有意思的细节RM恢复后所有正在运行的ApplicationMaster需要重新向新的RM注册并汇报状态。AM定期向RM发送心跳如果发现RM变了epoch编号变化就会重新注册。所以RM的HA切换对单个作业的影响很小最多是几分钟的暂停而不是作业失败。对于MapReduce作业来说一旦AM完成注册任务可以继续跑之前完成的Map任务结果依然有效。从实践来看在大规模集群中很多作业的失败并不是任务本身的问题而是RM或调度器层面的问题。YARN的HA设计可以说是MapReduce容错机制从Hadoop 1.x到2.x最重要的升级之一。它把作业失败可重跑提升到了框架本身可切换的层面这台机器终于不再有必死的单点了。6. 实际运维中的容错调优与踩坑记录6.1 Reduce输出原子性重复执行不一定重复写任务重试还有一个容易被忽视的问题Reduce任务成功提交输出之后AM收到了完成消息但还没来得及记录状态AM就挂了。恢复后的AM不知道这个Reduce任务已经完成了于是会重新调度这个Reduce任务。这时候如果Reduce任务的输出已经写入了最终结果目录那么第二次执行就会覆盖写入或报错。Hadoop通过OutputCommitter机制来解决这个问题。简单说每个Reduce任务先写到临时目录_temporary任务真正成功后才通过commit操作把临时文件移到最终目录。如果Reduce任务被重复执行第一次执行虽然写入了临时文件但没有commit成功因为AM挂了第二次执行会被识别为该任务已提交过或者继续在临时目录上操作最终只commit一次。这个机制在正常情况下运行得很好但在一些边界场景会出问题。比如如果你在作业中自己实现了OutputFormat却没有正确地实现OutputCommitter的原子提交逻辑就有可能在任务重试后产生重复数据或损坏数据。我做招聘数据清洗实训时就遇到过自定义OutputFormat没处理好commit逻辑导致作业在任务重试后输出文件里出现了重复行的情况。当时排查了很久最后确认是OutputCommitter的实现问题而不是Mapper或Reducer的bug。6.2 Shuffle期间的Disk故障中间结果丢失的恢复策略Reduce任务执行过程中需要从Map任务所在的节点拉取数据。如果某个Map输出文件所在的节点磁盘故障这些中间结果就丢了。怎么办框架的做法是如果该Map任务已经被标记为完成但输出不可读AM会判定该任务需要重新执行并调度它在另一个节点上重新计算。Hadoop里有个参数叫mapreduce.job.jvm.numtasks如果设置为1即JVM复用关闭每个任务单独起JVM那么重跑一个Map任务就是完全的冷启动如果设置了JVM复用重跑会更快一些但故障隔离性也更差。这个场景最容易暴露数据本地性和副本冗余的价值。如果HDFS副本数只剩1某个节点磁盘故障导致Map输入数据直接丢失那整个作业就彻底失败了。所以我在维护数据冗余时特别关注HDFS的健康状态——集群里每天多几块坏盘是常态但只要副本数还能维持在2以上MapReduce作业基本都能扛过去。6.3 参数调优建议清单结合排障经验我认为下面这些参数是与容错直接相关的值得在生产集群上仔细调mapreduce.map.maxattempts和mapreduce.reduce.maxattempts单个任务最大重试次数。偶发故障多的集群可以调到6~8但不要超过10否则毒任务会长时间拖死作业。mapreduce.task.timeout任务无进度超时阈值。默认10分钟偏保守计算密集任务建议调到3~5分钟IO密集任务可以保持10分钟。mapreduce.job.maxtaskfailures.per.tracker节点黑名单阈值。默认3如果集群老化严重建议调到5避免因为个别任务的问题误伤大量节点。mapreduce.map.speculative和mapreduce.reduce.speculative推测执行开关。如果作业有外部副作用或数据倾斜主动关掉Reduce推测执行是明智的。yarn.resourcemanager.am.max-attemptsAM最大重试次数。默认2对大作业建议调到4以上因为AM恢复的成本很低多点重试次数能显著提升作业成功率。mapreduce.reduce.shuffle.maxfetchfailuresReduce拉取Map输出的最大失败次数默认10。网络抖动严重的集群可以增加这个值避免shuffle阶段因为偶发网络问题把Reduce任务判定失败。这些参数不是越大越好也不是越小越好而是在快速失败和容忍偶发故障之间找一个平衡点。我个人建议先在测试集群里模拟节点宕机、磁盘故障等场景观察作业的恢复时间再根据生产环境的具体情况调整。6.4 一个和排序相关的容错案例分组排序作业为什么总是慢顺着热搜词里那个MapReduce排序—分组排序的实训案例多说两句。很多人做分组排序时Reducer拿到的数据并不是天然有序的要靠shuffle过程中的排序来保证同一个key的所有value是有序的。而在shuffle阶段如果某个Reduce任务的某个Map输出分区丢了比如对应节点宕机框架会重新调度那个Map任务来重新生成数据。但如果这个Map任务的输入数据本身就存在倾斜比如一个Block特别大或者里面包含大量相同的key重新执行的时间就会特别长。这时候Reduce任务可能已经被判定为失败并重试了一次而重试的Reduce任务又要重新拉取所有Map输出shuffle网络负载翻倍整个作业就陷入一种反复重试但永远跑不完的恶性循环。排查这种问题单靠调容错参数是不够的必须从数据分布入手比如对key做加盐处理或者调整Partitioner让数据分布更均匀。这个案例给我的教训是容错机制能防住节点故障这种硬故障但防不住数据问题和代码问题引发的软故障。前者靠框架解决后者靠开发者和运维者的经验解决。写在最后说实话我做MapReduce运维这几年真正让作业跑挂的大多数时候并不是框架的容错机制失效而是我们对容错机制的理解不够——要么没有配置好参数要么在代码里写了不满足幂等要求的逻辑要么在数据倾斜时还指望推测执行来救场。框架能做的是在硬件故障和偶发异常面前保证作业不崩而代码质量和数据质量才是作业能不能正确、高效跑完的最终决定因素。如果你正在学习MapReduce的容错机制我的建议很简单先去小集群上做一次断电实验把TaskTracker或NodeManager直接kill掉观察作业怎么恢复、日志怎么记录、哪些任务被重调度。这种实验做几次你对容错的理解会比看十篇文档都深。等你真正经历过一次几百个任务里有两三个节点宕机、作业依然按时跑完的场景就会明白MapReduce的容错设计到底值在哪了。
返回列表