Hadoop是Apache旗下的开源的分布式计算平台,为用户提供底层细节透明的分布式基础架构。它具有高可靠、高效、高扩展、高容错等特性,在业内得到广泛认可,同时hadoop也成为了大数据的代名词。本文将分析作为Hadoop的核心框架之一的MapReduce分布式计算框架。
分布式计算--MapReduce工作流程
MapReduce是一个用于大规模数据处理的编程模型和规范, 它采用了分而治之的思想, 基本形式有map和reduce两个处理阶段:主要思路是将大规模数据处理任务分为很多子任务, 并将子任务分配给若干个分布式的机器来并行完成批处理作业[1]。在工作中Map任务是生成<K,V>键值对,之后进行shuffle操作,Reduce机器将属于自己的分区数据“领走”,最后通过Reduce将最后的结果统计计算出来。整个处理过程的输出、输入都需要借助分布式文件系统进行存取。MapReduce具体执行阶段:
1)map任务的预处理
在初始阶段,使用InputFormat进行数据的预处理,检查输入的数据格式是否有误;之后使用InputSplit将大量的数据进行逻辑分割,记录每一个待处理数据的位置与长度;之后使用RecordReader处理已经被InputSplit逻辑划分好的数据,加载数据并且转化为键值对的形式,输出给Map任务。
2)Map端的shuffle操作
Shuffle操作就是将Map的输出数据进行重排输出给Reduce的过程,主要执行以下操作:
1.写入缓存
当Map得出结果时,不直接输出给磁盘,而是将K、V值序列化为字节数组,将其写入缓存中,当数据积累到一定量缓存发生溢出时,再一次性写至磁盘中。这种操作可以减轻I/O负担,降低开销。当然为了保证在发生溢写时,不影响新数据写至缓存中,就必须让缓存有多余的可用空间,这时就需要设置溢写比例<1.0,当达到该比例时就启动溢写操作,将数据刷入磁盘中。
2.溢写——分区、排序、合并
当发生溢写操作之前,缓存中首先会发生对输出<K,V>键值对进行分区操作。默认采用Hash函数对Key值进行Hash后再用Reduce任务数量进行取模即为hash(Key)mod(ReduceCount),这样就可以保证Reduce机器上的任务负载均衡,更加方便进行并行处理。在完成分区操作后,需要根据Key值进行内存排序,排序结束之后,还包含用户自定义的合并操作。其指将相同Key值的<K,V>键值对中Value值进行累加合并,从而减少写入磁盘中键值对的数量。
3.文件归并
每次溢写都会产生一个溢写文件,当Map任务结束之前,系统中会有大量的溢写文件,将这些溢写文件进行归并,即让相同Key值的键值对归并为一个新的键值对,形如<K1,<V1,V2,V3...>>,最后形成一个经过分区和排序后的大的溢写文件。
3)Reduce端的shuffle操作
在Reduce端的shuffle操作只需要从Map端领取数据,执行归并操作,最后传输给Reduce任务。具体操作如下:
1.“领取数据”
Map的输出结果都保存在自己的本地磁盘中,此时需要Reduce端将结果提取到自己的磁盘上。每个Reduce任务都会通过RPC向JobTracker询问Map任务是否已经结束,一旦Reduce收到JobTracker的通知,它就会到对应的Map机器上领取属于自己的分区数据。
2.归并数据
从Map端领取的数据会被存放在缓存中,如果缓存占满,就会溢写到磁盘中。当溢写过程启动时,会将具有相同Key值的键值对进行归并排序,溢写结束后会在磁盘中生成溢写文件。当所有Map端的数据都被领取后,多个溢写文件经过多轮会被归并排序形成多个大文件(多个大文件不会再继续归并为一个最终的大文件)。
3.将数据输入给Reduce任务
将多个大文件直接输入给Reduce任务,Reduce执行被定义的各种映射,输出最终结果,存储在分布式文件系统中。
实例分析:词频统计
一.基本步骤:
在执行程序时,会被部署到分布式集群中,其中有一台Master机器用于调度作业的执行,和多台worker机器用于执行map和reduce任务。系统将待处理的文件逻辑切分为多个数据分片,Master将不同分片分配给空闲的Map-worker进行处理。Map-worker执行任务生成<K,V>键值对存放进缓存中。当缓存数据刷入磁盘之前,会被划分为R个分区,Master记录好这几个分区的位置,之后通知Reduce-Worker领取对应的分区进行处理。Reduce-Worker先对领取的键值对进行排序,使得具有相同键值的键值对聚集在一起。之后开始执行Reduce任务,将每一个唯一Key执行Reduce函数,将结果输出到分布式文件系统中。
二.具体实例:

假设在集群中有3台机器执行Map任务,1台机器执行Reduce任务。文档中包含3行内容,每行由1个Map任务来处理,其中Map的输入是<K(文档的行号),V(每行的内容)>。Map任务结束后输出新的<K,V>键值对。之后经过Map端的shuffle后,将具有相同Key值的键值对进行归并操作得出形如<K1,<V1,V2...>>的键值对,并且将其排序。之后这些键值对会作为Reduce端的输入,由Reduce任务进行每个单词的计数,得出最终结果(这里不再给出Reduce端的shuffle操作)。
参考文献:
[1]李锐,王斌.
本处理中的MapReduce技术[J].中文信息学报,2012,26(04):9-20.
[2]MapReduce[J].Jeffrey Dean,Sanjay Ghemawat.Communications of the ACM. 2010 (1)
[3].A new algorithm for fast mining frequent itemsets using N-lists[J].Science China(Information Sciences),2012,55(09):2008-2030.