显示标签为“翻译”的博文。显示所有博文
显示标签为“翻译”的博文。显示所有博文

星期六, 八月 26, 2006

MapReduce:Simplified Data Processing on large Clusters-翻译版(下)

实现
拥有多个不同的MapReduce接口的实现是可能的。具体选择取决于环境。比如,一种实现适合于共享内存的机器,一种适合于NUMA(Non-Uniform Memory Access )多处理器,另外一种适合于大量的网络机器。
这节将描述在Google内部大量使用,适合于大量PC构成的集群系统这种计算环境的实现。在这个环境中:
  1. 双CPUx86机器,运行Linux,具有2-4G内存。
  2. 网络采用的是100M/s或1G/s的配置硬件,然实际的平均使用值大概为它们的一半。
  3. 采用由成百上千PC构成的集群系统,机器故障很经常。
  4. 存储采用的是直连每个机器的IDE硬盘。通过我们自己开发的GFS来管理磁盘上的数据。此文件系统使用复制策略使得自己在不可靠的硬件上具有很好的可用性和可靠性。
  5. 用户提交作业到调度系统。每个作业由一系列任务构成,这些任务被映射到集群范围内可用的机器上。
  • 执行概貌
通过把输入文件分割成M等份,Map调用能够在多台机器上分布执行。这M等分数据在不同机器上可以并行处理。通过使用分割函数(比如 hash(key) mod R)把中间形式的key空间分成R份,Reduce调用能够分布执行。分区数目和分割函数是由用户指定。
图1展示了整个MapReduce的操作流程。当用户程序调用MapReduce函数,将执行如下步骤(图上的数字号码与如下列表中的数字号码是一一对应的):
  1. MapReduce库首先把输入文件分成大小为16-64M之间块大小的M份(通过一个可选的用户参数控制),然后它在多台机器上启动程序副本。
  2. 其中存在于master上的一个副本是特殊的,其余的worker都将由master分配工作。
  3. 被分配map任务的worker将读取输入份中的内容。它从输入数据中解析key/value对,然后把这些对传给用户定义的Map函数。产生的中间key/value对将由内存缓存。
  4. 处在缓存中的对将被间隔性的写入磁盘,这些对将散步在由分割函数指定的R个区域中。这些缓存对在局部磁盘的具体位置被传回给master,然后master负责把这些信息转发给reduce worker。
  5. 当reduce worker从master得到位置信息,它将使用远程过程调用从map worker的局部磁盘中读取数据。读取完所有数据之后,reduce worker按照中间key进行排序,以使得相同值的key被分在一起。因为许多不同的key有可能映射到相同的reduce任务,所以排序是必须的。当 数据太大而不能将它们全部放入内存时,一种外部排序方法将被使用。
  6. reduce worker遍历排好序的中间数据,当它遇到一个唯一的中间形式key,它将把key和与它对应的中间value传给用户提供的Reduce函数。Reduce函数的输出将被附加到与这个reduce对应的最终输出文件中。
  7. 当所有的map和reduce任务被完成之后,master唤醒用户程序。于是MapReduce调用返回到用户代码。
当成功完成之后,输出结果将存在于R个输出文件中,每一个reduce任务对应一个。一般来说,用户没有必要合并这R个文件为一个。这些文件将被传给另一个MapReduce调用当作输入,或者能够处理输入被分割的其余形式的分布式程序。
  • master数据结构
master保存有多种数据结构。对于每个map和reduce任务,它存储有它们的状态(空闲,正在处理或者完成)和非空闲worker所在机器的标志。
关于中间文件的位置信息是通过master从map任务传递给reduce任务的。因此,对于每个完成的map任务,master存储有它所产生的R文件 的位置和大小信息。一旦map任务完成,master就更新这些信息,然后把这些信息传给正在处理中的reduce任务。
  • 容错
因为MapReduce库是用来在成百上千台机器上处理大量数据,所以库必须很好的进行容错处理。
  1. worker故障
master间隔性的ping worker。当尝试多次之后没有响应,master将认为worker已经出现故障。任何已经在wokrer上完成的任务将被设置为空闲状态,以使得它 可以重新被调度。同样,正在处理的map或reduce任务也将被设置为空闲状态,以便重新调度。
完成的map任务需要重新执行,那是因为它们的输出是存储在出现故障机器上的本地磁盘,而导致不可访问。完成的reduce任务输出结果是存储在全局文件系统而不存在这个问题。
当一个map任务首先在woker A上执行,然后在worker B上(因为A出现故障),所有执行reduce任务的worker将得到重新执行的通知。对于还没有从worker A上读取数据的reduce任务将从woker B上读取。
MapReduce对worker故障具有很强的适应能力。比如在一次MapReduce操作中,网络维护一次性造成80台机器变得不可抵达。此时MapReduce重新执行这些出现问题机器上的任务,以至最终完成任务。

参考:(1) 什么是MapReduce? Google的分布运算开发工具!

MapReduce:Simplified Data Processing on large Clusters-翻译版(上)

  • 概括
