以下是Hadoop的技术框架图,请选择相应的组件填入其中:
Hadoop概述
Hadoop是一个由Apache基金会所开发的分布式系统基础架构。主要解决海量数据的存储以及海量数据的分析计算,Hadoop核心概念是HDFS(分布式存储),以及MapReduce。通俗来将就是,当一个海量数据,如果你使用传统方法去处理,比如10TB的数据,那么数据库就不太够用,并且即使数据库够用,那么在有限的资源(也就是CPU、内存)中处理的数据量是极少 的,10TB数据可能需要运作很久,那么有了Hadoop,HDFS是分布式存储文件系统,它将10TB数据分块存储在不同服务器(节点)中,然后在每个服务器中处理相应文件,也就是 相当于n个服务器并行处理,并且每个服务器都不需要太高的配置。而这些处理的任务是由MapReduce分配的。
如果你不太清楚什么是并行:举个简单的例子。传统处理方法,假如要写10000个字,你一个人写字速度为一分钟100个,那么就需要100分钟才能写完。而采用并行处理,相当于有100同学一起写,并且分工明确,那么只需要一分钟,你们100个人就能写10000个字。 一个简单的案例就是WordCount,字数统计,假设有100G的文件,要统计里面所有单词出现的次数。
传统方法:
假设我们有个2G的机器。首先这100G我们不可能都加载到内存中,内存不够,所以要将100G文件依次读入内存中并操作,就产生类似下图的一个队列,每次读取文件的一部分然后再处理,然后处理时间。需要加载至内存的次数 = 100G/2G = 50次,假设每次处理字数统计所用时间为T(process),假设每次读取至内存所用时间为T(IO),假设将所有的50次统计结果汇总时间为T(wordcount),那么一共需要的时间为:50次 x [T(process) + T(IO)] + T(wordcount)
可以看出是非常耗时的

采用分布式并行计算后(Hadoop使用类似的方法,但是其多了reduce去汇总处理)

可以看到,同样将任务分割然后送到集群中不同的服务器(节点),每个服务器并行处理(可以理解为同时处理)。此时,依照上面传统方法的时间计算。共使用了1次
x [T(process)
+ T(IO)] + T(wordcount)
这里要注意,hadoop实际应用中,会将汇总处理也分布式并行处理,所以其实更加省时,这里先忽略,依然能看出两者相差的耗时。
看完例子你可能就差不多知道hadoop是干什么的了:
分布式存储:将文件分布式存储在多个服务器(节点)上。
分布式并行处理:通过编写处理任务的代码,将文件进行切割,然后将文件以及任务资源(比如jar包,也就是你写的处理逻辑)分发给多个服务器,每个服务器进行并行处理。
一、Hadoop介绍
1.1 Hadoop是什么
Hadoop是一个开源的软件,并且是可靠性的(reliable)、可扩展性的(scalable)、分布式计算的(distributed
computing)软件。
• Hadoop是一个框架,可以允许分布式处理大数据集(big data
sets),而且这个数据集是横跨在集群的机器上的(clusters of
computers)。也就是说数据可以分开地存储在集群中的每个机器
(节点)上,并且可以跨节点进行处理计算,它使用的是一种简单的编程模型(using
simple programming models)。
•Hadoop被设计成可以从单个服务器(single
servers)扩展到数千台机器的的集群上,每台机器都提供本地的存储和计算服务。也就是说当数据量小的时候,可以使用少一点的机器的集群,
而面对一个大量数据的情况时,现有的机器不足以支撑存储运算时,只需要再添加一些机器到集群中,就能解决问题。
•Hadoop并不是依赖硬件来提供高可用性,而是它自己被设计成可以检测和处理应用层的故障。在集群中每台机器上提供高可靠性服务,而这些机器可能会倾向于出现故障。
1.2 Hadoop能做什么
•Hadoop可以搭建大型数据仓库,PB级数据的存储、处理、分析、统计等业务。
•商业智能(BI)、可视化报表的产生等等。
•数据挖掘,从大量的数据中挖掘出有价值的结论等等。
1.3 Hadoop包括什么
Hadoop框架包含以下模块

• Hadoop Common:提供支撑其他模块的通用工具。
• Hadoop HDFS:为应用程序访问提供高吞吐量的分布式文件系统。
• Hadoop YARN:提供任务调度服务和集群资源管理的框架。
• Hadoop MapReduce:可以并行处理大数据集的编程计算模型。
• Hadoop Ozone:提供对象存储的功能。
• Hadoop Submarine :提供机器学习引擎。
二、HDFS数据存储原理
2.1 HDFS概述
HDFS,全称Hadoop Distributed File System,是一个分布式文件系统。
分布式文件系统(Distributed File
System)是指文件系统管理的物理存储资源不一定直接连接在本地节点上,而是通过计算机网络与节点相连。分布式文件系统的设计基于客户机/服务器
模式。一个典型的网络可能包括多个供多用户访问的服务器。

