大数据面试题——hadoop(hdfs、mapreduce、yarn)
文章目录
- Hadoop
- hadoop的常用配置文件有哪些
- 启动hadoop集群会分别启动哪些进程,各自的作用
- 简述java序列化和 hadoop自带序列化机制及其区别
- 请说下 HDFS 的组织架构
- 请说下 HDFS 读写流程
- NameNode 在启动的时候会做哪些操作
- Hadoop的HA的了解(High Availability高可用,HA)
- 在 NameNode HA 中,会出现脑裂问题吗?怎么解决脑裂
- 请说下 MapReduce组织架构
- MapReduce过程——map和reduce
- 环形缓冲区hadoop
- 请说下 MR 中 shuffle
- 在写 MR 时,什么情况下可以使用规约
- reduce函数如何知道从哪台机器获取map输出?
- MapReduce执行速度太慢,怎么办?
- yarn 集群的架构和工作原理
- yarn 的任务提交流程是怎样的
- yarn 的资源调度三种模型了解吗
- HDFS 在读取文件的时候, 如果其中一个块突然损坏了怎么办
- ------------------------------------------------
- HDFS 在上传文件的时候, 如果其中一个 DataNode 突然挂掉了怎么办
- SecondaryNameNode 了解吗,它的工作机制是怎样的
- Secondary NameNode 不能恢复 NameNode 的全部数据,那如何保证 NameNode 数据存储安全
- 8. 小文件过多会有什么危害, 如何避免
- 13. shuffle 阶段的数据压缩机制了解吗
Hadoop
hadoop 核心三套件
第一:存储——(HDFS )
第二:计算框架—— (MapReduce)
第三:资源调度框架—— (yarn)
hadoop的常用配置文件有哪些
-
hadoop-env.sh: 用于定义hadoop运行环境相关的配置信息,比如配置
JAVA_HOME环境变量、为hadoop的JVM指定特定的选项、指定日志文件所在的目录路径以及master和slave文件的位置等; -
core-site.xml: 用于定义系统级别的参数,如HDFS URL、Hadoop的临时目录以及用于rack-aware集群中的配置文件的配置等,此中的参数定义会覆盖core-default.xml文件中的默认配置;
-
hdfs-site.xml: HDFS的相关设定,如文件副本的个数、块大小及是否使用强制权限等,此中的参数定义会覆盖hdfs-default.xml文件中的默认配置;
-
mapred-site.xml:HDFS的相关设定,如reduce任务的默认个数、任务所能够使用内存的默认上下限等,此中的参数定义会覆盖mapred-default.xml文件中的默认配置;
启动hadoop集群会分别启动哪些进程,各自的作用
-
NameNode:
- 维护文件系统树及整棵树内所有的文件和目录。这些信息永久保存在本地磁盘的两个文件中:命名空间镜像文件、编辑日志文件
- 记录每个文件中各个块所在的数据节点信息,这些信息在内存中保存,每次启动系统时重建这些信息
- 负责响应客户端的 数据块位置请求 。也就是客户端想存数据,应该往哪些节点的哪些块存;客户端想取数据,应该到哪些节点取
- 接受记录在数据存取过程中,datanode节点报告过来的故障、损坏信息
-
SecondaryNameNode(非HA模式):
- 实现namenode容错的一种机制。定期合并编辑日志与命名空间镜像,当namenode挂掉时,可通过一定步骤进行上顶。(注意 并不是NameNode的备用节点)
-
DataNode:
- 根据需要存取并检索数据块
- 定期向namenode发送其存储的数据块列表
-
ResourceManager:
- 负责Job的调度,将一个任务与一个NodeManager相匹配。也就是将一个MapReduce之类的任务分配给一个从节点的NodeManager来执行。
-
NodeManager:
- 运行ResourceManager分配的任务,同时将任务进度向application master报告
-
JournalNode(HA下启用):
- 高可用情况下存放namenode的editlog文件
简述java序列化和 hadoop自带序列化机制及其区别
| Java基本类型 | Writable | 序列化后的长度 |
|---|---|---|
| boolean | BooleanWritable | 1 |
| byte | ByteWritable | 1 |
| int | IntWritableVIntWritable | 41~5 |
| float | FloatWritable | 4 |
| long | LongWritableVLongWritable | 81~9 |
| double | DoubleWritable | 8 |
-
什么是序列化?
将对象编码成一个字节流(序列化),以及从字节流中重新构建对象(反序列化)。、序列化的作用:
(1)一种持久化的格式:将将对象序列化后,把字节存到磁盘上;
(2)一种通信数据格式:进行网络传输,比如从一个虚拟机传输到另一个虚拟机;(3)一种拷贝和克隆的机制:用于深铂贝。
Hadoop中一般作用是前两种。 -
Java序列化
一个类只要实现了serializable(序列化接口),就可以进行序列化了。但是serializable只是一个标识,没有具体方法,只是判断一类是否可以序列化。
实现:创建一个 ObjectOutputStream对象,这个对象指示序列化写入的地方,然后调用writeObject()方法,进行序列化写入。
如果存在继承问题,父类实现序列化接口,则子类自动实现;若子类实现了序列化,父类没有实现,则父类需要一个无参构造器,子类将负责序列化父类的域。
序列化结果中包含了大量与类相关的信息导致结果膨胀,所以对于hadoop来说需要新的序列化机制。 -
Hadoop序列化
Hadoop采用序列化接口 Writable。可实现对基本数据类型进行序列化和自定义对象序列化。RawComparator 接口允许在数据流中比较大小,不用反序列化,节省开销,是在每个基本
类型封装器中通过静态内部类实现的,比如IntWritable内部类 Comparator 继承了WritableComparator类(RawComparator接口一个通用实现),可在数据流中比较大小。 -
java序列化机制和hadoop序列化比较。hadoop序列化框架相比较优势在于:
(1)绝对紧凑:序列化后数据更加紧凑,节省带宽。
(2)快速快:尽量避免减小序列化和反序列化的开销。
请说下 HDFS 的组织架构
-
Client:客户端
(1)切分文件。文件上传 HDFS 的时候,Client 将文件切分成一个一个的Block,然后进行存储
(2)与 NameNode 交互,获取文件的位置信息
(3)与 DataNode 交互,读取或者写入数据
(4)Client 提供一些命令来管理 HDFS,比如启动关闭 HDFS、访问 HDFS 目录及内容等 -
NameNode:名称节点,也称主节点,存储数据的元数据信息,不存储具体的数据
(1)管理 HDFS 的名称空间
(2)管理数据块(Block)映射信息
(3)配置副本策略
(4)处理客户端读写请求 -
DataNode:数据节点,也称从节点。NameNode 下达命令,DataNode 执行实际的操作,Block 默认是64MB(HDFS2.0改成了128MB),当客户端上传一个大文件时,HDFS 会自动将其切割成固定大小的 Block,为了保证数据可用性,每个 Block 会以多备份的形式存储,默认是3份。
(1)存储实际的数据块
(2)执行数据块的读 / 写操作 -
Secondary NameNode:并非 NameNode 的热备。当 NameNode 挂掉的时候,它并不能马上替换 NameNode 并提供服务
(1)定期合并 Fsimage 和 Edits,并推送给 NameNode
(2)辅助 NameNode,分担其工作量
(3) 在紧急情况下,可辅助恢复 NameNode
- fsimage:是内存命名空间元数据在外存的镜像文件;
- editlog:则是各种元数据操作的 write-ahead-log 文件,在体现到内存数据变化前首先会将操作记入 editlog 中,以防止数据丢失。
请说下 HDFS 读写流程