MapReduce既是一种编程模型,也是用来处理和生成大数据集的一种实现。用户指定的map函数操作key/value对以生成中间形式的key/value对;指定的reduce函数把具有相同中间形式key的value合并在一起。许多现实中的任务都可以用这种模型进行表述,在这篇文章中,将看到这方面的任务。
以MapReduce方式进行编写的程序可以自动并行化,运行在大规模集群系统上。运行系统替用户处理输入数据的分割,程序在一系列机器上的调度执行,故障处理,机器间的通讯。这样使得无并行和分布式经验的程序员利用此系统很容易的采用系统资源。
我们的MapReduce实现系统运行在用商用PC搭建的大规模集群系统上。它具有很好的可扩展性,比如可以对数TB的数据在成千上万机器上进行计算。程序员发现系统的可用性很好。在Google,已经有成百上千的基于MapReduce的应用程序。每一天,将有超过1千左右的MapReduce作业运行在Google的集群上。
  • 简单介绍
在过去的5年中,本文作者和其余在Google的开发者实现了用来处理大量原始数据(比如被抓取的文档,Web请求日志,等等)的各种特殊目的计算。比如反向索引,web文档图形结构的表示,基于每台主机被抓取页面数量的统计,在指定一天中最经常使用的查询,等等。大部分计算从原理上讲是很简单的。但是随着输入数据的加大,这些计算不得不分发在成百上千台计算机上进行计算,以便计算任务在一个可接受的时间限度内完成。这时本来很简单的应用程序不得不处理并行计算,分发数据,处理故障等问题而导致其变得晦涩难懂。
为了应付这种复杂性,我们设计了一种新的抽象。它可以使得用户表达自己的简单计算,而把并行化,容错,数据分发和负载平衡的细节全部隐藏在库中。此抽象的 灵感来自于Lisp和其它函数语言中的map和reduce操作。我们意识到大部分的计算都包括:对输入中的逻辑“记录”进行map操作,以获得一系列的 中间形式的key/value对;对共享同一key的中间值进行reduce操作,以并其合并。用户提供map和reduce操作,然后运用此函数模型, 可以使得我们非常容易进行大规模的并行计算,通过重新执行的方式来达到容错的目的。
这项工作最主要的贡献是提供了一个简单且强大的接口,通过这个接口实现,使得大规模计算得到自动并行与分布,在大规模集群系统上达到高性能。
第2节描述了基本的编程模型,另外给出了Google内部的一些实例;第3节描述了为集群系统量身定制的MapReduce接口实现;第4节阐述了一些对 现存编程接口模型十分有用的改进;第5节提供了基于各种任务,现有实现的性能测试;第6节说明了MapReduce在Google内部的使用,阐述了利用 它对现有索引系统进行重写所获得的一些经验;第7节讨论了相关及将来工作。
  • 编程模型
计算是把一系列key/value对当作输入,产生一系列key/value对。使用MapReduce库的用户用Map和Reduce表示计算过程。
用户提供的Map获得输入对,产生一系列的中间形式的key/value对。MapReduce库把中间形式key值为I的所有value集中在一起,然后把它们传给Reduce函数。
用户提供的Reduce函数,接受中间形式key值I和与它关联的一系列value,然后把这些值合并产生一个可能相对较小的value集合。一般来说,每个调用产生0个或1个输出值。中间形式的value是通过迭代器的形式提供给reduce函数。这样使得我们可以处理那些太大而不能一次性载入内存的value。

  1. 样例

  2. 考虑一个在大量文档中统计每个单词出现次数的问题。用户将会用如下类似的伪代码进行描述。
    map(String key,String value):
    //key: document name
    //value:document contents
    for each word w in value:
    EmitIntermediate(w,"1");

    reduce(String eky,Iterator values):
    //key: a word
    //values: a list of counts
    int result = 0;
    for each v in values:
    result += ParseInt(v);
    Emit(AsString(result));

    map函数产生一个单词和它出现次数的值(在这个样例中是1)。reduce函数对每个特定单词累加其出现次数。
    此外,用户代码提供输入与输出文件名,各种可选的可调整的参数以满足mapreduce specification对象需要。接着把对象传递给MapReduce函数,进行调用。此时,用户代码链接用C++实现的MapReduce库。附录A包含这个样例的完整程序。

  3. 类型

  4. 尽管前面的伪代码输入输出类型都是字符串,但从理论上说,map和reduce函数在类型处理上具有一定关联性。
    map (k1,v1) -> list(k2,v2)
    reduce (k2,list(v2)) ->list(v2)
    如:输入key/value与输出key/value是处在不同的域,中间的key/value与输出key/value是处在同一域。
    我们的C++实现都是把字符串作为函数的输入输出。这就需要用户代码在字符串与其它类型之间进行正确的转换。

  5. 更多实例

  6. 具体包括:分布grep,URL访问频率统计,web连接图反转,每台机器的词矢量,反向索引,分布排序。


参考:(1) 什么是MapReduce? Google的分布运算开发工具!