HDFS是一个设计可以运行在廉价硬件上的分布式文件系统。HDFS是一个高容错和可部署(deployed)在廉价机器上的系统,提供了对应用程序数据的高吞吐量访问,适用于具有
大型数据集的应用程序。
• 2.2 HDFS特点
•
商用硬件。硬件故障是常态,而不是异常。整个HDFS系统将由数百或数千个存储着文件数据片段的服务器组成。实际上它里面有非常巨大的组成部分,每一个组成部分都很可能出现故
障,这就意味着HDFS里的总是有一些部件是失效的,因此,故障的检测和自动快速恢复是HDFS一个很核心的设计目标。
•
流数据访问。运行在HDFS之上的应用程序必须流式地访问它们的数据集,它不是运行在普通文件系统之上的普通程序。HDFS被设计成适合批量处理的,而不是用户交互式的。重点是
在数据吞吐量,而不是数据访问的反应时间,POSIX的很多硬性需求对于HDFS应用都是非必须的,去掉POSIX一小部分关键语义可以获得更好的数据吞吐率。
•
大型数据集。运行在HDFS之上的程序有很大量的数据集。典型的HDFS文件大小是GB到TB的级别。所以,HDFS被调整成支持大文件。它应该提供很高的聚合数据带宽,一个集群中
支持数百个节点,一个集群中还应该支持千万级别的文件。
•
简单一致模型。HDFS应用需要一个一次写入多次读取的文件访问模型。一个文件一旦创建,写入和关闭都不需要改变除了追加和截断(truncate)。支持在文件的末端进行追加数据而不
支持在文件的任意位置进行修改。这个假设简化了数据一致性问题和支持高吞吐量的访问。一个Map/Reduce任务或者web爬虫(crawler)完美匹配了这个模型。
•
移动计算比移动数据便宜。在靠近计算数据所存储的位置来进行计算是最理想的状态,尤其是在数据集特别巨大的时候。这样消除了网络的拥堵,提高了系统的整体吞吐量。一个假定
就是迁移计算到离数据更近的位置比将数据移动到程序运行更近的位置要更好。HDFS提供了接口,来让程序将自己移动到离数据存储更近的位置。
•
在异构硬件和软件平台上的可移植性。HDFS被设计成可以简便地实现平台间的迁移,这将推动需要大数据集的应用更广泛地采用HDFS作为平台。
2.3 HDFS架构

•
NameNode:管理节点,维护着文件系统树及整个树内的所有文件和目录(元数据),同时负责客户端请求。
•
DataNode:文件系统的工作节点,根据需要存储和检索数据块,并且定期向NameNode发送他们所存储的块的列表。
• Blocks:HDFS中的存储单元,默认为128M,通常有多个备份,默认为3个。
2.4 HDFS数据读写


元数据有三种形式:内存、EditsLog、FsImage。
• 内存中保存的是最完整最新的元数据。
• EditsLog保存HDFS自最新的元数据检查点后的元数据变化的记录。
• FsImage保存最新的元数据检查点。
Checkpoint(检查点)指的是在NameNode启动时候,会先将fsimage中的文件系统元数据信息加载到内存,然后根据eidts中的记录将内存中的元数据同步至最新状态(这里读的是
journalnode中的editlog),将这个新版本的
FsImage 从内存中保存到本地磁盘上,然后删除旧的 Editlog。
fsimage存放上次checkpoint生成的文件系统元数据,Edits存放文件系统操作日志。checkpoint的过程,就是合并fsimage和Edits文件,然后生成最新的fsimage的过程。
2.6 HDFS副本存放

机架感知策略:
•
第一个复本放在运行客户端的节点上(如果客户端运行在集群之外,则在避免挑选存储太满或太忙的节点的情况下随机选择一个节点)。
• 第二个复本放在与第一个不同且随机选择的机架的节点上。
• 第三个复本与第二个复本放在同一个机架上,且随机选择另一个节点。
• 其它复本放在集群中随机选择的节点中,尽量避免在同一个机架上放太多复本。
三、MapReduce编程模型框架
3.1 MapReduce概述
•
定义:是谷歌开源的一种大数据并行计算编程模型,它降低了并行计算应用开发的门槛。
• 工作原理:利用一个输入key/value pair集合来产生一个输出的key/value
pair集合。
•
运行机制:以一种可靠的、容错的方式,在大型的商用硬件集群(数千个节点)上并行处理大量数据(多为TB级别的数据集)。
• 优点:简单容易使用、扩展性强、高容错性、可离线计算 PB
量级的数据。
• 缺点:实时计算、流式计算、有向图计算支持性不高。
3.2 工作原理

