Hadoop 面试
Hadoop 面试
简介
【简单】简介一下大数据技术生态?⭐
🎯 目标等级:L2 | ⏱ 建议用时:5 min | 🏷 标签:Hadoop / 大数据生态
💎 关键结论
大数据技术生态围绕数据的「采、存、算、查、管」展开:采集用 Flume/Sqoop,存储用 HDFS 与各类 NoSQL,计算用 MapReduce/Spark/Flink,查询分析用 Hive/Spark SQL/Flink SQL,资源管理用 YARN,协调用 ZooKeeper,调度用 Azkaban/Oozie,部署运维用 Ambari/Cloudera Manager。
⚡记忆卡片
- 口诀:采集存储加计算,查询资源协调调,部署运维收尾
- 关键词:Flume/Sqoop / HDFS / MapReduce/Spark/Flink / Hive / YARN / ZooKeeper
- 链路:采集(Flume/Sqoop)→ 存储(HDFS/NoSQL)→ 计算(MR/Spark/Flink)→ 查询(Hive/SQL)→ 资源管理(YARN)
📖 核心知识
大数据技术生态按职能分层:
- 数据采集:Flume、Sqoop、Logstash、Filebeat
- 分布式文件存储:Hadoop HDFS
- NoSql
- 文档数据库:Mongodb
- 列式数据库:HBase
- 搜索引擎:Solr、Elasticsearch
- 分布式计算
- 批处理:Hadoop MapReduce
- 流处理:Storm、Kafka
- 混合处理:Spark、Flink
- 查询分析:Hive、Spark SQL、Flink SQL、Pig、Phoenix
- 集群资源管理:Hadoop YARN
- 分布式协调:Zookeeper
- 任务调度:Azkaban、Oozie
- 集群部署和监控:Ambari、Cloudera Manager

🔀 发散问题
- Q:Hadoop 三大核心组件是什么,如何协作? → HDFS 负责存储、MapReduce 负责计算、YARN 负责资源调度;MapReduce 作业由 YARN 分配 Container 运行,数据存放在 HDFS 上,计算尽量调度到数据所在节点。
- Q:批处理与流处理的区别? → 批处理面向静态数据集一次性计算(MapReduce),流处理面向持续到达的事件流(Storm、Flink);Spark 微批、Flink 原生流处理可同时覆盖两类场景。
- Q:为什么大数据生态组件如此分散? → 采集、存储、计算、查询的负载特征差异很大,单一系统难以同时兼顾,因此按职能解耦、各司其职,再通过统一资源管理(YARN)与协调服务(ZooKeeper)集成。
【简单】什么是 HDFS?⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:5 min | 🏷 标签:Hadoop / HDFS
💎 关键结论
HDFS(Hadoop Distributed File System)是 Hadoop 的分布式文件系统,运行在廉价机器集群上,用于存储具有流数据访问模式的超大文件,对应用程序提供 PB 级存储容量,让用户像使用普通文件系统一样存储大规模文件数据。
⚡记忆卡片
- 口诀:大文件、流式读、廉价机器、PB 级
- 关键词:分布式文件系统 / 超大文件 / 廉价集群 / PB 级容量 / 流式访问
- 链路:大文件 → 切块分布存储 → 集群聚合为单一存储系统 → 应用按普通文件系统使用
📖 核心知识
HDFS 是 Hadoop Distributed File System 的缩写,即 Hadoop 的分布式文件系统。
HDFS 是一种用于存储具有流数据访问模式的超大文件的文件系统,它运行在廉价的机器集群上。
HDFS 的设计目标是管理数以千计的服务器、数以万计的磁盘,将这么大规模的服务器计算资源当作一个单一的存储系统进行管理,对应用程序提供 PB 级的存储容量,让应用程序像使用普通文件系统一样存储大规模的文件数据。
HDFS 是在一个大规模分布式服务器集群上,对数据分片后进行并行读写及冗余存储。因为 HDFS 可以部署在一个比较大的服务器集群上,集群中所有服务器的磁盘都可供 HDFS 使用,所以整个 HDFS 的存储空间可以达到 PB 级容量。
HDFS 的常见使用场景:
- 大数据存储 - HDFS 能够存储 PB 级甚至 EB 级的数据,适合存储日志数据、传感器数据、社交媒体数据等。
- 批处理与分析 - HDFS 是 Hadoop MapReduce 的默认存储系统,MapReduce 作业直接从 HDFS 读取数据并进行分布式计算。
- 数据仓库 - HDFS 可以作为数据仓库的底层存储,支持大规模数据的离线分析。
- 数据冷备 - 由于 HDFS 的高可靠和低成本,适用于存储访问频率较低的冷数据(如历史数据、备份数据)。
- 多媒体数据存储:HDFS 适合存储大规模的多媒体数据(如图像、视频、音频)。
🔀 发散问题
- Q:HDFS 与本地文件系统的主要区别? → HDFS 面向超大文件与高吞吐,数据分块分布存储在集群多节点上并带副本;本地文件系统面向单机小文件与低延迟随机访问。
- Q:HDFS 适合存大量小文件吗? → 不适合,海量小文件会占用 NameNode 大量内存来存元数据,详见本文档『HDFS 大量小文件会带来什么问题?如何解决?』。
- Q:HDFS 支持修改文件内容吗? → 不支持随机修改,一次写入后仅支持追加,详见本文档『HDFS 有什么特性(优缺点)?』。
【简单】HDFS 有什么特性(优缺点)?⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:5 min | 🏷 标签:Hadoop / HDFS
💎 关键结论
HDFS 的优点是高可用、易扩展、适合批处理、低成本;缺点是不适合低延迟访问、不适合大量小文件、不支持并发写入、不支持文件随机修改(仅支持追加)。一句话:为「大文件 + 高吞吐」而生,牺牲了低延迟与小文件友好性。
⚡记忆卡片
- 口诀:优点四高——高可用、易扩展、批处理、低成本;缺点四不——低延迟不行、小文件不行、并发写不行、随机改不行
- 关键词:高可用 / 易扩展 / 流式访问 / 低成本 / 低延迟短板 / 小文件短板 / 单写者 / 仅追加
- 链路:大文件高吞吐设计 → 副本保可用、分块保扩展 → 取舍掉低延迟与小文件场景
📖 核心知识
HDFS 的优点:
- 高可用 - 冗余数据副本,支持自动故障恢复;支持 NameNode HA、安全模式
- 易扩展 - 能够处理 10K 节点的规模;处理数据达到 GB、TB、甚至 PB 级别的数据;能够处理百万规模以上的文件数量,数量相当之大。
- 批处理 - 流式数据访问;数据位置暴露给计算框架
- 低成本 - HDFS 构建在廉价的商用机器上。
HDFS 的缺点:
- 不适合低延迟数据访问 - 适合高吞吐率的场景,就是在某一时间内写入大量的数据。但是它在低延时的情况下是不行的,比如毫秒级以内读取数据,它是很难做到的。
- 不适合大量小文件存储
- 存储大量小文件(这里的小文件是指小于 HDFS 系统的 Block 大小的文件,Block 默认 128MB)的话,它会占用 NameNode 大量的内存来存储文件、目录和块信息。这样是不可取的,因为 NameNode 的内存总是有限的。
- 磁盘寻道时间超过读取时间
- 不支持并发写入 - 一个文件同时只能有一个写入者
- 不支持文件随机修改 - 仅支持追加写入
🔀 发散问题
- Q:HDFS 为什么不支持并发写入? → 单写者模型让一致性控制大大简化,写入以 pipeline 方式一次成型多副本,避免多写者冲突,详见本文档『HDFS 的写数据流程是怎样的?』。
- Q:低延迟场景应该选什么? → HDFS 定位高吞吐,低延迟随机读写应考虑 HBase 等构建在其上的系统,或 KV 类存储。
- Q:HDFS 只支持追加,那如何"修改"数据? → 实践中通过重写文件(新写一份替代旧文件)实现更新,这也是数仓分层覆盖写的由来。
【简单】什么是 YARN?⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:5 min | 🏷 标签:Hadoop / YARN
💎 关键结论
YARN(Yet Another Resource Negotiator)是 Hadoop 的集群资源管理系统,负责统一的资源管理和调度;用户可以将各种服务框架部署在 YARN 上,由 YARN 统一管理和分配资源。它诞生于拆解 MRv1 JobTracker 的职责过载。
⚡记忆卡片
- 口诀:资源归 YARN,计算归框架
- 关键词:集群资源管理 / 资源调度 / JobTracker 拆分 / 框架无关
- 链路:MRv1 JobTracker 职责过重、单点故障 → 资源管理职责剥离 → YARN 统一资源调度 → 多框架共享集群
📖 核心知识
YARN(Yet Another Resource Negotiator,即另一种资源调度器) 是 Hadoop 的集群资源管理系统。YARN 负责资源管理和调度。用户可以将各种服务框架部署在 YARN 上,由 YARN 进行统一地管理和资源分配。
在 Hadoop 1.x 版本,MapReduce 中的 jobTracker 担负了太多的责任,接收任务是它,资源调度是它,监控 TaskTracker 运行情况还是它。这样实现的好处是比较简单,但相对的,就容易出现一些问题,比如常见的单点故障问题。要解决这些问题,只能将 jobTracker 进行拆分,将其中部分功能拆解出来。沿着这个思路,于是有了 YARN。
🔀 发散问题
- Q:YARN 的核心组件有哪些? → ResourceManager、NodeManager、ApplicationMaster、Container,详见本文档『YARN 有哪些核心组件?』。
- Q:YARN 只能跑 MapReduce 吗? → 不是,YARN 是通用资源管理层,MapReduce、Spark、Flink 等框架都可以部署在 YARN 上共享集群资源。
- Q:YARN 的任务提交流程是怎样的? → 详见本文档『YARN 是如何工作的?』。
【简单】什么是 MapReduce?⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:5 min | 🏷 标签:Hadoop / MapReduce
💎 关键结论
MapReduce 是 Hadoop 的分布式计算框架,通过「分而治之、移动计算而非移动数据」的思路,让用户以可靠、容错的方式在大型集群上并行处理海量数据(TB 级);整个计算过程围绕 <key, value> 对进行,只需编写 map 与 reduce 两个函数。
⚡记忆卡片
- 口诀:分而治之,移动计算不移动数据
- 关键词:分布式计算 / 键值对模型 / map-reduce 两阶段 / 移动计算
- 链路:输入拆块 → map 并行处理 → 排序分组 → reduce 聚合 → 结果落盘
📖 核心知识
MapReduce 是 Hadoop 项目中的分布式计算框架。它降低了分布式计算的门槛,可以让用户轻松编写程序,让其以可靠、容错的方式运行在大型集群上并行处理海量数据(TB 级)。
MapReduce 的设计思路是:
- 分而治之,并行计算
- 移动计算,而非移动数据
MapReduce 作业通过将输入的数据集拆分为独立的块,这些块由 map 任务以并行的方式处理。框架对 map 的输出进行排序,然后将其输入到 reduce 任务中。作业的输入和输出都存储在文件系统中。该框架负责调度任务、监控任务并重新执行失败的任务。
通常,计算节点和存储节点是相同的,即 MapReduce 框架和 HDFS 在同一组节点上运行。此配置允许框架在已存在数据的节点上有效地调度任务,从而在整个集群中实现非常高的聚合带宽。
MapReduce 框架由一个主 ResourceManager、每个集群节点一个工作程序 NodeManager 和每个应用程序的 MRAppMaster (YARN 组件) 组成。
MapReduce 框架仅对 <key、value> 对进行作,也就是说,框架将作业的输入视为一组 <key、value> 对,并生成一组 <key、value> 对作为作业的输出,可以想象是不同的类型。键和值类必须可由框架序列化,因此需要实现 Writable 接口。此外,关键类必须实现 WritableComparable 接口,以便于按框架进行排序。
MapReduce 作业的 Input 和 Output 类型:
(input) <k1, v1> -> map -> <k2, v2> -> combine -> <k2, v2> -> reduce -> <k3, v3> (output)🔀 发散问题
- Q:MapReduce 适合哪些场景,不适合哪些场景? → 适合数据统计(如网站 PV/UV)、搜索引擎构建索引、海量数据离线查询;不适合 OLAP 毫秒级响应、流计算(输入动态)与 DAG 多阶段作业(每阶段落盘、IO 开销大)。
- Q:MapReduce 有哪些核心组件? → Job、Mapper、Combiner、Reducer、Partitioner、InputFormat、OutputFormat,详见本文档『MapReduce 有哪些核心组件?』。
- Q:MapReduce 的完整工作流是怎样的? → 详见本文档『MapReduce 是如何工作的?』。
【简单】MapReduce 有什么特性(优缺点)?⭐
🎯 目标等级:L2 | ⏱ 建议用时:5 min | 🏷 标签:Hadoop / MapReduce
💎 关键结论
MapReduce 的核心特性:移动计算而非移动数据、扩展性近似线性、高可用、适合海量数据离线批处理、降低分布式编程门槛;短板是延迟高、不适合实时与迭代计算。
⚡记忆卡片
- 口诀:移动计算、线性扩展、离线批处理、门槛低
- 关键词:移动计算 / 线性扩展 / 高可用 / 离线批处理 / 编程门槛低
- 链路:节点增加 → 计算能力近似线性递增 → 海量数据离线处理
📖 核心知识
MapReduce 有以下特性:
- 移动计算,而非移动数据
- 良好的扩展性:计算能力随着节点数增加,近似线性递增
- 高可用
- 适合海量数据的离线批处理
- 降低了分布式编程的门槛
🔀 发散问题
- Q:MapReduce 的缺点是什么? → 作业启动与调度开销固定存在,不适合秒级响应的 OLAP、动态输入的流计算与多阶段 DAG 作业;详细边界见本文档『什么是 MapReduce?』。
- Q:移动计算为什么优于移动数据? → 大数据量下网络传输是瓶颈,把任务调度到数据所在节点可避免海量数据搬迁,集群内聚合带宽远高于跨节点传输。
- Q:MapReduce 作业慢通常如何调优? → 先看倾斜与小文件,再看 Shuffle 参数,详见本文档『MapReduce Shuffle 阶段的排序与合并细节是怎样的?』。
架构
【中等】HDFS 的架构是怎样设计的?⭐⭐⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:10 min | 🏷 标签:Hadoop / HDFS 架构
💎 关键结论
HDFS 采用主从架构:NameNode 管理命名空间与元数据(管「脑」),DataNode 负责实际数据的存储与读写(管「手」);文件按 Block 切块分布存储,每个 Block 多副本容错,配合层次化命名空间对外呈现为一个统一的文件系统。
⚡记忆卡片
- 口诀:主从分块加副本,命名空间一棵树;NN 管脑 DN 管手
- 关键词:主从架构 / 按块分区 / 数据副本 / 命名空间 / NameNode / DataNode
- 链路:文件 → 切 Block → 多副本分布到 DataNode → NameNode 统一管理元数据与寻址
📖 核心知识
HDFS 架构有以下几个核心要点:
- 主从架构
- 按块分区
- 数据副本
- 命名空间

