- 1 -
中国科技论文在线
改进型 mapreduce 框架的研究与设计
常 涛*
作者简介:常涛(1984-),男,硕士研究生,主要研究方向:云计算
(北京邮电大学计算机科学与技术学院,北京 100876)
摘要:随着网格计算被取代,云计算迎来了蓬勃的发展。Hadoop 作为开源云计算平台,得
到了国内外很多公司的青睐。相应的,作为 Hadoop的子项目和分布式并行处理框架,目前
基于 mapreduce的应用越来越多。随着应用的广泛性和多样性,其暴露处理来的不足和需要
改进之处越来越多。本文首先介绍了传统的 MapReduce 框架的处理流程,然后分析了处理
过程中可能出现的一些影响执行效率的问题,最后提供了框架的改进方案。经过试验表明,
改进后的框架在某些应用上确实能够提高作业的执行效率。
关键词: 并行处理;云计算;Hadoop;mapreduce;
中图分类号:TP311
research and design for improved mapreduce framework
CHANG Tao
(School of Computer Science and Technology, Beijing University of Posts and
Telecommunications,Beijing 100876, China)
Abstract: As replace the grid computing,cloud computing has a rapid an open source
cloud computing platform,Hadoop has been adopted by domestic and foreign
,as a sub-project of Hadoop and a distributed parallel processing
framework,there are more and more applications based on with the breadth and
diversity of the application, it exposes many places need to be this paper,we firstly
introduce the processing workflow of the traditional mapreduce framework,and secondly we analysis
some problems that maybe affect the efficiency of applications in the ,provides some
improvements for the framework.
Key words: parallel process;cloud computing;Hadoop; mapreduce;
0 引言
从 2003 年开始,许多大型互联网公司如 IBM 等就陆续提供了像买电买水一样通过一个
网络购买计算力和存储空间的服务,这种潮流在 2006-2007 年之间发生了质变,随着著名的
畅销书销售网站 Amazon 和 IBM、Google 先后推出了云计算服务,云计算作为网格的升级
实现产品已经在商业化的进程中阔步前行。
在云计算技术中,编程平台是一个非常重要的模块。Hadoop 平台是当今应用最为广泛
的开源云计算编程平台。
Hadoop 是 Apache 开源组织的一个分布式计算开源框架,在很多大型网站上都已经得到
了应用,如 Amazon,Facebook,Yahoo!,IBM 等等。本文针对传统的 MapReduce 框架进行
了改进,能够加快处理流程的处理时间。
随着 MapReduce 框架的大量应用,其也暴露出很多需要改进的地方,如对 map 和 reduce
进行的数据流化[1]和管道化,加快中间结果的传输以及 Doug Cutting 等人所做的一些扩展[2]
等。本文讨论了针对 map 操作后某些中间结果偏大的以至于影响到整个工作的处理时间和
- 2 -
中国科技论文在线
处理效率的问题以及相应的改进方案。Hadoop 针对这种情况会检测到有 ReduceTask 执行时
间过长时,会启动另一个 ReduceTask 作为备份执行,但本文提供了另一种改进方案。
1 传统 Mapreduce 框架处理流程
Mapreduce 概述
Hadoop 框架中核心设计是 MapReduce 和 HDFS(HadoopDistributedFileSystem)[3]。
MapReduce 是将作业分解为任务并对任务的结果进行汇总。"Map(映射)"和"Reduce(化
简)"的概念是从函数式编程语言借来的,还有从矢量编程语言借来的特性[4]。通过搭建
Hadoop 集群,可以方便的使用 MapReduce 进行分布式计算[5][6]。
Mapreduce 处理流程
Mapreduce 框架的微观处理流程在细节上的实现如下图[7]所示,
图 1传统 MapReduce 微观处理流程
普通的 MapReduce 处理流程需要经过如下五个阶段:
1)读入数据: key/value 对的记录格式数据。此前已经按照 map 任务的数目将输入文件
划分成相应的份数。
2)Map: 从每个记录里 extract something
map (in_key, in_value) -> list(out_key, intermediate_value)
处理 input key/value pair
输出中间结果 key/value pairs
3)Shuffle: 混排交换数据,把相同 key 的中间结果汇集到相同节点上
4)Reduce:
reduce (out_key, list(intermediate_value)) -> list(out_value)
归并某一个 key 的所有 values,进行计算
输出合并的计算结果 (usually just one)
5)输出结果
- 3 -
中国科技论文在线
2 传统 MapReduce 框架处理流程可能造成的问题
Mapreduce 的处理过程很多都会失败,经过分析,有很多原因。如 job 执行过程中机器
硬件发生故障或者错误;针对 job 分配的资源有限,譬如内存有限等;第三类有可能是中间
结果偏大,即会出现中间结果不均衡的问题,造成某个 ReduceTask 执行时间过长或者当掉。
数据不均衡可分为两类:
1)KEY 过于聚集,即不同 KEY 的个数虽多,但经过映射(如 HASH)后,过于聚集
在一起。采用两种办法相结合:
一是将 KEY 分散得足够大(如 HASH 桶数够多);
二是在 balance 的时候,进行重 HASH,将大的打成小的。
2)KEY 值单一,即不同 KEY 的个数少
对于这种情况,采用 HASH 再分散的方法无效,事先也无法分散得足够大,但处理的
方法也非常简单。对这种情况,按照大小进行横切即可,但这个时候一次 reduce 无法得到。
最终结果,至少需要连接两次 reduce,另外还需要增加 balance 接口,以方便区别是最
后一次 reduce,还是中间的 reduce。
3 改进型的 mapreduce 框架
体系结构
本文在进行系统架构设计时主要有以下两条原则:
1)可扩展,可配置。不同用户对某些特定功能具有不同的要求,因此,系统某些模块
在提供默认实现时应当提供用户扩展机制。
2)功能模块之间耦合性小,尽量使各模块功能在逻辑上相互独立。
针对以上功能需求,本文所涉及的改进型 MapReduce 框架体系结构如下图:
Data store 1 Data store n
map
(key 1,
values...)
(key 2,
values...)
(key 3,
values...)
map
(key 1,
values...)
(key 2,
values...)
(key 3,
values...)
Input key*value
pairs
Input key*value
pairs
== Barrier == : Aggregates intermediate values by output key
reduce reduce reduce
final key 1
values
final key 2
values
final key 3
values
...
balance
图 2改进型 MapReduce 框架系统架构图
- 4 -
中国科技论文在线
改进型的 MapReduce 框架系统架构如上图所示,整体上看,在 map 过程和 reduce 过程
间加上了一个 balance 过程。但这不代表改进后的框架一定会执行 balance 过程,只有在一
定条件的触发下,才会执行。
改进型 MapReduce 框架实际上在 MapReduce 上做两件事:保证 ReduceTask 均衡和控
制和 ReduceTask 大小。
改进方案
由于 map 的输入源自于 DFS,是相对静态的数据,所以各 MapTask 是均衡的,而且其
大小是已知和确定的。
reduce 不同于 map,它处理的是 map 后的数据,是动态产生的数据,只有在 map 完成
之后才能确定的数据(包括数据的分布和大小等)。数据在经过 map 之后,就映射到了某
个 ReduceTask,不能再更改,而且 ReduceTask 的个数也是不能修改的。
这会带来如下严重的问题:
1)reduce 端数据倾斜:比如有些 ReduceTask 需要处理较多数据,而另有一些只需要处
理较少数据,甚至还有些 ReduceTask 可能空转,没有任何数据需要处理。最严重时,可能
发生某个 ReduceTask 被分配了超出本地可用存储空间的数据量,或是超大数据量,需要特
长处理时间甚至需要处理的数据量超出了可用内存限制,造成溢出,直到整个 job 的失败;
2)reduce 端数据倾斜直接导致了 ReduceTask 不均衡;
3)并行 Job 困难。类似于操作系统进程调度,如果要并行,必然存在 Job 间的调度切
换,但由于 ReduceTask 需要处理的数据量可能很大,需要运行很长的时间,如果强制停止
ReduceTask,对于大的 ReduceTask 会浪费大量的已运行时间,甚至可能导致一个大的 Job
运行失败。因此,无法实现类似于进程的并行调度器。
4)针对改进型 MapReduce 框架,所作的其他改进,如在配置文件中增加阀值等字段;
更改 MapTask、ReduceTask 和 master 间心跳信息的数据格式等。
中间结果偏大的检测
本文提供的方法是在配置文件中,提供一个参数,作为检测中间结果是否偏大的阀值。
这个阀值的设计参考了将输入文件进行划分的设计方法。默认值的设定 64M,可以根据具
体应用的不同在执行 job 前灵活改变。系统检测到中间结果大小超过该阀值时即可认为需要
针对中间结果进行划分。
大数据块的切分函数
目前系统中只能提供 hash 函数对数据块进行在切分,接下来的工作可以从两方面进行
改进型:
这对数据类型的不同,提供针对 hash 函数在性能上的优化。
提供切分函数接口,切分函数由用户自己灵活实现。
Task 调度
在所有 Task 均衡,且其大小是可控的前提下,并行调度就可以仿照进程调度去做。我
们可以将 Task 当作一个运行时间片,由于其大小可以控制,所以只要大小适当,基本上就
可以控制其运行时长。
当一个 Task 运行完后,根据调度规则来决定下一个运行的 Task,下一个 Task 并不一
- 5 -
中国科技论文在线
定是同一个 Job。
4 结论
本文提出了一种经过改进的 MapReduce 框架,用来解决针对中间结果过大的情况。经
过基础实验对比验证,使作业的处理效率上有所提高。
致谢
感谢我的导师和实验室同学的帮助,特别感谢易剑前辈的慷慨帮助。
[参考文献] (References)
[1] UC Berkeley. MapReduce Online [OL].
[2] Doug Cutting. Scalable Computing with Hadoop. America: Yahoo Inc!, 2006.
[3] Apache Hadoop [OL].
[4] MapReduce [OL].
[5] White, T. . Hadoop: The Definitive Guide [M].America: O'Reilly Media,Inc,2009
[6] White, T. . Hadoop 权威指南(中文版)[M].曾大聃 周傲英.北京:清华大学出版社,2010
[7] J. Dean, S. Ghemawat. Mapreduce: Simplified data processing on large clusters[A]. In OSDI,2004.