HDFS 写流程

-
Client 调用
DistributedFileSystem对象的create方法,创建一个文件输出流(FSDataOutputStream)对象; -
通过
DistributedFileSystem对象与集群的NameNode进行一次RPC远程调用,namenode 检查该用户是否有上传权限,以及上传的文件是否在 hdfs 对应的目录下重名,
在 HDFS 的 Namespace 中创建一个文件条目(Entry),此时该条目没有任何的 Block,NameNode 会返回该数据每个块需要拷贝的 DataNode 地址信息; -
通过
FSDataOutputStream对象,开始向 DataNode 写入数据,数据首先被写入FSDataOutputStream对象内部的数据队列中,数据队列由 DataStreamer 使用,它通过选择合适的 DataNode 列表来存储副本,从而要求 NameNode 分配新的 block; -
namenode 收到请求之后,根据网络拓扑和机架感知以及副本机制进行文件分配,返回可用的 DataNode 的地址
DataStreamer将数据包以流式传输的方式传输到分配的第一个 DataNode 中,该数据流将数据包存储到第一个 DataNode 中并将其转发到第二个 DataNode 中,接着第二个 DataNode 节点会将数据包转发到第三个 DataNode 节点; -
DataNode 确认数据传输完成,最后由第一个 DataNode 通知 client 数据写入成功;
-
完成向文件写入数据,Client 在文件输出流(FSDataOutputStream)对象上调用
close方法,完成文件写入;
注:Hadoop 在设计时考虑到数据的安全与高效, 数据文件默认在 HDFS 上存放三份, 存储策略为本地一份,同机架内其它某一节点上一份, 不同机架的某一节点上一份
HDFS 读流程

-
Client 通过
DistributedFileSystem对象与集群的 NameNode 进行一次RPC远程调用,获取文件 block 位置信息; -
namenode 收到请求之后会检查用户权限以及是否有存在文件, NameNode 返回存储的每个块的
DataNode 列表;
NameNode 都会返回含有该 block 副本的 DataNode 地址; 这些返回的 DN 地址,会按照集群拓扑结构得出 DataNode 与客户端的距离,然后进行排序
排序两个规则:网络拓扑结构中距离 Client 近的排靠前;心跳机制中超时汇报的 DN 状态为 STALE,这样的排靠后 -
Client 选取排序靠前的 DataNode 来读取 block,如果客户端本身就是 DataNode, 那么将从本地直接获取数据 (短路读取特性)
-
Client 开始从 DataNode 并行读取数据;NameNode 只是返回 Client 请求包含块的 DataNode 地址,并不是返回请求块的数据
-
读取完一个 block 都会进行 checksum 验证,如果读取 DataNode 时出现错误,客户端会通知 NameNode,然后再从下一个拥有该 block 副本的 DataNode 继续读【一个块突然损坏】
-
最终读取来所有的 block 会合并成一个完整的最终文件
NameNode 在启动的时候会做哪些操作
NameNode 数据存储在内存和本地磁盘,本地磁盘数据存储在 fsimage镜像文件和 edits编辑日志文件
NameNode主要是用来保存HDFS的元数据信息,比如命名空间信息,块信息等。当它运行的时候,这些信息是存在内存中的。但是这些信息也可以持久化到磁盘上。
- fsimage - 它是在NameNode启动时对整个文件系统的快照
- edit logs - 它是在NameNode启动后,对文件系统的改动序列
-
首次启动 NameNode
1、格式化文件系统,为了生成fsimage镜像文件
2、启动 NameNode
(1)读取 fsimage 文件,将文件内容加载进内存
(2)等待 DataNade 注册与发送 block-report
3、启动 DataNode
(1)向 NameNode 注册
(2)发送 block-report
(3)检查 fsimage 中记录的块的数量和 block-report 中的块的总数是否相同
4、对文件系统进行操作(创建目录,上传文件,删除文件等)
此时内存中已经有文件系统改变的信息,但是磁盘中没有文件系统改变的信息,此时会将这些改变信息写入 edits 文件中,edits 文件中存储的是文件系统元数据改变的信息。 -
第二次启动 NameNode
1、读取 fsimage 和 edits 文件
2、将 fsimage 和 edits 文件合并成新的 fsimage 文件
3、创建新的 edits 文件,内容为空
4、启动 DataNode

