【转载】MapReduce:大型集群上的简化数据处理(MapReduce: Simplified Data Processing on Large Clusters)

MapReduce: Simplified Data Processing on Large Clusters

MapReduce:大规模集群上的简化数据处理

Jeffrey Dean, Sanjay Ghemawat
Google, Inc.


分类与主题描述

D.4.1 [操作系统]:进程管理——并发;分布式系统

通用术语

性能、设计、可靠性、算法

关键词

MapReduce、分布式处理、大规模集群


摘要

MapReduce是一种编程模型,也是一种用于处理和生成大规模数据集的相关实现。用户指定一个map函数来处理键值对,生成一组中间键值对,再指定一个reduce函数来合并所有关联到同一个中间键的中间值。

许多现实世界中的任务都可以用这个模型来表达,本文将展示这一点。

用这种函数式风格编写的程序可以自动在大规模商用机器集群上并行执行。运行时系统负责处理输入数据的分区、在一组机器上调度程序执行、处理机器故障,以及管理所需的机器间通信。这使得没有并行和分布式系统经验的程序员也能轻松利用大型分布式系统的资源。

我们的MapReduce实现运行在由数百台机器组成的大型集群上,具有高度的可扩展性:一个典型的MapReduce计算会在数千台机器上处理数TB的数据。程序员发现这个系统易于使用:目前已经实现了数百个MapReduce程序,每天在Google的集群上运行的MapReduce作业超过一千个。


1. 引言

在过去的五年里,作者和许多Google的工程师一起实现了数百个专用的大规模数据处理程序,处理了海量的原始数据,包括抓取的文档、网页请求日志等。我们从中计算得到了各种衍生数据,例如倒排索引、网页文档的各种表示形式、每个主机的抓取页面数量统计、按日期汇总的最频繁查询集合等。大多数此类数据处理在概念上都很简单。然而,由于输入数据量庞大,并且计算需要分布在数百台机器上才能在合理的时间内完成,因此它们最终都会变得复杂。为了处理并行计算、数据分发和错误处理等问题,原本简单的计算逻辑中会掺入大量复杂的代码。

为了应对这种复杂性,我们设计了一种新的抽象,它让我们只需表达简单的计算逻辑,而将并行化、容错、数据分发和负载均衡这些底层细节封装在一个库中。这种抽象的灵感来自于Lisp和许多其他函数式语言中常见的mapreduce原语。我们意识到,大多数计算都包含这样的操作:对输入中的每个逻辑记录应用一个map操作,以生成一组中间键值对,然后对所有共享同一个键的中间值应用一个reduce操作,以此来合并派生的数据。

通过让用户自定义mapreduce函数,我们的MapReduce运行时系统能够自动地、透明地将计算并行化,并在大规模机器集群上执行。系统的运行时系统提供的主要功能包括:自动并行化与分布、容错、调度以及机器间通信。

使用MapReduce抽象的好处有很多:

  1. 易于使用:程序员无需具备并行和分布式系统的经验就能轻松使用。
  2. 可扩展性强:每个MapReduce计算都可以扩展到数千台机器上运行。
  3. 容错性好:系统能够自动处理机器故障。

本文的主要贡献是提出了一个简单而强大的接口,它能够自动并行化和执行大规模数据处理计算,并且结合该接口的实现,在大型商用计算机集群上实现了高性能。

第2节描述基本的编程模型并给出一些示例。第3节描述MapReduce的实现,它针对我们的集群计算环境进行了定制。第4节描述我们发现对编程模型非常有用的几个改进。第5节给出了针对各种任务的MapReduce性能测量结果。第6节讨论了MapReduce在Google内部的使用情况,包括我们在将其用作生产索引系统重写基础时的一些经验。第7节介绍相关工作,第8节进行总结。


2. 编程模型

计算任务接收一组键值对作为输入,生成一组键值对作为输出。MapReduce库将计算过程表示为两个函数:MapReduce

  • Map:由用户编写,接收一个输入键值对,生成一组中间键值对。MapReduce库将所有关联到同一个中间键I的中间值聚合在一起,然后传递给Reduce函数。
  • Reduce:同样由用户编写,接收一个中间键I和该键对应的一组中间值。它将这些值合并,形成一个更小的值集合。通常,每次Reduce调用只会产生零个或一个输出值。中间值集合通过迭代器的方式提供给Reduce函数,这样我们就可以处理内存无法容纳的大规模数值集合。

2.1 示例

考虑统计大量文档中每个单词出现次数的问题。用户会写出类似下面的伪代码:

map(String key, String value):
  // key: document name
  // value: document contents
  for each word w in value:
    EmitIntermediate(w, "1");

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

Map函数为每个文档中的每个单词生成一个键值对(单词,出现次数1)。Reduce函数将同一个单词的所有出现次数累加起来,输出(单词,总次数)。

