跳转至

Hadoop

Hadoop是什么

Hadoop是Apache的开源的分布式系统基础架构,主要解决海量数据存储与计算的问题,主要包括HDFS分布式文件存储系统、MapReduce分布式计算框架和Yarn集群资源管理和任务调度框架

Hadoop的主要角色

  • NameNode: Hadoop中的主服务器, 管理文件系统的NameSpace命名空间, 对集群中存储的文件访问, 保存元数据,主要负责存储数据的元数据信息,不存储实际的数据块
  • SecondaryNameNode: 周期性的将NameNode镜像(fsImage)与操作日志(edit Log)合并, 防止edit Log过大, 并对元数据备份
  • DataNode:文件系统的工作节点, 根据客户端或者是namenode的调度存储和检索数据,并且定期向namenode发送他们所存储的块列表,存储实际的数据块

上面三个是HDFS的组件

  • ResourceManager: 负责集群中所有资源的统一管理和分配
  • NodeManager:管理集群中的单个节点, 以心跳的方式向ResourceManager汇报资源的使用情况(主要是cpu和内存), RM只接受NM的资源回报信息,不干涉其处理

HDFS是怎么管理元数据的

  • 元数据:文件的描述信息
  • 内存元数据:基于内存存储元数据,可以避免降低频繁的读写操作效率,元数据保存的比较完整
  • FsImage文件:磁盘元数据镜像文件,在NameNode工作目录中,NameNode作为全局的单点信息管理中心,但它不包含block所在的Datanode 信息
    • 块信息是启动时,由datanode向namenode汇报获得
  • edits文件:数据操作的日志文件,用于衔接内存元数据和fsimage之间的操作日志,可通过日志还原出元数据

元数据管理流程

  1. 首先,namenode 将元数据信息存储在内存中
  2. 然后,将元数据的所有更新操作不断顺序写入edits 文件
  3. 之后,为避免edits文件过大,导致重启的时候恢复元数据卡顿,secondary namenode在namenode重启的时候会首先将edits文件合并写入到磁盘的FsImage文件中

总结:

hdfs的元数据管理主要是通过namenode来进行, 他通过内存,fsimage,edits三者对元数据进行管理, 将元数据存放到内存中, 提高访问速度, 当有操作进来时,会把这条操作信息写入edits中保存, 通过snn利用edits去更新fsimage文件, 这样进行元数据的管理(先写入editlog在操作)

Secondary NameNode了解吗,它的工作机制是怎样的

Secondary Namenode 主要是用于edit logs 和 fsImage 的合并,edit logs 记录了 对namenode 元数据的增删改操作,fsimage 记录了最新的元数据检查点【不止有检查点,还包含了检查点之前的元数据】,Secondary Namenode在 namenode 重启的时候,会把edit logs 和fsimage进行合并,形成新的fsimage 文件。

  • 工作机制

  • Secondary NameNode 询问 NameNode 是否需要 checkpoint,并直接带回 NameNode的结果

  • Secondary NameNode 请求执行 checkpoint
  • NameNode 滚动正在写的 edits 日志
  • 将滚动前的编辑日志和镜像文件拷贝到 Secondary NameNode
  • Secondary NameNode 加载编辑日志edits和镜像文件fsimage到内存并合并
  • 生成新的镜像文件fsimage.chkpoint
  • 拷贝fsimage.chkpoint到 NameNode
  • NameNode 将 fsimage.chkpoint 重命名fsimage

所以如果 NameNode 中的元数据丢失,是可以从 Secondary NameNode 恢复一部分元数据信息的,但不是全部,因为 NameNode 正在写的 edits 日志还没有拷贝到 Secondary NameNode,这部分恢复不了

HDFS新加入节点需要做什么

  1. 修改namenode的hdfs-site.xml配置,增加dfs.hosts属性
  2. 刷新namenode
  3. 刷新resourceManager
  4. 编辑slaves文件,并添加新增节点的主机,更改完后,slaves文件不需要分发到其他机器上面去
  5. 启动datanode和nodemanager
  6. 对于新节点
  7. 修改mac地址以及IP地址
  8. 关闭防火墙,关闭selinux
  9. 配置ssh免密码登录
  10. 更改主机名,更改主机名与IP地址的映射
  11. 安装jdk
  12. 安装Hadoop