Hadoop的HA的了解(High Availability高可用,HA)
1. AvatarNode方案
Facebook提出在AvatarNode方案中使用Avatar Primary NameNode,Avatar Standby NameNode及NFS(网络文件系统(Network File System))配合通过共享的方式管理HDFS部分元数据EditLog。
Avatar Primary NameNode提供读写服务,并将Editlog写入到共享存储NFS上,Avatar Standby NameNode从共享存储NFS上读取Editlog数据并回放,这样尽可能保持与Primary之间状态一致。

但是AvatarNode这种方案并不是完美的。它的问题主要是共享存储NFS成为新的SPOF(单点故障(single point of failure)),必须保证其高可用;同时AvatarNode不具备自动Failover能力,一旦
- Avatar Primary NameNode出现故障,需要运维人员介入手动处理,
- 另外一点,昂贵的NFS设备引入与HDFS最初构建在“inexpensive commodity hardware”设计初衷多少有些出入。
2. QJM架构
QJM的思想最初来源于Paxos(帕克西)协议,摒弃了AvatarNode方案中的共享存储设备,改用多个JournalNode节点组成集群来管理和共享EditLog。
与Paxos协议类似,当NameNode向JournalNode请求读写时,要求至少大多数(Majority)成功返回才认为本次请求成功。
对于一个由2N+1台JournalNode组成的集群,可以容忍最多N台JournalNode节点挂掉。从这个角度来看,QJM相比AvatarNode方案具备了更强的HA能力。

在HA using QJM方案中,涉及到的核心组件包括:
Active NameNode(ANN):在HDFS集群中,对外提供读写服务的唯一Master节点。ANN将客户端请求过来的写操作通过EditLog写入共享存储系统(即JournalNode Cluster),为Standby NameNode及时同步数据提供支持;Standby NameNode(SBN):与ANN相互形成热备,SBN及时从共享存储系统中读取EditLog数据并更新内存,以保证当前状态尽可能与ANN同步。当前在整个HDFS集群中最多一台处于Active状态,最多一台处于Standby状态;JournalNode Cluster(JNs):ANN与SBN之间共享Editlog的一致性存储系统,是HDFS NameNode高可用的核心组件。借助JournalNode集群ANN可以尽可能及时同步元数据到SBN。其中ANN采用Push模式将EditLog写入JN,SBN通过Pull模式从JN拉取数据,整个过程中JN不主动进行数据交换;ZKFailoverController(ZKFC):ZKFailoverController以独立进程运行,对NameNode主备切换进行控制,正常情况ANN和SBN分别对应各自ZKFC进程。ZKFC主要功能:NameNode健康状况检测;借助Zookeeper实现NameNode自动选主;操作NameNode进行主从切换;Zookeeper(ZK):为ZKFC实现自动选主功能提供统一协调服务。
需要说明的是,在HA using QJM架构下,DataNode从仅向单个NameNode进行数据交互升级到同时向ANN和SBN进行数据交互,区别是仅执行ANN下发的指令,其他逻辑未发生大变化。
在 NameNode HA 中,会出现脑裂问题吗?怎么解决脑裂
假设 NameNode1 当前为 Active 状态,NameNode2 当前为 Standby 状态。如果某一时刻 NameNode1 对应的 ZKFailoverController 进程发生了 “假死” 现象,那么 Zookeeper 服务端会认为 NameNode1 挂掉了,根据前面的主备切换逻辑,NameNode2 会替代 NameNode1 进入 Active 状态。但是此时 NameNode1 可能仍然处于 Active 状态正常运行,这样 NameNode1 和 NameNode2 都处于 Active 状态,都可以对外提供服务。这种情况称为脑裂
脑裂对于 NameNode 这类对数据一致性要求非常高的系统来说是灾难性的,数据会发生错乱且无法恢复。Zookeeper 社区对这种问题的解决方法叫做 fencing,中文翻译为隔离,也就是想办法把旧的 Active NameNode 隔离起来,使它不能正常对外提供服务。
在进行 fencing 的时候,会执行以下的操作:
-
首先尝试调用这个旧 Active NameNode 的 HAServiceProtocol RPC 接口的 transitionToStandby 方法,看能不能把它转换为 Standby 状态。
-
如果 transitionToStandby 方法调用失败,那么就执行 Hadoop 配置文件之中预定义的隔离措施,Hadoop 目前主要提供两种隔离措施,通常会选择 sshfence:
(1) sshfence:通过 SSH 登录到目标机器上,执行命令 fuser 将对应的进程杀死
(2) shellfence:执行一个用户自定义的 shell 脚本来将对应的进程隔离
请说下 MapReduce组织架构