此外,用户还需要编写代码来指定输入和输出文件的位置,以及可选的调优参数。然后用户调用MapReduce函数,将用户定义的Map和Reduce函数传递给它,再加上配置参数。用户的代码与MapReduce库(用C++实现)链接在一起。附录A包含了这个示例的完整程序。

2.2 类型

尽管前面的示例使用字符串类型的输入输出,但从概念上讲,用户定义的Map和Reduce函数都关联着一组类型:

map    (k1, v1)       → list(k2, v2)
reduce (k2, list(v2)) → list(v2)

也就是说,输入键和值的类型(k1, v1)与输出键和值的类型不同。而中间键和值的类型(k2, v2)与输出键和值的类型相同。

我们的C++实现中,所有的键值对都以字符串的形式在用户函数之间传递,由用户代码负责在字符串和适当的类型之间进行转换。

2.3 更多示例

这里再举几个简单的示例,展示MapReduce模型的通用性:

  1. 分布式Grep:Map函数输出匹配模式的行。Reduce函数只是将中间数据原样复制到输出。
  2. URL访问频率统计:Map函数处理网页请求日志,输出(URL, 1)。Reduce函数将同一个URL的所有值累加,输出(URL, 总访问次数)。
  3. 反向网页链接图:Map函数输出每个目标URL对应的源URL,即(target, source)对。Reduce函数将同一个目标URL对应的所有源URL拼接成一个列表,输出(target, list(source))。
  4. 每主机词向量:词向量总结了一篇文档、一组文档或一个主机上所有网页中出现的最重要的单词,形式为(单词, 频率)对。Map函数为每个输入文档(主机名是从文档URL中提取的)输出(主机名, 词向量)。Reduce函数接收给定主机的所有词向量,将它们相加,丢弃低频词,然后输出最终的(主机名, 词向量)对。
  5. 倒排索引:Map函数解析每个文档,输出(单词, 文档ID)序列。Reduce函数接收给定单词的所有文档ID,对它们进行排序,并输出(单词, 列表(文档ID))。所有输出的集合形成一个简单的倒排索引,可以轻松跟踪每个单词在文档中的位置。
  6. 分布式排序:Map函数从每条记录中提取键,输出(键, 记录)对。Reduce函数原样输出所有键值对。这个计算依赖于我们在4.1节描述的分区机制和4.2节描述的排序保证。

3. 实现

MapReduce有很多不同的实现,适用于不同的场景。例如,一个小的共享内存实现、一个NUMA多处理器实现,以及一个更大的互联网连接机器集群实现。

本节描述的是针对Google广泛使用的计算环境定制的实现:通过以太网交换机连接的大量商用PC组成的集群。在我们的环境中:

  • 机器通常是双路x86处理器,运行Linux操作系统,内存为2-4GB。
  • 使用商用级别的IDE磁盘。一个分布式文件系统(GFS)用于管理存储在这些磁盘上的数据。GFS使用副本机制在不可靠的硬件之上提供可用性和可靠性。
  • 用户将作业提交到调度系统。每个作业由一组任务组成,由调度器映射到集群中一组可用机器上执行。

3.1 执行概述

通过自动将输入数据划分为M个分片(split),Map调用被分布到多台机器上执行。每个输入分片都可以由不同的机器并行处理。

Reduce调用则通过分区函数(例如hash(key) mod R)将中间键空间划分为R个分片,同样分布到多台机器上执行。分区数量R和分区函数由用户指定。

图1展示了MapReduce操作的完整流程。当用户程序调用MapReduce函数时,会发生以下一系列动作(图中的数字对应下面的步骤):

  1. 用户程序中的MapReduce库首先将输入文件划分为M个大小通常为16MB到64MB的分片(可由用户通过可选参数控制)。然后,它会在集群的机器上启动程序的多个副本。
  2. 这些副本中有一个是特殊的——主节点(master)。其余的都是工作节点(worker),由主节点分配任务。总共有M个Map任务和R个Reduce任务需要分配。主节点会为每个空闲的工作节点分配一个Map任务或Reduce任务。
  3. 被分配了Map任务的工作节点读取对应的输入分片。它从输入数据中解析出键值对,并将每个键值对传递给用户定义的Map函数。Map函数生成的中间键值对被缓存在内存中。
  4. 缓存的中间键值对会定期写入本地磁盘,通过分区函数划分为R个区域。这些缓冲数据在本地磁盘上的位置会被传回给主节点,主节点再将这些位置转发给执行Reduce任务的工作节点。
  5. 当Reduce任务的工作节点从主节点接收到这些位置信息后,它会使用远程过程调用(RPC)从Map任务工作节点的本地磁盘上读取缓冲的中间数据。当Reduce工作节点读取完所有中间数据后,它会根据中间键对数据进行排序,使得所有相同键的键值对都聚合在一起。之所以需要排序,是因为通常有许多不同的键会映射到同一个Reduce任务。如果数据量太大无法在内存中排序,则使用外部排序。
  6. Reduce任务的工作节点遍历排序后的中间数据,对于遇到的每个唯一的中间键,它会将该键和对应的中间值集合传递给用户定义的Reduce函数。Reduce函数的输出被追加到对应Reduce分区的输出文件中。
  7. 当所有的Map任务和Reduce任务都完成后,主节点唤醒用户程序。此时,用户程序中的MapReduce调用返回。