HDFS在数仓里的地位

HDFS是数仓底层用于数据存储的基础设施,它提供了一个可靠的且可拓展的存储平台。 数仓中的数据往往是结构化、半结构化或非结构化的,这些数据可以存储在 HDFS 中,以便后续的处理和分析。 HDFS与 Hadoop 生态系统中的许多工具和框架兼容,如 hive、hbase、spark等

namenode存在什么问题?怎么解决

  • namenode单点故障问题 在非HA架构中,NameNode是单点的,一旦NameNode故障,整个HDFS将不可用。 因此,考虑引入多个NameNode,设为主备NameNode(Active和Standby),高可用同时具有多个NN, 但同一时刻只有一个NN是active, 对外提供服务 ,其余都是standby状态,这些NN通过zookeeper进行管理。当active的NN挂了之后, 节点内的zkfc进程检测到会通知其他的NN的zkfc, 进行选举,等备用的NN选举成功后,会将挂掉的NN删除,同时激活自身的Active状态。高可用机制依赖于共享存储(如JournalNode)和故障切换机制,确保NameNode的高可用性。

  • namenode性能瓶颈问题 在早期HDFS架构中,单个NameNode负责管理整个文件系统的元数据,随着数据量增长,NameNode的内存和CPU成为瓶颈。 因为namenode将目录树存在内存中,而内存是有限的,一个namenode管理的文件是有限的, 但datanode可以通过添加节点认为是无限的, 所以就会出现能力不匹配。 于是提出引入多个独立的NameNode,每个NameNode管理文件系统的一部分命名空间(Namespace),这些Namenode之间相互独立,各自分工管理自己的命名空间,从而实现水平扩展,这就是HDFS的联邦(Federation)机制。 HDFS集群中的Datanode提供数据块的共享存储功能,每个Datanode都会向集群中所有的Namenode注册,且周期性地向所有的Namenode发送心跳和块汇报,然后执行Namenode通过响应发回的Namenode指令。

联邦和HA可以结合使用

hdfs中namenode如果挂掉了怎么办

  1. 当namenode发生故障宕机时, snn会保存所有的元数据信息, 在namenode重启时会将元数据发回给namenode
  2. 如果配置了HA, 当activate namenode宕机, standby namenode马上切换到active状态,消除单点故障问题
  3. 怎么知道发生宕机
  4. 利用zk, 心跳机制

HDFS的数据读写流程

  • 写数据hadoop fs -put a.tat /user/sl/

  • 首先客户端向namenode进行请求,然后namenode会检查该文件是否已经存在,如果不存在,就会允许客户端上传文件;

  • 然后客户端再次向namenode请求第一个block上传到哪几个datanode节点上。假设namenode返回了三个datanode节点
  • 那么客户端就会向datanode1请求上传数据,然后datanode1会继续调用datanode2,datanode2会继续调用datanode3,那么这个通信管道就建立起来了,紧接着dn3,dn2,dn1逐级应答客户端
  • 然后客户端就会向datanode1上传第一个block,以packet为单位(默认64k),datanode1收到后就会传给datanode2,datanode2传给datanode3
  • 当第一个block传输完成之后,客户端再次请求namenode上传第二个block。【串行写入数据块】

  • 读数据hadoop fs -get a.txt /opt/module/hadoop/data/

  • 首先客户端向namenode进行请求,然后namenode会检查文件是否存在,如果存在,就会返回该文件所在的datanode地址

  • 返回的datanode 地址会按照集群拓扑结构得出 datanode与客户端的距离并进行排序
  • 然后客户端会选择排序靠前的datanode来读取block,客户端会以packet为单位进行接收,先在本地进行缓存,然后写入目标文件中。【并行读取数据块】