分而治之。采用分布式并行计算,将计算任务进行拆分,由主节点下的各个子节点共同完成,最后汇总各子节点的计算结果,得出最终计算结果。
MapReduce任务通常将输入数据集分割成独立的块,由map任务以完全并行的方式处理。框架对map任务映射的输出进行排序,然后将这些输出输入到reduce任务中。
通常,作业的输入和输出都存储在文件系统中,框架负责调度任务,监视任务,并重新执行失败的任务。
MapReduce框架只对'key,value',对进行操作,也就是说,框架将作业的输入视为一组'key,value'对,并生成一组'key,value'对作为作业的输出,可以认为是不同类型的。
3.3 MapReduce执行步骤
•
第一阶段是把输入文件按照一定的标准分片(InputSplit),每个输入片的大小是固定的。默认情况下,输入片(InputSplit)的大小与数据块(Block)的大小
是相同的。如果数据块(Block)的大小是默认值64MB,输入文件有两个,一个是32MB,一个是72MB。那么小的文件是一个输入片,大文件会分为两个数据块,那么是两个输入片。一共产生三个输入片。每一个输入片由一个
Mapper进程处理。
•
第二阶段是对输入片中的记录按照一定的规则解析成键值对。有个默认规则是把每一行文本内容解析成键值对。“键”是每一行的起始位置(单位是字节),“值”是本行的文本内容
•
第三阶段是调用Mapper类中的map方法。第二阶段中解析出来的每一个键值对,调用一次map方法。如果有1000个键值对,就会调用1000次map方法。每一次调用map方法会输出零个或
者多个键值对。
•
第四阶段是按照一定的规则对第三阶段输出的键值对进行分区。分区是基于键进行的。比如我们的键表示省份(如北京、上海、山东等),那么就可以按照不同省份进行分区,同一个省份
的键值对划分到一个区中。默认是只有一个区。分区的数量就是Reducer任务运行的数量。默认只有一个Reducer任务。
•
第五阶段是对每个分区中的键值对进行排序。首先,按照键进行排序,对于键相同的键值对,按照值进行排序。比如三个键值对<2,2>、<1,3>、<2,1>,键和值分别是整数。那么排序后的
结果是<1,3>、<2,1>、<2,2>。如果有第六阶段,那么进入第六阶段。
•
第六阶段是对数据进行归约处理,也就是reduce处理,通常情况下的Combine过程,键相等的键值对会调用一次reduce方法,经过这一阶段,数据量会减少,归约后的数据输出到本地的linux文件中。
在MapReduec任务调度模型中,主要包含以下几个角色:
•
JobClient:接收客户端提交的程序jar包,并进行一系列的操作,具体会在下面进行详细分析。
•
JobTracker:负责接收JobTracker分配的作业,将Task分配到各个节点上去运行,并提供诸如监控工作节点状态及任务进度等管理功能,一个MapReduce集群只有一个JobTracker。
•
TaskTracker:负责监控任务的执行情况,并通过心跳连接周期性地向jobtracker汇报任务进度,资源使用量等。每台执行任务的节点都会有一个TaskTracker
,其使用“slot”等量划分本节点
上的资源量。“slot”代表计算资源(CPU、内存等)。一个Task
获取到一个slot 后才有机会运行,slot 分为Map slot和Reduce slot
两种,分别供MapTask 和Reduce Task 使用。TaskTracker
通过slot
数目(可配置参数)限定Task 的并发度。
• HDFS:用于提供数据存储服务和和作业间的文件共享。
四、YARN资源调度管理
4.1 概述
YARN全称是Yet Another Resource
Negotiatord。通用的资源管理系统,要申请资源统一经过YARN进行申请。为上层应用提供统一的资源管理和调度。不同计算框架可以共享同一个HDFS
集群上的数据,享受整体上的资源调度:Spark
on YARN 、 MapReduce on YAEN 、 Storm on YARN
...与其他计算框架共享集群资源,按资源需要分配,进而提高集群资源的利用率。
4.2 架构

Yarn 采用传统的 master-slave 架构模式,其主要由 4
种组件组成,它们的主要功能如下:
•
ResourceManager(RM):全局资源管理器,负责整个系统的资源管理和分配;
• ApplicationMaster(AM):负责应用程序(Application)的管理;
• NodeManager(NM):负责 slave 节点的资源管理和使用;
• Container(容器):对任务运行环境的一个抽象。
4.3 执行流程

• 客户端向RM中提交程序 。
• RM向NM中分配一个container,并在该container中启动AM 。
•
AM向RM注册,这样用户可以直接通过RM査看应用程序的运行状态(然后它将为各个任务申请资源,并监控它的运行状态,直到运行结束)
。
•
AM采用轮询的方式通过RPC协议向RM申请和领取资源,资源的协调通过异步完成。
• AM申请到资源后,便与对应的NM通信,要求它启动任务 。
•
NM为任务设置好运行环境(包括环境变量、JAR包、二进制程序等)后,将任务启动命令写到一个脚本中,并通过运行该脚本启动任务
。
•
各个任务通过某个RPC协议向AM汇报自己的状态和进度,以让AM随时掌握各个任务的运行状态,从而可以在任务失败时重新启动任务
。
• 应用程序运行完成后,AM向RM注销并关闭自己。