正常执行完成后,MapReduce计算的输出存放在R个输出文件中(每个Reduce任务对应一个,文件名由用户指定)。通常,用户不需要合并这R个输出文件——它们通常会作为另一个MapReduce计算的输入,或者在另一个可以处理分片输入的分布式应用中使用。

图1:执行流程概览

3.2 主节点数据结构

主节点维护着多种数据结构。对于每个Map任务和Reduce任务,它存储任务的当前状态(空闲、进行中、已完成),以及分配的工作机器(对于非空闲任务)。

主节点就像一个数据管道:Map任务完成时生成的中间文件区域的位置信息,通过主节点从Map任务传递给Reduce任务。因此,对于每个已完成的Map任务,主节点会存储该Map任务生成的R个中间文件区域的位置和大小。当Map任务完成时,这些信息会被增量式地推送给主节点。主节点再将这些信息转发给正在执行对应Reduce分区的、处于进行中的Reduce任务。

3.3 容错

由于MapReduce库被设计为可以使用数百台机器的协作来处理数TB级的数据,因此系统必须能够优雅地处理机器故障。

3.3.1 工作节点故障

主节点会周期性地向每个工作节点发送ping消息。如果在一定时间内没有收到工作节点的回复,主节点就会将该工作节点标记为失效。

任何由该失效工作节点完成的Map任务都会被重置回空闲状态,因此可以重新调度到其他工作节点上。同样,任何正在由失效工作节点执行的Map任务或Reduce任务也都会被重置为空闲状态。

已完成的Map任务之所以需要重新执行,是因为它们的输出存储在失效机器的本地磁盘上,因此无法访问。而已完成的Reduce任务不需要重新执行,因为它们的输出存储在全局文件系统中。

当一个Map任务因为工作节点A失效而被重新执行时,所有正在执行对应Reduce任务的工作节点都会收到重新执行的通知。任何还没有从工作节点A读取数据的Reduce任务,都会从新的工作节点读取数据。

MapReduce能够应对大规模工作节点失效的情况。例如,在一次大型MapReduce计算过程中,网络维护导致一组机器在几分钟内不可达。MapReduce主节点只是简单地重新执行了那些由不可达机器完成的工作,然后继续推进,最终成功完成了计算。

3.3.2 主节点故障

主节点会定期将它维护的数据结构写入检查点。如果主节点任务失效,可以从最后一个检查点启动一个新的主节点副本。然而,由于只有一个主节点,它的失效可能性很小。因此我们当前的实现中,如果主节点失效,就中止MapReduce计算。用户可以检查这种情况,并根据需要重试MapReduce操作。

3.3.3 故障场景下的语义

当用户提供的Map和Reduce函数是关于输入值的确定性函数时,我们的分布式实现与程序的无故障串行执行产生的输出相同。

我们通过依赖Map任务和Reduce任务输出的原子提交来实现这个特性。每个进行中的任务都会将其输出写入临时输出文件。一个Reduce任务生成一个输出文件,一个Map任务生成R个输出文件(每个Reduce任务对应一个)。当一个Map任务完成时,它会通知主节点,并附带R个临时文件的名称。如果主节点收到一个已经完成的Map任务的完成通知,它会忽略该通知。否则,它会在主数据结构中记录这R个文件的名称。

当一个Reduce任务完成时,它会将其临时输出文件原子地重命名为最终输出文件。如果同一个Reduce任务在多台机器上执行,多个重命名操作会并发地作用于同一个最终输出文件。我们依赖底层文件系统提供的原子重命名操作,来保证最终的文件系统状态只包含一个Reduce任务执行产生的数据。

确定性函数的情况下,我们的模型保证输出是等价的。对于非确定性函数,我们提供较弱但仍然合理的保证。当存在非确定性函数时,给定Reduce任务的输出可能等于针对该Reduce任务的某一个串行执行的输出,但不同的Reduce任务的输出可能对应不同的串行执行。

3.4 本地性

网络带宽是我们计算环境中相对稀缺的资源。我们通过将输入数据存储在集群中机器的本地磁盘上,并尽量将Map任务调度到包含对应输入数据副本的机器上执行,来节省网络带宽。

主节点在调度Map任务时,会考虑输入文件的位置信息。理想情况下,它会将Map任务分配到持有对应输入数据副本的机器上。或者,如果做不到,它会利用GFS的副本位置信息,将任务调度到离输入数据副本较近的机器上(例如,同一机架的机器)。当大型MapReduce计算在集群的大部分机器上运行时,大多数输入数据都在本地读取,几乎不消耗网络带宽。

3.5 任务粒度

