并行计算

MapReduce编程

MapReduce编程模型

起源:Google搜索。

编程模型抽象简洁,易于使用。运行于大规模分布式集群环境,自动实现分布式并行计算,具备容错和负载均衡,支持任务调度和状态监控。

1
2
3
4
5
6
7
// 输入和输出
sets of <key, values> paris

// 处理 <K, V> 键值对,产生中间结果
map (in_key, in_value) -> (out_key, intermediate_value) list
// 将所有相同key的中间结果值进行规约,产生规约后的结果
reduce (out_key, intermediate_value) list

map函数

将数据源中的记录(文本、数据库)作为map函数中的(key, value)对。

map()将生成一个或多个中间结果,以及一个与input相对应的output key

reduce函数

map操作结束后,所有与某指定out key相对应的中间结果组合为一个列表(list)。reduce()函数将这些中间结果组合为一个或多个对应于同一 output keyfinal value(实际上每一个output key通常只有一个final value)。

MapReduce逻辑过程

MapRduce并行化运行

  1. map()并行执行,不同的输入数据集生成不同的中间结果
  2. reduce()并行执行,分别处理不同的output key
  3. mapreduce的处理过程中不发生通信

瓶颈:当map处理全部结束后,reduce过程才能够开始

Hadoop/HDFS

Hadoop是一个分布式系统基础架构,由Apache基金会开发。用户可以在不了解分布式底层细节的情况下,开发分布式程序。充分利用集群的威力高速运算和存储。是对Google所提出的MapReduce、GFS、BigTable等模型的一个开源实现。

Hadoop集群逻辑结构

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的数据存储系统,是基于块的文件存储:

  1. 按块进行复制的形式放置,随机选择存储节点
  2. 副本的默认数目是3
  3. 默认的块的大小是64MB
  4. 减少元数据的量
  5. 有利于顺序读写(在磁盘上数据顺序存放)
  6. 适合于 MapReduce应用程序
  7. 依据Google File System的设计编写

HDFS的逻辑结构 HDFS的文件系统

特征:

  1. 存储规模大:大规模数据,大量存储节点,支持大文件。
  2. 高可靠:单个或者多个节点失效,对系统不会造成任何影响。
  3. 高可扩展:能简单加入更多服务器的方式便可服务更多的客户端。
  4. 为MapReduce优化:数据尽可能根据其本地局部性进行访问与计算。

使用场景:

  1. 大文件,顺序读:HDFS对顺序读进行了优化,但随机的访问负载较高。
  2. 数据支持一次写入,多次读取:不支持数据更新。
  3. 数据不进行本地缓存:文件很大,且顺序读没有局部性
  4. 任何一台服务器都有可能失效,需要通过大量的数据复制使得性能不会受到大的影响。

YARN

第一代Hadoop的缺陷:

  1. 单点故障:JobTracker只有一个。JobTracker负责接收来自各个TaskTracker节点的RPC请求,压力会很大,限制了集群的扩展;随着节点规模增 大之后,JobTracker就成为一个瓶颈
  2. 集群包含的节点超过 4,000 个时(其中每个节点可能是多核的),就会表现出一定的不可预测性
  3. 仅支持MapReduce计算框架
  4. 资源利用率低

第二代MapReduce框架

YARN就是第二代MapReduce的核心产物。

将JobTracker的两个主要功能(资源管理和作业调度/监控)分离

  1. 创建一个全局的ResourceManager(RM)和若干个针对应用程序的ApplicationMaster(AM)。
  2. 应用程序可以是传统的MapReduce作业或作业的DAG(有向无环图)。

RM

控制整个集群并管理应用程序向基础计算资源的分配,主要由两个组件构成:

  1. 调度器(Scheduler)
  2. 应用程序管理器(Applications Manager)

调度器不参与与具体应用程序相关的工作,不负责监控或者跟踪应用的执行状态等,也不负责重启任务,这些均由应用程序相关的ApplicationMaster完成。

调度器是一个可插拔的组件,用户可根据自己的需要设计新的调度器,YARN提供了多种直接可用的调度器,比如Fair Scheduler和Capacity Scheduler等。

AM

每个应用程序均包含一个AM,主要功能包括:

  1. 与RM调度器协商以获取资源(用Container表示);
  2. 将得到的任务进一步分配给内部的任务(资源的二次分配);
  3. 与NM通信以启动/停止任务;
  4. 监控所有任务运行状态,并在任务运行失败时重新为任务申请资源以重启任务。

NM

NM是每个节点上的资源和任务管理器,定时地向RM汇报本节点上的资源使用情况和各个Container的运行状态;接收并处理来自AM的Container启动/停止等各种请求。

Container

Container是YARN中的资源抽象,封装了某个节点上的多维度资源,如内存、CPU、磁盘、网络等。当AM向RM申请资源时,RM为AM返回的资源便是用Container表示。YARN会为每个任务分配一个Container,且该任务只能使用该Container中描述的资源。

YARN工作原理

网站总访客数:Loading

使用 Hugo 构建
主题 StackJimmy 设计