MapReduce主要有以下4个部分组成:Client、JobTracker、TaskTracker、Task
1. Client:
👉用户编写的MapReduce程序通过Client提交到JobTracker端。
👉用户可通过Client提供的一些接口查看作业运行状态。
2. JobTracker:
👉JobTracker负责资源监控和作业调度。
👉JobTracker 监控所有TaskTracker与Job的健康状况,一旦发现失败,就将相应的任务转移到其他节点。
👉JobTracker 会跟踪任务的执行进度、资源使用量等信息,并将这些信息告诉任务调度器(TaskScheduler),而调度器会在资源出现空闲时,选择合适的任务去使用这些资源。
3. TaskTracker:
👉 TaskTracker 会周期性地通过“心跳”将本节点上资源的使用情况和任务的运行进度汇报给JobTracker,同时接收JobTracker 发送过来的命令并执行相应的操作(如启动新任务、杀死任务等)。
👉TaskTracker 使用“slot”等量划分本节点上的资源量(CPU、内存等)。一个Task 获取到一个slot 后才有机会运行,而Hadoop调度器的作用就是将各个TaskTracker上的空闲slot分配给Task使用。slot 分为Map slot 和Reduce slot 两种,分别供MapTask 和Reduce Task 使用。
4. Task:
👉 Task 分为Map Task和Reduce Task 两种,均由TaskTracker 启动。
MapReduce过程——map和reduce
map端
- 由程序内的
InputFormat(默认实现类TextInputFormat)来读取外部数据,它会调用RecordReader(它的成员变量)的read()方法来读取,返回k-v键值对。 - 读取的k,v键值对传送给
map()方法,作为其入参来执行角户定义的 map逻辑。 context.write方法被调用时,outputCollector组件会将map()方法的输出结果写入到环形缓决区内。- 环形缓冲区其实就是一个数组,后端不断接受数据的同时,前端数据不断被溢出,长度用完后读取的新数据再从前端开始覆盖。这个缓冲区默认大小100M,可以通过
MR.SORT.MB配置。 spiller组件会从环形缓冲区溢出文件,这过程会按照定义的partitioner分区(默认是hashpartition),并且按照 key.compareTo进行排序(底层主要用快排和外部排序),若有combiner也会执行 combiner。spiller 的不断工作,会不断澄出许多小文件。这些文件仍在map task所处机器上。- 小文件执行
merge(合并),行程分区且区内有序的大文件(归并排序,会再一次调用combiner)。 - Reduce会根据自己的分区,去所有 map task中,从文件读取对应的数据。
reduce端
1.reduce task通过网络向map task获取某一分区的数据。
2.通过 GroupingComparator()分辨同一组的数据,把他们发送给reduce(k,.iterator)方法-(这里多个数据合成一组时,只取其中一个key,取得是第一个)。
3.调用context.write()方法,会让 OutPurFormat调用RecodeWriter的 write()方法将处理结果写入到数据仓库中。写出的只有一个分区的文件数据
环形缓冲区hadoop

底层就是一个字节数组:数组前面记录关于KV的索引位置,数组后面记录KV数据。首尾相接构成一个环形的缓冲区,中间是赤道。用于数据spil溢出处理。
缓冲区采用典型单生产者消费者模型。
MapOutputBuffer的 collect方法和MapOutputBuffer.Buffer的write方法作为生产者
spillThread 线程是消费者,其间同步是通过可重入互斥锁spillLock和 spillLock 上的两个条件变量(spillDone和 spillReady)实现的。
kvoffsets主要由三个变量控制,kvstart、kvend、kvindex。开始时kvstart=kvend,kvindex指向待写入位置,当写入一条数据后,kvindex向后移动一位,当kvoffsets内存使用率超过io.sort.spill.percent(默认80%)后,数据开始溢出到磁盘。
请说下 MR 中 shuffle
- 作业启动:开发者通过控制台启动作业;
- 作业初始化:这里主要是切分数据、创建作业和提交作业,与第三步紧密相联;
- 作业/任务调度:
👉对于1.0版的Hadoop来说就是JobTracker来负责任务调度,
👉对于2.0版的Hadoop来说就是Yarn中的Resource Manager负责整个系统的资源管理与分配 - Map任务;
- Shuffle;
- Reduce任务;
- 作业完成:通知开发者任务完成。
而这其中最主要的MapReduce过程,主要是第4、5、6步三部分
- Map:数据输入,做初步的处理,输出形式的中间结果;
- Shuffle:按照partition、key对中间结果进行排序合并,输出给reduce线程;
- Reduce:对相同key的输入进行最终的处理,并将结果写入到文件中。

map_shuffle
实际包含了输入(input)过程、切分(partition)过程、溢写spill过程(sort和combine过程)、merge过程。