如前所述,我们将Map阶段划分为M个分片,Reduce阶段划分为R个分片。理想情况下,M和R都应该远大于集群中的机器数量。每个工作节点执行多个不同的任务可以提高负载均衡,并且当工作节点失效时,也可以加快恢复速度:它完成的所有Map任务都可以分散到所有其他工作节点上重新执行。

在我们的实现中,M和R的大小是有实际限制的,因为主节点必须做出O(M+R)次调度决策,并且内存中要保存O(M*R)个状态(不过这个状态占用的内存很小)。

此外,R通常受限于用户,因为每个Reduce任务的输出最终都会成为一个独立的输出文件。实际上,我们通常选择M的大小,使得每个单独的任务大约处理16MB到64MB的输入数据(这样可以使本地性优化达到最佳效果),而R的值则设置为用户期望的输出文件数量。我们通常使用的M和R的比例是,当集群中有几千台机器时,M=20000,R=5000。

3.6 备份任务

导致MapReduce总执行时间变长的常见原因之一是“掉队者”(straggler):集群中某台机器完成最后几个Map或Reduce任务的速度异常缓慢。掉队者可能由多种原因导致。例如,一台磁盘有问题的机器可能会因为频繁的纠错而将读取性能从30MB/s降低到几MB/s。再比如,集群调度系统可能将其他任务调度到了这台机器上,导致MapReduce代码的执行速度因为CPU、内存、本地磁盘或网络带宽的竞争而变慢。我们最近遇到的一个问题是,机器初始化代码中的一个bug导致处理器缓存被禁用,这使机器的计算性能下降了一个数量级。

我们有一个通用的机制来缓解掉队者导致的问题。当一个MapReduce操作接近完成时,主节点会将正在进行中的任务调度到其他空闲机器上执行备份。无论主副本还是备份副本完成,任务都会被标记为已完成。我们对这个机制进行了微调:它通常只会使计算使用的机器资源增加几个百分点,但却能显著缩短大型MapReduce计算的总执行时间。例如,在排序程序中,当启用备份任务机制时,完成时间缩短了44%。


4. 改进

基本的编程模型简单易用,但我们发现还可以进行一些补充和改进,使其适用于更广泛的问题。本节将描述这些改进。

4.1 分区函数

MapReduce的用户指定Reduce任务的数量R,以及中间键到Reduce任务的分区函数。默认的分区函数使用哈希取模(hash(key) mod R),这通常能产生非常均衡的分区。

然而,有时候用户希望按键的其他属性进行分区。例如,有时输出的键是URL,用户希望同一个主机的所有条目最终都在同一个输出文件中。为了支持这种情况,用户可以提供自定义的分区函数。例如,使用hash(Hostname(urlkey)) mod R作为分区函数,可以让所有来自同一个主机的URL都被划分到同一个Reduce任务中。

4.2 排序保证

我们保证,在给定的Reduce分区内,中间键值对是按照键的升序排列的。这种排序保证使得在分区内生成排序后的输出变得容易,当输出文件需要支持按键随机访问查找时,或者当用户发现排序后的输出更容易处理时,这一点非常有用。

4.3 合并器函数

在某些情况下,每个Map任务都会产生大量的重复键,而用户定义的Reduce函数满足交换律和结合律。这就允许使用**合并器(Combiner)**函数。

合并器函数在每个执行Map任务的工作节点上本地执行。通常,合并器函数的实现与Reduce函数的实现是一样的,区别只在于MapReduce库如何处理函数的输出。Reduce函数的输出被写入最终输出文件。而合并器函数的输出被写入将要发送给Reduce任务的中间数据。

例如,在单词统计的例子中,由于单词出现的频率服从齐普夫定律,每个Map任务都会生成大量形如(the, 1)的键值对。这些数据会通过网络发送给单个Reduce任务,然后由Reduce函数将所有值加起来得到一个数字。通过使用合并器,每个Map任务都会先在本地将重复键的值合并,从而大大减少需要通过网络传输的数据量。

4.4 输入输出类型

MapReduce库支持多种不同格式的输入数据。例如,文本模式输入将每一行视为一个键值对:键是行在文件中的偏移量,值是行的内容。另一种常见的格式是以键排序的序列,其中每个键值对都按顺序存储。每种输入类型的实现都知道如何将输入划分为适合单独处理的分片。

用户也可以通过实现一个简单的读取器接口来添加新的输入类型。不过,用户并不总是需要用键值对的方式来访问输入数据,读取器接口的实现有很大的灵活性。

类似地,我们支持多种输出格式来生成不同格式的数据。用户也可以添加自定义的输出格式。

4.5 副作用

在某些情况下,用户发现生成辅助文件或Map和Reduce任务的其他副作用会很方便。我们依赖程序员来检测这种副作用,并使其具有原子性和幂等性。通常,应用会先将结果写入临时文件,然后在计算完全完成后原子地重命名这些文件。

