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

资讯详情

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

MapReduce高可用架构详解:NameNode与ResourceManager故障转移实践

MapReduce高可用架构详解:NameNode与ResourceManager故障转移实践 1. 大数据MapReduce高可用架构的整体设计思路1.1 为什么高可用是MapReduce绕不开的坎大数据集群跑着跑着Master节点突然挂掉整个MapReduce作业全部报废这是最让人崩溃的场景。尤其是生产环境里任务一跑就是几个小时一个NameNode宕机直接让所有客户端卡死数据写不进去也读不出来运维群里瞬间炸锅。我在一线维护过好几套Hadoop集群坦白讲高可用这件事不是加分项而是保命项。很多人误以为MapReduce天生就能容错因为任务失败会自动重试。但这个容错只针对TaskTracker/NodeManager层面的计算节点挂了如果是资源调度的主节点或者存储元数据的主节点挂了整个集群就瘫痪了。MapReduce的运行依赖两个关键角色ResourceManager负责分配计算资源NameNode负责管理文件元数据。这两个节点一旦单点故障后面所有作业都得排队等人工干预严重时还会丢元数据数据目录直接不可用。所以高可用架构的核心目标就两条第一主节点故障后能自动切换业务无感知或秒级感知第二数据不丢状态不丢作业可以继续跑。围绕这两条主流方案就是基于ZooKeeper的Active/Standby模式配合共享存储或JournalNode来同步元数据。这也是我今天想重点拆解的内容为什么要这样设计关键配置项怎么选踩过哪些坑。1.2 高可用架构方案的选型逻辑Hadoop生态里高可用方案不是只有一种得根据集群规模和业务场景来选。第一种是单机伪分布式就是学习用的NameNode和ResourceManager都在一台机器上谈不上高可用。第二种是HA模式用两台机器一主一备通过ZooKeeper协调自动故障转移。第三种是联邦模式多个NameNode分管不同目录但这主要是为了解决横向扩展问题并不完全等价于高可用。第四种是跨机房容灾比如用BackupNode或快照同步成本高一般公司用不上。实际生产环境最常用的就是第二种两个NameNode 两个ResourceManager 一个三节点的ZooKeeper集群。这套架构的好处是比联邦简单ZooKeeper本身也是成熟组件运维知识通用。而且MapReduce的HA方案和HDFS的HA是配套的ResourceManager的HA实现可以复用ZooKeeper的协调机制部署一次解决两个主节点的高可用问题。选型时还要注意版本Hadoop 2.x之后才正式支持NameNode HA和ResourceManager HA如果你还在用老掉牙的1.x版本那还是先升级再说。另外如果集群规模特别大比如上千节点可以考虑用YARN Federation但那是后话一般公司到不了那个量级。1.3 高可用架构的整体拓扑与角色划分一个标准的MapReduce高可用集群角色划分是这样的Active NameNode负责处理所有客户端读写请求维护元数据。Standby NameNode实时同步Active的元数据随时准备接管。JournalNode至少3个存储EditLogNameNode写日志的共享存储。ZooKeeper集群至少3台负责自动故障转移时的主节点选举。Active ResourceManager负责资源调度接收作业提交。Standby ResourceManager通过ZooKeeper感知Active状态故障时自动切换。注意Active和Standby的NameNode之间不是简单的热备Standby是实时同步EditLog的内存里有完整的元数据镜像所以切换时间在几十秒级别。而ResourceManager的HA则稍有不同两个RM通过ZooKeeper抢占Active状态但作业的运行状态和资源分配信息需要通过ZK来持久化和恢复。整体拓扑可以画成一张图但这里不用mermaid我用文字描述客户端请求经过ZooKeeper找到Active NameNode的地址NameNode写入EditLog到JournalNode集群Standby NameNode从JournalNode拉取日志并应用到内存镜像ResourceManager的HA依赖ZK的分布式锁来实现主备抢占。这套架构把控制流和数据流分开控制流走ZooKeeper和JournalNode数据流直接走DataNode避免了大量元数据通信拥塞在单一节点上。2. 核心组件与高可用实现原理2.1 NameNode高可用EditLog的共享存储与自动故障转移NameNode的高可用核心在于EditLog的共享。Active NameNode每做一次元数据修改比如创建文件、删除目录都会把操作记录写入EditLog。普通模式下EditLog写在本地磁盘故障了没法恢复。HA模式下EditLog必须同时写到多个JournalNode上只要大多数一般是2/3写入成功就算提交成功。这样设计的好处是什么我举个例子Active节点写入EditLog是写本地 写JournalNode客户端发起的元数据操作要等FSNamesystem向JournalNode flush成功后才返回成功。如果Active节点挂了Standby节点已经通过JournalNode同步到了最新的日志它可以在内存里重放这些日志把元数据恢复到最新状态。关键是有没有丢日志的问题这取决于客户端是否收到了成功响应。只要Active节点确认了写成功日志就一定在JournalNode上所以不会丢。自动故障转移则依赖ZooKeeper。具体流程是Active NameNode和Standby NameNode都以临时节点的形式注册在ZK上Active会持有ZK的一个分布式锁。Active故障后临时节点自动消失Standby通过ZK的watcher机制感知到锁被释放然后发起选举把自己变成Active。这里有个关键设计两个节点必须能同时访问共享存储JournalNode但同一时刻只能有一个Active否则会出脑裂。防止脑裂的机制叫fencing简单的说当节点被切换后旧Active会被强制隔离比如kill掉进程或者用SSH到原节点执行fuser命令。2.2 ResourceManager高可用基于ZooKeeper的自动切换ResourceManager的HA原理和NameNode类似但实现细节不同。RM的HA使用ActiveStandbyElector这个类它会在ZooKeeper上创建一个临时节点来抢Active状态。抢到锁的那个RM就是Active另一个是Standby。RM是无状态的不对RM本身维护了每个Application的资源请求、Container分配、NodeManager状态等大量运行时数据。为了实现HA在RM发生转移后所有的ApplicationMaster重新向新的RM发送心跳RM会从ZK中恢复之前持久化的应用状态。这里需要注意RM的HA不是零中断现有作业会经历一次恢复期新RM会把所有未完成的Application重新调度这个过程可能让正在跑的MapReduce任务重新申请Container但任务不会丢。还要注意一个点YARN的HA并不需要共享存储。RM把状态写到ZooKeeperApplicationMaster也向ZK写入自己的状态。所以配置上不需要JournalNode但ZK的性能会影响切换时间。ResourceManager HA的配置项在yarn-site.xml中需要设置yarn.resourcemanager.ha.enabled为true配置多个RM节点并指定ZK地址。切换后客户端需要通过yarn.resourcemanager.ha.rm-ids来找到Active节点所以推荐配置yarn-site.xml里的yarn.resourcemanager.cluster-id和yarn.resourcemanager.zk-state-store.address。2.3 MapReduce作业在HA架构下的容错机制MapReduce作业本身有两个层面的容错一个是任务级别的一个是应用级别的。任务级别MapTask或ReduceTask执行失败时ApplicationMaster会将其重新调度到另一台NodeManager上默认重试次数为4次。这种容错不需要HA系统参与是YARN框架自带的。应用级别当ResourceManager故障转移后ApplicationMaster的信息从ZK恢复YARN会尽量保留已运行的Container让作业继续跑。但这里有个现实情况如果ApplicationMaster本身也挂了作业就会失败。YARN可以对ApplicationMaster设置重启次数通过yarn.resourcemanager.am.max-attempts来配置默认是2次。你说这个够不够在HA场景下建议将AM重启次数设置到4次因为RM切换的那几十秒AM可能因为无法连接RM而自杀多点余地总是好的。MapReduce作业还需要考虑数据本地性问题。在HA架构下如果NameNode切换DataNode的块报告会延迟调度器可能暂时无法精确感知数据位置任务会退化为随机分配造成网络开销。这块不用太焦虑等DataNode重新向新的NameNode注册后很快恢复。如果频繁切换就要排查是不是元数据同步有延迟或者ZK选举配置不合理。3. 实操搭建一套高可用MapReduce集群的详细步骤3.1 环境规划与前置条件动手搭建前先把机器和软件版本定好。我用的是Hadoop 3.3.6JDK 8ZooKeeper 3.7.1三台物理机或虚机配置如下主机名IP角色node110.0.0.11NameNode, ResourceManager, ZooKeeper, JournalNodenode210.0.0.12NameNode, ResourceManager, ZooKeeper, JournalNodenode310.0.0.13DataNode, NodeManager, ZooKeeper, JournalNode这里我用了3个节点其中node1和node2互为主备node3作为仲裁和存储节点。注意最少需要3个JournalNode才能容忍1个故障所以我在3台都部署了JournalNode。如果机器不够也可以用单节点JournalNode但那就不叫高可用了。部署前需要做几件事配置主机名和/etc/hosts保证所有节点免密SSH登录关闭防火墙或开放必要端口设置NTP时间同步。这些基本功不说了直接跳过容易被忽略但很重要必须确认机器时钟偏差不要超过ZooKeeper默认的sessionTimeout默认40秒否则ZooKeeper会频繁断开连接。3.2 ZooKeeper集群的搭建与验证先把ZooKeeper装好因为它是一切协调的基础。下载解压后在conf目录下创建zoo.cfgtickTime2000 initLimit10 syncLimit5 dataDir/data/zookeeper clientPort2181 server.1node1:2888:3888 server.2node2:2888:3888 server.3node3:2888:3888然后在dataDir目录下创建myid文件分别写入1、2、3。启动后用zkServer.sh status查看状态应该有一个leader两个follower。如果都显示standalone说明配置没生效检查myid文件路径和端口是否被占用。到这里注意ZooKeeper集群的网络分区容忍度是超过一半也就是说3个节点中最多挂1个如果挂2个整个ZK集群就不可用了HA也就崩了。所以生产环境推荐至少5个ZK节点能容忍2个故障。我演示环境用3个但你在生产千万别这么做。3.3 HDFS HA的配置文件详解HDFS HA的核心配置在hdfs-site.xml里。我把关键项列出来property namedfs.nameservices/name valuemycluster/value /property property namedfs.ha.namenodes.mycluster/name valuenn1,nn2/value /property property namedfs.namenode.rpc-address.mycluster.nn1/name valuenode1:8020/value /property property namedfs.namenode.rpc-address.mycluster.nn2/name valuenode2:8020/value /property property namedfs.namenode.http-address.mycluster.nn1/name valuenode1:9870/value /property property namedfs.namenode.http-address.mycluster.nn2/name valuenode2:9870/value /property property namedfs.namenode.shared.edits.dir/name valueqjournal://node1:8485;node2:8485;node3:8485/mycluster/value /property property namedfs.client.failover.proxy.provider.mycluster/name valueorg.apache.hadoop.hdfs.server.namenode.ha.ConfiguredFailoverProxyProvider/value /property property namedfs.ha.automatic-failover.enabled/name valuetrue/value /property property namedfs.ha.fencing.methods/name valuesshfence/value /property property namedfs.ha.fencing.ssh.connect-timeout/name value30000/value /property /xml这里面最难理解的是dfs.namenode.shared.edits.dir它指定了JournalNode的地址。注意qjournal前缀后面是三个节点的IPC端口默认8485最后是nameservice ID。如果你用的是NFS共享存储这里就换成NFS的路径但NFS本身也有单点问题我建议用JournalNode。还有一个很重要的点dfs.ha.fencing.methods。如果设置成sshfence那么当ZK触发切换时旧Active节点会被SSH连接并执行fuser -k 8020/tcp默认命令来杀掉进程。这里需要配置免密SSH而且要确保fuser命令在PATH路径下。如果不配置fencing脑裂风险会大增两个NameNode同时写EditLog数据直接损坏。这是我实际踩过的坑别省这一步。3.4 YARN HA的配置文件详解ResourceManager的HA配置在yarn-site.xmlproperty nameyarn.resourcemanager.ha.enabled/name valuetrue/value /property property nameyarn.resourcemanager.cluster-id/name valuemycluster/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.webapp.address.rm1/name valuenode1:8088/value /property property nameyarn.resourcemanager.webapp.address.rm2/name valuenode2:8088/value /property property nameyarn.resourcemanager.zk-address/name valuenode1:2181,node2:2181,node3:2181/value /property property nameyarn.resourcemanager.recovery.enabled/name valuetrue/value /property property nameyarn.resourcemanager.store.class/name valueorg.apache.hadoop.yarn.server.resourcemanager.recovery.ZKRMStateStore/value /property property nameyarn.resourcemanager.zk-state-store.address/name valuenode1:2181,node2:2181,node3:2181/value /property /xml注意yarn.resourcemanager.recovery.enabled必须设置为true这样RM才能把应用状态持久化到ZK。如果不设置切换后所有历史作业都会丢失。核心参数是yarn.resourcemanager.store.class默认是ZooKeeper实现这个不需要改动但你需要确认ZK地址写对。还有一个容易被忽略的yarn.resourcemanager.ha.enabled和yarn.resourcemanager.cluster-id必须和HDFS的nameservice区分开不要写成一样的值否则在ZK上路径会冲突。3.5 启动顺序与状态检查搭建完成后启动顺序非常关键。我第一次部署时直接start-dfs.sh结果ZK和JournalNode还没就绪NameNode一直注册失败日志刷屏。正确的顺序是启动ZooKeeper3台机器分别执行zkServer.sh start启动JournalNode在core-site.xml配置了journalnode后执行hdfs --daemon start journalnode在所有节点上启动格式化HDFS元数据只在node1上执行hdfs namenode -format在node1上启动NameNodehdfs --daemon start namenode在node2上执行hdfs namenode -bootstrapStandby把元数据从node1拉取过来初始化HA状态在node1上执行hdfs zkfc -formatZK启动HDFSstart-dfs.sh启动YARNstart-yarn.sh启动后用hdfs haadmin -getAllServiceState来查看NameNode状态应该一个active一个standby。用yarn rmadmin -getServiceState rm1查看YARN状态。如果状态都是standby说明ZK中锁的初始化有问题检查zkfc进程是否启动。初始化HA状态那步很容易忘我干过几次结果两个NameNode都认为自己是standby客户端根本无法连接。执行hdfs zkfc -formatZK会清空ZK上的HA状态等于重新初始化选举环境所以这个命令在首次部署和后续故障恢复中都很重要。4. 常见故障与排查技巧实录4.1 故障转移失败ZooKeeper节点没起来现象主NameNode宕机后集群没有任何反应所有客户端都卡死但Standby节点没有接管。排查步骤先看Standby节点上zkfc的日志hadoop-*-zkfc-node2.log。最常见的原因是zkfc进程因为ZooKeeper连接超时退出了。这时候检查ZooKeeper的存活状态zkServer.sh status。如果ZK节点没问题再看JournalNode是否正常因为Standby只有从JN拉取到日志才算ready。用jps查看每个节点的进程JournalNode如果没起来EditLog就无法同步。解决方案按顺序重新启动ZK、JN、NameNode、zkfc。如果仍然不行手动执行hdfs haadmin -transitionToActive --forceactive nn1强制切换。注意forceactive参数在没有其他Active节点时才用否则可能出现双Active。4.2 脑裂两个NameNode同时变Active现象客户端报错Name node is in safe mode或者Operation category READ is not supported同时发现两个NameNode都显示active。原因fencing配置失效。我之前把dfs.ha.fencing.methods设置成shell(true)Shell命令没有正确执行旧节点进程没被杀掉。另一个原因是SSH免密失效了fencing连接不上旧节点。预防措施fencing方法建议用sshfence并定期验证免密SSH是否正常。还可以设置脚本让fencing执行更强的操作比如关闭物理机网卡或者重启。这听着很暴力但在生产环境确实是有效手段。4.3 RM切换后MapReduce作业心跳丢失现象RM从rm1切到rm2后正在跑的MapReduce作业全部失败日志显示ApplicationMaster连续N次心跳超时。原因分析RM切换期间AM和RM的TCP连接断开AM需要重新发现新的RM地址。YARN中AM可以通过yarn.resourcemanager.ha.admin.address来获取新RM地址但如果客户端或AM在RM切换后的长段时间内反复重试最终会失败。解决方案调大AM的重试次数和间隔。在mapred-site.xml中设置yarn.resourcemanager.am.attempts和yarn.am.liveness-monitor.expiry-interval-ms。常用配置是am.max-attempts4liveness监测时间10分钟。同时确保客户端引用的yarn-site.xml是最新的特别是yarn.resourcemanager.ha.rm-ids配置要包含所有RM节点否则AM只向一个RM地址发心跳一旦切换就失联。4.4 性能问题故障切换后集群变慢这种情况通常是原来的Active节点恢复后重新成为Active但内存中的元数据是几小时前的快照需要花很长时间重新读取DataNode的块报告。表现为HDFS吞吐量下降MapReduce作业卡在等待数据块阶段。我喜欢在NameNode切换后先手工触发块报告通过hdfs dfsadmin -report检查DataNode状态然后用hdfs dfsadmin -triggerBlockReport命令让所有DataNode立刻上报。这会加重网络负载但能加速恢复。另一个办法是在切换后临时提高namenode的service handler数即在hdfs-site.xml中调大dfs.namenode.service.handler.count但需要重启NameNode不是紧急手段。5. 高可用集群的参数调优与运维建议5.1 关键的ZooKeeper参数ZooKeeper的性能直接影响故障转移时间。tickTime默认2000毫秒是心跳的基本单位。initLimit是follower在启动时连接leader并同步数据的最大时间单位是tickTime。如果你的集群跨机房网络延迟高建议把initLimit和syncLimit调大。我一般设置initLimit20syncLimit10给网络抖动留出余量。另外sessionTimeout默认是40000毫秒由tickTime*20计算得出如果ZK节点频繁断开可以把sessionTimeout调高到60000但这会让故障检测变慢有利有弊。我的经验是在高可用切换速度要求不苛刻比如10秒内的场景调大sessionTimeout比调小更稳因为很多断开是瞬时抖动。5.2 JournalNode的数量与吞吐量JournalNode写入是同步的每个EditLog要写多个JN写入性能受影响。如果集群的元数据操作非常频繁比如每秒上千次mkdir、put操作建议把JN数量从3增加到5因为3个JN要等2个响应5个JN要等3个响应在吞吐量上差距不大但容错性提高很多。还有一点JournalNode最好禁掉swap分区因为JN的写入延迟会直接影响客户端操作延迟。如果JN因为swap变慢客户端所有写操作都会变慢这是连锁反应。5.3 监控与告警HA集群没有监控就是裸奔。推荐至少监控以下指标ZooKeeper节点的连接数、延迟和leader状态。NameNode的Active/Standby状态切换次数。JournalNode的写入延迟和磁盘空间。RM的Active/Standby状态以及ZK中存储状态的size。我用过自定义脚本配合Prometheus和Grafana来做监控zookeeper和hadoop都有现成的exporter不用完全自己造轮子。告警规则就两条状态异常立即告警切换次数在短时间内超过2次立即通知值班。5.4 备份与降级策略高可用不是万无一失的。JournalNode集群如果全部挂掉NameNode无法写EditLogHA名存实亡。所以我还建议做每日自动备份用hdfs dfsadmin -fetchImage定期抓取元数据镜像存到独立存储。JournalNode本身也有冗余但物理上别和NameNode放在同一台机器虽然有例外但风险大。另外要制定降级策略。如果ZK集群彻底损坏可以临时把hdfs-site.xml里的dfs.ha.automatic-failover.enabled设为false手动运维主备。操作步骤在Standby上执行hdfs haadmin -transitionToStandby然后手动指定Active。这种方式不优雅但能保住数据。6. 总结之后还想多说两句最后分享一个小经验高可用架构配置完成后一定要定期做故障演练。每个月挑个深夜手动kill掉Active NameNode的进程模拟一次真正的主备切换检查整个流程是否顺畅日志是否有异常。很多问题不演练根本发现不了比如SSH密钥过期、fencing命令路径变更、ZK节点磁盘爆了等等。我经历过的两次大事故都是因为演练做少了。另外如果你是刚接触大数据的新手建议先在一台机器上把HA原理和配置项走通再扩展到三台真实机器。纯看文档和实操完全是两回事尤其是JournalNode和ZK的启动顺序多踩几次坑就记住了。以后遇到集群故障你会感谢今天愿意花时间打基础的自己。这套基于ZooKeeper和JournalNode的高可用方案是Hadoop生态里久经考验的成熟设计。虽然现在Spark、Flink很流行但MapReduce依然是很多企业的数据批处理底座学会把它架构设计成高可用对整个大数据平台的稳定运行很有意义。
返回列表