mapredeuce的过程(详细)

  1. map阶段
  2. 首先通过InputFormat把输入目录下的文件进行逻辑切片,切片默认大小为一个block的大小,并且每一个切片对应一个Map任务,由一个maptask来处理
  3. maptask再将切片中的数据解析成<k,v>的键值对,k表示偏移量,v表示一行内容
  4. 然后调用Mapper类中的map方法对每一行内容进行处理,解析为<k,v>的键值对。在wordCount案例中,k表示单词,v表示数字1
  5. shuffle阶段
  6. map端shuffle
    1. 将map后的<k,v>写入环形缓冲区【默认大小为100M】,一半写元数据信息(key的起始位置,value的起始位置,value的长度,partition号),一半写<k,v>数据
    2. 等缓冲区写入到达阈值【默认是80%】(或map任务结束)的时候,就要进行spill溢写操作
    3. 溢写之前数据会通过 Partitioner按key进行分区,决定数据发送到哪个Reducer
      1. 默认的分区算法根据key的hashcode对reduce task的个数取模得到分区号
    4. 然后数据先按照分区号进行排序,再按照key进行排序
    5. 溢写时,如果设置了Combiner则会先在分区将数据进行一次聚合操作,再按照分区号从小到大依次写入磁盘临时文件,
      1. (可选)使用Combiner进行预聚合,将有相同Key的Value 合并起来, 减少溢写到磁盘的数据量【只能在累加、最大值使用,不能在求平均值的时候使用】
    6. 然后将每次溢写生成的临时文件merge成一个大文件【保证一个maptask只对应一个文件】同时生成对应的索引文件【包含分区号等信息】
  7. reduce端shuffle:merge
    1. copy:reduce会从完成任务的maptask节点上将同一分区的各个maptask的结果拉取到内存中,如果放不下,就会溢写到磁盘上
    2. merge:然后对内存和磁盘上的数据进行merge操作【可以保证是相同key的数据】
      1. Merge有3种形式,分别是内存到内存,内存到磁盘,磁盘到磁盘。默认情况下第一种形式不启用,第二种Merge方式一直在运行(spill阶段)直到结束,然后启用第三种磁盘到磁盘的Merge方式生成最终的文件
    3. sort:对merge的数据进行排序
  8. reduce阶段
  9. Reduce 任务通过索引文件快速定位到对应分区的偏移量范围或直接从 part-r-xxxxx 文件中读取指定分区的数据块
  10. key相同的数据会调用一次reduce方法,每次调用产生一个键值对,最后将这些键值对写入到HDFS文件中

mapredeuce的过程(简略)

  1. map阶段
  2. 先对输入文件切片
  3. 分别调用map方法
  4. 输出KV键值对
  5. shuffle阶段
  6. map端shuffle
    • 将map输出写入到默认100M的环形缓冲区,一半写元数据,另一半写实际数据
    • 到80%的时候开始溢写到磁盘
    • 溢写前先 对key按照分区进行快速排序,然后才溢写到文件中
    • 溢写到文件中后,进行 merge 归并排序
  7. reduce端shuffle
    • reduce 会拉取同一分区各个 mapTask 的结果到内存中,放不下就溢写到磁盘
    • 然后对内存和磁盘上的数据进行 merge 归并排序
  8. reduce阶段
  9. key 相同的数据会调用一次 reduce 方法,每次调用产生一个键值对
  10. 最后将这些键值对写入到 HDFS 文件