我们不支持单个任务生成的多个输出文件的原子两阶段提交。因此,对于需要跨多个文件具有一致性语义的输出,应该由同一个Reduce任务生成。

4.6 跳过损坏记录

有时,用户代码中的bug会导致Map或Reduce函数在处理特定记录时崩溃。这种情况会干扰MapReduce操作的完成。通常的做法是修复bug,但有时这并不可行——bug可能存在于第三方库中,而我们没有源代码。此外,有时忽略少数有问题的记录是可以接受的,例如在对大型数据集进行统计分析时。

我们提供了一种可选的执行模式,在这种模式下,MapReduce库会检测到导致崩溃的记录,然后跳过这些记录以继续执行。

每个工作进程都安装了一个信号处理器,用于捕获段错误和总线错误。在调用用户定义的Map或Reduce操作之前,MapReduce库会将当前处理的记录序号保存在全局变量中。如果用户代码触发了信号,信号处理器就会向主节点发送一个包含序号的”最后一口气”UDP包。当主节点看到某个特定记录多次失败时,它就会标记该记录需要跳过,并在下次重新执行相应的Map或Reduce任务时跳过它。

4.7 本地执行

调试Map和Reduce函数的bug可能会很棘手,因为实际的计算是在分布式系统中执行的,通常有几千台机器,而且工作分配是由主节点动态决定的。为了帮助调试、性能分析和小规模测试,我们开发了MapReduce库的另一个实现,它可以在本地机器上串行执行MapReduce计算的所有工作。这样用户就可以很容易地使用他们熟悉的调试工具(如gdb)来调试MapReduce程序。

4.8 状态信息

主节点运行一个内置的HTTP服务器,并显示一组状态页面,用户可以通过这些页面监控计算的进度。状态页面显示计算是否正在进行、已经完成了多少、还有多少没完成,以及处理的字节数、输入字节数、中间数据字节数、输出字节数、处理速率等。页面还包含指向每个任务的标准错误和标准输出的链接,用户可以用这些来调试计算。

进度数据可以用来估计计算还需要多长时间完成,以及是否应该向计算中添加更多资源。这些页面也可以用来找出为什么计算比预期慢。

此外,顶级状态页面显示了哪些工作节点失效了,以及它们失效时正在运行哪些Map和Reduce任务。这对于诊断用户代码中的bug非常有用。

4.9 计数器

MapReduce库提供了一个计数器工具,用于统计各种事件的发生次数。例如,用户可能想要统计处理的单词总数或索引的德语文档数量。

要使用这个工具,用户代码中创建一个命名的计数器对象,然后在Map和/或Reduce函数中递增计数器。例如:

Counter* uppercase;
uppercase = GetCounter("uppercase");

map(String key, String value):
  for each word w in value:
    if (IsCapitalized(w)):
      uppercase->Increment();
    EmitIntermediate(w, "1");

各个工作节点上的计数器值会定期传递给主节点(附带在心跳响应中)。主节点将来自成功的Map和Reduce任务的计数器值聚合起来,并在MapReduce操作完成时返回给用户代码。聚合后的计数器值也会显示在主节点的状态页面上,这样用户就可以看到实时的计算进度。

主节点在聚合计数器值时,会消除重复执行的Map或Reduce任务产生的重复计数。重复执行可能由备份任务机制或工作节点失效导致的重新执行引起。

有些计数器是由MapReduce库自动维护的,例如处理的输入键值对数量和生成的输出键值对数量。

用户发现计数器对于检查MapReduce操作的完整性非常有用。例如,在某些MapReduce操作中,用户可能希望确保生成的输出键值对数量恰好等于处理的输入键值对数量,或者处理的德语文档数量在总文档数量中所占比例在合理范围内。


5. 性能

在本节中,我们通过在一个大型机器集群上运行的几个程序来衡量MapReduce的性能。一些程序在整个集群上进行数据洗牌,另一些则进行大量的计算。

5.1 集群配置

所有程序都在一个由大约1800台机器组成的集群上运行。每台机器有两个2GHz的Intel Xeon处理器,支持超线程,4GB内存,两个160GB的IDE磁盘,以及一个千兆以太网链路。这些机器部署在两层的树形以太网交换机中,根交换机的总带宽大约为100-200Gbps。所有机器都在同一个托管设施中,因此任意两台机器之间的往返时间都小于1毫秒。

尽管有4GB内存,但大约1-1.5GB内存被集群中运行的其他任务占用了。程序在中午左右执行,此时机器大多空闲。

输入数据、中间数据和输出数据都存储在GFS中。我们使用的GFS副本因子为3(即每个文件有三个副本)。

5.2 查找(Grep)

查找程序扫描大约10¹⁰个100字节的记录(总共约1TB的数据),查找出现次数相对较少的三字符模式(模式在92337个记录中出现)。Map函数将匹配的记录输出。Reduce函数是恒等函数,只是将中间数据原样复制到输出。

