3.1.6 MapReduce应用编程

更新于 2026年10月10日 版权声明
3.1.6 MapReduce应用编程

单词计数是最简单也是最能体现MapReduce思想的程序之一,可称为MapReduce版“Hello World”。单词计数的主要功能是统计一系列文本文件中每个单词出现的次数。通过单词计数实例来阐述采用MapReduce解决实际问题的基本思路和具体实现过程。

1)设计思路

首先,检查单词计数是否可以使用MapReduce进行处理。因为在单词计数程序任务中,不同单词出现的次数之间不存在相关性,是相互独立的,所以,可以把不同的单词分发给不同的机器进行并行处理。因此,可以采用MapReduce来实现单词计数的统计任务。

其次,确定MapReduce程序的设计思路。把文件内容分解成许多个单词,然后把所有相同的单词聚集到一起,计算出每个单词的出现次数。

最后,确定MapReduce程序的执行过程。把一个大的文件切分成许多个分片,将每个分片输入到不同结点上形成不同的Map任务。每个Map任务分别负责完成从不同的文件块中解析出所有的单词。

Map函数的输入采用<key,value>方式,用文件的行号作为key,文件的一行作为value。Map函数的输出以单词作为key,1作为value,即<单词,1>表示该单词出现了1次。

Map阶段结束以后,会输出许多<单词,1>形式的中间结果,然后Sort会把这些中间结果进行排序,并把同一单词的出现次数合并成一个列表,得到<key,list(value)>形式。例如,<Hello,<1,1,1,1,1>>就表明Hello单词在5个地方出现过。

如果使用Combine,那么Combine会把每个单词的list(value)值进行合并,得到<key,value>形式,例如,<Hello,5>表明Hello单词出现过5次。

在Partition阶段,框架会把Combine的结果分发给不同的Reduce任务。Reduce任务接收到所有分配给自己的中间结果以后,就开始执行汇总计算工作,计算得到每个单词出现的次数并把结果输入到HDFS中。

2)处理过程

下面通过一个实例对单词计数进行更详细的讲解。

第一步 将文件拆分成多个分片。

该实例把文件拆分成两个分片,每个分片包含两行内容。在该作业中,有两个执行Map任务的结点和一个执行Reduce任务的结点。每个分片分配给一个Map结点,并将文件按行分割形成<key,value>对,如图3⁃6所示。这一步由MapReduce框架自动完成,其中key的值为行号。

图示

图3⁃6 分割过程

第二步 将分割好的<key,value>对交给用户定义的Map函数进行处理,生成新的<key,value>对,如图3⁃7所示。

图示

图3⁃7 执行Map函数

第三步 在实际应用中,每个输入分片在经过Map函数分解以后都会生成大量类似<Hello,1>的中间结果,为了减少网络传输开销,框架会把Map方法输出的<key,value>对按照key值进行排序,并执行Combine过程,将key值相同的value值累加,得到Map的最终输出结果,如图3⁃8所示。

图示

图3⁃8 Map端排序及Combine过程

第四步 Reduce先对从Map端接收的数据进行排序,再交由用户自定义的Reduce方法进行处理,得到新的<key,value>对,并作为结果输出,如图3⁃9所示。

图示

图3⁃9 Reduce端排序及输出结果

3)程序实现

通过前面的分析,已经清楚了单词计数程序的设计思路和处理过程,接下来介绍如何编写基本的MapReduce程序实现单词计数。

第一步 任务准备。

单词计数(Word Count)的任务是对一组输入文档中的单词进行分别计数。假设文件的量比较大,每个文档又包含大量的单词,则无法使用传统的线性程序进行处理,而这类问题正是MapReduce可以发挥优势的地方。

首先,在本地创建3个文件:file00l、file002和file003,文件的具体内容见表3⁃1。

表3⁃1 单词计数输入文件

图示

再使用HDFS命令创建一个input文件目录。

hadoop fs ⁃mkdir input

然后,把file001、file002和file003上传到HDFS中的input目录下。

编写MapReduce程序的第一个任务就是编写Map程序。在单词计数任务中,Map需要完成的任务就是把输入的文本数据按单词进行拆分,然后以特定的键值对的形式进行输出。

第二步 编写Map程序。

Hadoop MapReduce框架已经在类Mapper中实现了Map任务的基本功能。为了实现Map任务,开发者只需要继承类Mapper,并实现该类的Map函数。