输入
- map task只读取split分片,split与block(hdfs的最小存储单位,默认为64MB)可能是一对一也能是一对多,但是对于一个split只会对应一个文件的一个block或多个block,不允许一个split对应多个文件的多个block;
- 这里切分和输入数据的时会涉及到InputFormat的文件切分算法和host选择算法。
文件切分算法,主要用于确定InputSplit的个数以及每个InputSplit对应的数据段。FileInputFormat以文件为单位切分生成InputSplit,对于每个文件,由以下三个属性值决定其对应的InputSplit的个数:
- goalSize: 它是根据用户期望的InputSplit数目计算出来的,即totalSize/numSplits。其中,totalSize为文件的总大小;numSplits为用户设定的Map Task个数,默认情况下是1;
- minSize:InputSplit的最小值,由配置参数
mapred.min.split.size确定,默认是1; - blockSize:文件在hdfs中存储的block大小,不同文件可能不同,默认是64MB。
这三个参数共同决定InputSplit的最终大小,计算方法如下:
splitSize=max{minSize, min{gogalSize,blockSize}}
FileInputFormat的host选择算法参考《Hadoop技术内幕-深入解析MapReduce架构设计与实现原理》的p50.
Partition
- 作用:将map的结果发送到相应的reduce端,总的partition的数目等于reducer的数量。
- 实现功能:
- map输出的是key/value对,决定于当前的mapper的part交给哪个reduce的方法是:mapreduce提供的Partitioner接口,对key进行hash后,再以reducetask数量取模,然后到指定的job上(HashPartitioner,可以通过
job.setPartitionerClass(MyPartition.class)自定义)。 - 然后将数据写入到内存缓冲区,缓冲区的作用是批量收集map结果,减少磁盘IO的影响。key/value对以及Partition的结果都会被写入缓冲区。在写入之前,key与value值都会被序列化成字节数组。
- map输出的是key/value对,决定于当前的mapper的part交给哪个reduce的方法是:mapreduce提供的Partitioner接口,对key进行hash后,再以reducetask数量取模,然后到指定的job上(HashPartitioner,可以通过
- 要求:负载均衡,效率;
spill(溢写):sort & combiner
- 作用:把内存缓冲区中的数据写入到本地磁盘,在写入本地磁盘时先按照partition、再按照key进行排序(
quick sort); - 注意:
- 这个spill是由另外单独的线程来完成,不影响往缓冲区写map结果的线程;
- 内存缓冲区默认大小限制为100MB,它有个溢写比例(
spill.percent),默认为0.8,当缓冲区的数据达到阈值时,溢写线程就会启动,先锁定这80MB的内存,执行溢写过程,maptask的输出结果还可以往剩下的20MB内存中写,互不影响。然后再重新利用这块缓冲区,因此Map的内存缓冲区又叫做环形缓冲区(两个指针的方向不会变,下面会详述); - 在将数据写入磁盘之前,先要对要写入磁盘的数据进行一次排序操作,先按
<key,value,partition>中的partition分区号排序,然后再按key排序,这个就是sort操作,最后溢出的小文件是分区的,且同一个分区内是保证key有序的;
combine:执行combine操作要求开发者必须在程序中设置了combine(程序中通过job.setCombinerClass(myCombine.class)自定义combine操作)。
- 程序中有两个阶段可能会执行combine操作:
- map输出数据根据分区排序完成后,在写入文件之前会执行一次combine操作(前提是作业中设置了这个操作);
- 如果map输出比较大,溢出文件个数大于3(此值可以通过属性
min.num.spills.for.combine配置)时,在merge的过程(多个spill文件合并为一个大文件)中还会执行combine操作;
- combine主要是把形如
<aa,1>,<aa,2>这样的key值相同的数据进行计算,计算规则与reduce一致,比如:当前计算是求key对应的值求和,则combine操作后得到<aa,3>这样的结果。 - 注意事项:不是每种作业都可以做combine操作的,只有满足以下条件才可以:
- reduce的输入输出类型都一样,因为combine本质上就是reduce操作;
- 计算逻辑上,combine操作后不会影响计算结果,像求和就不会影响;
merge
- merge过程:当map很大时,每次溢写会产生一个spill_file,这样会有多个spill_file,而最终的一个map task输出只有一个文件,因此,最终的结果输出之前会对多个中间过程进行多次溢写文件(spill_file)的合并,此过程就是merge过程。也即是,待Map Task任务的所有数据都处理完后,会对任务产生的所有中间数据文件做一次合并操作,以确保一个Map Task最终只生成一个中间数据文件。
- 注意:
- 如果生成的文件太多,可能会执行多次合并,每次最多能合并的文件数默认为10,可以通过属性
min.num.spills.for.combine配置; - 多个溢出文件合并时,会进行一次排序,排序算法是多路归并排序;
- 是否还需要做combine操作,一是看是否设置了combine,二是看溢出的文件数是否大于等于3;
- 最终生成的文件格式与单个溢出文件一致,也是按分区顺序存储,并且输出文件会有一个对应的索引文件,记录每个分区数据的起始位置,长度以及压缩长度,这个索引文件名叫做
file.out.index。
- 如果生成的文件太多,可能会执行多次合并,每次最多能合并的文件数默认为10,可以通过属性
reduce_shuffle
在 reduce task 之前,不断拉取当前 job 里每个 maptask 的最终结果,然后对从不同地方拉取过来的数据不断地做 merge ,也最终形成一个文件作为 reduce task 的输入文件。

copy过程
- 作用:拉取数据;
- 过程:Reduce进程启动一些数据copy线程(
Fetcher),通过HTTP方式请求map task所在的TaskTracker获取map task的输出文件。因为这时map task早已结束,这些文件就归TaskTracker管理在本地磁盘中。 - 默认情况下,当整个MapReduce作业的所有已执行完成的Map Task任务数超过Map Task总数的5%后,JobTracker便会开始调度执行Reduce Task任务。然后Reduce Task任务默认启动
mapred.reduce.parallel.copies(默认为5)个MapOutputCopier线程到已完成的Map Task任务节点上分别copy一份属于自己的数据。 这些copy的数据会首先保存的内存缓冲区中,当内冲缓冲区的使用率达到一定阀值后,则写到磁盘上。
内存缓冲区
- 这个内存缓冲区大小的控制就不像map那样可以通过
io.sort.mb来设定了,而是通过另外一个参数来设置:mapred.job.shuffle.input.buffer.percent(default 0.7), 这个参数其实是一个百分比,意思是说,shuffile在reduce内存中的数据最多使用内存量为:0.7 ×maxHeap of reduce task。 - 如果该reduce task的最大heap使用量(通常通过
mapred.child.java.opts来设置,比如设置为-Xmx1024m)的一定比例用来缓存数据。默认情况下,reduce会使用其heapsize的70%来在内存中缓存数据。如果reduce的heap由于业务原因调整的比较大,相应的缓存大小也会变大,这也是为什么reduce用来做缓存的参数是一个百分比,而不是一个固定的值了。
merge过程
- Copy过来的数据会先放入内存缓冲区中,这里的缓冲区大小要比 map 端的更为灵活,它基于 JVM 的
heap size设置,因为 Shuffle 阶段 Reducer 不运行,所以应该把绝大部分的内存都给 Shuffle 用。 - 这里需要强调的是,merge 有三种形式:1)内存到内存 2)内存到磁盘 3)磁盘到磁盘。默认情况下第一种形式是不启用的。当内存中的数据量到达一定阈值,就启动内存到磁盘的 merge(图中的第一个merge,之所以进行merge是因为reduce端在从多个map端copy数据的时候,并没有进行sort,只是把它们加载到内存,当达到阈值写入磁盘时,需要进行merge) 。这和map端的很类似,这实际上就是溢写的过程,在这个过程中如果你设置有Combiner,它也是会启用的,然后在磁盘中生成了众多的溢写文件,这种merge方式一直在运行,直到没有 map 端的数据时才结束,然后才会启动第三种磁盘到磁盘的 merge (图中的第二个merge)方式生成最终的那个文件。
- 在远程copy数据的同时,Reduce Task在后台启动了两个后台线程对内存和磁盘上的数据文件做合并操作,以防止内存使用过多或磁盘生的文件过多。
reducer的输入文件
- merge的最后会生成一个文件,大多数情况下存在于磁盘中,但是需要将其放入内存中。当reducer 输入文件已定,整个 Shuffle 阶段才算结束。然后就是 Reducer 执行,把结果放到 HDFS 上。
在写 MR 时,什么情况下可以使用规约
(1)Map端溢出spill文件后,讲行分区排序后,会执行combiner 函数。
(2)多个spill文件合并成一个大文件时,会使用。
(3)在Reduce任务端,复制map输出,在合弁后溢出到磁盘时,也会执行combiner函数。
reduce函数如何知道从哪台机器获取map输出?
Map和Reduce中间有个MRAppMaster负责协调。MRAppMaster负责任务调度和监控,map结束会通知MRAppMaster,Reduce也会定期去询问MRAppMaster以便获取 map输出的位置。
MapReduce执行速度太慢,怎么办?
1、自定义分区函数,让key值较为均匀的分布在 Reduce 上。
2、对map端进行压缩处理。
3、使用combiner函数
yarn 集群的架构和工作原理
Apache YARN (Yet Another Resource Negotiator) 是 hadoop 2.0 引入的集群资源管理系统。用户可以将各种服务框架部署在 YARN 上,由 YARN 进行统一地管理和资源分配。
组成:Client、ResourceManager、NodeManager、ApplicationMaster、Container
- Client
向RM提交任务
杀死任务 - ResourceManager
👉ResourceManager通常在独立的机器上以后台进程的形式运行,它是整个 集群资源的主要协调者和管理者 。
👉 负责给用户提交的所有应用程序分配资源 ,(调度器:Scheduler)它根据应用程序优先级、队列容量、ACLs、数据位置等信息,做出决策,然后以共享的、安全的、多租户的方式制定分配策略,调度集群资源。 - NodeManager
👉NodeManager是 YARN 集群中的每个具体 节点的管理者 。
👉 主要 负责该节点内所有容器的生命周期的管理,监视资源和跟踪节点健康 。具体如下:- 启动时向
ResourceManager注册并定时发送心跳消息,等待ResourceManager的指令; - 维护
Container的生命周期,监控Container的资源使用情况; - 管理任务运行时的相关依赖,根据
ApplicationMaster的需要,在启动Container之前将需要的程序及其依赖拷贝到本地。
- 启动时向
- ApplicationMaster
👉 在用户提交一个应用程序时,YARN 会启动一个轻量级的 进程ApplicationMaster。
👉ApplicationMaster负责协调来自ResourceManager的资源,并通过NodeManager监视容器内资源的使用情况,同时还负责任务的监控与容错。具体如下:
- 根据应用的运行状态来决定动态计算资源需求;
- 向
ResourceManager申请资源,监控申请的资源的使用情况; - 跟踪任务状态和进度,报告资源的使用情况和应用的进度信息;
- 负责任务的容错。
- Container
👉Container是 YARN 中的 资源抽象 ,它封装了某个节点上的多维度资源,如内存、CPU、磁盘、网络等。
👉 当 AM 向 RM 申请资源时,RM 为 AM 返回的资源是用Container表示的。
👉 YARN 会为每个任务分配一个Container,该任务只能使用该Container中描述的资源。ApplicationMaster可在Container内运行任何类型的任务。例如,MapReduce ApplicationMaster请求一个容器来启动 map 或 reduce 任务
yarn 的任务提交流程是怎样的