输入数据被分成大约64MB的分片(M=15000),因此整个计算分布在大约1746台机器上执行。所有输入数据都在本地磁盘上,因此几乎不消耗网络带宽。

整个计算在150秒内完成,包括启动开销。启动开销大约是1分钟,主要是将程序分发到所有工作机器上的时间。

5.3 排序

排序程序对1TB的数据进行排序。排序程序使用不到50行的用户代码:三行Map函数从每行中提取一个10字节的键,然后输出(键, 行)对。标准的恒等Reduce函数作为排序操作的Reduce函数。

我们使用默认的分区函数(键的哈希),它会产生均匀分布的键。更一般地说,对于排序,我们会使用一个分区函数,它根据键的采样值来分割键空间。

和查找程序一样,输入数据被分成64MB的分片(M=15000)。我们使用4000个Reduce任务,因此会生成4000个输出文件。

图2展示了排序程序执行过程中数据通过系统的速率。左上方的图显示了输入数据被读取的速率。数据读取速率峰值大约为13GB/s,当所有Map任务完成时,速率迅速降为0。请注意,读取速率比排序的总吞吐量高,因为排序花费了大约一半的时间在网络上传输中间数据。

中间的图显示了中间数据从Map任务通过网络发送到Reduce任务的速率。第一个Reduce任务在Map任务开始后大约20秒开始启动。图中的第一个低谷是因为第一批Reduce任务完成了,而我们还在等待更多的Map任务完成,以生成更多的数据来排序。

右下方的图显示了Reduce任务将排序后的数据写入最终输出文件的速率。在第一个排序阶段结束和写入开始之间有一个延迟,因为机器正忙于对中间数据进行排序。写入速率峰值大约为2-4GB/s。

图2:排序程序的数据传输时间

整个排序过程从开始到结束总共花费了891秒。这与TeraSort基准测试的最佳结果非常接近。例如,使用1024台Iron机器(每台比我们的机器稍快)的TeraSort最佳记录是1057秒。

有几件事值得注意:

  1. 输入速率远高于中间数据传输速率和输出速率,这是因为我们的本地性优化——大多数输入数据都是从本地磁盘读取的,绕过了相对较慢的网络。
  2. 中间数据传输速率远高于输出速率。这是因为输出数据需要写入两个副本(我们使用GFS,副本因子为2,以实现可靠性和可用性)。输出数据写入两个副本,因此消耗的磁盘写入带宽大约是中间数据读取消耗的磁盘读取带宽的两倍。
  3. Map阶段的持续时间比Reduce阶段长。这部分是因为Map阶段处理的输入数据更多(输入数据加上中间数据的开销)。

5.4 反向网页链接图

反向链接图程序使用大约20亿个网页(大约140GB的数据)作为输入,生成反向链接图——每个网页对应一个包含所有链接到它的网页的列表。程序代码不到50行。Map函数输出(目标URL, 源URL)对。Reduce函数将给定目标URL的所有源URL拼接成一个列表,并输出(目标URL, 列表(源URL))对。

这个计算的输出大约是同样大小的另一个数据集(大约140GB),被存储为2000个文件(R=2000)。整个计算耗时1500秒。

5.5 每主机词向量

词向量计算为每个主机生成一个词向量,输入是大约90000个HTML文档(大约18GB)。输出是大约300000个(主机名, 词向量)对,总共大约2GB。

Map函数为每个文档生成一个(主机名, 词向量)对。Reduce函数接收给定主机的所有词向量,将它们相加,丢弃低频词,然后输出最终的(主机名, 词向量)对。

整个计算耗时234秒。

5.6 排序(重复测试)

我们重新运行了排序程序,这次使用了R=1500个Reduce任务和M=15000个Map任务。输出数据的副本因子为2,因此输出数据量大约是2TB(而不是1TB)。

整个计算耗时1007秒。

5.7 备份任务的影响

为了衡量备份任务的效果,我们在禁用备份任务的情况下重新运行了排序程序。图3显示了正常执行和禁用备份任务执行的进度。

上方的曲线显示了正常执行的进度。x轴是时间(秒),y轴是已完成的任务百分比(Map任务和Reduce任务分别显示)。可以看到,Map任务在大约800秒时几乎全部完成,但剩下最后几个任务花了很长时间才完成。Reduce任务在大约600秒时开始启动,在大约1300秒时几乎全部完成。整个计算在1283秒时完成(不包括将输出文件提交到GFS的时间)。

下方的曲线显示了禁用备份任务时的执行进度。可以看到,计算在进行了1700秒后,仍然有几个Reduce任务没有完成。在1744秒时,我们中止了计算,因为进度已经停滞不前了。

掉队者导致了巨大的性能差异。启用备份任务后,排序程序的执行时间缩短了44%。

图3:1500GB排序的归一化执行时间

5.8 机器故障

为了测试故障恢复能力,我们在排序程序执行过程中故意杀死了200个工作进程中的一组。集群有足够的备用容量,可以立即替换这些机器。