为实现单词计数的Map任务,首先为类Mapper设定好输入类型和输出类型。这里,Map函数的输入是<key,value>形式,其中,key是输入文件中一行的行号,value是该行号对应的一行内容。所以,Map函数的输入类型为<LongWritable,Text>。Map函数的功能为完成文本分割工作,Map函数的输出也是<key,value>形式,其中,key是单词,value为该单词出现的次数。所以,Map函数的输出类型为<Text,LongWritable>。

以下是单词计数程序的Map任务的实现代码。(https://www.daowen.com)

在上述代码中,实现Map任务的类为CoreMapper。该类将需要输出的两个变量one和label进行初始化。变量one的初始值直接设置为1,表示某个单词在文本中出现过。

Map函数的前两个参数是函数的输入参数,value为Text类型,是指每次读入文本的一行,key为Object类型,是指输入的行数据在文本中的行号。首先,StringTokenizer类机器方法将value变量中文本的一行文字进行拆分,拆分后的单词放在tokenizer列表中。然后,程序通过循环对每一个单词进行处理,把单词放在label中,把one作为单词计数。在函数的整个执行过程中,one的值一直是1。在该实例中,key没有被明显地使用到。context是Map函数的一种输出方式,通过使用该变量,可以直接将中间结果存储在其中。

Map任务执行后,3个文件的输出结果见表3⁃2。

表3⁃2 单词计数Map任务输出结果

图示

第三步 编写Reduce程序。

编写MapReduce程序的第二个任务就是编写Reduce程序。在单词计数任务中,Reduce需要完成的任务就是把输入结果中的数字序列进行求和,从而得到每个单词的出现次数。

在执行完Map函数之后,会进入Map端的合并阶段,在这个阶段中,MapReduce框架会自动将Map阶段的输出结果进行排序和分区,然后再分发给相应的Reduce任务去处理。经过Map端合并阶段后的输出结果见表3⁃3。

表3⁃3 单词计数Map端合并输出结果

图示

Reduce端接收各个Map端发来的数据后,会进行合并,即把同一个key,也就是同一单词的键值对进行合并,形成<key, <V1, V2, ...Vn>>形式的输出。经过Reduce端合并后的输出结果见表3⁃4。

表3⁃4 单词计数Reduce端合并输出结果

图示

Reduce阶段需要对上述数据进行处理,从而得到每个单词的出现次数。从Reduce函数的输入已经可以理解Reduce函数需要完成的工作,就是对输入数据value中的数字序列进行求和。以下是单词计数程序的Reduce任务的实现代码。

与Map任务实现相似,Reduce任务也是继承Hadoop提供的类Reducer并实现其接口。 Reduce函数的输入、输出类型与Map函数的输出类型本质上是相同的。在Reduce函数的开始部分,首先设置sum参数用来记录每个单词的出现次数,然后遍历value列表,并对其中的数字进行累加,最终就可以得到每个单词总的出现次数。在输出的时候,仍然使用context类型的变量存储信息。当Reduce阶段结束时,就可以得到最终需要的结果(表3⁃5)。

表3⁃5 单词计数Reduce任务输出结果

图示

第四步 编写main函数。

为了使用CoreMapper和CoreReducer类进行真正的数据处理,还需要在main函数中通过Job类设置Hadoop MapReduce程序运行时的环境变量,以下是具体代码。

代码首先检查参数是否正确,如果不正确就提醒用户。随后,通过Job类设置环境参数,并设置程序的类、Mapper类和Reducer类。然后,设置程序的输出类型,也就是Reduce函数的输出结果<key,value>中key和value各自的类型。最后,根据程序运行时的参数,设置输入、输出文件路径。

第五步 运行程序。

在运行代码前,需要把当前工作目录设置为/user/local/Hadoop。编译WordCount程序需要以下3个jar,为了简便起见,把这3个jar添加到CLASSPATH中。

使用JDK包中的工具对代码进行编译。

$ javac WordCount.java

编译之后,在文件目录下可以发现有3个“.class”文件,这是Java的可执行文件,将它们打包并命名为wordcount.jar。

$ jar -cvf wordcount.jar *.class

这样就得到了单词计数程序的jar包。在运行程序之前,需要启动Hadoop系统,包括启动HDFS和MapReduce。然后,就可以运行程序了。

$ ./bin/Hadoop jar wordcount.jar WordCount input output

最后,可以运行下面的命令,查看结果。

$ ./bin/Hadoop fs -cat output/*

4)核心代码包

编写MapReduce程序需要引用Hadoop的以下几个核心组件包,它们实现了Hadoop MapReduce框架。

这些核心组件包的基本功能描述见表3⁃6。

表3⁃6 Hadoop MapReduce核心组件包的基本功能

图示

↑上一章 ↓下一章
关注公众号获取验证码
复制内容需要验证码(7.99元/天)