- 客户端
client向yarn集群提交作业 , 首先①向ResourceManager申请分配资源 Resource Manager会为作业分配一个Container(Application manager),Container里面运行这(Application Manager)Resource Manager会找一个对应的NodeManager通信②,要求NodeManager在这个container上启动应用程序Application Master③Application Master向Resource Manager申请资源④(采用轮询的方式通过RPC协议),Resource scheduler将资源封装发给Application master④,Application Master将获取到的资源分配给各个Node Manager,并监控运行情况⑤Node Manage得到任务和资源开始执行作业⑥- 再细分作业的话可以分为 先执行
Map Task,结束后在执行Reduce Task最后再将结果返回給Application Master等依次往上层递交⑦
yarn 的资源调度三种模型了解吗
在 Yarn 中有三种调度器可以选择:
FIFO Scheduler
Capacity Scheduler
Fair Scheduler
apache 版本的 hadoop 默认使用的是 capacity scheduler 调度方式。
CDH 版本的默认使用的是 fair scheduler 调度方式
FIFO Scheduler(先来先服务):
FIFO Scheduler 把应用按提交的顺序排成一个队列,这是一个先进先出队列,在进行资源分配的时候,先给队列中最头上的应用进行分配资源,待最头上的应用需求满足后再给下一个分配,以此类推。
FIFO Scheduler 是最简单也是最容易理解的调度器,也不需要任何配置,但它并不适用于共享集群。大的应用可能会占用所有集群资源,这就导致其它应用被阻塞,比如有个大任务在执行,占用了全部的资源,再提交一个小任务,则此小任务会一直被阻塞。
Capacity Scheduler(能力调度器):
对于 Capacity 调度器,有一个专门的队列用来运行小任务,但是为小任务专门设置一个队列会预先占用一定的集群资源,这就导致大任务的执行时间会落后于使用 FIFO 调度器时的时间。
Fair Scheduler(公平调度器):
在 Fair 调度器中,我们不需要预先占用一定的系统资源,Fair 调度器会为所有运行的 job 动态的调整系统资源。
比如:当第一个大 job 提交时,只有这一个 job 在运行,此时它获得了所有集群资源;当第二个小任务提交后,Fair 调度器会分配一半资源给这个小任务,让这两个任务公平的共享集群资源。
需要注意的是,在 Fair 调度器中,从第二个任务提交到获得资源会有一定的延迟,因为它需要等待第一个任务释放占用的 Container。小任务执行完成之后也会释放自己占用的资源,大任务又获得了全部的系统资源。最终的效果就是 Fair 调度器即得到了高的资源利用率又能保证小任务及时完成。
引用
https://hiszm.blog.csdn.net/article/details/109060789
https://blog.csdn.net/xh16319/article/details/31375197
https://cshihong.github.io/2018/05/11/MapReduce%E6%8A%80%E6%9C%AF%E5%8E%9F%E7%90%86/
https://andr-robot.github.io/Hadoop%E4%B8%ADMapReduce%E6%89%A7%E8%A1%8C%E6%B5%81%E7%A8%8B%E8%AF%A6%E8%A7%A3/
HDFS 在读取文件的时候, 如果其中一个块突然损坏了怎么办
客户端读取完 DataNode 上的块之后会进行 checksum (总和检验码)验证,也就是把客户端读取到本地的块与 HDFS 上的原始块进行校验,如果发现校验结果不一致,客户端会通知 NameNode,然后再从下一个拥有该 block 副本的 DataNode 继续读
------------------------------------------------
HDFS 在上传文件的时候, 如果其中一个 DataNode 突然挂掉了怎么办
客户端上传文件时与 DataNode 建立 pipeline 管道,
- 管道正向是客户端向 DataNode 发送的数据包,package
- 管道反向是 DataNode 向客户端发送 ack 确认,也就是正确接收到数据包之后发送一个已确认接收到的应答,当 DataNode 突然挂掉了,
客户端接收不到这个 DataNode 发送的 ack 确认,客户端会通知 NameNode,NameNode检查该块的副本与规定的不符,NameNode 会通知 DataNode 去复制副本,并将挂掉的 DataNode 作下线处理,不再让它参与文件上传与下载。
SecondaryNameNode 了解吗,它的工作机制是怎样的
只有在NameNode重启时,edit logs才会合并到fsimage文件中,从而得到一个文件系统的最新快照。但是在产品集群中NameNode是很少重启的,这也意味着当NameNode运行了很长时间后,edit logs文件会变得很大。在这种情况下就会出现下面一些问题:
- edit logs文件会变的很大,怎么去管理这个文件是一个挑战。
- NameNode的重启会花费很长时间,因为有很多改动在edit logs中要合并到fsimage文件上。
- 如果NameNode挂掉了,那我们就丢失了很多改动因为此时的fsimage文件非常旧。在这个情况下丢失的改动不会很多, 因为丢失的改动应该是还在内存中但是没有写到edit logs的这部分。
因此为了克服这个问题,我们需要一个易于管理的机制来帮助我们减小edit logs文件的大小和得到一个最新的fsimage文件,这样也会减小在NameNode上的压力。这跟Windows的恢复点是非常像的,Windows的恢复点机制允许我们对OS进行快照,这样当系统发生问题时,我们能够回滚到最新的一次恢复点上。
Secondary NameNode 是合并 NameNode 的 edit logs 到 fsimage 文件中;