杀死200个工作进程导致计算时间增加了大约5%。这比我们预期的要多,因为被杀死的工作进程中有一些恰好在执行接近完成的任务,这些任务必须重新执行,导致了额外的延迟。


6. 经验

MapReduce于2003年2月部署完成。在撰写本文时,它已经被Google内部的各个领域的大量程序员使用了大约一年半。我们已经实现了超过300个不同的MapReduce程序,每天在Google的集群上运行的MapReduce作业超过一千个。

表1展示了2004年8月期间运行的MapReduce作业的汇总统计数据。

指标 数值
作业总数 29,423
平均作业完成时间 634秒
消耗的机器天数 79,186天
读取的输入数据总量 3,288 TB
生成的中间数据总量 758 TB
写入的输出数据总量 193 TB
每个作业平均使用工作节点数 157台
每个作业平均工作节点故障数 1.2台
每个作业平均Map任务数 3,351个
每个作业平均Reduce任务数 55个
不同的Map实现数量 395个
不同的Reduce实现数量 269个
不同的Map/Reduce组合数量 426个

表1:2004年8月运行的MapReduce作业统计

6.1 大规模索引

MapReduce最令人兴奋的用途之一是重写了Google的网页搜索服务所使用的索引系统。索引系统将大量的网页文档(通过爬虫系统获取)作为输入,生成用于搜索查询的倒排索引。

索引代码是作为一系列MapReduce操作来实现的。我们使用MapReduce带来了几个好处:

  1. 索引代码更简单、更小、更容易理解,因为处理容错、分发和并行化的细节都隐藏在MapReduce库中了。例如,使用MapReduce前后,与计算相关的代码行数减少了大约70%。
  2. 性能已经足够好,因此我们可以用更简单、计算成本更高的算法来替代复杂的、为效率而调整的算法。
  3. 操作变得更加容易。大多数由机器故障、机器速度慢以及网络瞬态问题导致的问题都由MapReduce库自动处理,无需操作人员干预。

6.2 重构的经验

在使用MapReduce的过程中,我们学到了一些东西:

  1. 限制编程模型是一件好事:它让并行化和分布式计算变得容易,并且使计算具有容错性。
  2. 网络带宽是稀缺资源:我们的许多优化都是为了减少通过网络发送的数据量:本地性优化允许我们从本地磁盘读取数据,从而节省网络带宽;合并器函数也大大减少了中间数据的网络传输量。
  3. 可以用简单的工具来完成许多事情:我们惊讶地发现,使用MapReduce可以轻松解决如此多的问题。
  4. 代码重用非常容易:MapReduce库使我们能够快速试验新的算法,并在大规模数据上运行它们。

7. 相关工作

许多系统都提供了用于并行编程的受限编程模型,并自动实现并行化。例如,MPI提供了消息传递原语,使程序员能够开发并行程序,但与MapReduce相比,它需要程序员掌握更多的并行编程知识。

我们的分区机制与一些并行数据库系统中使用的水平分区类似。我们的排序保证与有序分区相结合,类似于有序分区并行数据库中的功能。

我们的本地性优化与主动磁盘(Active Disks)的工作原理类似,通过将计算推送到数据所在的位置来减少网络流量。

我们的备份任务机制与TACC系统中使用的”reaper”进程类似,用于检测掉队者并处理它们。不同之处在于,我们的备份任务机制可以应对由硬件故障、软件bug或资源竞争导致的掉队者。

River系统提供了一个编程模型,其中数据流通过进程图流动。与MapReduce不同,River的重点是在大型集群上实现均衡的数据流,而不是通过抽象来简化分布式计算。

还有许多系统实现了特定类型的并行计算。例如,FASTSort是一个外部排序算法,可以在大型机器集群上高效运行。与这些专用系统相比,MapReduce适用于更广泛的问题。


8. 结论

MapReduce编程模型在Google已经成功应用于许多不同的领域。我们将其成功归因于几个原因。

首先,模型易于使用,即使对于没有并行和分布式系统经验的程序员来说也是如此,因为它隐藏了并行化、容错、局部性优化和负载均衡的细节。

其次,各种各样的问题都可以用MapReduce来表达。例如,MapReduce被用于生成Google网页搜索服务的数据、排序、数据挖掘、机器学习以及许多其他系统。

第三,我们已经能够将计算部署到由数千台机器组成的大型集群上,并且能够很好地扩展。这种实现能够很好地应对机器故障,易于管理,并且具有很高的性能。

根据我们使用MapReduce的经验,我们相信它将成为处理大规模数据计算的重要工具。


参考文献

[1] R. Arpaci-Dusseau, E. Anderson, N. Treuhaft, D. Culler, J. Hellerstein, D. Patterson, and K. Yelick. Cluster I/O with River: Making the fast case rare. In Proceedings of the 1997 ACM/IEEE Conference on Supercomputing, pages 10–22. ACM Press, November 1997.