MapReduce中WordCount的键值对变化

  1. map输入切片阶段:
  2. Key:行首字节偏移量(LongWritable)。
  3. Value:行内容(Text
  4. map切片解析阶段:
  5. Key:单词(Text)。
  6. Value:计数(IntWritable)。
  7. shuffle阶段:
  8. Key:单词(Text)。
  9. Value:对map得到的键值对累计计数(IntWritable
  10. reduce输入阶段:
  11. Key:单词(Text)。
  12. Value:值的迭代器(Iterable<IntWritable>),包含所有 Map 任务对该单词的计数。
  13. reduce处理输出阶段:
  14. Key:单词(Text)。
  15. Value:总计数(IntWritable)。

hadoop中的序列化和反序列化,为什么不用java的

  • 将数据结构转换为二进制格式在网络中传输
  • java的序列化太过重量级, 有很多附带信息, 不便在网络上高效传输
  • hadoop的序列化更加轻量化

MapTask的个数和ReduceTask的个数由什么决定

  • MapTask和block的块数有关, 有几个切片就有几个maptask
  • reduceTask的并行度与分区个数有关
  • 自定义分区可以改变reduce个数

map端为什么要排序?

因为reduce阶段需要分组,目的将key相同的放在一起进行规约,而如果全部都在reduce阶段进行sort排序(内部排序)就太消耗内存,恰好因为map阶段的输出是溢写到磁盘,那么理论上只要磁盘够大,在磁盘中进行排序可以对任意数据量分组,将相同key的数据组织在一起建立索引,reduce拉取的时候可以连续访问,提高效率,这就是为了通过外部排序降低内存的使用量。即map端排序(shuffle阶段)可以减轻reduce端排序的压力。

map端排序使用了两种算法:hashmap和sort

map端输出的文件组织形式是什么样的?

简而言之就是多个溢写文件合并后的大文件

具体形式如下:

在本地磁盘生成一个输出目录(如 _temporary/task_2023.../output),包含以下文件:

  1. 数据文件:part-r-xxxxx(如 part-r-00000
part-r-00000:
  [分区0数据] apple 1, banana 1, apple 1, ...  # 分区内按键排序
  [分区1数据] hadoop 1, mapreduce 1, hadoop 1, ...
  [分区2数据] world 1, count 1, world 1, ...
  1. 索引文件:part-r-xxxxx.index
part-r-00000.index:
  partition=0, start=0, end=1024
  partition=1, start=1025, end=2048
  partition=2, start=2049, end=3072
  1. 元数据文件(可选)

(2)reduce怎么知道从哪里下载map输出的文件

当Map任务执行结束,会向JobTracker(Hadoop 1.x)或ApplicationMaster(YARN/Hadoop 2.x+)报告输出文件的位置信息(如文件路径、大小、所在节点等),reduce中的一个线程会定期询问ApplicationMaster以便获取map输出的位置,JobTracker/ApplicationMaster会为每个Reduce任务分配需处理的分区,并告知其对应的Map输出文件位置(即哪些TaskTracker节点上有它所需的数据)

※join原理

  • Reduce端Join(传统MapReduce Join):通过MapReduce的Shuffle机制将关联数据汇聚到同一Reduce任务中完成关联计算,但是由于Reduce端需处理全量数据,易导致数据倾斜(如热点Key)。
  • map 阶段:对来自不同表的数据打标签,然后用连接字段作为key,其余部分和标签作为value,最后进行输出
  • shuffle 阶段:根据key的值进行hash,这样就可以将key相同的送入一 个reduce 中
  • reduce 阶段:同一个key的数据会调用一次reduce方法,就 是对来自不同表的数据进行join(笛卡尔积)
  • Map端Join(Map Join):将小表全量加载到Map任务的内存中,直接在Map阶段完成关联计算,避免Shuffle和Reduce。适合一张表十分小、一张表很大的场景,好处是增加Map端业务,减少Reduce端数据的压力,尽可能的减少数据倾斜
  • 原理:将小表复制多份,让每个map task内存中存在一份(比如存放到 HashMap 中),然后只扫描大表。对于大表中的每一条记录 key/value, 在HashMap中查找是否有相同的key的记录,如果有,则join连接后输出即可
  • 具体办法:
    • 使用DistributedCache 缓存文件到Task运行节点
    • 在mapper的setup方法中,将文件读取到缓存集合中

(2)如果map输出太多小文件该如何调优

  1. 输入端:使用combineinputformat合并小文件
  2. map端:提高环形缓冲区的大小、减少IO次数、开启Combiner
  3. 调整并行度:合理设置Reduce任务数,避免过多小文件或数据倾斜

为什么要有环形缓冲区

  • 环形缓冲区不需要重新申请新的内存且内存是连续的, 没有内存碎片, 始终用的都是这个内存, 规避了Full GC导致的问题, 不用频繁申请内存
  • 另外呢,环形缓冲区同时做了两件事情:1、排序;2、索引。在这里一次排序,将无序的数据变为有序,写磁盘的时候顺序写,读数据的时候顺序读,效率高非常多!

(2)环形缓冲区及其阈值高低的影响

环形缓冲区底层就是一个数组,默认大小是100M。数组中存放着<key,value>数据,以及关于<key,value>的元数据信息,每个<key,value>对应一个元数据。元数据由4个int组成,第一个int存放value的起始位置,第二个int存放key的起始位置,第三个int存放partition,第四个int存放value的长度。

<key,value>数据和元数据在环形缓冲区中的存储是由equator(赤道)分隔的,<``key,value>按照索引递增的方向存储,元数据则按照索引递减的方向存储。将数组抽象为一个环形结构之后,以equator为界,<key,value>数据顺时针存储,元数据逆时针存储。

环形缓冲区有阈值的目的是

  1. 避免内存溢出导致程序崩溃
  2. 减少频繁磁盘I/O
  3. 防止阻塞生产者,不需要等消费组处理,生产者可以继续发送

  4. 高阈值的影响

  5. 优点:
    • 减少溢写次数:缓冲区更满时才溢写,降低磁盘 I/O 频率。
    • 提高内存利用率:充分利用内存暂存数据,减少合并(Merge)阶段的压力。
  6. 缺点:
    • 内存压力风险:若数据生成速度突增,缓冲区可能快速填满,导致频繁溢写或内存溢出(OOM)。
    • 合并效率下降:单次溢写的文件较大,合并阶段需要更多内存和时间。
  7. 低阈值的影响
  8. 优点:
    • 降低内存压力:缓冲区更快溢写,减少内存占用峰值,避免 OOM。
    • 提高合并效率:溢写文件较小,合并阶段处理更灵活(如并行合并)。
  9. 缺点:
    • 增加溢写次数:频繁的小文件溢写导致更多磁盘 I/O 操作。
    • 降低内存利用率:缓冲区未满即溢写,可能浪费内存资源。

哪个阶段最费时间,环形缓冲区的调优以及什么时候需要调

shuffle阶段的排序和溢写磁盘 时间花费最大

参数 默认值 调优建议
mapreduce.task.io.sort.mb 100MB 根据Map输出数据量调整,如设为200-500MB(需监控内存使用)。
mapreduce.map.sort.spill.percent 0.8 若溢写频繁,可适当降低(如0.7),减少单次溢写数据量。
mapreduce.shuffle.file.buffer.percent 0.1 控制合并文件时使用的内存比例,避免合并阶段OOM。
  • 何时需要调优环形缓冲区?
  • 原则上说,缓冲区越大,磁盘 io 的次数越少,执行速度就越快
  • 触发条件:
    • Map任务日志中出现频繁溢写(如Spill 0Spill 1等)。
    • 任务执行时间中Shuffle阶段占比过高(通过JobHistory分析)。
    • 集群内存资源充足,但CPU或磁盘利用率低。
  • 调优目标:
    • 减少溢写次数,降低磁盘IO压力。
    • 平衡内存使用与GC开销。

副本存储策略

  • 第一个放在client所在的节点上, client不在集群上则随机选择一个
  • 第二个在另一个机架的节点上
  • 第三个在第二个副本所在机架的随机节点上

HDFS底层存储原理

  • hdfs将要存储的大文件进行切割(客户端进行), 存放在datanode中, 元数据存放在namenode
  • 文件有多个副本
  • HDFS可以做冷热存储分离,即固态存储热数据, 机械存储冷数据

hdfs如何保证数据不丢失

  • 副本机制:默认有三个副本,数据块存在不同的datanode上
  • 心跳机制:datanode会定时发送心跳给namenode, 超过10分钟 + 30秒会被认为该datanode不可用, namenode会安排其他的datanode复制不可用的datanode上的副本 , 保证副本数量是设置的数量
  • 安全模式:每个DN会定时汇报自己所有块的信息(一小时), NameNode会计算损坏率, 损坏率低于99%, 则进入安全模式(安全模式下,客户端不能读或修改数据块的信息), namenode会标记损坏的数据块为不可用,此时副本不足3个, 安排其他datanode去复制一个好的副本,复制完成,删除损坏的副本

hdfs如何保证数据一致性

  • 校验和: hdfs会为文件生成一个校验和, 校验和文件会和文件本身保存在同一空间中, 传输数据时datanode会将数据与校验数据进行检测, 检验结果不同, 则文件出错, 数据块就是无效, datanode不会发送, 客户端会从其他副本读取
  • 副本机制
  • 副本修复:副本数量不足,启动副本修复

hdfs如何保证数据高可用

  • 磁盘故障
  • 由于有副本机制, DateNode检测到磁盘出问题会向NameNode请求从其他节点复制故障所丢失的数据, 或者在汇报块信息namenode检测到故障, 或者在客户端读取的时候检测到故障,namenode都会..., 保证副本数量是设置的数量
  • DataNode故障容错
  • 心跳机制
  • NameNode故障容错
  • HA高可用, 部署两台NameNode
  • 机架感知
  • 客户端重试机制

写入数据流程中,一个DataNode挂掉了怎么办?

客户端上传文件时与DataNode建立pipeline管道, 正向发送packet, 反向发送ack, 当客户端发完收不到ack, 说明DataNode可能挂了, 会去通知namenode, namenode去检查该块的副本与规定的数量是否一致, 不符合则会通知DataNode去复制副本, 将挂掉的DataNode作下线处理, 不让它参加文件的上传与下载

namenode返回更新后的datanode列表给客户端, 客户端重建新pipeline, 继续上传未上传的文件

Hdfs小文件过多带来的危害以及解决办法

  • 危害
  • 存储大量的小文件,会占用namenode大量的内存来存储元数据信息
  • 在计算的时候,每个小文件需要一个maptask进行处理,浪费资源
  • 读取的时候,寻址时间超过读取时间 【文件以block为单位被写入hdfs,默认情况下一个block会被放在三台机器上。所以写入速度取决于内存,硬盘带宽以及网络带宽
  • 解决方法
  • 先对小文件进行合并之后再上传到hdfs
  • Hadoop Archive: 能够高效的将小文件放入hdfs中,并将多个小文件打包为一个HAR文件,从而减少namenode的内存使用
  • 在计算前,采用combineinputformat的切片方式,可以将多个小文件放到一个切片中进行计算
    1. map阶段读入 采用CombineTextInputFormat提高效率
    2. 若是map阶段输出太多小文件,则在shuffle阶段后写入磁盘前开启combiner进行预聚合减少小文件产出
  • SequenceFile二进制文件: key为文件名, value是文件内容 , 用于合并小文件, 不能追加写
  • 开启uber模式,实现JVM的重用,也就是说让多个task共用一个jvm, 这样就不必为每一个task开启一个jvm

Hadoop1.0 2.0 3.0 区别

  • hadoop1.0 和hadoop2.0 的区别

  • 新增了YARN框架,1.0的时候,MapReduce 既负责资源调度又负责计 算,到了2.0,资源调度就交给了yarn框架

  • 新增了HDFS 高可用机制(HA),通过配置Active和Standby两个NameNode,实现在集群中对NameNode的热备,解决了1.0存在的单点故障问题

  • hadoop2.0 和hadoop3.0 的区别

  • hadoop3.0 要求的最低java版本为jdk1.8

  • hadoop3.0 的时候支持HDFS的纠删码机制,作用就是节省存储空间【普通副本机制假设需要3倍存储空间而这种机制只需1.4倍即可】
  • hadoop3.0 的 MapReduce 进行了优化,性能提高了30%
  • hadoop3.0 支持两个以上的 namenode,也就是可以设置一个Active和多个Standby

yarn是什么,有哪些优缺点

yarn是hadoop2.x引入的新一代资源管理器, 用于管理hadoop集群的资源和作业的调度(从mapreduce抽取出来的功能)

  • 优点有:
  • 支持多种计算框架, 扩展了Hadoop的应用场景
  • 可以根据不同作业的资源需求动态分配资源, 提高了集群的利用率和灵活性
  • 允许多个作业同时运行, 提高了集群的并发度
  • YARN本身的资源占用很小, 对集群的负载影响较小
  • 缺点有:
  • yarn需要额外的配置和管理, 增加了部署和维护的复杂度
  • yarn对内存和cpu的要求较高, 可能需要更多的硬件资源支持

yarn的架构组件

  • ResourceManager: 全局资源管理器处理客户端请求; 监控NodeManager;
  • Nodemanager: 单节点上管理者: 定期汇报给RM资源情况,对Container进行管理
  • ApplicationMaster: 每个应用程序的管理者,由RM启动, 向RM申请资源, 任务监控
  • Container: 资源抽象单位,封装多维度资源,一个节点可以运行多个

yarn的任务提交流程

  1. 首先客户端Client提交任务到RM(Resource Manager)上,同时客户端会向RM申请一个application
  2. RM返回一个唯一的ApplicationID和HDFS临时路径,客户端将作业的JAR包、配置文件、输入分片等资源上传至该路径。
  3. 然后客户端就会提交任务运行需要的资源到对应路径上
  4. 客户端的资源提交完毕后,就会向RM申请Appmaster(Application Master)

MapReduce Application Master(mrAppmaster) 是 Application Master的一种具体实现

  1. RM 会将用户的请求初始化成一个Task,放入调度队列中
  2. 接着RM选择一个NM(NodeManager)(动态选择)领取task任务并且创建任务容器Container和启动Appmaster
  3. 然后Appmaster会向RM申请运行MapTask的资源,RM将MapTask分配给多个NM(分配NM),NM创建任务容器然后AppMaster发送程序启动脚本给分配的NM,分别启动MapTask

假设有两个数据切片,RM 就会将MapTask任务分配给两个NodeManager,这两个NodeManager分别领取MapTask任务并创建任务容器Container;Appmaster向这两个NM发送程序启动脚本,分别启动MapTask

  1. Appmaster等待所有MapTask运行完毕后,再次向RM申请容器, 发送启动脚本给NM来启动ReduceTask
  2. NM定时汇报给AM,AM监控任务执行情况,失败重试,客户端也会轮询执行情况
  3. 程序运行完毕后,Appmaster会向RM申请注销自己

yarn调度算法

  • 先进先出调度器(FIFO)
  • 容量调度器
  • 将集群资源划分为多个队列,每个队列配置固定的资源配额(如CPU、内存),队列内部采用FIFO策略空闲资源可临时借给其他队列,需归还时触发资源回收

    • 资源分配策略:

    • 队列选择:优先选择资源使用率最低的队列(深度优先算法)。

    • 作业选择:按提交时间和优先级分配资源。
    • 容器分配:优先满足数据本地性(同一节点 > 同一机架 > 跨机架)
    • 公平调度器
    • 动态平衡资源分配,确保所有作业在时间维度上公平共享资源。

常见的压缩算法

  • Snappy:无需安装,不支持切分,压缩后和文本处理一样
  • 适用:Mapper输出数据较大时,作为中间数据的压缩格式,或者作为一个MR到另一个MR的中间数据
  • LZO:需要安装,支持切分,压缩后需要建立索引
  • 适用:压缩后体积还超过块大小的,单个文件越大,优点越明显
  • GZip:无需安装,压缩率高,无需处理
  • 缺点:不支持切片
  • 适用:压缩后一个块大小内
  • Bzip2:压缩率极高,无需安装
  • 缺点:速度慢
  • 适用:对速度要求低,对压缩率要求高;或存储后使用较少

评论