上面的图片展示了Secondary NameNode是怎么工作的。
- 首先,请求执行定时
checkpoint,它定时到NameNode去获取edit logs,并更新到fsimage.chkpoint上(Secondary NameNode自己的fsimage) - 一旦它有了新的
fsimage.chkpoint文件,它将其拷贝回NameNode中,并且重命名为fsimage。 - NameNode在下次重启时会使用这个新的fsimage文件,从而减少重启的时间。
Secondary NameNode的整个目的是在HDFS中提供一个检查点。它只是NameNode的一个助手节点。这也是它在社区内被认为是检查点节点的原因。
现在,我们明白了Secondary NameNode所做的不过是在文件系统中设置一个检查点来帮助NameNode更好的工作。
它不是要取代掉NameNode也不是NameNode的备份。
所以从现在起,让我们养成一个习惯,称呼它为检查点节点吧。
所以如果 NameNode 中的元数据丢失,是可以从 Secondary NameNode 恢复一部分元数据信息的,但不是全部,因为 NameNode 正在写的 edits 日志还没有拷贝到 Secondary NameNode,这部分恢复不了
Secondary NameNode 不能恢复 NameNode 的全部数据,那如何保证 NameNode 数据存储安全
Import Checkpoint(恢复数据)
如果主节点namenode挂掉了,硬盘数据需要时间恢复或者不能恢复了,现在又想立刻恢复HDFS,这个时候就可以import checkpoint。步骤如下:
- 准备原来机器一样的机器,包括配置和文件
- 创建一个空的文件夹,该文件夹就是配置文件中dfs.name.dir所指向的文件夹。
- 拷贝你的secondary NameNode checkpoint出来的文件,到某个文件夹,该文件夹为fs.checkpoint.dir指向的文件夹(例如:/home/hadadm/clusterdir/tmp/dfs/namesecondary)
- 执行命令bin/hadoop namenode –importCheckpoint
- 这样NameNode会读取checkpoint文件,保存到dfs.name.dir。但是如果你的dfs.name.dir包含合法的 fsimage,是会执行失败的。因为NameNode会检查fs.checkpoint.dir目录下镜像的一致性,但是不会去改动它。
一般建议给maste配置多台机器,让namesecondary与namenode不在同一台机器上值得推荐的是,你要注意备份你的dfs.name.dir和 ${hadoop.tmp.dir}/dfs/namesecondary
这个问题就要说 NameNode 的高可用了,即 NameNode HA
一个 NameNode 有单点故障的问题,那就配置双 NameNode,配置有两个关键点,一是必须要保证这两个 NN 的元数据信息必须要同步的,二是一个 NN 挂掉之后另一个要立马补上。
-
元数据信息同步在 HA 方案中采用的是 “共享存储”。每次写文件时,需要将日志同步写入共享存储,这个步骤成功才能认定写文件成功。然后备份节点定期从共享存储同步日志,以便进行主备切换。
-
监控 NN 状态采用 zookeeper,两个 NN 节点的状态存放在 ZK 中,另外两个 NN 节点分别有一个进程监控程序,实施读取 ZK 中有 NN 的状态,来判断当前的 NN 是不是已经 down 机。如果 standby 的 NN 节点的 ZKFC 发现主节点已经挂掉,那么就会强制给原本的 active NN 节点发送强制关闭请求,之后将备用的 NN 设置为 active。
可以进行解释下:NameNode 共享存储方案有很多,比如 Linux HA, VMware FT, QJM 等,目前社区已经把由 Clouderea 公司实现的基于 QJM(Quorum Journal Manager)的方案合并到 HDFS 的 trunk 之中并且作为默认的共享存储实现
基于 QJM 的共享存储系统主要用于保存 EditLog,并不保存 FSImage 文件。FSImage 文件还是在 NameNode 的本地磁盘上。QJM 共享存储的基本思想来自于 Paxos 算法,采用多个称为 JournalNode 的节点组成的 JournalNode 集群来存储 EditLog。每个 JournalNode 保存同样的 EditLog 副本。每次 NameNode 写 EditLog 的时候,除了向本地磁盘写入 EditLog 之外,也会并行地向 JournalNode 集群之中的每一个 JournalNode 发送写请求,只要大多数 (majority) 的 JournalNode 节点返回成功就认为向 JournalNode 集群写入 EditLog 成功。如果有 2N+1 台 JournalNode,那么根据大多数的原则,最多可以容忍有 N 台 JournalNode 节点挂掉
8. 小文件过多会有什么危害, 如何避免
Hadoop 上大量 HDFS 元数据信息存储在 NameNode 内存中, 因此过多的小文件必定会压垮 NameNode 的内存
每个元数据对象约占 150byte,所以如果有 1 千万个小文件,每个文件占用一个 block,则 NameNode 大约需要 2G 空间。如果存储 1 亿个文件,则 NameNode 需要 20G 空间
显而易见的解决这个问题的方法就是合并小文件, 可以选择在客户端上传时执行一定的策略先合并, 或者是使用 Hadoop 的 CombineFileInputFormat<K,V> 实现小文件的合并
-
文件是由许许多多的records组成的,那么可以通过调用HDFS的sync()方法(和append方法结合使用)来解 决。或者,可以通过些一个程序来专门合并这些小文件(see Nathan Marz’s post about a tool called the Consolidator which does exactly this).
-
就需要某种形式的容器来通过某种方式来group这些file。
- HAR files:一个HAR文件是通过hadoop的archive命令来创建,而这个命令实 际上也是 运行了一个MapReduce任务来将小文件打包成HAR.对大文件里面的文件在做一个索引。 这样通过一个二级索引,HAR 就可以避免 保存太多的小文件