[2] T. Anderson, D. Culler, and D. Patterson. A case for NOW (Networks of Workstations). IEEE Micro, 15(1):54–64, February 1995.

[3] D. Bitton, D. J. DeWitt, D. K. Hsiao, and J. Menon. A taxonomy of parallel sorting. ACM Computing Surveys, 16(3):287–318, September 1984.

[4] J. S. Chase, D. E. Irwin, L. E. Grit, J. D. Moore, and S. E. Sprenkle. Dynamic virtual clusters in a grid site manager. In Proceedings of the 12th International Symposium on High Performance Distributed Computing, pages 181–194. IEEE Computer Society, June 2003.

[5] D. DeWitt and J. Gray. Parallel database systems: The future of high performance database systems. Communications of the ACM, 35(6):85–98, June 1992.

[6] S. Ghemawat, H. Gobioff, and S.-T. Leung. The Google File System. In Proceedings of the 19th ACM Symposium on Operating Systems Principles, pages 29–43. ACM Press, October 2003.

[7] J. Gray. The transaction concept: Virtues and limitations. In Proceedings of the 7th International Conference on Very Large Data Bases, pages 144–154. IEEE Computer Society, September 1981.

[8] J. Hennessy and D. Patterson. Computer Architecture: A Quantitative Approach. Morgan Kaufmann Publishers, 1996.

[9] D. Joseph and D. Grunwald. Prefetching using Markov predictors. In Proceedings of the 24th International Symposium on Computer Architecture, pages 252–263. ACM Press, June 1997.

[10] L. Lamport. The part-time parliament. ACM Transactions on Computer Systems, 16(2):133–169, May 1998.

[11] M. Litzkow, M. Livny, and M. Mutka. Condor – a hunter of idle workstations. In Proceedings of the 8th International Conference of Distributed Computing Systems, pages 104–111. IEEE Computer Society, June 1988.

[12] C. Nyberg, C. Koester, and J. Gray. Sorting 100 Terabytes with TeraSort. Technical report, NCR Corp., 1997.

[13] D. Patterson, G. Gibson, and R. Katz. A case for redundant arrays of inexpensive disks (RAID). In Proceedings of the 1988 ACM SIGMOD International Conference on Management of Data, pages 109–116. ACM Press, June 1988.

[14] M. Stumm and S. Zhou. Algorithms implementing distributed shared memory. Proceedings of the IEEE, 78(12):1955–1968, December 1990.

[15] S. Tuecke, C. Czajkowski, K. Foster, J. Frey, S. Graham, T. Kesselman, D. Martin, P. Vanderbilt, and D. Welch. Open Grid Services Infrastructure (OGSI) Version 1.0. Technical report, Global Grid Forum, June 2003.

[16] R. van Renesse, K. Birman, and S. Maffeis. Horus: A flexible group communication system. Communications of the ACM, 39(4):76–83, April 1996.

[17] J. Wilkes, R. Golding, C. Staelin, and T. Sullivan. The HP AutoRAID hierarchical storage system. ACM Transactions on Computer Systems, 14(1):108–136, February 1996.


附录A:词频统计程序

#include "mapreduce/mapreduce.h"

// User's map function
class WordCounter : public Mapper {
 public:
  virtual void Map(const MapInput& input) {
    const string& text = input.value();
    const int n = text.size();
    for (int i = 0; i < n; ) {
      // Skip past leading whitespace
      while ((i < n) && isspace(text[i]))
        i++;

      // Find word end
      int start = i;
      while ((i < n) && !isspace(text[i]))
        i++;
      if (start < i)
        Emit(text.substr(start, i-start), "1");
    }
  }
};

REGISTER_MAPPER(WordCounter);

// User's reduce function
class Adder : public Reducer {
 public:
  virtual void Reduce(ReduceInput* input) {
    int64 value = 0;

    // Iterate over all entries with the
    // same key and add the values
    while (!input->done()) {
      value += StringToInt(input->value());
      input->NextValue();
    }

    // Emit sum for input->key()
    Emit(IntToString(value));
  }
};

REGISTER_REDUCER(Adder);

int main(int argc, char** argv) {
  ParseCommandLineFlags(argc, argv);

  MapReduceSpecification spec;

  // Store list of input files in "spec"
  for (int i = 1; i < argc; i++) {
    spec.add_input(argv[i]);
  }
  spec.set_output("wordcounts");
  spec.set_mapper_class("WordCounter");
  spec.set_reducer_class("Adder");

  // Optional: start 8 copies of reduce tasks
  // (defaults to 1 copy)
  spec.set_reduce_tasks(8);

  // Now run the mapreduce computation
  MapReduceResult result;
  if (!MapReduce(spec, &result)) {
    fprintf(stderr, "MapReduce failed: %s\n",
            result.error_string().c_str());
    return 1;
  }
  return 0;
}

Leave a Reply

Your email address will not be published. Required fields are marked *

*