MapReduce编程模型
起源:Google搜索。
编程模型抽象简洁,易于使用。运行于大规模分布式集群环境,自动实现分布式并行计算,具备容错和负载均衡,支持任务调度和状态监控。
|
|
map函数
将数据源中的记录(文本、数据库)作为map函数中的(key, value)对。
map()将生成一个或多个中间结果,以及一个与input相对应的output key。
reduce函数
map操作结束后,所有与某指定out key相对应的中间结果组合为一个列表(list)。reduce()函数将这些中间结果组合为一个或多个对应于同一
output key的final value(实际上每一个output key通常只有一个final value)。

MapRduce并行化运行
map()并行执行,不同的输入数据集生成不同的中间结果reduce()并行执行,分别处理不同的output keymap和reduce的处理过程中不发生通信
瓶颈:当
map处理全部结束后,reduce过程才能够开始
Hadoop/HDFS
Hadoop是一个分布式系统基础架构,由Apache基金会开发。用户可以在不了解分布式底层细节的情况下,开发分布式程序。充分利用集群的威力高速运算和存储。是对Google所提出的MapReduce、GFS、BigTable等模型的一个开源实现。

Map函数处理
当map开始产生输出时,并不是简单的把数据写到磁盘,因为频繁的磁盘操作会导致性能严重下降。它的处理过程更复杂,数据首先是写到本地内存中的一个缓冲区,并进行预排序,以提升效率。
每个map任务都有一个用来写入输出数据的循环内存缓冲区(默认大小是100MB)。当缓冲区中的数据量达到一个特定阀值(默认是0.80),系统将会启动一个后台线程把缓冲区中的内容spill到磁盘。在spill过程中,map的输出将会继续写入到缓冲区,但如果缓冲区已满,map就会被阻塞直到spill完成。
每当内存中的数据达到spill阀值的时候,都会产生一个新的spill文件,所以在map任务写完它的最后一个输出记录时,可能会有多个spill文件。在map任务完成前,所有的spill文件将会被归并排序为一个索引文件和数据文件。当spill文件归并完毕后,map将删除所有的临时spill文件,并告知TaskTracker任务已完成。
如果设定了Combiner,将在排序输出的基础上运行。Combiner就是一个Mini Reducer,它在执行map任务的节点本身运行,先对map的输出做一次简单reduce,使得map的输出更紧凑,更少的数据会被写入磁盘和传送到Reducer。
Reduce函数处理
map的输出文件放置在运行map任务的TaskTracker的本地磁盘上,它是运行reduce任务的TaskTracker所需要的输入数据。reduce任务的输入数据分布在集群内的多个map任务的输出中,map任务可能会在不同的时间内完成,只要有其中的一个map任务完成,reduce任务就开始拷贝它的输出。这个阶段称之为拷贝阶段。
reduce任务拥有多个拷贝线程,可以并行的获取map输出。线程数默认是5。
拷贝来的数据叠加在磁盘上,有一个后台线程会将它们归并为更大的排序文件,节省后期归并的时间。当所有的map输出都被拷贝后,reduce任务进入归并排序阶段,对所有的map输出进行归并排序,这个工作可能会重复多次。
假设有50个
map输出(可能有保存在内存中),并且归并因子是10,则最终需要5次归并。每次归并会把10个文件归并为一个,最终生成5个中间文件。之后,系统不再把5个中间文件归并成一个文件,而是排序后直接提交给reduce函数,省去向磁盘写数据这一步。
HDFS
HDFS是Hadoop的数据存储系统,是基于块的文件存储:
- 按块进行复制的形式放置,随机选择存储节点
- 副本的默认数目是3
- 默认的块的大小是64MB
- 减少元数据的量
- 有利于顺序读写(在磁盘上数据顺序存放)
- 适合于 MapReduce应用程序
- 依据Google File System的设计编写

特征:
- 存储规模大:大规模数据,大量存储节点,支持大文件。
- 高可靠:单个或者多个节点失效,对系统不会造成任何影响。
- 高可扩展:能简单加入更多服务器的方式便可服务更多的客户端。
- 为MapReduce优化:数据尽可能根据其本地局部性进行访问与计算。
使用场景:
- 大文件,顺序读:HDFS对顺序读进行了优化,但随机的访问负载较高。
- 数据支持一次写入,多次读取:不支持数据更新。
- 数据不进行本地缓存:文件很大,且顺序读没有局部性
- 任何一台服务器都有可能失效,需要通过大量的数据复制使得性能不会受到大的影响。
YARN
第一代Hadoop的缺陷:
- 单点故障:
JobTracker只有一个。JobTracker负责接收来自各个TaskTracker节点的RPC请求,压力会很大,限制了集群的扩展;随着节点规模增 大之后,JobTracker就成为一个瓶颈 - 集群包含的节点超过 4,000 个时(其中每个节点可能是多核的),就会表现出一定的不可预测性
- 仅支持
MapReduce计算框架 - 资源利用率低

YARN就是第二代MapReduce的核心产物。
将JobTracker的两个主要功能(资源管理和作业调度/监控)分离
- 创建一个全局的
ResourceManager(RM)和若干个针对应用程序的ApplicationMaster(AM)。 - 应用程序可以是传统的
MapReduce作业或作业的DAG(有向无环图)。
RM
控制整个集群并管理应用程序向基础计算资源的分配,主要由两个组件构成:
- 调度器(Scheduler)
- 应用程序管理器(Applications Manager)
调度器不参与与具体应用程序相关的工作,不负责监控或者跟踪应用的执行状态等,也不负责重启任务,这些均由应用程序相关的ApplicationMaster完成。
调度器是一个可插拔的组件,用户可根据自己的需要设计新的调度器,YARN提供了多种直接可用的调度器,比如Fair Scheduler和Capacity Scheduler等。
AM
每个应用程序均包含一个AM,主要功能包括:
- 与RM调度器协商以获取资源(用Container表示);
- 将得到的任务进一步分配给内部的任务(资源的二次分配);
- 与NM通信以启动/停止任务;
- 监控所有任务运行状态,并在任务运行失败时重新为任务申请资源以重启任务。
NM
NM是每个节点上的资源和任务管理器,定时地向RM汇报本节点上的资源使用情况和各个Container的运行状态;接收并处理来自AM的Container启动/停止等各种请求。
Container
Container是YARN中的资源抽象,封装了某个节点上的多维度资源,如内存、CPU、磁盘、网络等。当AM向RM申请资源时,RM为AM返回的资源便是用Container表示。YARN会为每个任务分配一个Container,且该任务只能使用该Container中描述的资源。