- Sequence Files:使用filename作为key,并且file contents作为value。

以上三种方法虽然能够解决小文件的问题,但是这些方法都有局限:**HAR Files 和 Sequence Files 一旦创建,之后都不支持修改,**所以这是对读场景很友好的;
使用 HBase 需要引入外部系统,维护成本很高。
HFS的出现解决了需要在HDFS中存储海量小文件,同时也要存储一些大文件的混合的场景。简单来说,就是在HBase表中,需要存放大量的小文件(10MB以下),同时又需要存放一些比较大的文件(10MB以上。)
HBase MOB:
MOB数据(即100KB到10MB大小的数据)直接以HFile的格式存储在文件系统上(例如HDFS文件系统),然后把这个文件的地址信息及大小信息作为value存储在普通HBase的store上,通过攻击集中管理这些文件。这样就可以大大降低HBase的compation和split频率,提升性能。
13. shuffle 阶段的数据压缩机制了解吗
在 shuffle 阶段,可以看到数据通过大量的拷贝,从 map 阶段输出的数据,都要通过网络拷贝,发送到 reduce 阶段,这一过程中,涉及到大量的网络 IO,如果数据能够进行压缩,那么数据的发送量就会少得多。
hadoop 当中支持的压缩算法:
gzip、bzip2、LZO、LZ4、Snappy,这几种压缩算法综合压缩和解压缩的速率,谷歌的 Snappy 是最优的,一般都选择 Snappy 压缩。
谷歌出品,必属精品
更多推荐

所有评论(0)