(1)HDFS 主从架构
HDFS 采用 master/slave 架构。一个 HDFS 集群是由一个 NameNode 和一定数目的 DataNode 组成。NameNode 是一个中心服务器,负责管理文件系统的命名空间 (namespace) 以及客户端对文件的访问。集群中的 DataNode 一般是一个节点一个,负责管理它所在节点上的存储。HDFS 暴露了文件系统的命名空间,用户能够以文件的形式在上面存储数据。从内部看,一个文件其实被分成一个或多个数据块,这些块存储在一组 DataNode 上。NameNode 执行文件系统的命名空间操作,比如打开、关闭、重命名文件或目录。它也负责确定数据块到具体 DataNode 节点的映射。DataNode 负责处理文件系统客户端的读写请求。在 NameNode 的统一调度下进行数据块的创建、删除和复制。
- NameNode - 负责 HDFS 集群的管理、协调。具体来说,主要有以下职责:
- 管理命名空间 - 执行有关命名空间的操作,例如打开,关闭、重命名文件和目录等。
- 管理元数据 - 维护文件的位置、所有者、权限、数据块等。
- 管理 Block 副本策略 - 默认 3 个副本
- 客户端读写请求寻址
- DataNode:负责提供来自文件系统客户端的读写请求,执行块的创建,删除等操作。具体来说,主要有以下职责:
- 执行客户端发送的读写操作
- 存储 Block 和数据校验和
- 定期向 NameNode 发送心跳以续活
- 定期向 NameNode 上报 Block 信息
(2)按块分区
HDFS 将文件数据分割成若干数据块(Block),每个 DataNode 存储一部分数据块,这样文件就分布存储在整个 HDFS 服务器集群中。
将大文件分割成 Block 的主要目的是为了优化网络传输和数据处理的效率。这种分割机制使得文件的不同部分可以并行处理,大大提高了数据处理的速度。
HDFS Block 有以下要点:
- Block 是 HDFS 最小存储单元
- 文件写入 HDFS 会被切分成若干个 Block
- Block 大小固定,默认为 128MB,可通过
dfs.blocksize参数修改 - 若一个 Block 的大小小于设定值,不会占用整个块空间
- 默认情况下每个 Block 有 3 个副本
这实际上是典型的分布式分区思想,使得 HDFS 具备了扩展能力。
(3)数据复制
HDFS 被设计成能够在一个大集群中跨机器可靠地存储超大文件。它将每个文件存储成一系列的数据块,除了最后一个,所有的数据块都是同样大小的。为了容错,文件的所有数据块都会有副本。每个文件的数据块大小和副本系数都是可配置的。应用程序可以指定某个文件的副本数目。副本系数可以在文件创建的时候指定,也可以在之后改变。HDFS 中的文件都是一次性写入的,并且严格要求在任何时候只能有一个写入者。
NameNode 全权管理数据块的复制,它周期性地从集群中的每个 DataNode 接收心跳信号和块状态报告 (Blockreport)。接收到心跳信号意味着该 DataNode 节点工作正常。块状态报告包含了一个该 DataNode 上所有数据块的列表。
(4)命名空间
HDFS 支持传统的层次型文件组织结构。用户或者应用程序可以创建目录,然后将文件保存在这些目录里。文件系统命名空间的层次结构和大多数现有的文件系统类似:用户可以创建、删除、移动或重命名文件。HDFS 不支持用户磁盘配额和访问权限控制,也不支持硬链接和软链接。但是 HDFS 架构并不妨碍实现这些特性。
NameNode 负责维护文件系统的命名空间,任何对文件系统命名空间或属性的修改都将被 NameNode 记录下来。应用程序可以设置 HDFS 保存的文件的副本数目。文件副本的数目称为文件的副本系数,这个信息也是由 NameNode 保存的。
🔬 扩展知识
扩展知识
- 【L3】机架感知与副本放置 - NameNode 通过机架感知(rack awareness)为每个 DataNode 维护机架拓扑信息。默认 3 副本的放置策略为:第一副本放客户端所在节点(远程客户端则随机选择),第二副本放另一机架的节点,第三副本放与第二副本同机架的不同节点。这一布局在可靠性(跨机架容灾)与写入成本(跨机架写只发生一次)之间做了折中;读路径则依据该拓扑优先返回距离客户端最近的副本,以降低延迟、节省带宽。
- 【L3】元数据内存化与架构权衡 - NameNode 将全部元数据(目录树、文件属性、Block 映射)常驻内存,量化估算约为每个文件/目录/Block 对象 150 字节,1 亿个对象约吃掉 15GB 堆内存——这直接决定了单 NameNode 可管理的文件数上限。由此衍生出两条演进路线:
- NameNode HA(主备) - 解决的是可用性(Active 故障时 Standby 秒级接管),不解决元数据规模问题,适合文件数可控(亿级对象以内)的集群。
- HDFS Federation - 多个 NameNode 分管不同命名空间、共享 DataNode,实现命名空间水平扩展,适合超大命名空间或多业务线隔离场景,代价是块池管理与挂载视图(RBF)的运维复杂度上升。
- 【L4】生产踩坑:元数据膨胀引发 GC 风暴 - 曾遇到 NameNode 频繁出现 20~30 秒的 Full GC,期间 DataNode 心跳积压,GC 结束后 NameNode 集中处理积压心跳,误判一批 DataNode 抖动并触发大规模块复制,网络带宽瞬间被打满,形成恶性循环。排查 GC 日志后定位根因:元数据对象数超 12 亿而堆内存仅 64GB。修复手段:扩堆 + 换用 G1 收集器 + 小文件归档治理,GC 停顿降到秒级以内。
- 【L4】实战场景:小文件元数据膨胀应急 - 场景:NameNode 堆内存从 60% 一路涨到 88%,每天增长约 1%,新接入埋点日志业务每 5 分钟落一批小文件、单日新增 3000 万个文件。先算账——按每对象约 150 字节,3000 万个文件加对应 Block 对象约需 9GB 堆内存/天;立即要求上游改按小时攒批写入(单文件从 5MB 提升到 500MB 级别),存量小文件用 HAR 或 SequenceFile 离线合并,必要时临时扩堆争取窗口期。长期方案:接入层强制攒批 + 最小文件大小准入;计算侧用
CombineFileInputFormat合并输入;冷数据定期 HAR 归档;配置目录配额(hdfs dfsadmin -setSpaceQuota)防失控;上线小文件巡检自动合并作业。权衡:攒批牺牲时效性(5 分钟变小时级),若要求实时可见可改落 Kafka 再批量入湖。
🔀 发散问题
- Q:NameNode 为什么把所有元数据放在内存里,而不是落盘存储? → 元数据访问极其频繁——每次 open 文件、定位块都要查询——内存化才能支撑数万级 QPS 的低延迟响应;若落盘,单次寻址就变成一次磁盘 IO,集群整体吞吐会掉几个数量级。代价是文件数受堆内存上限制约,这正是小文件问题和 Federation 出现的根源。
- Q:NameNode 发生一次 30 秒的 Full GC,集群会发生什么? → DataNode 心跳超时判定阈值通常以分钟计(默认 10 分钟),不会立刻被判死;但停顿期间 NameNode 不处理心跳、不做副本决策,恢复后可能集中触发块复制与副本删除,引发网络和磁盘风暴。GC 期间的写请求会超时,客户端需重试或重建 pipeline。
- Q:文件数逼近单机内存上限时,选 Federation 还是垂直扩容 NameNode? → 垂直扩容受单机内存与 GC 效率约束,超大堆的 Full GC 代价很高,属于短期止血;Federation 用多 NameNode 分管命名空间水平扩展,适合长期。生产中通常先做小文件治理压缩对象数,再评估是否上 Federation。
【中等】HDFS 使用 NameNode 的好处?⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:10 min | 🏷 标签:Hadoop / HDFS 架构
💎 关键结论
NameNode 集中管理元数据,让 HDFS 获得一致的文件视图、快速的数据定位、简单的扩展与容错编排;代价是它自身成为中心瓶颈与单点,需要靠主备 HA 来保障可用性。
⚡记忆卡片
- 口诀:元数据集中管,数据分散存
- 关键词:中心化元数据 / 易扩展 / 快速寻址 / 容错调度 / 单点瓶颈
- 链路:客户端问 NameNode 拿位置 → 直连 DataNode 读写 → NameNode 监控状态并调度副本补齐
📖 核心知识
HDFS 使用 NameNode 的好处主要体现在以下几个方面:
- 中心化的元数据管理 - NameNode 在 HDFS 中负责存储整个文件系统的元数据,包括文件和目录的结构、每个文件的数据块信息及其在 DataNode 上的位置等。这种中心化的管理,使得文件系统的组织和管理变得更加简洁高效,并且可以确保整个文件系统的一致性。
- 易扩展 - 由于实际的数据存储在 DataNode 上,而 NameNode 只存储元数据,这样的架构设计使得 HDFS 可以轻松扩展到处理 PB 级别甚至更大规模的数据集。
- 快速的文件访问:用户或应用程序在访问文件时,首先与 NameNode 交互以获得数据块的位置信息,然后直接从 DataNode 读取数据。这种方式可以快速定位数据,提高文件访问的效率。
- 容错和恢复机制:NameNode 可以监控 DataNode 的状态,实现系统的容错。在 DataNode 发生故障时,NameNode 可以指导其它 DataNode 复制丢失的数据块,保证数据的可靠性。
- 简化数据管理:NameNode 的存在简化了数据的管理和维护。例如,在进行数据备份、系统升级或扩展时,管理员只需要关注 NameNode 上的元数据,而不是每个节点上存储的实际数据。
然而,由于 NameNode 是中心节点,它也成为了系统的一个潜在瓶颈和单点故障。因此,HDFS 后来引入了主备 NameNode 机制来保证 NameNode 自身的可用性。
🔬 扩展知识
扩展知识
- 【L3】中心化的代价 - 全量元数据常驻 NameNode 内存,文件总数受限于 NameNode 内存规模;命名空间继续水平扩展需要 Federation,见本文档『什么是 HDFS Federation?』。
- 【L3】与 DataNode 的分工 - NameNode 只管元数据路径,不经过数据流量;读写数据时客户端直连 DataNode,因此 NameNode 不会成为数据传输的带宽瓶颈,只可能成为元数据请求与内存的瓶颈。
- 【L4】单点治理 - 生产集群用 Active/Standby + QJM 解决 NameNode 自身可用性,把单点故障转化为自动切换,见本文档『HDFS 如何实现高可用?』。
🔀 发散问题
- Q:NameNode 的单点问题如何解决? → 通过 Active/Standby 主备架构 + 共享编辑日志(QJM)实现 HA,详见本文档『HDFS 如何实现高可用?』。
- Q:NameNode 与 SecondaryNameNode 是什么关系? → SecondaryNameNode 只是辅助合并元数据、不是热备,详见本文档『NameNode 与 SecondaryNameNode 的区别与联系?』。
- Q:为什么客户端不直接访问 DataNode? → 块的位置与映射关系只有 NameNode 掌握,客户端必须先向 NameNode 寻址,才能找到数据所在的 DataNode。
【中等】HDFS 使用 Block 的好处?⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:10 min | 🏷 标签:Hadoop / HDFS 架构
💎 关键结论
HDFS 把文件切分为固定大小的 Block 分布存储,换来的是容错(块级副本)、并行处理、传输效率、易扩展与负载均衡五个维度的收益,Block 是 HDFS 分布与容错的基本单元。
⚡记忆卡片
- 口诀:切块存、副本存、并行算
- 关键词:块级副本 / 并行处理 / 传输效率 / 易扩展 / 负载均衡
- 链路:大文件 → 切 Block → 块级多副本分布到多 DataNode → 并行读写与故障恢复
📖 核心知识
HDFS 采用文件分块(Block)进行存储管理,主要是基于以下几个原因:
- 提高可靠性和容错性 - 通过将文件分成多个块,并在不同的 DataNode 上存储这些块的副本,HDFS 可以提高数据的可靠性。即使某些 DataNode 出现故障,其他节点上的副本仍然可以用于数据恢复。
- 提高数据处理效率:在处理大规模数据集时,将大文件分割成小块可以提高数据处理的效率。这样,可以并行地在多个节点上处理不同的块,从而加速数据处理和分析。
- 提高网络传输效率:分块存储还有利于网络传输。当处理或传输一个大文件的部分数据时,只需处理或传输相关的几个块,而不是整个文件,这减少了网络传输负担。
- 易于扩展:分块机制使得 HDFS 易于扩展。可以简单地通过增加更多的 DataNode 来扩大存储容量和处理能力,而不需要对现有的数据块进行任何修改。
- 负载均衡:分块存储还有助于在集群中实现负载均衡。不同的数据块可以分布在不同的节点上,从而均衡各个节点的存储和处理负载。
🔬 扩展知识
扩展知识
- 【L3】Block 是副本与恢复的基本单位 - 副本以块为单位管理,故障后的副本补齐、再均衡(balancer)都以块为粒度编排,见本文档『HDFS 的副本机制是怎样的?』。
- 【L3】块抽象屏蔽存储细节 - 分块后文件可以跨磁盘、跨节点存储,单个文件不再受单机磁盘容量限制,用户无需关心数据的物理位置。
- 【L4】小文件问题的根源不在块存储 - 小于 Block 的文件不会占满整块磁盘空间,但每个文件仍要在 NameNode 占用元数据对象,瓶颈在 NameNode 内存而非磁盘,见本文档『HDFS 大量小文件会带来什么问题?如何解决?』。
🔀 发散问题
- Q:Block 为什么默认设计得比较大(128MB)? → 大 Block 减少元数据量、降低寻址开销,且让一次传输的时间远大于寻道时间,充分发挥顺序读写的高吞吐优势。
- Q:Block 大小可以修改吗? → 可以,通过
dfs.blocksize全局修改,也可以在创建文件时指定。 - Q:小文件与 Block 的关系是什么? → 小于 Block 的文件不会占满整块,但每个文件仍要在 NameNode 占用元数据对象,海量小文件仍是问题,详见本文档『HDFS 大量小文件会带来什么问题?如何解决?』。
【中等】NameNode 与 SecondaryNameNode 的区别与联系?⭐⭐⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:10 min | 🏷 标签:Hadoop / HDFS 元数据
💎 关键结论
NameNode 是 HDFS 的主节点,实时管理命名空间与元数据;SecondaryNameNode 只是辅助节点,定期拉取 EditLog 与 FsImage 做 checkpoint 合并、缩短 NameNode 重启时间。它不是 NameNode 的备份/热备,NameNode 故障时无法接管服务。
⚡记忆卡片
- 口诀:NN 管实时,SNN 管合并;SNN 不是备胎
- 关键词:命名空间管理 / 辅助节点 / checkpoint / FsImage 合并 / 非热备
- 链路:NameNode 写 EditLog → SNN 定期拉取 → 与 FsImage 合并 → 新 FsImage 送回 NameNode
📖 核心知识
NameNode 和 SecondaryNameNode 的区别:
- NameNode 是 HDFS 的主要节点,负责管理文件系统的命名空间。它维护着整个文件系统的目录和文件结构,以及所有文件的元数据,包括文件的数据块(block)信息、数据块的位置等。
- SecondaryNameNode 是 NameNode 的辅助节点。
- SecondaryNameNode 不是 NameNode 的备份,不能在 NameNode 故障时接管其功能。
- HDFS 在运行过程中,所有的事务(如文件创建、删除等)都会首先记录在 NameNode 的内存和 EditLog 中。SecondaryNameNode 定期从 NameNode 获取这些日志文件,与文件系统的命名空间镜像(FsImage)合并,然后把新的 FsImage 送回给 NameNode,以帮助减少 NameNode 的内存压力。
NameNode 和 SecondaryNameNode 的联系:
- 共同目标:二者共同目的是维护 HDFS 的稳定和高效运作。NameNode 作为核心,负责实时的元数据管理;而 SecondaryNameNode 辅助 NameNode,通过定期处理 FsImage 和 EditLog,减轻 NameNode 的负担。
- 数据交互:SecondaryNameNode 的工作依赖于与 NameNode 的交互,从 NameNode 获取元数据的状态和编辑日志。
🔬 扩展知识
扩展知识
- 【L3】Checkpoint 触发时机与量化参数 - 由
dfs.namenode.checkpoint.period(默认 3600 秒)与dfs.namenode.checkpoint.txns(默认 100 万条事务)两个条件任一满足即触发。这意味着极端情况下 SecondaryNameNode 上的 FsImage 落后主 NameNode 最近 1 小时的变更。 - 【L3】失效场景(为什么它不是热备)
- SecondaryNameNode 不接收任何写请求、不维护实时元数据状态,NameNode 宕机时它无法接管服务,只能用于辅助恢复元数据,且恢复期间集群不可用。
- 它的 FsImage 仅反映最近一次 checkpoint 的状态,之后的 EditLog 增量仍在 NameNode 本地;若 NameNode 本地磁盘同时损坏,checkpoint 之后的元数据就会丢失。
- 真正的热备是 Standby NameNode + 共享编辑日志(QJM) 的 HA 架构,秒级切换;HDFS 2.x 之后生产环境基本以 HA 取代 SecondaryNameNode。
- 【L3】SecondaryNameNode 的真实价值 - 缩短 NameNode 重启时间。没有 checkpoint,NameNode 重启需回放全量 EditLog,千万级事务的回放可达数十分钟;SecondaryNameNode 提前完成 FsImage 与 EditLog 的合并,NameNode 重启只需加载最近镜像加少量增量。
- 【L4】生产踩坑:checkpoint 监控缺失 - SecondaryNameNode 所在机器故障一周无人察觉(未配 checkpoint 监控),期间 EditLog 持续累积。后来 NameNode 升级重启,需回放 8000 万条事务,启动耗时近 1 小时,远超预期维护窗口。教训:必须为 checkpoint 间隔和 SecondaryNameNode 存活状态配置告警。
- 【L4】实战场景:NameNode 整机报废后的元数据恢复 - 老集群未配 HA,NameNode 所在机器整机报废(磁盘也损坏),只有 SecondaryNameNode。恢复步骤:在备用机器部署新 NameNode,把 SNN 上最近一次 checkpoint 的 FsImage 与 EditLog 拷过来,用
hdfs namenode -importCheckpoint导入,启动进入安全模式,等 DataNode 心跳与块报告重建块位置映射,再用hdfs fsck校验缺块。损失边界:最后一次 checkpoint 之后的元数据变更(最长约 1 小时窗口)无法恢复;数据块还在 DataNode,但失去文件到块的映射。防复发:短期异地定时备份 FsImage/EditLog,中期上 HA(Standby + QJM)才是根本解法。
🔀 发散问题
- Q:checkpoint 为什么能缩短 NameNode 重启时间?没有它会怎样? → NameNode 重启 = 加载 FsImage + 回放 EditLog。没有 checkpoint,EditLog 只增不减,重启回放时间随历史事务数线性增长;checkpoint 定期把 EditLog 合并进 FsImage 并截断日志,把回放范围缩小到一个 checkpoint 周期内。
- Q:SecondaryNameNode 挂掉之后,NameNode 还能正常工作吗? → 短期可以——NameNode 内存中元数据照常维护。风险在于 EditLog 持续累积不被合并,一旦 NameNode 重启就要回放更长的日志,重启时间被大幅拉长;且缺少一份异地镜像,元数据抗灾能力下降。
- Q:既然有 SecondaryNameNode,为什么还要引入 HA? → SecondaryNameNode 只解决“元数据合并”,不解决“可用性”:它无法接管服务,且 checkpoint 有滞后。HA 的 Standby NameNode 实时同步 edits、支持自动故障切换,才满足生产集群分钟级甚至秒级恢复的要求。
【中等】什么是 FsImage 和 EditLog?⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:10 min | 🏷 标签:Hadoop / HDFS 元数据
💎 关键结论
FsImage 是元数据的全量快照(静态),EditLog 是自上次快照以来的增量变更日志(动态实时);NameNode 重启时通过 FsImage + 回放 EditLog 重建最新元数据,checkpoint 定期把两者合并成新的 FsImage。
⚡记忆卡片
- 口诀:镜像存全量,日志记增量,重启两相合
- 关键词:FsImage 全量快照 / EditLog 增量日志 / checkpoint 合并 / 重启重建
- 链路:元数据变更 → 实时写 EditLog → checkpoint 合并进 FsImage → 重启时 FsImage + 回放 EditLog
📖 核心知识
HDFS 中,FsImage和EditLog是两个关键的文件,用于存储和管理文件系统的元数据。它们的主要区别如下:
FsImage(文件系统镜像)
- 内容:
FsImage包含 HDFS 元数据的完整快照,例如文件系统的目录树、文件和目录的属性等。 - 静态性:它是在特定时间点上的静态快照。一旦创建,除非进行新的快照操作,否则内容不会改变。
- 使用场景:在 NameNode 启动时使用,用于加载文件系统的最初状态。此外,在进行系统备份时也会生成新的
FsImage。 - 更新频率:不是实时更新的。通常在系统进行 checkpoint 操作时才会更新。
EditLog(编辑日志)
- 内容:
EditLog记录了自上一个FsImage快照以来所有对文件系统所做的增量更改。这些更改包括文件和目录的创建、删除、重命名等操作。 - 动态性:它是一个动态更新的日志文件。每次对文件系统进行更改时,这个更改就会记录在
EditLog中。 - 使用场景:用于记录所有的文件系统更改操作。在 NameNode 重启时,
FsImage将与EditLog结合使用,以重建文件系统的最新状态。 - 更新频率:实时更新。每次对文件系统的更改都会迅速反映在
EditLog中。
结合使用
在 HDFS 中,FsImage和EditLog一起工作,以确保文件系统的元数据既能够被可靠地存储,又能够反映最新的更改。定期进行 checkpoint 操作(由 Secondary NameNode 或 Standby NameNode 执行)会将EditLog中的更改应用到FsImage中,创建一个新的、更新的快照。这样可以保证在系统重启或恢复时,可以快速加载最新的文件系统状态。
🔬 扩展知识
扩展知识
- 【L3】EditLog 在 HA 中的容错 - HA 环境下 EditLog 通过 QJM 写入 JournalNode 集群(过半写入即成功),避免单点磁盘损坏导致元数据丢失,见本文档『HDFS 如何实现高可用?』。
- 【L3】checkpoint 的执行者与价值 - checkpoint 由 SecondaryNameNode 或 HA 中的 Standby NameNode 执行,把重启回放范围截断到最近一个 checkpoint 周期内,见本文档『NameNode 与 SecondaryNameNode 的区别与联系?』。
- 【L4】元数据备份的本质 - 元数据容灾本质是「备份 FsImage + 保留 edits」,恢复时用备份镜像叠加增量重建命名空间;因此生产中两者都应定期异地备份。
🔀 发散问题
- Q:checkpoint 由谁执行? → 由 SecondaryNameNode 或 HA 架构中的 Standby NameNode 执行,详见本文档『NameNode 与 SecondaryNameNode 的区别与联系?』。
- Q:EditLog 丢了会怎样? → checkpoint 之后的元数据变更无法恢复,因此 HA 架构用 QJM 把 edits 同步到多个 JournalNode 保障元数据可靠,详见本文档『HDFS 如何实现高可用?』。
- Q:为什么不用 EditLog 直接代替 FsImage? → 回放全量日志耗时随历史事务线性增长,必须靠 FsImage 把回放范围截断到最近一个 checkpoint 周期内。
【中等】HDFS 大量小文件会带来什么问题?如何解决?⭐⭐⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:10 min | 🏷 标签:Hadoop / HDFS 小文件
💎 关键结论
小文件问题的本质是 NameNode 内存与计算效率的双重挑战:海量小文件耗尽 NameNode 元数据内存、拖低读写吞吐、产生海量 Map 任务;治理核心就四个字——合并归大(源头攒批、HAR/SequenceFile 合并、CombineFileInputFormat)。
⚡记忆卡片
- 口诀:小文件三宗罪——吃内存、慢读写、炸任务;治理就靠合并归大
- 关键词:NameNode 内存压力 / 寻址开销 / Map 任务膨胀 / HAR / SequenceFile / CombineFileInputFormat
- 链路:源头攒批 → 存量 HAR/SequenceFile 合并 → 计算侧 CombineFileInputFormat → 配额与巡检防复发
📖 核心知识
问题:
- NameNode 内存压力大 - NameNode 将每个文件/目录/Block 的元数据(约 150 字节/对象)存储在内存中,海量小文件会耗尽 NameNode 内存,直接限制集群可管理的文件总数。
- 读写效率低 - 读取小文件时磁盘寻道时间远超数据传输时间,吞吐率急剧下降。
- 计算开销大 - 默认一个 Block 对应一个 Map 任务,海量小文件会产生海量 Map 任务,调度开销远超实际计算。
解决方案:
- 源头治理 - 应用层尽量合并产出,按天/按小时分区聚合,避免碎片化写入。
- HAR 归档(Hadoop Archive) - 将多个小文件打包成一个 HAR 文件,减少 NameNode 元数据数量(客户端通过 har:// 协议访问)。
- 合并存储 - 使用 SequenceFile / MapFile 将小文件合并为大文件的键值对存储,配合压缩效果更佳。
- 输入层合并 - 使用
CombineFileInputFormat将多个小文件组合成一个 InputSplit,减少 Map 任务数。 - 输出层控制 - 合理设置 Reduce 任务数,避免输出大量小文件;必要时追加一轮合并作业。
一句话总结:小文件问题的本质是 NameNode 内存与计算效率的双重挑战,治理核心就四个字——合并归大。
🔬 扩展知识
扩展知识
- 【L3】方案权衡:三种主流治理手段各有适用边界:
| 方案 | 原理 | 适用边界 | 局限 |
|---|---|---|---|
| HAR 归档 | 多文件打包为少量 HAR 文件,元数据从 NameNode 转移到 HAR 内部索引 | 只读冷数据归档 | 不支持追加与修改,归档后访问多一层解析 |
| SequenceFile 合并 | 以文件名为 key、内容为 value 的 KV 大文件,支持 Block 级压缩 | 后续用 MR/Spark 批处理读取 | 不能当普通文件单独访问,需改写下游读取逻辑 |
| HDFS Federation | 多 NameNode 分担元数据,提高文件数上限 | 业务隔离 + 海量文件刚需 | 不解决寻址效率问题,小文件读性能依旧差 |
注意:HAR 和 SequenceFile 是“减少对象数”,Federation 是“分摊对象数”,前者治本,后者扩容。
- 【L3】失效场景 -
CombineFileInputFormat只缓解 Map 任务数膨胀,不缓解 NameNode 内存压力——文件元数据一个不少;对仍在持续写入小文件的业务,一次性合并作业只是快照,很快又会劣化,必须从源头治理。 - 【L4】生产踩坑:小文件拖垮 ETL - 某 ETL 作业从 40 分钟劣化到 5 小时。现象:Map 阶段耗时占比 95%,任务数从 2000 暴涨到 40 万。排查发现上游新业务按分钟级滚动生成日志文件,单文件仅几十 KB。根因:默认一个文件一个 split,调度器启动 40 万任务的开销远超计算本身。修复:接入层改按小时聚合写入 + 历史小文件 SequenceFile 合并 +
CombineFileInputFormat三管齐下,作业回到 45 分钟。 - 【L4】实战场景:小文件应急止血 - 场景:NameNode 内存使用率持续上涨到 90%,新接入日志业务单日 5000 万个平均 200KB 的小文件。应急:先算账——5000 万文件 + 对应 Block 约 1 亿对象 × 150 字节 ≈ 15GB 堆内存/天,协调上游把滚动周期从分钟级改为小时级,存量跑一轮 SequenceFile 合并并删除原文件,临时调大 NameNode 堆内存观察 GC。长期:接入层强制攒批 + 最小文件大小校验;存量目录设
setSpaceQuota与文件数配额;建立小文件巡检与自动合并机制。权衡:合并作业安排低峰期;实时性要求高的链路改走 Kafka。
🔀 发散问题
- Q:为什么“1 亿个小文件”对 NameNode 是致命的,而 1 亿个 Block 的大文件却没问题? → 两者对象数量级看似相同,但大文件的 Block 是顺序写入、元数据规整;小文件场景下每个文件还要额外占一个文件对象 + 目录项,对象数翻倍,且创建/删除操作 QPS 远高于大文件场景,NameNode 内存与 RPC 线程双重承压。
- Q:HAR 归档后,对下游 MapReduce 作业有什么影响? → HAR 对上层透明(
har://协议访问),但归档过程本身是一轮 MR 作业、消耗资源;HAR 不支持追加与修改,只适合冷数据;读取时多一层索引解析,对吞吐影响很小。若下游需要频繁单独访问文件,HAR 就不合适,应选 SequenceFile。 - Q:如果业务就是需要实时写小文件(如分钟级日志滚动),怎么办? → 思路是“写入不合并、读取时合并、定期归档”:写入侧攒批到小时级再滚动;计算侧用
CombineFileInputFormat;读取时效要求高的改落 Kafka/HBase;定期跑归档作业把超过 N 天的小文件合成大文件。源头治理永远优先于事后补救。
【困难】什么是 HDFS Federation?⭐⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:15 min | 🏷 标签:Hadoop / HDFS Federation
💎 关键结论
HDFS Federation 用“多个 NameNode 分管多个命名空间、共享底层 DataNode”的方式,突破单一 NameNode 的内存、性能与隔离性三大瓶颈,实现命名空间的水平扩展;生产中常用 RBF(Router-Based Federation)向客户端提供统一挂载视图。
⚡记忆卡片
- 口诀:多个 NN 分管命名空间,DN 共享不变,块池隔离
- 关键词:命名空间水平扩展 / Namespace Volume / Block Pool / 共享 DataNode / RBF
- 链路:单 NN 内存/吞吐瓶颈 → 多 NN 分管命名空间 → 共享 DataNode 存储层 → 块池隔离 + RBF 统一视图
📖 核心知识
背景:单一 NameNode 架构存在三大瓶颈——内存受限(文件数受限于单节点内存)、性能受限(所有请求打到单个 NameNode)、隔离性差(一个长任务可能影响整个集群)。HDFS Federation 通过引入多个 NameNode 实现命名空间的水平扩展。
核心概念:
- Namespace Volume - 每个 NameNode 管理一个独立的命名空间卷,各卷之间完全隔离。
- Block Pool - 每个 Namespace Volume 对应一个块池,DataNode 需要向所有 NameNode 注册,并存储所有块池的数据块;块 ID 全局唯一,不同块池之间不会冲突。
- DataNode 不变 - DataNode 仍是共享的存储层,向每个 NameNode 发送心跳和块报告。
收益:
- 命名空间可随节点数水平扩展,突破单机内存瓶颈
- 业务之间相互隔离,单个 NameNode 故障只影响其管辖的命名空间
- 聚合吞吐量随 NameNode 数量线性提升
实际生产中常用 RBF(Router-Based Federation),通过 Router 组件向客户端屏蔽多个 NameNode,提供统一挂载视图。
一句话总结:Federation 用“多个 NameNode 分管多个命名空间、共享底层 DataNode”的方式,解开了 HDFS 的单机元数据天花板。
🔬 扩展知识
扩展知识
- 【L3】与 HA 的分工 - HA(主备)解决的是可用性,不解决元数据规模;Federation 解决命名空间扩展与业务隔离,两者可叠加——每个 Namespace Volume 内部仍可配自己的 HA。选型边界:文件数可控的集群优先 HA + 小文件治理;超大命名空间或多业务线隔离才需要 Federation。
- 【L3】代价与运维复杂度 - DataNode 需向所有 NameNode 注册并存储所有块池数据,块池管理与挂载视图(RBF)的运维复杂度上升;集群均衡、配额、权限策略都要按多命名空间重新规划。
- 【L4】与单点治理的关系 - Federation 分摊的是元数据对象数,并不提升单个小文件的读性能;若问题根源是小文件,优先做合并治理而不是盲目上 Federation,见本文档『HDFS 大量小文件会带来什么问题?如何解决?』。
🏭 实战场景
实战场景
- 场景:多业务线共用集群的隔离治理 - 现象:单集群承载多条业务线,某业务批量脚本创建海量文件,NameNode 请求队列积压,拖慢全部业务的元数据操作。处置:按业务线划分 Namespace Volume,故障与配额相互隔离,客户端通过 RBF 获得统一挂载视图。经验:隔离诉求(而非单纯容量)往往是上 Federation 的真实驱动力。
- 场景:NameNode 内存瓶颈的分级应对 - 文件数持续增长、NameNode 内存与 GC 压力上升时,先做小文件治理压缩元数据对象数;仍不解决再规划 Federation;垂直扩容(加大内存)只作短期止血,超大堆的 Full GC 代价会随内存增长而恶化。
⚠️ 常见误区
常见误区
- ❌ "NameNode 内存不够就直接上 Federation" → 应先做小文件治理与垂直扩容;Federation 带来多命名空间的运维复杂度(块池、挂载视图、配额重新规划),文件数真正超出单机承载时才值得上。
- ❌ "Federation 后每个 NameNode 管自己的 DataNode" → DataNode 仍是共享存储层,要向所有 NameNode 注册、发心跳和块报告,并存储所有块池的数据块。
- ❌ "Federation 能顺便解决可用性问题" → Federation 解决的是规模与隔离,可用性要靠每个命名空间内部各自配置 HA,两者叠加使用。
🔀 发散问题
- Q:Federation 后 DataNode 需要改造吗? → 不需要改变存储角色,DataNode 仍是共享存储层,但要向每个 NameNode 注册、发送心跳和块报告,并存储所有块池的数据块。
- Q:不同 NameNode 的块会冲突吗? → 不会,每个 Namespace Volume 对应独立块池,块 ID 全局唯一,不同块池之间不会冲突。
- Q:Federation 与 NameNode 垂直扩容怎么选? → 垂直扩容受单机内存与 GC 效率约束,是短期止血;Federation 水平扩展命名空间,是长期方案,但运维成本更高。
【简单】YARN 有哪些核心组件?⭐⭐⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:5 min | 🏷 标签:Hadoop / YARN
💎 关键结论
YARN 四大核心组件:ResourceManager(中央资源调度,内含 Scheduler 与 ApplicationManager)、NodeManager(单机代理,管 Container)、ApplicationMaster(每应用一个,负责任务拆分与资源申请)、Container(资源抽象单位)。
⚡记忆卡片
- 口诀:RM 管全局,NM 管单机,AM 管应用,Container 装资源
- 关键词:ResourceManager / NodeManager / ApplicationMaster / Container / Scheduler / ApplicationManager
- 链路:Client 提交 → RM 分配 AM Container → AM 申请任务 Container → NM 启动并监控任务
📖 核心知识

YARN 有以下核心组件:
- ResourceManager - ResourceManager 是管理资源和安排在 YARN 上运行的中央调度器。整个系统有且只有一个 ResourceManager,因为号令发布都来自一处,因此不存在调度不一致的情况(很多分布式系统都是通过经典的一主多从模式来解决一致性问题的)。它也包含了两个主要的子组件:
- 定时调度器(Scheduler) - 从本质上来说,定时调度器就是一种策略,或者说一种算法。当 Client 提交一个任务的时候,它会根据所需要的资源以及当前集群的资源状况进行分配。注意,它只负责向应用程序分配资源,并不做监控以及应用程序的状态跟踪。
- 应用管理器(ApplicationManager) - 应用管理器就是负责管理 Client 提交的应用。上面不是说到定时调度器(Scheduler)不对用户提交的程序监控嘛,其实啊,监控应用的工作正是由应用管理器(ApplicationManager)完成的。
- NodeManager - NodeManager 是 ResourceManager 在每台机器的上代理,负责容器的管理,并监控他们的资源使用情况(cpu、内存、磁盘及网络等),以及向 ResourceManager/Scheduler 提供这些资源使用报告。
- ApplicationMaster - 每当 Client 提交一个 Application 时候,就会新建一个 ApplicationMaster 。由这个 ApplicationMaster 去与 ResourceManager 申请容器资源,获得资源后会将要运行的程序发送到容器上启动,然后进行分布式计算。这么设计的原因在于,数据量大的时候,移动数据成本太高,耗时太久,改为移动计算代价较小。
- Container -
Container是 YARN 对资源的抽象,它封装了某个节点上的多维度资源,如内存、CPU、磁盘、网络等。当 AM 向 RM 申请资源时,RM 为 AM 返回的资源是用Container表示的。- YARN 会为每个任务分配一个
Container,该任务只能使用该Container中描述的资源。 ApplicationMaster可在Container内运行任何类型的任务。例如,MapReduce ApplicationMaster请求一个容器来启动 map 或 reduce 任务,而Giraph ApplicationMaster请求一个容器来运行 Giraph 任务。- 容器由 NodeManager 启动和管理,并被它所监控。
- 容器被 ResourceManager 所调度。
- YARN 会为每个任务分配一个
🔬 扩展知识
扩展知识
- 【L3】AM 分布式化(设计权衡) - 把任务拆分、重试、状态监控从中心 RM 下放到每个应用自己的 AM,RM 只做纯资源分配,这是 YARN 相对 MRv1 JobTracker 的核心进步;代价是 AM 成为单个应用的单点,AM 挂掉该应用需整体重来,只能靠
yarn.resourcemanager.am.max-attempts(默认 2)有限重试兜底。 - 【L3】RM 单点(失效场景) - 整个系统有且只有一个 ResourceManager 意味着它是集群级单点,生产环境必须配置 ResourceManager HA(Active/Standby + ZooKeeper 选举),见本文档『YARN 如何实现高可用?』。
- 【L3】Container 资源粒度(设计权衡) - Container 最小分配粒度由
yarn.scheduler.minimum-allocation-mb(默认 1024MB)决定,申请值会向上对齐到该粒度的整数倍;粒度过粗造成内存碎片浪费,过细则调度开销上升。 - 【L4】生产踩坑:Container 频繁被杀 - 作业频繁因 Container “running beyond memory limits” 被杀,根因是资源申请未算上堆外内存与 YARN overhead;教训:资源申请必须与 JVM 参数联动核算。
🔀 发散问题
- Q:ResourceManager 的 Scheduler 为什么不做任务监控,而是交给 AM? → 这是职责分离的设计:Scheduler 保持无状态的纯分配角色,决策路径最短、吞吐最高、易于替换调度策略;任务级监控下放给每个应用的 AM,避免中心 RM 成为监控瓶颈。MRv1 JobTracker 正是因为把这些职责搅在一起而崩溃的。
- Q:Container 和 JVM 进程是什么关系? → Container 是资源抽象单位(内存 + vCore 等),NodeManager 按 Container 规格启动对应进程(map/reduce 任务、AM、Spark Executor 等),并通过内存监控强制限额;一个 Container 通常对应一个 JVM 进程,超额即被 kill。
- Q:NodeManager 挂掉后,其上运行的 Container 会怎样? → RM 检测到心跳超时后将该节点移出可用列表,其上所有 Container 被认定失败,AM 收到通知后在其他节点重新申请 Container 重跑失败任务;NodeManager 恢复时不会自动恢复之前的 Container,需重新申请。
- Q:作业长时间停在 ACCEPTED 状态如何排查? → 先看队列资源水位(大概率被其他大作业占满),再查集群资源碎片化(单节点剩余不满足 Container 申请粒度),最后看 AM Container 是否反复失败耗尽重试;修复靠队列保底容量与弹性借用、调小 Container 规格、必要时抢占。
【简单】MapReduce 有哪些核心组件?⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:5 min | 🏷 标签:Hadoop / MapReduce
💎 关键结论
MapReduce 核心组件:Job(作业配置)、Mapper、Combiner(本地聚合)、Reducer(含 shuffle/sort/reduce 三子阶段)、Partitioner(分区)、InputFormat/OutputFormat(输入输出规范)。
⚡记忆卡片
- 口诀:Job 统领全局,Map 拆、Partition 分、Combine 聚、Reduce 合
- 关键词:Job / Mapper / Combiner / Reducer / Partitioner / InputFormat / OutputFormat
- 链路:InputFormat 切分输入 → Mapper 映射 → Combiner 本地聚合 → Partitioner 分区 → Reducer 归并 → OutputFormat 落盘
📖 核心知识
MapReduce 有以下核心组件:
- Job - Job 表示 MapReduce 作业配置。
Job通常用于指定Mapper、combiner(如果有)、Partitioner、Reducer、InputFormat、OutputFormat实现。 - Mapper - Mapper 负责将输入键值对映射到一组中间键值对。转换的中间记录不需要与输入记录具有相同的类型。一个给定的输入键值对可能映射到零个或多个输出键值对。
- Combiner -
combiner是map运算后的可选操作,它实际上是一个本地化的reduce操作。它执行中间输出的本地聚合,这有助于减少从Mapper传输到Reducer的数据量。 - Reducer - Reducer 将共享一个 key 的一组中间值归并为一个小的数值集。Reducer 有 3 个主要子阶段:shuffle,sort 和 reduce。
- shuffle - Reducer 的输入就是 mapper 的排序输出。在这个阶段,框架通过 HTTP 获取所有 mapper 输出的相关分区。
- sort - 在这个阶段中,框架将按照 key (因为不同 mapper 的输出中可能会有相同的 key) 对 Reducer 的输入进行分组。shuffle 和 sort 两个阶段是同时发生的。
- reduce - 对按键分组的数据进行聚合统计。
- Partitioner - Partitioner 负责控制 map 中间输出结果的键的分区。
- 键(或者键的子集)用于产生分区,通常通过一个散列函数。
- 分区总数与作业的 reduce 任务数是一样的。因此,它控制中间输出结果(也就是这条记录)的键发送给 m 个 reduce 任务中的哪一个来进行 reduce 操作。
- InputFormat - InputFormat 描述 MapReduce 作业的输入规范。MapReduce 框架依赖作业的 InputFormat 来完成以下工作:
- 确认作业的输入规范。
- 把输入文件分割成多个逻辑的 InputSplit 实例,然后将每个实例分配给一个单独的 Mapper。InputSplit 表示要由单个
Mapper处理的数据。 - 提供 RecordReader 的实现。RecordReader 从
InputSplit中读取<key, value>对,并提供给Mapper实现进行处理。
- OutputFormat - OutputFormat 描述 MapReduce 作业的输出规范。MapReduce 框架依赖作业的 OutputFormat 来完成以下工作:
- 确认作业的输出规范,例如检查输出路径是否已经存在。
- 提供 RecordWriter 实现。RecordWriter 将输出
<key, value>对到文件系统。
🔀 发散问题
- Q:Combiner 什么时候不能用? → 聚合不满足结合律时不能用(如求平均值),否则结果错误;可改写为输出(sum, count)再由 Reduce 相除。
- Q:InputSplit 与 Block 是什么关系? → InputSplit 是逻辑切片、Block 是物理存储单元,默认 split 大小与 Block 对齐,但两者并非强绑定。
- Q:Shuffle 阶段的内部细节是怎样的? → 详见本文档『MapReduce Shuffle 阶段的排序与合并细节是怎样的?』。
工作流
【中等】HDFS 的写数据流程是怎样的?⭐⭐⭐⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:10 min | 🏷 标签:Hadoop / HDFS 工作流
💎 关键结论
HDFS 写流程四步:按 Block 大小分割数据 → 通过 NameNode 寻址 DataNode → 按 pipeline 逐跳写数据、ack 反向回传 → 写完后通知 NameNode;核心语义是「全管道 ack 确认、单写者租约、故障重建管道重发」。
⚡记忆卡片
- 口诀:切块寻址管道写,全员 ack 才算数
- 关键词:Block 分割 / NameNode 寻址 / pipeline / ack 确认队列 / 写租约
- 链路:create 建文件 → NameNode 分配 DataNode 列表 → DataStreamer 管道写 packet → 全 ack 后出队 → close 通知 NameNode
📖 核心知识
HDFS 写数据流程大致为:
- 按 Block 大小分割数据
- 通过 NameNode 寻址 DataNode
- 向 DataNode 写数据
- 完成后通知 NameNode
扩展:下面的漫画生动的展示了 HDFS 的写入流程,图片引用自博客:翻译经典 HDFS 原理讲解漫画
HDFS 写数据的源码流程:

- 客户端通过对
DistributedFileSystem对象调用create()函数来新建文件。 - 分布式文件系统对 NameNode 创建一个 RPC 调用,在文件系统的命名空间中新建一个文件。
- NameNode 对新建文件进行检查无误后,分布式文件系统返回给客户端一个
FSDataOutputStream对象,FSDataOutputStream对象封装一个DFSoutPutstream对象,负责处理 NameNode 和 DataNode 之间的通信,客户端开始写入数据。 FSDataOutputStream将数据分成一个一个的数据包,写入内部数据队列,DataStreamer 负责将数据包依次流式传输到由一组 DataNode 构成的管道中。DFSOutputStream维护着确认队列来等待 DataNode 收到确认回执,收到管道中所有 DataNode 确认后,数据包从确认队列删除。- 客户端完成数据的写入,调用
close()方法关闭传输通道。 - NameNode 确认完成。
🔬 扩展知识
扩展知识
- 【L3】packet 与 chunk - 数据以 packet(默认 64KB)为单位进入 pipeline,每个 packet 由若干 chunk(512 字节数据 + 4 字节校验和)组成。
- 【L3】pipeline 构建 - DataStreamer 先向 NameNode 申请 DataNode 列表(按机架感知策略),构成 pipeline,数据在管道中逐跳转发,ack 沿管道反向回传。
- 【L3】写失败恢复语义 - pipeline 中某个 DataNode 故障时:客户端把它移出 pipeline,未确认的数据包重新入队;客户端向 NameNode 申请新节点补入 pipeline,满足最小副本数(
dfs.namenode.replication.min,默认 1)后继续写入,文件关闭时再补齐全部副本。故障节点上未确认的块会被 NameNode 删除,不会出现部分可见的脏块。 - 【L3】lease(租约) - 文件打开时客户端持有写租约,HDFS 保证同一文件只有一个写入者;客户端宕机后租约超时(默认 1 小时),NameNode 会自动关闭文件并将已确认的块纳入命名空间,避免悬挂的写会话。
- 【L3】失效场景 - 若 pipeline 中靠后的节点故障,靠前节点已写入的数据以已收到的 ack 为准;若所有 DataNode 都写失败或 NameNode 拒绝分配新节点,写入直接抛异常,客户端必须自己重试,HDFS 不保证写操作自动幂等重放。
- 【L4】生产踩坑:慢盘拖垮 pipeline - 某导入作业频繁报写超时,重试后文件块数对不上、下游校验失败。排查发现某台 DataNode 数据盘老化,写延迟从毫秒级劣化到秒级,拖慢整条 pipeline 的 ack 链路,客户端超时后不断重建管道。修复:配置慢盘检测(disk balancer / 磁盘健康检查)并替换故障节点后恢复。教训:pipeline 写入的性能取决于最慢的那个节点。
- 【L4】实战场景:DataNode 掉线时如何保证不丢数据 - 场景:Flume 实时写 HDFS,凌晨一台 DataNode 掉线、作业报写超时。不丢的语义边界:未收齐 ack 的 packet 会重发到新 pipeline,已确认的 packet 至少有一份完整副本;掉线节点上未完成的块由 NameNode 判定缺失并触发副本补齐;风险在 Flume 若按时间滚动且未正确
hflush,缓冲区数据可能随进程异常丢失。排查:看客户端日志 pipeline 重建记录、确认 NameNode 及时把该 DataNode 标记为 dead(心跳超时默认 10 分钟,可用dfs.namenode.heartbeat.recheck-interval调短)、检查欠复制告警。优化:合理 roll 策略与hflush周期、关键数据提高副本数、开启磁盘健康检查。权衡:hflush频率越高可见性越好但吞吐下降、小文件风险上升;副本数提高以存储成本换安全,核心链路 3 副本起步。
🔀 发散问题
- Q:为什么 HDFS 写采用 pipeline 逐跳转发,而不是客户端并行写三副本? → pipeline 让客户端出口带宽只需一份,副本复制由 DataNode 之间级联完成,充分利用集群内网带宽,且 ack 链天然有序;并行写三副本会让客户端带宽放大 3 倍,跨机房场景成本极高,还要自行处理三个写流的一致性。
- Q:ack 确认队列(ack queue)在故障恢复中起什么作用? → 已发送但未确认的 packet 保存在确认队列中,pipeline 上所有 DataNode ack 后才移除。任何节点 ack 超时,客户端重建 pipeline,未确认的 packet 重新入队重发,这正是写入不丢数据的机制基础。
- Q:如果 DataNode 收到数据但宕机了、没来得及发 ack,会不会出现重复数据? → 客户端会重发该 packet,目标节点重启后块内容由最终确认的完整块为准,未确认的半成品块在租约超时后被 NameNode 删除;校验和机制保证块内容完整,不会出现重复或半块数据。
【中等】HDFS 的读数据流程是怎样的?⭐⭐⭐⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:10 min | 🏷 标签:Hadoop / HDFS 工作流
💎 关键结论
HDFS 读流程三步:客户端向 NameNode 查文件块位置 → NameNode 返回各块副本的 DataNode 列表(按距离排序)→ 客户端就近直连 DataNode 逐块读取;读路径不需要 pipeline,是点对点的单副本拉取。
⚡记忆卡片
- 口诀:问 NN 拿位置,就近读、逐块连
- 关键词:NameNode 寻址 / 副本位置列表 / 就近读 / 校验和 / 短路读
- 链路:open 打开文件 → NameNode 返回块位置 → 连最近 DataNode 读 → 块读完换下一块最优节点 → close
📖 核心知识
HDFS 读数据流程大致为:
- 客户端向 NameNode 查询文件信息
- NameNode 返回相关信息
- 该文件的所有数据块
- 每个数据块对应的 DataNode(按距离客户端的远近排序)
- 客户端向 DataNode 读数据

HDFS 读数据的源码流程:

- 客户端调用
FileSystem对象的open()方法在 HDFS 中打开要读取的文件。 - HDFS 通过使用 RPC(远程过程调用)来调用 NameNode,确定文件起始块(Block)的位置。
DistributedFileSystem类返回一个支持文件定位的输入流FSDataInputStream对象,FSDataInputStream对象接着封装DFSInputStream对象(存储着文件起始几个块的 DataNode 地址),客户端对这个输入流调用read()方法。DFSInputStream连接距离最近的 DataNode,通过反复调用read方法,将数据从 DataNode 传输到客户端。- 到达块的末端时,
DFSInputStream关闭与该 DataNode 的连接,寻找下一个块的最佳 DataNode。 - 客户端完成读取,对
FSDataInputStream调用close()方法关闭连接。
🔬 扩展知识
扩展知识
- 【L3】副本定位 - NameNode 返回每个块的全部副本位置,并按网络拓扑距离排序;客户端优先读本节点 > 同机架 > 跨机架的副本。
- 【L3】校验和失败处理 - 读取时每 512 字节验证一次校验和,不匹配则抛 ChecksumException,客户端向 NameNode 报告该副本损坏,并自动切换其他副本重读;NameNode 随后标记 corrupt 副本并触发重新复制替换。
- 【L3】短路读(short-circuit read) - 客户端与 DataNode 同机时,可通过
dfs.client.read.shortcircuit直接读本地块文件,绕过 DataNode 网络栈,显著降低本地读延迟,HBase 等低延迟场景依赖此特性。 - 【L3】hedged read - 通过
dfs.client.hedged.read.threshold.millis开启:第一个副本读超过阈值未返回时,并发读另一副本,谁先返回用谁,专治个别慢节点拖长尾。 - 【L3】失效场景 - 若某个块的所有副本校验和都不匹配,说明数据真损坏,只能依赖外部备份恢复;若集群拓扑脚本配置错误,“最近副本”可能实际在远端机房,读流量大量跨机房,带宽被打满而排查方向往往被误导。
- 【L4】生产踩坑:机架感知失效导致跨机架读 - 机房搬迁后集群读带宽突然打满、作业普遍变慢,而写入正常。排查网络流量发现几乎所有读都跨机架进行。根因:机架感知拓扑脚本未随新机房更新,NameNode 对所有节点的距离判定失真,“就近读”全部退化。修复:更新拓扑脚本并验证后恢复。教训:集群拓扑变更后必须验证机架感知配置。
- 【L4】实战场景:P99 延迟飙升定位 - 场景:HBase 读 HDFS 的在线查询系统 P99 从 30ms 涨到 2 秒,磁盘/CPU 指标正常。均值正常 + P99 飙高是典型长尾特征:先看读请求是否集中在少数 DataNode,再查慢盘/GC/网络拥塞,验证机架感知。应急:开 hedged read 对冲长尾、疑似慢节点 decommission、
hdfs balancer均衡热点块;长期建立磁盘巡检、读延迟分位数告警与定期均衡。权衡:hedged read 用额外带宽换尾延迟,吞吐型作业更应优先做节点治理与数据均衡。
🔀 发散问题
- Q:读路径为什么不需要 pipeline?与写路径的本质差异是什么? → 读是点对点的单副本拉取,不存在多节点级联与确认需求;客户端依据 NameNode 返回的副本位置列表直连最近的 DataNode,读完一块断开再连下一块的最优节点。pipeline 是写路径为了一次带宽、多处复制而设计的。
- Q:客户端如何缓存块位置?缓存会不会导致读到过期位置? →
DFSInputStream会缓存已获取的块位置列表以减少 NameNode RPC;若按缓存位置连接失败或读到损坏副本,会重新向 NameNode 拉取最新位置再重试,过期缓存由失败驱动刷新,不会造成数据错误。 - Q:hedged read 适合什么场景?代价是什么? → 适合对 P99 延迟敏感的在线读取(如 HBase 读 HFile);代价是读流量成倍增加,吞吐型批处理作业通常不开启,否则集群带宽会被对冲读白白消耗。
【中等】MapReduce 是如何工作的?⭐⭐⭐⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:10 min | 🏷 标签:Hadoop / MapReduce 工作流
💎 关键结论
MapReduce 作业分 map 与 reduce 两阶段,每阶段都以键值对为输入输出;框架负责 splitting 与 shuffling,开发者只写 mapping 与 reducing 两个函数;作业由 YARN 分配 Container 启动 MRAppMaster 统一调度,失败任务自动重试。
⚡记忆卡片
- 口诀:拆行 map 发单词,shuffle 同 key 聚一筐,reduce 求和出结果
- 关键词:input / splitting / mapping / shuffling / reducing / MRAppMaster
- 链路:读取输入 → 按行拆分 K1/V1 → map 并行产出 K2/V2 → shuffle 同 key 汇聚 → reduce 归约输出
📖 核心知识
MapReduce 任务过程分为两个处理阶段:map 阶段和 reduce 阶段。每阶段都以键值对作为输入和输出,其类型由程序员来选择。程序员还需要写两个函数:map 函数和 reduce 函数。
以词频统计为例,其工作流再细分一下,可以划分为以下阶段:
- input - 读取文本文件;
- splitting - 将文件按照行进行拆分,此时得到的
K1行数,V1表示对应行的文本内容; - mapping - 并行将每一行按照空格进行拆分,拆分得到的
List(K2,V2),其中K2代表每一个单词,由于是做词频统计,所以V2的值为 1,代表出现 1 次; - shuffling - 由于
Mapping操作可能是在不同的机器上并行处理的,所以需要通过shuffling将相同key值的数据分发到同一个节点上去合并,这样才能统计出最终的结果,此时得到K2为每一个单词,List(V2)为可迭代集合,V2就是 Mapping 中的 V2; - reducing - 这里的案例是统计单词出现的总次数,所以
Reducing对List(V2)进行归约求和操作,最终输出。
MapReduce 编程模型中 splitting 和 shuffing 操作都是由框架实现的,需要我们自己编程实现的只有 mapping 和 reducing,这也就是 MapReduce 这个称呼的来源。

全流程骨架补充:作业提交后由 YARN 分配 Container 启动 MRAppMaster,由它负责切分 InputSplit(默认 128MB,与 HDFS Block 对齐)、调度 map/reduce 任务、跟踪进度;任务失败自动重试(默认 4 次),重试仍慢则交给推测执行兜底。Shuffle 阶段的环形缓冲、溢写、归并等微观机制见本文档『MapReduce Shuffle 阶段的排序与合并细节是怎样的?』。
🔬 扩展知识
扩展知识
- 【L3】失效场景(MR 的适用边界)
- DAG/多阶段作业 - 阶段间数据必须落盘 HDFS,多阶段作业的中间结果 IO 与副本开销成倍放大,性能远不如 Spark 等内存框架。
- 迭代计算 - 每轮迭代都要完整经历落盘、再读取,机器学习训练类负载完全不适合。
- 低延迟查询 - 作业启动、调度、Shuffle 建连开销固定存在,秒级响应的场景应交给 MPP 或 OLAP 引擎。
- 【L4】生产踩坑:小文件引发任务数爆炸 - 某统计作业平时 40 分钟,某天突然跑了 9 小时。排查发现实际计算只需 20 分钟,其余时间全耗在任务调度上——上游产出 40 万个几十 KB 的小文件,导致 40 万个 Map 任务,每个任务的启动与 Shuffle 建连都是固定开销。用
CombineFileInputFormat合并输入、推动上游合并产出后,作业回到 25 分钟。教训:任务数不是越多越好,调度开销本身有成本。 - 【L4】实战场景:2TB 作业从 8 小时优化到 1 小时 - 输入侧:检查小文件导致的 Map 任务数爆炸(CombineFileInputFormat 合并)、split 与 Block 对齐保数据本地性、输入压缩(如 Snappy);计算侧:启用 Combiner 本地聚合、用 Counter 查 Reduce 输入分布治倾斜、Shuffle 调优;输出侧:控制 Reduce 数避免海量小文件、输出压缩;资源侧:调大任务内存减少溢写、确认队列资源充足。权衡:Combiner 仅适用可结合聚合;加盐打散多一轮作业,轻度倾斜不值得;优先做收益最大的环节——通常倾斜治理与 Combiner 的收益远超参数微调。
🔀 发散问题
- Q:为什么 MapReduce 不适合迭代式计算(如机器学习训练)? → 每轮迭代的输出是下一轮的输入,MR 每轮间必须落盘 HDFS(序列化、三副本写、再读取),IO 开销巨大;Spark 用内存 RDD 把中间结果留在内存,同类负载性能可高出数个数量级。
- Q:Map 输出为什么不直接写 HDFS? → Map 输出是中间结果,需按分区重组供 Reduce 拉取,且生命周期很短;写 HDFS 要付出三副本复制与 NameNode 元数据开销,得不偿失。中间结果写任务节点本地磁盘、由 Reduce 通过 HTTP 拉取,是吞吐与开销的折中。
- Q:Reduce 任务数如何决定?太多或太少分别有什么问题? → 经验目标是让每个 Reduce 处理约 1~2 个 HDFS 块大小(128~256MB)的数据;太少则单任务数据量大、易内存溢出且并行度不足,太多则输出大量小文件、Shuffle 建连开销上升,还会拖垮 NameNode。
【中等】MapReduce Shuffle 阶段的排序与合并细节是怎样的?⭐⭐⭐⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:10 min | 🏷 标签:Hadoop / MapReduce Shuffle
💎 关键结论
Shuffle 是 Map 与 Reduce 之间的桥梁,排序和合并贯穿 Map 端与 Reduce 端两侧:一句话总结——Shuffle 的本质是“排序定序、合并降量”,它是 MapReduce 性能的胜负手。
⚡记忆卡片
- 口诀:Map 端环形缓冲溢写排序,Reduce 端拉取归并分组
- 关键词:环形缓冲区 / 分区快排 / Combiner / 归并排序 / HTTP 拉取 / GroupingComparator
- 链路:Map 输出进环形缓冲 → 分区内快排 + Combiner → 溢写归并 → Reduce HTTP 拉取 → 归并分组 → reduce 函数
📖 核心知识
Map 端(写侧):
- 环形缓冲区 - Map 输出先写入内存环形缓冲区(默认 100MB),达到阈值(默认 80%)触发溢写。
- 分区与排序 - 溢写前先按 Partitioner 分区,每个分区内对 key 做快速排序。
- Combiner - 若设置了 Combiner,溢写前做本地聚合,减少写盘数据量。
- 归并 - 多次溢写产生多个有序小文件,最终通过归并排序合并为一个分区有序的磁盘文件(溢写文件过多时先做中间归并)。
Reduce 端(读侧):
- 复制(Copy) - 通过 HTTP 并行拉取各 Map 输出中属于本分区的数据;数据量小放内存,大则落盘。
- 归并(Merge) - 边拉取边做归并排序,保持 key 有序;落盘文件过多时进行多轮归并。
- 分组(Sort/Group) - 最终归并时通过 GroupingComparator 将相同 key(或等价 key)的数据聚成一组,交给 reduce 函数处理。
整个过程中,“排序”用归并贯穿始终,“合并”用 Combiner 和归并双管齐下,目的都是压缩跨网络传输的数据量。
一句话总结:Shuffle 的本质是“排序定序、合并降量”,它是 MapReduce 性能的胜负手。
🔬 扩展知识
扩展知识
- 【L3】为什么是环形缓冲区 - 记录数据从一端写入、索引从另一端写入,两者相向增长、相遇即触发溢写,避免普通缓冲区溢写时的整体拷贝与搬移;溢写时只对索引排序,因为排序的是元数据(key + 指针)而非记录本身,CPU 与内存开销大幅下降。
- 【L3】关键参数 -
mapreduce.task.io.sort.mb(默认 100MB)控制缓冲区大小;mapreduce.map.sort.spill.percent(默认 0.8)控制溢写阈值;io.sort.factor(默认 10)控制归并路数;Reduce 端拉取并行度由mapreduce.reduce.shuffle.parallelcopies(默认 5)控制。 - 【L3】中间归并 - 溢写文件数超过阈值时会分轮归并,避免最终归并路数过大导致磁盘随机 IO 爆炸。
- 【L3】失效场景
- 单条记录超过缓冲区剩余空间时,该记录直接落盘,绕过缓冲与 Combiner,性能急剧下降。
- Combiner 对不可结合的聚合(如平均值)会产生错误结果,此时必须禁用或改写为可结合形式。
- Reduce 端拉取阶段若内存不足,会出现 "Too many fetch failures",任务反复重试甚至作业失败。
- 【L4】生产踩坑:Reduce 端拉取雪崩 - 某作业 Reduce 阶段大量报 "Too many fetch failures",任务重试 4 次后作业失败。现象:Map 数高达 5 万个,每个 Reduce 需从 5 万个 Map 拉取数据。排查发现 Reduce 堆内存仅 1GB,拉取缓冲频繁落盘、归并队列积压导致 HTTP 拉取超时。修复:调大 Reduce 堆内存与
mapreduce.reduce.shuffle.input.buffer.percent、控制 Map 端输出压缩,故障消失。教训:Reduce 端内存要在拉取缓冲与归并缓冲之间合理分配,Map 数过多时要优先压源头。 - 【L4】实战场景:500GB Shuffle 量调优 - 场景:作业 Shuffle 数据量 500GB,Reduce 端频繁报 merge 失败、磁盘 IO 打满。诊断:先看 Counter 确认 Reduce 输入分布排除倾斜,再看 Map 溢写次数与 Reduce 落盘文件数定位瓶颈侧。Map 端:调大
mapreduce.task.io.sort.mb减少溢写轮次、优化 Combiner、Map 输出压缩(Snappy);Reduce 端:调大堆内存与拉取缓冲占比提高内存归并、控制io.sort.factor归并路数。根源治理:若是倾斜导致单分区过大,先打散再谈调参。权衡:内存给 Shuffle 多了留给业务逻辑的就少,Shuffle 优化的本质是在内存、磁盘、网络三者间做预算分配,必须按作业数据量实测。
🔀 发散问题
- Q:环形缓冲区相比普通缓冲区,到底省在哪里? → 省在拷贝与排序开销:数据与索引相向写入,溢写只需对索引做快排后顺序落盘,记录本身一次写入、零搬移;普通缓冲区溢写前要先整体排序再搬移写入,对高频触发的 Shuffle 来说差距显著。
- Q:什么场景下不能用 Combiner?给一个例子。 → 聚合运算不满足结合律时不能用:求平均值就是典型——局部平均后再求平均在样本量不等时结果错误;正确做法是 Combiner 输出(sum, count)两个值,在 Reduce 端相除。不能改写就必须放弃 Combiner。
- Q:Reduce 端为什么要边拉取边归并,而不是等所有 Map 输出到齐再统一归并? → 受内存与磁盘约束:等到齐再归并,磁盘峰值占用等于本分区全部 Shuffle 数据,且归并路数受
io.sort.factor限制,文件数过多要多轮归并。边拉边归并尽早压缩文件数,把磁盘与内存峰值控制在可接受范围。
【中等】什么是 MapReduce 推测执行?⭐⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:10 min | 🏷 标签:Hadoop / MapReduce 容错
💎 关键结论
推测执行是“备份赛跑、择优录取”的兜底机制:框架发现某任务明显落后于同类任务时,在另一个健康节点启动备份任务,谁先完成采用谁的结果;它能治慢节点,治不了数据倾斜。
⚡记忆卡片
- 口诀:发现落后开备份,谁快用谁
- 关键词:落后任务(straggler)/ 备份任务 / 赛跑择优 / 慢节点兜底
- 链路:统计任务进度得分 → 识别落后任务 → 健康节点启动备份 → 先完成者生效,另一个被杀
📖 核心知识
背景:集群由廉价机器组成,硬件性能不均衡,个别节点可能因负载高、磁盘老化等原因成为“慢节点”,拖慢整个作业(作业的完成时间取决于最慢的任务)。
机制:
- 框架周期性统计各任务的进度得分(progress score),发现某任务明显落后于同类任务的平均水平时,判定其为“落后任务(straggler)”。
- 在另一个健康节点上启动一个备份任务,处理相同的输入数据,与原任务赛跑。
- 谁先完成采用谁的结果,另一个任务被杀掉。
配置:通过 mapreduce.map.speculative 和 mapreduce.reduce.speculative 开关控制(默认开启)。
注意事项:
- 推测执行会消耗额外集群资源,资源紧张时可关闭。
- 如果任务慢的根因是数据倾斜而非节点问题,推测执行只会浪费资源,应从数据层面解决。
一句话总结:推测执行是“备份赛跑、择优录取”的兜底机制,能治慢节点,治不了数据倾斜。
🔬 扩展知识
扩展知识
- 【L3】生效边界:慢节点 vs 数据倾斜 - 推测执行的备份任务跑在健康节点上,解决的是“节点慢”;而倾斜任务的慢来自数据量本身,备份任务要处理同样多的数据、同样慢,只会浪费一倍资源。判断技巧:若绝大多数任务很快、只剩个别长尾,先看是分散在不同节点(慢节点)还是输入数据量异常大(倾斜)。
- 【L3】资源代价与开关策略 - 备份任务占用额外 Container 与集群资源,资源紧张的集群可关闭;map 与 reduce 可分别通过
mapreduce.map.speculative、mapreduce.reduce.speculative控制。
🔀 发散问题
- Q:为什么推测执行治不了数据倾斜? → 倾斜任务的慢来自数据量本身,备份任务处理同样多的数据依然慢;倾斜要用过滤、打散、广播等手段,详见本文档『MapReduce 中如何定位和解决数据倾斜?』。
- Q:备份任务和原任务会同时写输出吗,会冲突吗? → 两者处理相同输入、各自独立计算,框架只采纳先完成者的结果并杀掉另一个,不会产生两份生效的输出。
- Q:任务失败和任务慢是一回事吗? → 不是:失败任务由框架直接重试(失败重跑),推测执行针对的是仍在运行但明显落后的任务(备份赛跑),两者是不同的容错机制。
【困难】MapReduce 中如何定位和解决数据倾斜?⭐⭐⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:15 min | 🏷 标签:Hadoop / MapReduce 调优
💎 关键结论
数据倾斜的本质是“个别 key 背了全组的活”:大部分 Reduce 很快完成、个别卡在 99%;解法三板斧——过滤脏数据、加盐两阶段聚合打散、小表广播 MapJoin。
⚡记忆卡片
- 口诀:先定位再治理,过滤、打散、广播三板斧
- 关键词:长尾 99% / 热点 key / 脏数据 / 两阶段聚合 / MapJoin / Counter
- 链路:Counter/Web UI 定位热点 key → 过滤脏数据 → 加盐打散两阶段聚合 → 验证各 Reduce 输入分布
📖 核心知识
定位:作业大部分 Reduce 很快完成,个别 Reduce 长时间卡在 99%;或日志显示某 Reduce 处理的数据量远超其他 Reduce。
常见原因:
- 某些 key 的记录数远超其他 key(如热门商品、头部用户)
- 大量空值/无效值(null、空串)未过滤,全部落入同一 key
- 分区/分组键选取不当,导致数据分布不均
解决方案:
- 过滤脏数据 - 对 null、空串、异常 key 单独处理或提前过滤,避免其混入聚合。
- Combiner 本地聚合 - 在 Map 端先聚合,大幅减少传入 Reduce 的数据量,对求和/计数类场景效果显著。
- 两阶段聚合(打散 + 汇总) - 第一阶段给倾斜 key 加随机前缀,打散到多个 Reduce 做局部聚合;第二阶段去掉前缀做全局汇总。
- Map 端 Join 替代 Reduce 端 Join - 大表 Join 小表导致倾斜时,将小表广播到内存做 MapJoin,绕开 Shuffle。
- 调整并行度 - 增加 Reduce 数量或自定义 Partitioner,让数据分布更均匀。
一句话总结:数据倾斜的本质是“个别 key 背了全组的活”,解法就是过滤、打散、广播三板斧。
🔬 扩展知识
扩展知识
- 【L3】定位手段(量化)
- 作业 Web UI 看各 Reduce 进度与输入字节数,倾斜 Reduce 的输入通常是其他的数十倍。
- 用自定义 Counter 输出 map 输出 key 分布的 Top N,直接定位热点 key。
- 经验判断:若 99% 任务完成后最后一个跑数小时,基本可断定倾斜而非慢节点(慢节点用推测执行可缓解,倾斜不能)。
- 【L3】方案权衡:
| 手段 | 适用边界 | 代价 |
|---|---|---|
| Combiner | 可结合聚合(sum/count/max) | 不可结合聚合(平均值)会算错 |
| 两阶段聚合(加盐) | 单 key 倾斜、聚合类作业 | 多一轮作业 + 一次中间落盘 |
| 自定义 Partitioner | 分区不均但 key 分布可控 | 需维护映射规则,业务变更要同步改 |
| Map 端 Join | 大表 Join 小表(小表能装进内存) | 小表过大时内存装不下即失效 |
- 【L3】失效场景 - 若倾斜来自业务语义本身(如全局汇总只有一个 key),加盐、调 Partitioner 都无效,必须改写计算逻辑(如分组合并);Hive 层的
mapjoin、skew join参数适用于 Hive 作业,属 SQL 层手段,归 Hive 专题讨论。
:::
🏭 实战场景
实战场景
- 案例:脏数据引发倾斜 - 一次订单统计作业,99% 的 Reduce 十分钟完成,最后一个跑了 6 小时。通过 Counter 定位到 null 和空串 user_id 占全部记录的 90%,全部落入同一分区。修复:空值单独过滤走独立统计逻辑 + 对剩余数据加随机前缀打散,作业从 6 小时降到 12 分钟。教训:脏数据是倾斜的第一大来源,先过滤再谈打散。
- 场景:头部店铺热点治理 - 大促日终统计各店铺成交额,头部 1% 店铺贡献 60% 交易量,作业长尾 5 小时。先用 Counter 输出店铺维度记录数 Top N,确认是业务热点还是脏数据(空店铺 ID);脏数据走过滤/单独统计,真实热点用两阶段聚合(店铺 ID 加随机前缀拆到多个 Reduce 局部求和,再去前缀全局汇总)。验证:对比各 Reduce 输入记录数分位数,P99 与中位数差距应缩到 3 倍以内。权衡:加盐多一轮作业与中间结果存储成本,轻度倾斜用过滤 + 调并行度即可,重度倾斜才值得上两阶段聚合。
⚠️ 常见误区
常见误区
- ❌ "开推测执行就能缓解数据倾斜" → 推测执行解决的是「节点慢」;倾斜任务的慢来自数据量本身,备份任务处理同样多的数据依然慢,只会浪费一倍资源。
- ❌ "加盐两阶段聚合是倾斜的万能药" → 它只对单 key 倾斜的聚合类作业有效;全局汇总只有一个 key 时必须改写计算逻辑,不可结合的聚合(如平均值)也不能直接加盐求和。
- ❌ "增加 Reduce 数量就能解决倾斜" → 若倾斜来自热点 key,同一 key 的记录无论多少分区都会进同一个 Reduce;调并行度只能改善分区级不均,治不了 key 级倾斜。
🔀 发散问题
- Q:两阶段聚合为什么通常要两个 MR 作业,而不是一个作业内完成? → 第一阶段加盐后同一原始 key 被拆到多个 Reduce,输出的是“盐 + key”的局部结果;第二阶段必须去盐后再做全局汇总,两次分组键不同,单次 Shuffle 无法同时满足两次分组,所以需要两轮作业串联。
- Q:为什么推测执行治不了数据倾斜? → 推测执行解决的是“节点慢”,备份任务跑在健康节点上;而倾斜任务的慢来自数据量本身,备份任务要处理同样多的数据,同样慢,只会浪费一倍资源。
- Q:空值 key 除了过滤,还有什么处理方式? → 若空值也需要参与统计,可把空值替换成随机值打散到各分区,汇总时再把随机前缀归并回空值;或者把空值拆到独立 Reduce 单独计算后合并结果,避免其淹没正常分区。
【中等】YARN 是如何工作的?⭐⭐⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:10 min | 🏷 标签:Hadoop / YARN 工作流
💎 关键结论
YARN 任务提交五步:Client 向 RM 申请运行 AM → RM 通过 NM 分配第一个 Container 启动 AM → AM 拆分任务并向 RM 申请任务 Container → AM 与 NM 通信把任务分发到 Container 运行 → 任务向 AM 汇报心跳,完成后 AM 向 RM 注销释放资源;核心设计是「RM 只管资源、AM 管应用」。
⚡记忆卡片
- 口诀:申请起 AM,AM 要容器,任务进容器,完事注销
- 关键词:Client 提交 / ApplicationMaster / Container 申请 / NodeManager 启动 / 心跳汇报
- 链路:Client → RM 分配 AM Container → AM 拆分任务申请 Container → NM 启动任务 → 心跳汇报 → AM 注销释放资源
📖 核心知识

这张图简单地标明了提交一个程序所经历的流程,接下来我们来具体说说每一步的过程。
- Client 向 ResourceManager 申请运行一个 Application 进程,这里我们假设是一个 MapReduce 作业。
- ResourceManager 向 NodeManager 通信,为该 Application 进程分配第一个容器。并在这个容器中运行这个应用程序对应的 ApplicationMaster。
- ApplicationMaster 启动以后,对作业(也就是 Application) 进行拆分,拆分 task 出来,这些 task 可以运行在一个或多个容器中。然后向 ResourceManager 申请要运行程序的容器,并定时向 ResourceManager 发送心跳。
- 申请到容器后,ApplicationMaster 会去和容器对应的 NodeManager 通信,而后将作业分发到对应的 NodeManager 中的容器去运行,这里会将拆分后的 MapReduce 进行分发,对应容器中运行的可能是 Map 任务,也可能是 Reduce 任务。
- 容器中运行的任务会向 ApplicationMaster 发送心跳,汇报自身情况。当程序运行完成后, ApplicationMaster 再向 ResourceManager 注销并释放容器资源。
🔬 扩展知识
扩展知识
- 【L3】资源请求协议 - AM 通过心跳携带资源请求(ResourceRequest:节点 + 资源规格 + 数量),RM 返回分配结果,属于异步协商模式;因此资源分配有秒级延迟,并非即时响应。
- 【L3】AM 失败处理 - AM 心跳超时后 RM 判定其失败,整个 Application 在剩余重试次数内重新调度(
yarn.resourcemanager.am.max-attempts默认 2),已运行的任务随 Container 回收,任务级进度不保留。 - 【L3】RM 故障 - 未配置 RM HA 时,RM 宕机导致所有应用无法提交、AM 无法续心跳,集群级停摆。
- 【L3】Container 超限 - 任务实际内存超过 Container 限额,NodeManager 会直接将其杀死,AM 收到失败通知后决定是否重试。
- 【L3】量化感受 - AM 向 RM 的心跳默认约 1 秒一次;AM Container 启动通常在数秒级;从提交到首个 map 任务跑起来的端到端耗时正常在 10~30 秒内,超过则需排查队列资源或 AM 启动失败。
- 【L4】生产踩坑:AM 内存不足引发整体重跑 - 某夜间批处理作业偶发整体失败重跑,损失数小时。现象:AM Container 反复被 NodeManager 以内存超限杀掉,耗尽 2 次重试机会。排查发现 AM 申请的 1.5GB 内存不足以覆盖堆 + 堆外开销,且发生在作业任务数最多、AM 内存峰值的时刻。修复:AM 内存调至 3GB 并统一“申请资源 > JVM 堆 + overhead”的配置规范后稳定。
- 【L4】实战场景:关键作业停在 ACCEPTED - 场景:凌晨批处理高峰,关键报表作业提交后半小时没拿到 AM Container。排查:
yarn application -status确认状态,RM Web UI 看队列使用率与待分配,大概率队列被低优先级大作业占满或资源碎片化。应急:杀掉/降级占资源的低优先级作业,或临时移到空闲队列。长期:按业务线划分队列配最小保障容量、关键作业配优先级与可抢占属性、建立高峰容量规划与错峰调度。权衡:抢占保护高优先级但被抢方任务白跑,需配套准入控制与容量规划。
🔀 发散问题
- Q:为什么 AM 要通过心跳“顺便”请求资源,而不是单独发一个资源申请 RPC? → 心跳捎带请求减少 RPC 次数,也让 RM 能以 AM 存活状态为前提分配资源——AM 死了就不必再分配;代价是分配是异步的,有秒级延迟。这是吞吐与实时性的折中。
- Q:AM 挂掉后,已经跑完的 map 任务进度还在吗? → 不在。AM 承载应用的任务状态与调度信息,AM 失败后整个 Application 按重试策略从头再来,已完成的 map 任务也要重跑,这正是 AM 重试次数默认只给 2 次的原因——重跑代价很高。
- Q:RM 和 AM 都挂了,哪个对集群影响更大? → RM 影响全局:所有应用无法提交、无法续心跳,属于集群级故障,必须靠 RM HA 保障;AM 只影响单个应用,影响面是作业级的。两者的容灾设计也因此不同:RM 靠 ZooKeeper 选举,AM 靠有限重试。
【中等】YARN 有哪些资源调度器?它们有什么区别?⭐⭐⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:10 min | 🏷 标签:Hadoop / YARN 调度
💎 关键结论
YARN 三种调度器:FIFO(单队列先来先服务,队头阻塞)、Capacity(多队列保底容量、弹性借用,Hadoop 默认)、Fair(多队列公平均分、支持抢占,CDH 默认);一句话:FIFO 适合独占,Capacity 保底优先,Fair 均分为上——多租户生产环境基本二选一。
⚡记忆卡片
- 口诀:FIFO 排队、Capacity 保底、Fair 均分
- 关键词:FIFO Scheduler / Capacity Scheduler / Fair Scheduler / 队列保底 / 弹性借用 / 抢占 / DRF
- 链路:多租户需求 → 多队列划分 → 保底容量 + 弹性借用 → 抢占纠正失衡
📖 核心知识
YARN 提供三种资源调度器,通过 yarn.resourcemanager.scheduler.class 配置:
| 维度 | FIFO Scheduler | Capacity Scheduler | Fair Scheduler |
|---|---|---|---|
| 队列模型 | 单队列,先来先服务 | 多队列,按容量划分,队列内 FIFO | 多队列,资源共享 |
| 分配策略 | 按提交顺序独占资源 | 各队列保底容量,空闲资源可弹性借用 | 按“公平”原则均分资源(DRF 主导资源公平) |
| 抢占 | 不支持 | 不支持 | 支持(可抢占超额资源) |
| 适用场景 | 小型专用集群 | 多租户生产集群(Hadoop 默认) | 多租户共享集群(CDH 默认) |
Capacity Scheduler(容量调度器):管理员预先为每个队列分配固定容量百分比(如 A 队列 40%、B 队列 60%),队列内部按 FIFO 排序;某队列空闲时资源可被其他队列临时借用,但不会饿死保底容量。优点是容量可控、隔离性好。
Fair Scheduler(公平调度器):目标是让所有应用随时间推移获得均等的资源份额。新任务提交后很快获得资源,随着任务增多资源被均摊;支持基于 DRF(主导资源公平)的多维资源公平,并支持抢占机制纠正资源失衡。
一句话总结:FIFO 适合独占,Capacity 保底优先,Fair 均分为上——多租户生产环境基本二选一。
🔬 扩展知识
扩展知识
- 【L3】失效场景与选型边界
- FIFO - 一个长时间的大作业就能让后面所有小作业无限期排队(队头阻塞),多租户环境基本不可用。
- Capacity - 借用资源的归还不是即时的:需等借用方任务自然结束才能收回,紧急时只能依赖抢占;队列划分过细则资源碎片化,过粗则隔离失效。
- Fair - 抢占会直接杀掉被抢方的 Container,长任务被抢后白跑,频繁抢占会造成集群吞吐抖动;抢占需配置延迟与最小保障,防止弱势队列被反复抢成饿死。
- 【L3】量化示例 - 典型生产配置——root 下划分 etl(50%)、adhoc(20%)、realtime(30%)三个队列,队列内
maximum-capacity限制弹性上限;Fair Scheduler 用yarn.scheduler.fair.preemption.cluster-utilization-threshold(默认 0.8)控制抢占触发水位。 - 【L4】生产踩坑:队列规划缺失 - 集群把所有作业都扔在默认的 root.default 队列,某晚批处理大作业把队列打满,次日凌晨实时报表链路排队超 1 小时。现象:实时任务 ACCEPTED 堆积。排查发现队列规划缺失 + FIFO 语义下大作业阻塞小作业。修复:拆分 etl/adhoc/realtime 队列、配容量保底与弹性借用、实时队列设高优先级 + 可抢占,延迟恢复到分钟级。教训:队列规划是 YARN 治理的第一课,比任何参数调优都重要。
- 【L4】实战场景:白天查询 + 夜间 ETL 调度设计 - 场景:白天交互式查询要求分钟级响应,夜间 ETL 要求最大吞吐。架构:Capacity/Fair 多队列 + 抢占(如 interactive 30%、etl 50%、adhoc 20%,均允许弹性借用);白天 interactive 高优先级、可抢占 etl 超额 Container;夜间 etl 借用 interactive 空闲资源拉满吞吐。护栏:为抢占配延迟与队列最小保障防饿死,对单作业设资源上限防独占。权衡:保障越强夜间利用率越低,需基于监控数据持续调整队列比例。
🔀 发散问题
- Q:为什么 Capacity Scheduler 是 Hadoop 官方默认,而 Fair 是 CDH 默认? → Capacity 以容量百分比做硬隔离,保底资源可审计、可承诺,适合按部门核算的多租户场景;Fair 以均分与抢占见长,交互式查询体验好,契合 Cloudera 面向查询负载的产品定位。选型看组织是按预算划分资源(Capacity)还是追求整体公平(Fair)。
- Q:Capacity 调度器中,借出去的资源如何归还? → 借用是软约束:不会主动杀掉借用方任务来即时归还,而是在借用方 Container 结束后,后续分配优先满足原队列的保底需求;若急需归还,需依赖抢占机制(新版本 Capacity Scheduler 也支持抢占)。
- Q:为什么 FIFO 在生产集群几乎绝迹? → 单队列按提交顺序执行,没有优先级与保障概念,一个长作业阻塞全部后续作业,且无抢占纠正手段;只适合单一用途的专用集群,多租户场景下 FIFO 是事故之源。
复制
复制主要指通过网络在多台机器上保存相同数据的副本。
复制数据,可能出于各种各样的原因:
- 提高可用性 - 当部分组件出现位障,系统依然可以继续工作,系统依然可以继续工作。
- 降低访问延迟 - 使数据在地理位置上更接近用户。
- 提高读吞吐量 - 扩展至多台机器以同时提供数据访问服务。
所有分布式系统都需要支持复制。
【中等】HDFS 的副本机制是怎样的?⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:10 min | 🏷 标签:Hadoop / HDFS 复制
💎 关键结论
HDFS 以 Block 为单位存副本,默认 3 副本:第一副本放客户端所在节点、第二副本放另一机架、第三副本放与第二副本同机架的不同节点;读时优先就近副本,NameNode 通过心跳与块报告全权管理复制。
⚡记忆卡片
- 口诀:块级三副本,跨机架保平安,读时就近选
- 关键词:块级副本 / dfs.replication / 机架感知 / 副本放置策略 / 就近读
- 链路:文件切 Block → 每 Block 按策略放 3 副本 → NameNode 心跳/块报告监控 → 缺副本时调度补齐
📖 核心知识
基于块的副本
由于 Hadoop 被设计运行在廉价的机器上,这意味着硬件是不可靠的,为了保证容错性,HDFS 提供了副本机制。HDFS 将文件分解为若干 Block,Block 是 HDFS 最小存储单元,每个 Block 有多个副本。
HDFS 的默认副本数为 3,更多的副本意味着更高的数据安全性,但同时也会带来更高的额外开销(存储成本和带宽成本)。3 个副本是在保障数据可靠性和系统成本之间的一个较好的平衡点。
副本数可以通过以下方式修改:
- 在 HDFS 的配置文件 hdfs-site.xml 中,有一个名为
dfs.replication的属性,可以设置全局的默认副本数。修改这个值后,需要重启 HDFS 使配置生效。 - 针对单个文件或目录修改副本数:如果只想改变某个特定文件或目录的副本数,而不影响整个系统的默认设置,可以使用 HDFS 的命令行工具。例如,使用命令
hdfs dfs -setrep -w <副本数> <文件/目录路径>来修改特定文件或目录的副本数。

NameNode 全权管理数据块的复制,它周期性地从集群中的每个 DataNode 接收心跳信号和块状态报告 (BlockReport)。接收到心跳信号意味着该 DataNode 节点工作正常。块状态报告包含了一个该 DataNode 上所有数据块的列表。

副本分布策略
副本分布策略是 HDFS 可靠性和性能的关键。优化的副本存放策略是 HDFS 区分于其他大部分分布式文件系统的重要特性。HDFS 采用一种称为机架感知 (rack-aware) 的策略来改进数据的可靠性、可用性和网络带宽的利用率。大型 HDFS 实例一般运行在跨越多个机架的计算机组成的集群上,不同机架上的两台机器之间的通信需要经过交换机。在大多数情况下,同一个机架内的两台机器间的带宽会比不同机架的两台机器间的带宽大。
通过一个机架感知的过程,NameNode 可以确定每个 DataNode 所属的机架 id。一个简单但没有优化的策略就是将副本存放在不同的机架上。这样可以有效防止当整个机架失效时数据的丢失,并且允许读数据的时候充分利用多个机架的带宽。这种策略设置可以将副本均匀分布在集群中,有利于当组件失效情况下的负载均衡。但是,因为这种策略的一个写操作需要传输数据块到多个机架,这增加了写的代价。
HDFS 默认的副本数为 3,此时 HDFS 的副本分布策略是:
- 副本 1 - 放在 Client 所在节点;对于远程 Client,系统会随机选择节点
- 副本 2 - 放在不同机架的节点上
- 副本 3 - 放在与第二个副本同一机架的不同节点上
- 副本 N - 在满足以下条件的节点中随机选择
- 每个节点只存储一份副本
- 每个机架最多存储两份副本
- 优选 - 同等条件下优先选择空闲节点。
- 如果某个 DataNode 节点上的空闲空间低于特定的临界点,按照均衡策略系统就会自动地将数据从这个 DataNode 移动到其他空闲的 DataNode。
副本选择
为了降低整体的带宽消耗和读取延时,HDFS 会尽量让客户端程序读取离它最近的副本。如果在客户端程序的同一个机架上有一个副本,那么就读取该副本。如果一个 HDFS 集群跨越多个数据中心,那么客户端也将首先读本地数据中心的副本。
为了最大限度地减少带宽消耗和读取延迟,HDFS 在执行读取请求时,优先读取距离读取器最近的副本。如果在与读取器节点相同的机架上存在副本,则优先选择该副本。如果 HDFS 群集跨越多个数据中心,则优先选择本地数据中心上的副本。

🔬 扩展知识
扩展知识
- 【L3】机架感知的取舍 - 副本跨机架存放能防止整机架失效、读时可利用多机架带宽,但写需要跨机架传输、代价上升;默认 3 副本「本机架两份 + 跨机架一份」正是可靠性与写成本的平衡点。
- 【L3】副本存放约束 - 每节点最多一份、每机架最多两份、同等条件优先选空闲节点,这些约束共同避免副本集中在同一故障域,并在存储水位失衡时自动迁移数据。
- 【L4】大规模欠复制的恢复代价 - 多节点同时故障后副本补齐集中爆发,会占用大量网络带宽;生产中可对复制带宽限流,避免冲击在线读写业务。
🔀 发散问题
- Q:为什么默认是 3 副本而不是更多? → 3 副本在可靠性与成本间取得平衡:本机架两副本 + 跨机架一副本,既能容忍单机架整体故障,又不会让存储与写带宽成本无限放大;关键数据可按需调高。
- Q:副本不足时 HDFS 怎么处理? → NameNode 通过心跳与块报告发现欠复制块后,调度其他 DataNode 复制补齐,详见本文档『DataNode 故障如何处理?』。
- Q:读时如何选择副本? → 按网络拓扑距离优先读本节点 > 同机架 > 跨机架的副本,详见本文档『HDFS 的读数据流程是怎样的?』。
【中等】HDFS 如何保证数据一致性?⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:10 min | 🏷 标签:Hadoop / HDFS 一致性
💎 关键结论
HDFS 不提供数据库式的强一致性,而是靠五套机制保障大规模场景下的有效性与健壮性:NameNode 中心化元数据、块复制、写操作与复制的原子性、客户端一致性协议、心跳/校验和的定期检查与错误恢复。
⚡记忆卡片
- 口诀:中心管元数据,副本保数据,写一次不修改,心跳加校验
- 关键词:中心化元数据 / 块复制 / 原子写入 / 客户端协议 / 心跳与校验和
- 链路:写入完成 → 副本分散存储 → 心跳监控 + 校验和验证 → 异常时副本替换/重复制
📖 核心知识
HDFS 的数据一致性主要依赖以下机制来保证:
- NameNode 的中心化管理 - NameNode 在 HDFS 中负责存储整个文件系统的元数据,包括文件和目录的结构、每个文件的数据块信息及其在 DataNode 上的位置等。这种中心化的管理,使得文件系统的组织和管理变得更加简洁高效,并且可以确保整个文件系统的一致性。
- 数据块的复制(Replication) - HDFS 采用副本来保证数据的可靠性。一旦数据写入完成,副本就会分散存储在不同的 DataNodes 上。尽管这种方法不是强一致性模型,但通过足够数量的副本和及时的副本替换策略,HDFS 能够提供较高水平的数据一致性和可靠性。
- 写入和复制的原子性保证 - 在 HDFS 中,文件一旦创建,其内容就不能被更新,只能被追加或重写。这种方式简化了并发控制,因为写操作在文件级别上是原子的。在复制数据块时,HDFS 保证原子性复制,即一个数据块的所有副本在任何时间点上都是相同的。如果复制过程中出现错误,那么不完整的副本会被删除,系统会重新尝试复制直到成功。
- 客户端的一致性协议 - 客户端在与 HDFS 交互时,遵循特定的协议。例如,客户端在完成文件写入之后,需要向 NameNode 通知,以确保 NameNode 更新文件的元数据。这样可以保证 NameNode 的元数据与实际存储的数据保持一致。
- 定期检查和错误恢复
- 心跳和健康检查 - DataNodes 定期向 NameNode 发送心跳和 Block 健康状况报告。NameNode 利用这些信息来检查和维护系统的整体一致性。例如,如果某个 DataNode 失败,NameNode 会重新组织数据块的副本。
- 校验 - HDFS 在存储和传输数据时,会计算数据的校验和。在读取数据时,会验证这些校验和,确保数据的完整性。
通过这些机制,HDFS 确保了系统中的数据在正常操作和故障情况下的一致性和可靠性。虽然 HDFS 不提供像传统数据库那样的强一致性保证,但它的设计和实现确保了在大规模数据处理场景中的有效性和健壮性。
🔬 扩展知识
扩展知识
- 【L3】写入可见性边界 - 数据可见性以已确认的写入为准:pipeline 中所有 DataNode ack 后数据包才算成功;未完成 ack 的数据不对外可见且会重发,完整语义见本文档『HDFS 的写数据流程是怎样的?』。
- 【L3】校验和失败的修复闭环 - 读时校验失败 → 客户端报告 corrupt 并切换其他副本重读 → NameNode 标记损坏副本并触发重新复制替换,形成「检测—隔离—替换」闭环。
- 【L4】设计哲学:用不可换一致 - 一次写入、仅追加、不支持随机修改的模型让并发控制极大简化,这是 HDFS 能在大规模场景兼顾高吞吐与一致性的前提。
🔀 发散问题
- Q:HDFS 是强一致性系统吗? → 不是传统意义上的强一致性:它靠副本、原子写与校验和提供高水平可靠性,但读到的可见性以已确认写入为准,与数据库 ACID 语义不同。
- Q:写租约在保证一致性中起什么作用? → 租约保证同一文件同一时刻只有一个写入者,避免多写者冲突,详见本文档『HDFS 的写数据流程是怎样的?』。
- Q:校验和验证失败会怎样? → 客户端报告损坏副本并自动切换其他副本重读,NameNode 标记 corrupt 副本并触发重新复制替换。
容错
【中等】HDFS 有哪些故障类型?如何检测故障?⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:10 min | 🏷 标签:Hadoop / HDFS 容错
💎 关键结论
HDFS 三类常见故障:节点故障(心跳超时检测)、通信故障(ACK 机制检测)、数据损坏(校验和检测);检测手段分别是心跳、ACK 回执与 CheckSum 校验。
⚡记忆卡片
- 口诀:节点看心跳,通信看 ACK,数据看校验和
- 关键词:节点故障 / 通信故障 / 数据损坏 / 心跳超时 / CheckSum
- 链路:心跳/ACK/校验和检测 → 定位故障类型 → 切副本或重复制恢复
📖 核心知识
HDFS 常见故障及检测方法:
- 节点故障
- DataNode 每 3 秒向 NameNode 发送心跳
- 超时未收到心跳,NameNode 判定 DataNode 宕机
- 通信故障
- 客户端请求 DataNode 会收到 ACK
- 数据损坏
- 磁盘介质在存储过程中受环境或者老化影响,其存储的数据可能会出现错乱。HDFS 的应对措施是,对于存储在 DataNode 上的数据块,计算并存储校验和(CheckSum)。在读取数据的时候,重新计算读取出来的数据的校验和,如果校验不正确就抛出异常,应用程序捕获异常后就到其他 DataNode 上读取备份数据。
- 如果 DataNode 监测到本机的某块磁盘损坏,就将该块磁盘上存储的所有 BlockID 报告给 NameNode,NameNode 检查这些数据块还在哪些 DataNode 上有备份,通知相应的 DataNode 服务器将对应的数据块复制到其他服务器上,以保证数据块的备份数满足要求。


🔬 扩展知识
扩展知识
- 【L3】三类检测手段的分工 - 心跳解决「节点是否存活」、ACK 解决「链路是否可用」、校验和解决「数据是否正确」,三者分别覆盖节点、通信、数据三层故障,互为补充。
- 【L3】坏盘主动上报的价值 - DataNode 检测到本机磁盘损坏后主动上报其上所有 BlockID,不必等读取时才暴露问题,NameNode 可立即调度副本补齐,缩短数据欠复制窗口。
- 【L4】监控建议 - 生产中应对 DataNode 心跳丢失数、欠复制块数、损坏块数设置告警,在故障演变成业务可见问题之前介入。
🔀 发散问题
- Q:心跳超时后 NameNode 会做什么? → 判定 DataNode 宕机,查找其上数据块的其他副本并调度补齐,详见本文档『DataNode 故障如何处理?』。
- Q:读写过程中遇到故障怎么处理? → 写跳过故障节点、读切换其他副本,详见本文档『HDFS 读写故障如何处理?』。
- Q:校验和是在什么时候计算的? → 存储与传输时计算并保存校验和,读取时重新计算比对,不一致则抛异常并改读其他副本。
【中等】HDFS 读写故障如何处理?⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:10 min | 🏷 标签:Hadoop / HDFS 容错
💎 关键结论
写故障靠 ACK 检测:客户端收不到 ACK 就判定节点宕机、跳过该节点,副本不足的块信息通知 NameNode 补齐;读故障靠 NameNode 寻址全部副本,某节点宕机就改读其他节点。
⚡记忆卡片
- 口诀:写看 ACK 跳节点,读靠寻址换节点
- 关键词:ACK 检测 / 跳过故障节点 / 欠复制上报 / 多副本寻址 / 切换读
- 链路:写:数据包 → 无 ACK → 跳过节点 → 通知 NameNode 补副本;读:寻址全部副本 → 故障节点不可用 → 读其他节点
📖 核心知识
写入故障处理
- 写入数据通过数据包传输
- DataNode 接收数据后,返回 ACK
- 如果客户端没有收到 ACK,就判定 DataNode 宕机,跳过节点
- 没有充分备份的数据块信息通知到 NameNode
读取故障处理
- 读数据先要通过 NameNode 寻址该数据块的所有 DataNode
- 如果某 DataNode 宕机,则读取其他节点

🔬 扩展知识
扩展知识
- 【L3】跳过写与最小副本数的关系 - 写时跳过故障节点后,只要满足最小副本数即可继续写入,剩余副本在文件关闭时补齐;完整恢复语义见本文档『HDFS 的写数据流程是怎样的?』。
- 【L3】读的降级与失效边界 - 读故障切换的前提是至少还有一个副本可用;全部副本失效时块不可读,只能靠副本策略与机架感知提前预防,或依赖外部备份恢复。
- 【L4】客户端的重试责任 - HDFS 不保证写操作自动幂等重放,写异常场景下客户端需自行重试并校验结果,避免重复写入或数据缺失。
🔀 发散问题
- Q:写时跳过故障节点后,副本数不足怎么办? → 客户端把欠复制信息通知 NameNode,由 NameNode 调度其他 DataNode 复制补齐,满足最小副本数后继续写入、关闭时补齐全部副本。
- Q:读时所有副本都不可用怎么办? → 说明数据可能真正丢失,只能依赖外部备份恢复;这也是多副本与机架感知策略要预防的场景。
- Q:写故障的详细恢复语义是怎样的? → 详见本文档『HDFS 的写数据流程是怎样的?』中的 pipeline 故障恢复部分。
【中等】DataNode 故障如何处理?⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:10 min | 🏷 标签:Hadoop / HDFS 容错
💎 关键结论
DataNode 每 3 秒向 NameNode 发心跳,超时未发则被判宕机;NameNode 随即查找其上数据块及剩余副本,通知其他 DataNode 复制补齐,保证副本数符合配置,数据不丢。
⚡记忆卡片
- 口诀:心跳断、判宕机,找副本、补一份
- 关键词:心跳超时 / 宕机判定 / 块定位 / 副本补齐
- 链路:心跳超时 → NameNode 判定宕机 → 查找块及副本分布 → 通知其他 DataNode 复制补齐
📖 核心知识
DataNode 每 3 秒会向 NameNode 发送心跳消息,以证明自身正常工作。如果 DataNode 超时未发送心跳,NameNode 就会认为该 DataNode 已经宕机。
NameNode 会立即查找该 DataNode 上存储的数据块有哪些,以及这些数据块还存储在哪些其他 DataNode 上。
随后,NameNode 通知这些 DataNode 再复制一份数据块到其他 DataNode 上,保证 HDFS 存储的数据块副本数符合配置数。即使再出现服务器宕机,也不会丢失数据。

🔬 扩展知识
扩展知识
- 【L3】宕机判定后的恢复链路 - 判定宕机 → 定位其上欠复制块 → 调度健康 DataNode 复制补齐;补齐期间数据仍可读(依赖剩余副本),新写入不再路由到该节点。
- 【L3】多余副本的清理 - DataNode 携旧副本重新上线后,NameNode 对比块报告,保留有效副本、调度清理多余副本,保证副本数不超配置、不多占存储。
- 【L4】大规模故障的处置顺序 - 多节点同时故障时欠复制块暴增,应优先补齐低于最小副本数的块(不可再丢的高风险块),其余按带宽窗口渐进补齐。
🔀 发散问题
- Q:DataNode 恢复上线后会发生什么? → 重新向 NameNode 注册、上报块报告;其上原有副本仍有效,多余的副本会由 NameNode 调度清理,不会重复占用配额。
- Q:副本补齐会影响集群性能吗? → 大规模块复制会占用网络带宽,生产中可通过限流控制复制带宽,避免冲击在线业务。
- Q:如何提前发现慢节点而不是等宕机? → 结合磁盘健康检查与读写延迟监控提前摘除慢盘,避免其拖慢 pipeline,见本文档『HDFS 的写数据流程是怎样的?』。
【中等】NameNode 故障如何处理?⭐⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:10 min | 🏷 标签:Hadoop / HDFS 容错
💎 关键结论
NameNode 是整个 HDFS 的核心,所有的文件路径和数据块存储信息都保存在 NameNode,它一旦故障整个集群都无法使用;因此通过 Active/Standby 主备架构实现故障转移:Active 宕机后 Standby 快速升级为新的 Active,元数据靠 edits 编辑日志持续同步,并由 QJM 共享存储保障可靠。
⚡记忆卡片
- 口诀:主备双 NN,edits 同步,QJM 过半写,ZK 选主
- 关键词:Active/Standby / edits / FsImage / QJM / JournalNode / ZooKeeper 选举
- 链路:Active 故障 → Standby 确认元数据同步完成 → 升级为 Active → 对外服务
📖 核心知识
NameNode 是整个 HDFS 的核心,记录着 HDFS 文件分配表信息,所有的文件路径和数据块存储信息都保存在 NameNode,如果 NameNode 故障,整个 HDFS 系统集群都无法使用。如果 NameNode 上记录的数据丢失,整个集群所有 DataNode 存储的数据也就没用了。
NameNode 通过主备架构实现故障转移:
- Active NameNode - 是正在工作的 NameNode;
- Standby NameNode - 是备份的 NameNode。
Active NameNode 宕机后,Standby NameNode 快速升级为新的 Active NameNode。Standby NameNode 周期性同步 edits 编辑日志,定期合并 FsImage 与 edits 到本地磁盘。
注:Hadoop 3.0 允许配置多个 Standby NameNode。
元数据文件:
- edits(编辑日志文件) - 保存了自最新检查点(Checkpoint)之后的所有文件更新操作。
- FsImage(元数据检查点镜像文件) - 保存了文件系统中所有的目录和文件信息,如:某个目录下有哪些子目录和文件,以及文件名、文件副本数、文件由哪些 Block 组成等。
Active NameNode 内存中有一份最新的元数据(= FsImage + edits)。
Standby NameNode 在检查点定期将内存中的元数据保存到 FsImage 文件中。
🔬 扩展知识
扩展知识
- 【L3】利用 QJM 实现元数据高可用 - QJM(Quorum Journal Manager)基于 Paxos 算法:只要保证 Quorum(法定人数)数量的操作成功,就认为这是一次最终成功的操作。QJM 共享存储系统:部署奇数(2N+1)个 JournalNode 负责存储 edits 编辑日志;写 edits 时只要超过半数(N+1)的 JournalNode 返回成功,就代表本次写入成功;最多可容忍 N 个 JournalNode 宕机。
- 【L3】Active 节点选举 - 利用 ZooKeeper 实现 Active 节点选举,切换流程详见本文档『NameNode 如何实现主备切换?』。
- 【L4】故障处理的优先级 - 无 HA 的集群 NameNode 故障只能靠重启并回放 edits 恢复,停机时间与元数据规模正相关;生产集群应优先部署 HA,把「故障处理」转化为「自动切换」,详见本文档『HDFS 如何实现高可用?』。
🔀 发散问题
- Q:Standby 升级为 Active 前为什么必须确认元数据完全同步? → 否则新 Active 的命名空间视图落后于实际写入,可能造成客户端已确认的元数据操作丢失,还会与 DataNode 块报告对不上。
- Q:为什么 JournalNode 要部署奇数个? → 过半写入原则下,2N+1 个节点最多容忍 N 个宕机;偶数部署并不比少一个节点多换来容错能力,纯属浪费。
- Q:NameNode 元数据彻底丢失怎么办? → 只能靠 JournalNode 上的 edits 与备份的 FsImage 重建命名空间;因此生产中 FsImage 必须定期异地备份。
【简单】HDFS 安全模式有什么作用?⭐
🎯 目标等级:L2 | ⏱ 建议用时:5 min | 🏷 标签:Hadoop / HDFS 容错
💎 关键结论
NameNode 启动时进入安全模式,检查数据块的健康状况和副本数量;只有足够数量的数据块可用后,才退出安全模式开始正常对外服务。
⚡记忆卡片
- 口诀:启动先进安全模式,块够数才开工
- 关键词:启动阶段 / 块健康检查 / 副本数量检查 / 退出条件
- 链路:NameNode 启动 → 进入安全模式 → 检查块健康与副本数 → 达标后退出、正常服务
📖 核心知识
在启动过程中,NameNode 进入安全模式。在这个模式下,它会检查数据块的健康状况和副本数量。只有在足够数量的数据块可用时,NameNode 才会退出安全模式,开始正常的操作。
🔀 发散问题
- Q:安全模式下集群能读写吗? → 安全模式主要是只读自检阶段,写操作受限,直到退出安全模式才恢复正常服务。
- Q:为什么启动时需要安全模式? → NameNode 重启后需要等 DataNode 心跳与块报告上报,重建块位置映射并确认副本达标,避免在数据视图不完整时提供服务。
- Q:长时间卡在安全模式怎么办? → 通常意味着大量 DataNode 未上线或块报告异常,应先检查 DataNode 存活与网络连通性,而不是强退安全模式。
HA
【困难】HDFS 如何实现高可用?⭐⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:15 min | 🏷 标签:Hadoop / HDFS HA
💎 关键结论
HDFS 通过 Active/Standby 双 NameNode 互备 + 共享存储(QJM/NFS)同步元数据 + ZKFailoverController 借助 ZooKeeper 自动选举切换来实现高可用:Active 把 EditLog 写入 JournalNode 集群(过半写入即成功),Standby 持续拉取同步,故障时确认元数据完全同步后升主对外服务。
⚡记忆卡片
- 口诀:双 NN 互备,QJM 过半写,ZKFC 盯健康,ZK 管选主
- 关键词:Active/Standby / ZKFailoverController / ZooKeeper / QJM / JournalNode / DataNode 双向上报
- 链路:Active 写 EditLog 到 JournalNode → Standby 定时拉取同步 → 故障时 ZKFC 检测 + ZK 选举 → Standby 确认同步完成升 Active
📖 核心知识
HDFS 高可用架构如下:

HDFS 高可用架构主要由以下组件所构成:
- Active NameNode 和 Standby NameNode:两台 NameNode 形成互备,一台处于 Active 状态,为主 NameNode,另外一台处于 Standby 状态,为备 NameNode,只有主 NameNode 才能对外提供读写服务。
- 主备切换控制器 ZKFailoverController:ZKFailoverController 作为独立的进程运行,对 NameNode 的主备切换进行总体控制。ZKFailoverController 能及时检测到 NameNode 的健康状况,在主 NameNode 故障时借助 Zookeeper 实现自动的主备选举和切换,当然 NameNode 目前也支持不依赖于 Zookeeper 的手动主备切换。
- Zookeeper 集群:为主备切换控制器提供主备选举支持。
- 共享存储系统:共享存储系统是实现 NameNode 的高可用最为关键的部分,共享存储系统保存了 NameNode 在运行过程中所产生的 HDFS 的元数据。主 NameNode 和 NameNode 通过共享存储系统实现元数据同步。在进行主备切换的时候,新的主 NameNode 在确认元数据完全同步之后才能继续对外提供服务。
- DataNode 节点:除了通过共享存储系统共享 HDFS 的元数据信息之外,主 NameNode 和备 NameNode 还需要共享 HDFS 的数据块和 DataNode 之间的映射关系。DataNode 会同时向主 NameNode 和备 NameNode 上报数据块的位置信息。
目前 Hadoop 支持使用 Quorum Journal Manager (QJM) 或 Network File System (NFS) 作为共享的存储系统,这里以 QJM 集群为例进行说明:Active NameNode 首先把 EditLog 提交到 JournalNode 集群,然后 Standby NameNode 再从 JournalNode 集群定时同步 EditLog,当 Active NameNode 宕机后, Standby NameNode 在确认元数据完全同步之后就可以对外提供服务。
需要说明的是向 JournalNode 集群写入 EditLog 是遵循 “过半写入则成功” 的策略,所以你至少要有 3 个 JournalNode 节点,当然你也可以继续增加节点数量,但是应该保证节点总数是奇数。同时如果有 2N+1 台 JournalNode,那么根据过半写的原则,最多可以容忍有 N 台 JournalNode 节点挂掉。

🔬 扩展知识
扩展知识
- 【L3】Standby 如何快速接管 - Standby 平时通过 JournalNode 持续追平 EditLog,并周期性做 checkpoint 生成 FsImage;切换时只需回放剩余少量 edits,才能做到快速接管;若 Standby 长期落后,切换时间会显著拉长甚至拒绝切换。
- 【L3】DataNode 双向上报 - DataNode 同时向主备两个 NameNode 上报块位置,使 Standby 无需冷启动重建块映射,切换后立即掌握数据分布。
- 【L4】QJM 与 NFS 的选型 - QJM 基于多 JournalNode 过半写,无额外单点;NFS 依赖外部共享存储设备,存在单点与运维依赖;生产一般选 QJM。
- 【L4】与 Federation 的关系 - HA 解决可用性,不解决元数据规模;超大命名空间可叠加 Federation,每个命名空间内部各自配置 HA,见本文档『什么是 HDFS Federation?』。
🏭 实战场景
实战场景
- 场景:Active NameNode 所在节点宕机的自动切换 - 现象:Active NameNode 进程所在机器整体宕机,ZooKeeper 会话超时后 ZKFC 触发选举,Standby 确认 edits 追平后升为 Active,客户端短暂重试后恢复;期间 DataNode 上的块数据不受影响,只是元数据操作短暂阻塞。经验:切换耗时主要取决于 Standby 的元数据追赶进度,日常必须监控 checkpoint 与 Standby 同步健康度。
- 场景:计划内手动切换演练 - 集群升级或迁移前,先手动把 Active 转为 Standby、再把目标节点转为 Active,避免非预期故障触发切换;同时验证 fencing 配置有效,这是上线 HA 后的必做演练。
⚠️ 常见误区
常见误区
- ❌ "配了 HA 就永远不会停服、不会丢数据" → HA 保障的是 NameNode 故障后的快速接管,切换窗口内元数据操作仍会短暂不可用;数据可靠性仍靠副本机制保障,客户端也需要重试配合。
- ❌ "JournalNode 随便部署几台就行" → 必须奇数台且至少 3 台,过半写入原则下 2N+1 台最多容忍 N 台宕机;偶数台并不提升容错能力。
- ❌ "Standby 只是备份,平时不用管" → Standby 的同步进度、checkpoint 健康度直接决定切换速度,长期落后会导致切换耗时剧增甚至失败,必须纳入监控。
🔀 发散问题
- Q:HDFS HA 与 HDFS Federation 能同时使用吗? → 可以,Federation 的每个命名空间(Namespace Volume)内部都可以独立配置自己的 Active/Standby HA,两者分别解决扩展性与可用性问题,详见本文档『什么是 HDFS Federation?』。
- Q:为什么不能只定时拷贝 FsImage 给 Standby? → FsImage 只是某个检查点的快照,中间还有大量 edits 变更;只拷快照会丢失最新变更,所以必须实时同步 EditLog 并在切换前确认完全同步。
- Q:主备切换的具体流程是怎样的? → 详见本文档『NameNode 如何实现主备切换?』。
【困难】NameNode 如何实现主备切换?⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:12 min | 🏷 标签:Hadoop / HDFS HA
💎 关键结论
NameNode 主备切换由 ZKFailoverController(ZKFC)总体控制:HealthMonitor 定时检测健康状态 → 状态变化回调 ZKFC → ActiveStandbyElector 与 ZooKeeper 交互完成选举 → 回调 ZKFC 通过 HAServiceProtocol 把 NameNode 转为 Active/Standby;选举结果以 ZooKeeper 上的临时锁节点为凭据。
⚡记忆卡片
- 口诀:HM 探病,ZKFC 决策,Elector 抢锁,RPC 翻牌
- 关键词:HealthMonitor / ZKFailoverController / ActiveStandbyElector / ActiveStandbyElectorLock / HAServiceProtocol
- 链路:HealthMonitor 检测异常 → ZKFC 判断需切换 → ActiveStandbyElector 在 ZK 抢锁选举 → 回调 ZKFC → RPC 转 Active/Standby
📖 核心知识
NameNode 实现主备切换的流程下图所示:

工作流程说明:
- HealthMonitor 初始化完成之后会启动内部的线程来定时调用对应 NameNode 的 HAServiceProtocol RPC 接口的方法,对 NameNode 的健康状态进行检测。
- HealthMonitor 如果检测到 NameNode 的健康状态发生变化,会回调 ZKFailoverController 注册的相应方法进行处理。
- 如果 ZKFailoverController 判断需要进行主备切换,会首先使用 ActiveStandbyElector 来进行自动的主备选举。
- ActiveStandbyElector 与 Zookeeper 进行交互完成自动的主备选举。
- ActiveStandbyElector 在主备选举完成后,会回调 ZKFailoverController 的相应方法来通知当前的 NameNode 成为主 NameNode 或备 NameNode。
- ZKFailoverController 调用对应 NameNode 的 HAServiceProtocol RPC 接口的方法将 NameNode 转换为 Active 状态或 Standby 状态。
主备选举过程:
NameNode 在选举成功后,会在 zk 上创建了一个 /hadoop-ha/${dfs.nameservices}/ActiveStandbyElectorLock 节点,而没有选举成功的备 NameNode 会监控这个节点,通过 Watcher 来监听这个节点的状态变化事件,ZKFC 的 ActiveStandbyElector 主要关注这个节点的 NodeDeleted 事件(这部分实现跟 Kafka 中 Controller 的选举一样)。
如果 Active NameNode 对应的 HealthMonitor 检测到 NameNode 的状态异常时, ZKFailoverController 会主动删除当前在 Zookeeper 上建立的临时节点 /hadoop-ha/${dfs.nameservices}/ActiveStandbyElectorLock,这样处于 Standby 状态的 NameNode 的 ActiveStandbyElector 注册的监听器就会收到这个节点的 NodeDeleted 事件。收到这个事件之后,会马上再次进入到创建 /hadoop-ha/${dfs.nameservices}/ActiveStandbyElectorLock 节点的流程,如果创建成功,这个本来处于 Standby 状态的 NameNode 就选举为主 NameNode 并随后开始切换为 Active 状态。
当然,如果是 Active 状态的 NameNode 所在的机器整个宕掉的话,那么根据 Zookeeper 的临时节点特性,/hadoop-ha/${dfs.nameservices}/ActiveStandbyElectorLock 节点会自动被删除,从而也会自动进行一次主备切换。
🔬 扩展知识
扩展知识
- 【L3】六段链路的职责分工 - HealthMonitor 只负责探测,ZKFC 负责决策与编排,ActiveStandbyElector 只负责选举,状态翻转最终靠 HAServiceProtocol RPC 完成;各环节职责单一,便于独立测试与替换。
- 【L3】与 Kafka Controller 选举的相似性 - 抢锁 + Watcher 监听 NodeDeleted 的实现与 Kafka Controller 选举一致,是基于 ZooKeeper 选主的经典模式。
- 【L4】切换失败的常见根因 - ZooKeeper 会话抖动、备节点元数据未同步完成、fencing 执行失败都会导致切换中止或延迟,见本文档『如何应对 HDFS 脑裂问题?』。
🔀 发散问题
- Q:支持手动切换吗? → 支持,NameNode 也支持不依赖 ZooKeeper 的手动主备切换,适合计划内运维场景。
- Q:为什么用临时节点做锁? → Active 所在机器宕机后会话断开、临时节点自动删除,备节点的 Watcher 收到 NodeDeleted 事件后立即发起新一轮选举,天然实现故障感知。
- Q:切换过程中如何避免脑裂? → 靠 fencing 隔离旧主,详见本文档『如何应对 HDFS 脑裂问题?』。
【困难】如何应对 HDFS 脑裂问题?⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:12 min | 🏷 标签:Hadoop / HDFS HA
💎 关键结论
脑裂指旧 Active 假死(如长时间 GC)导致备节点升主后,集群同时存在两个「主」;HDFS 用隔离(Fencing)应对:在共享存储、DataNode、Client 三处保证只有一个 NameNode 生效,新主上任前先通过 ActiveBreadCrumb 发现旧主,先劝退、失败再执行 sshfence/shellfence,fencing 成功后才转为 Active。
⚡记忆卡片
- 口诀:假死生双主,上任先清旧,fence 成功才转正
- 关键词:脑裂 / Fencing / ActiveBreadCrumb / transitionToStandby / sshfence / shellfence
- 链路:旧 Active 假死失联 → 备节点当选 → 读 ActiveBreadCrumb 发现旧主 → 先 transitionToStandby、失败再 fence → 成功后 becomeActive
📖 核心知识
在实际中,NameNode 可能会出现这种情况,NameNode 在垃圾回收(GC)时,可能会在长时间内整个系统无响应,因此,也就无法向 zk 写入心跳信息,这样的话可能会导致临时节点掉线,备 NameNode 会切换到 Active 状态,这种情况,可能会导致整个集群会有同时有两个 NameNode,这就是脑裂问题。
脑裂问题的解决方案是隔离(Fencing),主要是在以下三处采用隔离措施:
- 第三方共享存储:任一时刻,只有一个 NN 可以写入;
- DataNode:需要保证只有一个 NN 发出与管理数据副本有关的删除命令;
- Client:需要保证同一时刻只有一个 NN 能够对 Client 的请求发出正确的响应。
关于这个问题目前解决方案的实现如下:
- ActiveStandbyElector 为了实现隔离,会在成功创建 Zookeeper 节点
hadoop-ha/${dfs.nameservices}/ActiveStandbyElectorLock从而成为 Active NameNode 之后,创建另外一个路径为/hadoop-ha/${dfs.nameservices}/ActiveBreadCrumb的持久节点,这个节点里面保存了这个 Active NameNode 的地址信息; - Active NameNode 的 ActiveStandbyElector 在正常的状态下关闭 Zookeeper Session 的时候,会一起删除这个持久节点;
- 但如果 ActiveStandbyElector 在异常的状态下 Zookeeper Session 关闭 (比如前述的 Zookeeper 假死),那么由于
/hadoop-ha/${dfs.nameservices}/ActiveBreadCrumb是持久节点,会一直保留下来,后面当另一个 NameNode 选主成功之后,会注意到上一个 Active NameNode 遗留下来的这个节点,从而会回调 ZKFailoverController 的方法对旧的 Active NameNode 进行 fencing。
在进行隔离的时候,会执行以下的操作:
首先尝试调用这个旧 Active NameNode 的 HAServiceProtocol RPC 接口的 transitionToStandby 方法,看能不能把它转换为 Standby 状态; 如果 transitionToStandby 方法调用失败,那么就执行 Hadoop 配置文件之中预定义的隔离措施。
Hadoop 目前主要提供两种隔离措施,通常会选择第一种:sshfence:通过 SSH 登录到目标机器上,执行命令 fuser 将对应的进程杀死; shellfence:执行一个用户自定义的 shell 脚本来将对应的进程隔离。 只有在成功地执行完成 fencing 之后,选主成功的 ActiveStandbyElector 才会回调 ZKFailoverController 的 becomeActive 方法将对应的 NameNode 转换为 Active 状态,开始对外提供服务。
🔬 扩展知识
扩展知识
- 【L3】NameNode 选举机制与 Kafka Controller 的对照 - NameNode 选举的实现机制与 Kafka 的 Controller 类似:Kafka 给 Broker 发送的请求中都会携带 controller epoch 信息,如果 Broker 发现当前请求的 epoch 小于缓存中的值,就证明这是来自旧 Controller 的请求,就会拒绝这个请求;但异常情况下(Broker 先收到旧 Controller 的请求)并没有完美方案,新 Controller 选出后会向全局所有 Broker 发送 metadata 请求收敛 epoch,出现问题的概率非常小,且即使出现由于 Kafka 的高可靠架构影响也很有限。
- 【L3】epoch 通用做法 - 通过标识每次选举的版本号,并以最新版本选举结果为准,是分布式选举避免脑裂的常见做法。在其他分布式系统中,epoch 可能会被称为 term、version 等。
- 【L4】假死根因治理 - GC 假死是触发脑裂场景的典型诱因,生产上应监控 NameNode GC 时长并优化堆配置,从源头减少误切换。
🔀 发散问题
- Q:为什么 fencing 必须成功才允许转正? → 若旧主仍存活且未被隔离,两个 Active 可能同时下发副本管理等命令或响应客户端,造成元数据与数据不一致。
- Q:sshfence 和 shellfence 怎么选? → 通常选 sshfence(SSH 登录目标机执行 fuser 杀死对应进程);环境不便时用 shellfence 执行自定义脚本隔离进程。
- Q:脑裂会导致数据丢失吗? → 防护到位时不会,fencing 保证单写者;风险在切换窗口内客户端可能收到旧主的过期响应,客户端需重试并以新主为准。
【困难】YARN 如何实现高可用?⭐⭐
🎯 目标等级:L2 | ⏱ 建议用时:10 min | 🏷 标签:Hadoop / YARN HA
💎 关键结论
YARN ResourceManager 的 HA 与 HDFS NameNode 类似,采用 Active/Standby 双 RM;但 RM 不像 NameNode 有海量元数据需要维护,它的状态信息可以直接写到 ZooKeeper 上,并依赖 ZooKeeper 进行主备选举,无需额外的共享存储(如 JournalNode)。
⚡记忆卡片
- 口诀:双 RM 互备,状态存 ZK,选举靠 ZK
- 关键词:Active/Standby ResourceManager / ZooKeeper / 状态存储 / 主备选举
- 链路:Active RM 把状态写入 ZooKeeper → 故障时 ZK 选举 → 新 Active 从 ZK 恢复状态 → 接管调度
📖 核心知识
YARN ResourceManager 的高可用与 HDFS NameNode 的高可用类似,但是 ResourceManager 不像 NameNode,没有那么多的元数据信息需要维护,所以它的状态信息可以直接写到 Zookeeper 上,并依赖 Zookeeper 来进行主备选举。

🔬 扩展知识
扩展知识
- 【L3】与 HDFS HA 的对比 - HDFS 需要 QJM 这类共享编辑日志来同步海量元数据;RM 维护的是调度状态而非全量文件元数据,体量小,直接以 ZooKeeper 作为状态存储与选举协调器,架构更轻量。
- 【L4】切换期间的作业影响 - RM 切换时正在运行的 Container 状态需要重建,部分任务可能重跑;上层框架(如 MapReduce AM 重试机制)可吸收这种抖动,生产中应把 RM HA 纳入集群基线配置。
🔀 发散问题
- Q:为什么 RM 不需要 JournalNode 这样的共享存储? → RM 维护的是调度状态而非全量文件元数据,体量小,直接写入 ZooKeeper 即可满足同步与选举需求。
- Q:RM HA 切换后已提交的作业还在吗? → 作业状态持久化在 ZooKeeper 中,新 Active 恢复后继续管理;正在执行的任务可能因状态重建而重跑。
- Q:与 HDFS HA 的切换机制有什么共同点? → 都采用 Active/Standby + 协调服务选举的模式,HDFS 用 ZooKeeper + ZKFC,YARN 直接用 ZooKeeper,详见本文档『HDFS 如何实现高可用?』。


