分布式学习1.1

分布式系统学习1.1

本篇是第一篇分布式系统学习(一)的补充,在分布式系统学习(一)中更多是从分布式系统八股的层面入手,更偏向实际应用(或许),而本篇采用MIT6.824课程中的体系,使用更学院派的方式来了解分布式系统。

同时这门课的授课老师是 Robert Morris 教授,世界上第一个蠕虫病毒 Morris 病毒就是出自他之手。

并发与分布式系统

我们为什么需要分布式系统?

首先,不像是我们把主力系统从Windows切换成Linux一样,我们换系统的动机可能只是因为我们觉得Linux比Windows更好用,但Windows系统其实完全可以胜任我们的开发工作,甚至在某些方面比Linux更方便,也就是说我们通常并不是出于无法在Windows上完成我们的工作转而使用Linux,但分布式系统却是因为单机系统已经无法满足需求,而不得不采用的架构,甚至因此我们不得不忍受分布式系统带来的更多的复杂性(用windows和linux举例可能不是很合适,但我暂时想不到更贴切的比喻)

所以在实际业务中采用分布式系统往往是不得已而为之,如果在设计系统时单机系统已经完全可以满足你的需求时,请多思考如何用单机系统实现,因为分布式系统会让解决同一个问题变得复杂很多

现在假设我们有一个单机系统,由一个简单的Web服务器和数据库组成的Web应用对外提供服务,最开始的时候这个单机系统可以完全满足我们的需求,每天稳定有几百人访问你的网站,但机缘巧合下某一天你的网站突然火了,有几万甚至十几万人访问你的网站,可以遇见的是你的服务器撑不住了,现在你有两个选择,重构你的代码,通过极致优化来降低占用,以及加更多的服务器,很显然正常都会选择第二个方案,因为重构服务付出很多的时间成本,所以你选择了第二个方案,问题很快解决了,你不断地购买服务器,把服务放在多台机器上跑,水多了加面,面多了加水(btw,恭喜你搭建了一个集群),所有的用户最终都需要看到相同的数据,所以所有的Web服务器都与后端数据库通信,单台Web服务器对数据库不会造成很大的压力,但随着Web服务器数量的增加,数据库成为了你系统的瓶颈,这时增加更多的服务器已经无济于事了,这时你就不得不考虑做一些重构,采用分布式系统(或许你可以将一个数据库拆分成多个数据库来提升性能,但这个工作量太大了),在取舍下,你不得不采用分布式系统。

所以你将搭建你的第一个分布式系统——分布式数据库/分布式存储系统(分布式数据库和分布式存储不一样,分布式存储和分布式存储系统不一样)

建立分布式系统将要面临的问题

并发的问题,计算机网络的问题等等,在分布式系统(一)中已经详细描述过了

分布式系统要解决的问题/满足的特性

可扩展性(Salability)

回顾我们为什么需要分布式系统中,我们默认你编写的Web服务是具备可扩展性的,这样你才可以通过增加服务器的方式来提升性能,所以构建分布式系统的前提就是你的服务需要满足可扩展性

可用性(Availability)

在分布式系统领域中有很重要的CAP理论,即可用性,一致性,和分区容错,但由于在实际业务中由于网络通信的问题,分区是不可避免的,所以只能在可用性和一致性中做取舍,具体CAP理论将会在下一篇博文中具体描述

如果你只使用一台计算机来构建你的系统,那么你的系统大概率是可靠的。然而如果你通过数千台计算机构建你的系统,那么即使每台计算机可以稳定运行一年,对于1000台计算机也意味着平均每天会有3台计算机故障。

所以在分布式系统中你不得不考虑的是容错,因为错误总会发生,不可避免,所以在系统设计的时候就要考虑在发生错误时保证系统依然可用,即保证可用性。在6.824课程中也提到:

同时,因为我们需要为第三方应用开发人员提供方便的抽象接口,我们的确也需要构建这样一种基础架构,它能够尽可能多的对应用开发人员屏蔽和掩盖错误。这样,应用开发人员就不需要处理各种各样的可能发生的错误。

这一点也体现在课前推荐我们阅读的论文MapReduce: Simplified Data Processing on Large Clusters中,在本篇博文的最后我将会贴出这篇论文的原文和对应的中文翻译

在可用性之下还有一种要求更低的容错是自我可恢复性。即一旦出了问题,服务会停止工作,然后工作人员来修,等修好之后系统继续提供服务。要实现自我可恢复性,首先要保证数据不丢失,从而使修复后提供和修复前完全相同的服务。我们知道程序运行时的数据是存储在内存中的,意外断电就会丢失,所以一个可恢复的系统需要将数据进行持久化,存在硬盘里。对于具备可用性的系统,我们也需要提供可恢复性,因为具备可用性的系统的容错能力是有限的,如果故障太多,系统也会停止响应,我们希望在修好之后依旧提供相同的服务。

要实现可用性,我们通常有两个工具

一个是通过持久化数据到硬盘中,另一个就是存副本。

但可预见的是在多台服务器上,副本的同步就成了问题,我们需要做一系列工作来保证不同服务器上存储的副本数据是一致的(不过这里的一致性只是实际业务中很小的一个例子,无法代表全部的数据一致性要求)

一致性(Consistency)

所以我们要满足一致性,我们现在有一个分布式系统提供服务,我们可能有不同的服务器分布在跨洲的地理范围,我们需要保证用户在访问相隔几千公里的两台服务器时得到的数据是相同的

加入我们现在有一个服务,部署在服务器A和服务器B上,当我们向发送请求给服务器A将一个值由20改为21,但服务器B由于卡顿等原因导致其数值仍是20,这时候如果有用户通过get读取,他可能获取21,也可能获取20.

这时候我们就需要想办法来让我们的系统满足一致性,一致性分为强一致性和弱一致性(也叫最终一致性)

强一致性要求用户获取的数据必须都是最新一次更新,银行系统经常有这样的要求

弱一致性不保证用户能获取到最新的更新,但可以保证用户能获取到数据。

因为要满足强一致性的代价非常高,整个分布式系统需要做大量的通信才能实现强一致性,在所有数据同步一致前用户就需要一直等待,这对一些性能要求比较高的系统来说是不可接受的

所以在大多数业务中常常采用最终一致性,比如常见的评论点赞,你获取到的可以不是最新的点赞数

MapReduce

接下来就是重头戏 MapReduce

MapReduce是由Google设计,开发和使用的一个系统,相关的论文在2004年发表。Google当时面临的问题是,他们需要在TB级别的数据上进行大量的计算。比如说,为所有网页创建索引,分析整个互联网的连接路径并得出最重要或者最权威的网页。如你所知,在当时,整个互联网的数据也有数十TB。构建索引基本上等同于对整个数据做排序,而排序比较费时。如果用一台计算机对整个互联网数据进行排序,要花费多长时间呢?可能要几周,几个月,甚至几年。所以,当时Google非常希望能将对大量数据的大量运算并行跑在几千台计算机上,这样才能快速完成计算。对Google来说,购买大量的计算机是没问题的,这样Google的工程师就不用花大量时间来看报纸来等他们的大型计算任务完成。所以,有段时间,Google买了大量的计算机,并让它的聪明的工程师在这些计算机上编写分布式软件,这样工程师们可以将手头的问题分包到大量计算机上去完成,管理这些运算,并将数据取回。

如果你只雇佣熟练的分布式系统专家作为工程师,尽管可能会有些浪费,也是可以的。但是Google想雇用的是各方面有特长的人,不一定是想把所有时间都花在编写分布式软件上的工程师。所以Google需要一种框架,可以让它的工程师能够进行任意的数据分析,例如排序,网络索引器,链接分析器以及任何的运算。工程师只需要实现应用程序的核心,就能将应用程序运行在数千台计算机上,而不用考虑如何将运算工作分发到数千台计算机,如何组织这些计算机,如何移动数据,如何处理故障等等这些细节。所以,当时Google需要一种框架,使得普通工程师也可以很容易的完成并运行大规模的分布式运算。这就是MapReduce出现的背景。

MapReduce的思想是,应用程序设计人员和分布式运算的使用者,只需要写简单的Map函数和Reduce函数,而不需要知道任何有关分布式的事情,MapReduce框架会处理剩下的事情。

这里理解起来可能有点抽象,简单来说我们现在有大量的数据需要处理,传统的方式我们如果用一台机器处理大量的数据,那会等到天荒地老,如果我们用成千上万的服务器把任务分散来跑或许会很快,但我们同样需要用大量的代码来处理分散数据,并发处理,处理完数据收集以及处理故障等问题。

所以Google开发了MapReduce,使得程序员只需要编写Map和Reduce两个函数,就可以完成上面分发数据,并行计算,收集数据等等问题,同时使得故障处理变得简单。

原文是这样描述的:

Over the past five years, the authors and many others at Google have implemented hundreds of special-purpose computations that process large amounts of raw data, such as crawled documents, web request logs, etc., to compute various kinds of derived data, such as inverted indices, various representations of the graph structure of web documents, summaries of the number of pages crawled per host, the set of most frequent queries in a given day, etc. Most such computations are conceptu- ally straightforward. However, the input data is usually large and the computations have to be distributed across hundreds or thousands of machines in order to finish in a reasonable amount of time. The issues of how to par- allelize the computation, distribute the data, and handle failures conspire to obscure the original simple compu- tation with large amounts of complex code to deal with these issues. As a reaction to this complexity, we designed a new abstraction that allows us to express the simple computa- tions we were trying to perform but hides the messy de- tails of parallelization, fault-tolerance, data distribution and load balancing in a library. Our abstraction is in- spired by the map and reduce primitives present in Lisp and many other functional languages. We realized that most of our computations involved applying a map op- eration to each logical “record” in our input in order to compute a set of intermediate key/value pairs, and then applying a reduce operation to all the values that shared the same key, in order to combine the derived data ap- propriately. Our use of a functional model with user- specified map and reduce operations allows us to paral- lelize large computations easily and to use re-execution as the primary mechanism for fault tolerance. The major contributions of this work are a simple and powerful interface that enables automatic parallelization and distribution of large-scale computations, combined with an implementation of this interface that achieves high performance on large clusters of commodity PCs.

在过去五年中,本文作者及 Google 的许多同事实现了数百个专用计算程序,用于处理大量原始数据(如抓取的网页文档、Web 请求日志等),以计算各种派生数据(如倒排索引、网页文档图结构的各种表示、每个主机被抓取页面数量的汇总、某一天最频繁的查询集合等)。大多数此类计算在概念上都很简单。然而,输入数据通常非常庞大,计算必须分布在数百甚至数千台机器上才能在合理的时间内完成。如何并行化计算、分发数据以及处理故障等问题,使得原本简单的计算被大量复杂的代码所掩盖。

为了应对这种复杂性,我们设计了一种新的抽象,使我们能够表达想要执行的简单计算,同时将并行化、容错、数据分发和负载均衡等繁琐细节隐藏在库中。我们的抽象受到 Lisp 及许多其他函数式语言中 map 和 reduce 原语的启发。我们意识到,我们的大多数计算都涉及对输入中的每条逻辑"记录"应用一个 map 操作以计算一组中间键/值对,然后对所有共享相同键的值应用一个 reduce 操作,以适当地合并派生数据。我们使用带有用户指定 map 和 reduce 操作的函数式模型,使我们能够轻松地并行化大规模计算,并将重新执行作为容错的主要机制。

本工作的主要贡献是一个简单而强大的接口,能够实现大规模计算的自动并行化和分布式执行,以及该接口的一个实现,该实现在大规模普通 PC 集群上实现了高性能。

现在我们来关注其具体实现

抽象来看,MapReduce假设有一些输入,这些输入被分割成大量的不同的文件或者数据块。所以我们假设现在有输入文件1,输入文件2和输入文件3,这些输入可能是从网上抓取的网页,更可能是包含了大量网页的文件。

MapReduce启动时,会查找Map函数。之后,MapReduce框架会为每个输入文件运行Map函数。这里很明显有一些可以并行运算的地方,比如可以并行运行多个只关注输入和输出的Map函数。

Map函数以文件作为输入,文件又是整个输入数据的一部分。Map函数的输出是一个key-value对的列表。

假如我们有一个数单词的程序,我们被提供一个文件,里面有大量的混乱且可能重复的单词,我们需要做的工作是统计文件中所有单词及他们出现的次数

在这个例子中,Map函数会输出key-value对,其中key是单词,而value是1。Map函数会将输入中的每个单词拆分,并输出一个key-value对。最后需要对所有的key-value进行计数,以获得最终的输出。所以,假设输入文件1包含了单词a和单词b,Map函数的输出将会是key=a,value=1和key=b,value=1。第二个Map函数只从输入文件2看到了b,那么输出将会是key=b,value=1。第三个输入文件有一个a和一个c。

image

我们对所有的输入文件都运行了Map函数,并得到了论文中称为中间输出(intermediate output),也就是每个Map函数输出的key-value对。

运算的第二阶段是运行Reduce函数。MapReduce框架会收集所有Map函数输出的每一个单词的统计。比如说,MapReduce框架会先收集每一个Map函数输出的key为a的key-value对。收集了之后,会将它们提交给Reduce函数。

之后会收集所有的b。这里的收集是真正意义上的收集,因为b是由不同计算机上的不同Map函数生成,所以不仅仅是数据从一台计算机移动到另一台(如果Map只在一台计算机的一个实例里,可以直接通过一个RPC将数据从Map移到Reduce)。我们收集所有的b,并将它们提交给另一个Reduce函数。这个Reduce函数的入参是所有的key为b的key-value对。对c也是一样。所以,MapReduce框架会为所有Map函数输出的每一个key,调用一次Reduce函数。

2

在我们这个简单的单词计数器的例子中,Reduce函数只需要统计传入参数的长度,甚至都不用查看传入参数的具体内容,因为每一个传入参数代表对单词加1,而我们只需要统计个数。最后,每个Reduce都输出与其关联的单词和这个单词的数量。所以第一个Reduce输出a=2,第二个Reduce输出b=2,第三个Reduce输出c=1。

这就是一个典型的MapReduce Job。从整体来看,为了保证完整性,有一些术语要介绍一下:

  • Job。整个MapReduce计算称为Job。
  • Task。每一次MapReduce调用称为Task。

所以,对于一个完整的MapReduce Job,它由一些Map Task和一些Reduce Task组成。所以这是一个单词计数器的例子,它解释了MapReduce的基本工作方式。

在论文中的伪代码实例:

map(String key, String value):
  // key: 文档名称
  // value: 文档内容
  for each word w in value:
    EmitIntermediate(w, "1");

reduce(String key, Iterator values):
  // key: 一个单词
  // values: 计数列表
  int result = 0;
  for each v in values:
    result += ParseInt(v);
  Emit(AsString(result));

参考

课程翻译Lecture 01

原始论文

原文:MapReduce: Simplified Data Processing on Large Clusters

中文翻译:

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

Jeffrey Dean 和 Sanjay Ghemawat jeff@google.com, sanjay@google.com Google 公司


摘要

MapReduce 是一种用于处理和生成大规模数据集的编程模型及其相关实现。用户指定一个 map 函数,用于处理键/值对以生成一组中间键/值对;再指定一个 reduce 函数,用于合并所有具有相同中间键的中间值。如本文所示,许多现实世界的任务都可以用这种模型来表达。

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

我们的 MapReduce 实现运行在由大量普通商用机器组成的大规模集群上,并且具有高度的可扩展性:一次典型的 MapReduce 计算可以在数千台机器上处理数 TB 的数据。程序员发现该系统非常易于使用:已经实现了数百个 MapReduce 程序,在 Google 的集群上每天执行超过一千个 MapReduce 作业。


1 引言

在过去五年中,本文作者及 Google 的许多同事实现了数百个专用计算程序,用于处理大量原始数据(如抓取的网页文档、Web 请求日志等),以计算各种派生数据(如倒排索引、网页文档图结构的各种表示、每个主机被抓取页面数量的汇总、某一天最频繁的查询集合等)。大多数此类计算在概念上都很简单。然而,输入数据通常非常庞大,计算必须分布在数百甚至数千台机器上才能在合理的时间内完成。如何并行化计算、分发数据以及处理故障等问题,使得原本简单的计算被大量复杂的代码所掩盖。

为了应对这种复杂性,我们设计了一种新的抽象,使我们能够表达想要执行的简单计算,同时将并行化、容错、数据分发和负载均衡等繁琐细节隐藏在库中。我们的抽象受到 Lisp 及许多其他函数式语言中 map 和 reduce 原语的启发。我们意识到,我们的大多数计算都涉及对输入中的每条逻辑"记录"应用一个 map 操作以计算一组中间键/值对,然后对所有共享相同键的值应用一个 reduce 操作,以适当地合并派生数据。我们使用带有用户指定 map 和 reduce 操作的函数式模型,使我们能够轻松地并行化大规模计算,并将重新执行作为容错的主要机制。

本工作的主要贡献是一个简单而强大的接口,能够实现大规模计算的自动并行化和分布式执行,以及该接口的一个实现,该实现在大规模普通 PC 集群上实现了高性能。

第 2 节描述基本编程模型并给出若干示例。第 3 节描述针对我们基于集群的计算环境量身定制的 MapReduce 接口实现。第 4 节描述我们发现有用的编程模型的若干改进。第 5 节给出我们实现的各种任务的性能测量数据。第 6 节探讨 MapReduce 在 Google 内部的使用情况,包括我们将其作为生产索引系统重写基础的经验。第 7 节讨论相关工作和未来工作。


2 编程模型

计算接受一组输入键/值对,并产生一组输出键/值对。MapReduce 库的用户将计算表达为两个函数:Map 和 Reduce。

Map 由用户编写,接受一个输入对并产生一组中间键/值对。MapReduce 库将所有与相同中间键 I 关联的中间值归组在一起,并将它们传递给 Reduce 函数。

Reduce 函数也由用户编写,接受一个中间键 I 和该键的一组值。它将这些值合并在一起,形成一个可能更小的值集合。通常每次 Reduce 调用只产生零个或一个输出值。中间值通过一个迭代器提供给用户的 reduce 函数。这使我们能够处理太大而无法放入内存的值列表。

2.1 示例

考虑统计大量文档集合中每个单词出现次数的问题。用户将编写类似于以下伪代码的代码:

map(String key, String value):
  // key: 文档名称
  // value: 文档内容
  for each word w in value:
    EmitIntermediate(w, "1");

reduce(String key, Iterator values):
  // key: 一个单词
  // values: 计数列表
  int result = 0;
  for each v in values:
    result += ParseInt(v);
  Emit(AsString(result));

map 函数输出每个单词及其相关的出现计数(在这个简单示例中仅为"1")。reduce 函数将针对特定单词输出的所有计数求和。

此外,用户还需编写代码来填充一个 mapreduce 规范对象,其中包含输入和输出文件的名称以及可选的调优参数。然后用户调用 MapReduce 函数,将规范对象传递给它。用户的代码与 MapReduce 库(用 C++ 实现)链接在一起。附录 A 包含此示例的完整程序文本。

2.2 类型

尽管前面的伪代码是以字符串输入和输出的形式编写的,但从概念上讲,用户提供的 map 和 reduce 函数具有关联的类型:

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

即,输入的键和值与输出的键和值来自不同的域。此外,中间键和值与输出键和值来自相同的域。

我们的 C++ 实现向用户定义的函数传递字符串并从中接收字符串,由用户代码负责在字符串和适当类型之间进行转换。

2.3 更多示例

以下是一些可以轻松地表达为 MapReduce 计算的有趣程序的简单示例。

分布式 Grep: map 函数在匹配给定模式时输出该行。reduce 函数是一个恒等函数,只是将提供的中间数据复制到输出。

URL 访问频率计数: map 函数处理网页请求日志并输出〈URL, 1〉。reduce 函数将相同 URL 的所有值相加,并输出〈URL, 总计数〉对。

反向 Web 链接图: map 函数为名为 source 的页面中找到的每个指向目标 URL 的链接输出〈target, source〉对。reduce 函数将与给定目标 URL 关联的所有源 URL 列表连接起来,并输出对:〈target, list(source)〉。

每个主机的词向量: 词向量以〈word, frequency〉对的列表形式总结文档或一组文档中出现的最重要单词。map 函数为每个输入文档输出一个〈hostname, term vector〉对(其中 hostname 从文档的 URL 中提取)。reduce 函数接收给定主机的所有每文档词向量。它将这些词向量相加,丢弃不常见的词,然后输出最终的〈hostname, term vector〉对。

倒排索引: map 函数解析每个文档,并输出一系列〈word, document ID〉对。reduce 函数接受给定单词的所有对,对相应的文档 ID 进行排序,并输出〈word, list(document ID)〉对。所有输出对的集合构成一个简单的倒排索引。可以轻松地扩展此计算以跟踪单词位置。

分布式排序: map 函数从每条记录中提取键,并输出〈key, record〉对。reduce 函数原样输出所有对。此计算依赖于第 4.1 节描述的分区功能和第 4.2 节描述的排序特性。


3 实现

MapReduce 接口有许多不同的实现方式。正确的选择取决于环境。例如,一种实现可能适合小型共享内存机器,另一种适合大型 NUMA 多处理器,还有一种适合更大的联网机器集合。

本节描述针对 Google 广泛使用的计算环境的实现:通过交换式以太网连接的大规模普通 PC 集群。在我们的环境中:

(1) 机器通常是运行 Linux 的双处理器 x86 系统,每台机器有 2-4 GB 内存。

(2) 使用普通商用网络硬件——通常在机器级别为 100 兆位/秒或 1 千兆位/秒,但总体对分带宽要低得多。

(3) 集群由数百或数千台机器组成,因此机器故障很常见。

(4) 存储由直接连接到各台机器的廉价 IDE 磁盘提供。使用内部开发的分布式文件系统来管理存储在这些磁盘上的数据。该文件系统使用复制技术在不可靠的硬件之上提供可用性和可靠性。

(5) 用户向调度系统提交作业。每个作业由一组任务组成,由调度器映射到集群内的一组可用机器。

3.1 执行概述

Map 调用通过自动将输入数据分区为 M 个分片来分布在多台机器上。输入分片可以由不同的机器并行处理。Reduce 调用通过使用分区函数(例如 hash(key) mod R)将中间键空间分区为 R 个部分来分布。分区数(R)和分区函数由用户指定。

图 1 展示了我们实现中 MapReduce 操作的整体流程。当用户程序调用 MapReduce 函数时,会发生以下一系列操作(图 1 中的编号标签对应下面列表中的编号):

  1. 用户程序中的 MapReduce 库首先将输入文件分割为 M 个片段,通常每个片段为 16 兆字节到 64 兆字节(MB)(用户可通过可选参数控制)。然后它在集群的机器上启动程序的多个副本。

  2. 程序副本中有一个是特殊的——master。其余的是由 master 分配工作的 worker。有 M 个 map 任务和 R 个 reduce 任务需要分配。master 选择空闲的 worker 并为每个 worker 分配一个 map 任务或一个 reduce 任务。

  3. 被分配了 map 任务的 worker 读取相应输入分片的内容。它从输入数据中解析出键/值对,并将每对传递给用户定义的 Map 函数。由 Map 函数产生的中间键/值对被缓存在内存中。

  4. 缓冲的对被定期写入本地磁盘,由分区函数分区为 R 个区域。这些缓冲对在本地磁盘上的位置被传回给 master,master 负责将这些位置转发给 reduce worker。

  5. 当 reduce worker 被 master 通知这些位置时,它使用远程过程调用从 map worker 的本地磁盘读取缓冲数据。当 reduce worker 读取了所有中间数据后,它按中间键排序,使相同键的所有出现归组在一起。排序是必要的,因为通常许多不同的键映射到同一个 reduce 任务。如果中间数据量太大而无法放入内存,则使用外部排序。

  6. reduce worker 遍历排序后的中间数据,对于遇到的每个唯一中间键,它将该键和相应的中间值集合传递给用户的 Reduce 函数。Reduce 函数的输出被追加到该 reduce 分区的最终输出文件中。

  7. 当所有 map 任务和 reduce 任务都完成后,master 唤醒用户程序。此时,用户程序中的 MapReduce 调用返回到用户代码。

成功完成后,mapreduce 执行的输出可在 R 个输出文件中获得(每个 reduce 任务一个,文件名由用户指定)。通常,用户不需要将这 R 个输出文件合并为一个文件——他们通常将这些文件作为另一个 MapReduce 调用的输入,或在另一个能够处理分区为多个文件的输入的分布式应用程序中使用它们。

3.2 Master 数据结构

master 维护若干数据结构。对于每个 map 任务和 reduce 任务,它存储状态(空闲、进行中或已完成)以及 worker 机器的标识(对于非空闲任务)。

master 是中间文件区域位置从 map 任务传播到 reduce 任务的管道。因此,对于每个已完成的 map 任务,master 存储该 map 任务产生的 R 个中间文件区域的位置和大小。随着 map 任务的完成,会接收到这些位置和大小信息的更新。这些信息被增量地推送给正在执行 reduce 任务的 worker。

3.3 容错

由于 MapReduce 库旨在帮助使用数百或数千台机器处理非常大量的数据,因此该库必须能够优雅地容忍机器故障。

Worker 故障

master 定期 ping 每个 worker。如果在一定时间内没有收到某个 worker 的响应,master 将该 worker 标记为失败。该 worker 已完成的任何 map 任务都被重置为初始的空闲状态,因此可以在其他 worker 上重新调度。类似地,在故障 worker 上正在执行的任何 map 任务或 reduce 任务也被重置为空闲状态,并可以重新调度。

已完成的 map 任务在故障时需要重新执行,因为它们的输出存储在故障机器的本地磁盘上,因此不可访问。已完成的 reduce 任务不需要重新执行,因为它们的输出存储在全局文件系统中。

当一个 map 任务首先由 worker A 执行,然后由 worker B 执行(因为 A 失败了)时,所有执行 reduce 任务的 worker 都会被通知重新执行。任何尚未从 worker A 读取数据的 reduce 任务将从 worker B 读取数据。

MapReduce 对大规模 worker 故障具有弹性。例如,在一次 MapReduce 操作期间,正在运行的集群上的网络维护导致每次有 80 台机器组在几分钟内变得不可达。MapReduce master 只是重新执行了不可达 worker 机器完成的工作,并继续取得进展,最终完成了 MapReduce 操作。

Master 故障

可以很容易地让 master 定期写入上述 master 数据结构的检查点。如果 master 任务死亡,可以从最后检查点状态启动新的副本。然而,由于只有一个 master,其故障不太可能发生;因此我们当前的实现在 master 故障时中止 MapReduce 计算。客户端可以检查此条件,并根据需要重试 MapReduce 操作。

故障存在时的语义

当用户提供的 map 和 reduce 操作是其输入值的确定性函数时,我们的分布式实现产生的输出与整个程序的无故障顺序执行产生的输出相同。

我们依赖 map 和 reduce 任务输出的原子提交来实现此属性。每个进行中的任务将其输出写入私有临时文件。一个 reduce 任务产生一个这样的文件,一个 map 任务产生 R 个这样的文件(每个 reduce 任务一个)。当 map 任务完成时,worker 向 master 发送一条消息,其中包含 R 个临时文件的名称。如果 master 收到已完成 map 任务的完成消息,则忽略该消息。否则,它将 R 个文件的名称记录在 master 数据结构中。

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

我们绝大多数的 map 和 reduce 操作都是确定性的,在这种情况下,我们的语义等同于顺序执行,这使得程序员很容易推断其程序的行为。当 map 和/或 reduce 操作是非确定性的时,我们提供较弱但仍然合理的语义。在非确定性操作存在的情况下,特定 reduce 任务 R1 的输出等同于非确定性程序顺序执行产生的 R1 的输出。然而,不同 reduce 任务 R2 的输出可能对应于非确定性程序的不同顺序执行产生的 R2 的输出。

考虑 map 任务 M 和 reduce 任务 R1 和 R2。令 e(Ri) 为 Ri 已提交的执行(恰好有一次这样的执行)。较弱的语义产生是因为 e(R1) 可能读取了 M 的一次执行产生的输出,而 e(R2) 可能读取了 M 的另一次执行产生的输出。

3.4 本地性

在我们的计算环境中,网络带宽是相对稀缺的资源。我们利用输入数据(由 GFS 管理)存储在构成集群的机器的本地磁盘上这一事实来节约网络带宽。GFS 将每个文件划分为 64 MB 的块,并在不同的机器上存储每个块的多个副本(通常 3 个副本)。MapReduce master 考虑输入文件的位置信息,并尝试在包含相应输入数据副本的机器上调度 map 任务。如果失败,它尝试在该任务输入数据副本附近调度 map 任务(例如,在与包含数据的机器在同一网络交换机上的 worker 机器上)。当在集群中相当大比例的 worker 上运行大型 MapReduce 操作时,大多数输入数据是本地读取的,不消耗网络带宽。

3.5 任务粒度

如上所述,我们将 map 阶段细分为 M 个部分,将 reduce 阶段细分为 R 个部分。理想情况下,M 和 R 应远大于 worker 机器的数量。让每个 worker 执行许多不同的任务可以改善动态负载均衡,并在 worker 故障时加速恢复:它已完成的许多 map 任务可以分散到所有其他 worker 机器上。

在我们的实现中,M 和 R 的大小有实际限制,因为 master 必须做出 O(M+R) 个调度决策,并在内存中保持 O(MR) 的状态(如上所述)。(然而,内存使用的常数因子很小:O(MR) 部分的状态大约由每个 map 任务/reduce 任务对约一个字节的数据组成。)

此外,R 通常受到用户的约束,因为每个 reduce 任务的输出最终在单独的输出文件中。在实践中,我们倾向于选择 M 使得每个单独任务大约处理 16 MB 到 64 MB 的输入数据(以使上述本地性优化最有效),并使 R 为预期使用的 worker 机器数量的一个小倍数。我们经常执行 M=200,000 和 R=5,000 的 MapReduce 计算,使用 2,000 台 worker 机器。

3.6 备份任务

延长 MapReduce 操作总时间的一个常见原因是"落后者":一台机器花费异常长的时间来完成计算中最后几个 map 或 reduce 任务之一。落后者可能由多种原因引起。例如,一台磁盘有问题的机器可能经历频繁的可纠正错误,使其读取性能从 30 MB/s 降至 1 MB/s。集群调度系统可能在机器上调度了其他任务,由于 CPU、内存、本地磁盘或网络带宽的竞争,导致其执行 MapReduce 代码更慢。我们最近遇到的一个问题是机器初始化代码中的一个 bug 导致处理器缓存被禁用:受影响机器上的计算速度降低了一百倍以上。

我们有一个通用机制来缓解落后者问题。当 MapReduce 操作接近完成时,master 调度剩余进行中任务的备份执行。当主执行或备份执行中的任何一个完成时,该任务即被标记为已完成。我们调整了此机制,使其通常仅将操作使用的计算资源增加不超过几个百分点。

我们发现这显著减少了完成大型 MapReduce 操作的时间。例如,第 5.3 节描述的排序程序在禁用备份任务机制时完成时间增加了 44%。


4 改进

虽然仅通过编写 Map 和 Reduce 函数提供的基本功能对大多数需求来说已经足够,但我们发现一些扩展很有用。这些在本节中描述。

4.1 分区函数

MapReduce 的用户指定他们期望的 reduce 任务/输出文件的数量(R)。数据使用中间键上的分区函数分区到这些任务中。提供了一个默认的分区函数,使用哈希(例如"hash(key) mod R")。这往往产生相当均衡的分区。然而,在某些情况下,按键的某些其他函数来分区数据是有用的。例如,有时输出键是 URL,我们希望单个主机的所有条目最终在同一个输出文件中。为了支持这种情况,MapReduce 库的用户可以提供特殊的分区函数。例如,使用"hash(Hostname(urlkey)) mod R"作为分区函数会使来自同一主机的所有 URL 最终在同一个输出文件中。

4.2 排序保证

我们保证在给定的分区内,中间键/值对按键的递增顺序处理。此排序保证使得为每个分区生成排序的输出文件变得容易,这在输出文件格式需要支持按键的高效随机访问查找时很有用,或者当用户发现数据已排序很方便时。

4.3 Combiner 函数

在某些情况下,每个 map 任务产生的中间键有大量重复,并且用户指定的 Reduce 函数是可交换和可结合的。一个很好的例子是第 2.1 节中的单词计数示例。由于单词频率往往遵循 Zipf 分布,每个 map 任务将产生数百或数千个形如 <the, 1> 的记录。所有这些计数都将通过网络发送到单个 reduce 任务,然后由 Reduce 函数相加以产生一个数字。我们允许用户指定一个可选的 Combiner 函数,在数据通过网络发送之前对其进行部分合并。

Combiner 函数在每台执行 map 任务的机器上执行。通常使用相同的代码来实现 combiner 和 reduce 函数。reduce 函数和 combiner 函数之间的唯一区别是 MapReduce 库如何处理函数的输出。reduce 函数的输出写入最终输出文件。combiner 函数的输出写入将发送到 reduce 任务的中间文件。

部分合并显著加速了某些类别的 MapReduce 操作。附录 A 包含使用 combiner 的示例。

4.4 输入和输出类型

MapReduce 库支持以多种不同格式读取输入数据。例如,“text"模式的输入将每行视为一个键/值对:键是文件中的偏移量,值是行的内容。另一种常见的支持格式存储按键排序的键/值对序列。每种输入类型实现都知道如何将自己分割为有意义的范围以作为单独的 map 任务处理(例如,text 模式的范围分割确保范围分割仅发生在行边界)。用户可以通过提供简单 reader 接口的实现来添加对新输入类型的支持,尽管大多数用户只使用少数预定义输入类型之一。

reader 不一定需要提供从文件读取的数据。例如,可以很容易地定义一个从数据库或从内存中映射的数据结构读取记录的 reader。

类似地,我们支持一组输出类型以不同格式生成数据,用户代码可以很容易地添加对新输出类型的支持。

4.5 副作用

在某些情况下,MapReduce 的用户发现从他们的 map 和/或 reduce 操作产生辅助文件作为额外输出很方便。我们依赖应用程序编写者使这些副作用具有原子性和幂等性。通常应用程序写入临时文件,并在文件完全生成后原子地重命名。

我们不支持单个任务产生的多个输出文件的原子两阶段提交。因此,产生具有跨文件一致性要求的多个输出文件的任务应该是确定性的。此限制在实践中从未成为问题。

4.6 跳过坏记录

有时用户代码中的 bug 会导致 Map 或 Reduce 函数在某些记录上确定性地崩溃。这样的 bug 会阻止 MapReduce 操作完成。通常的做法是修复 bug,但有时这不可行;也许 bug 在源代码不可用的第三方库中。此外,有时忽略少数记录是可以接受的,例如在对大数据集进行统计分析时。我们提供了一种可选的执行模式,其中 MapReduce 库检测哪些记录导致确定性崩溃并跳过这些记录以继续推进。

每个 worker 进程安装一个信号处理程序来捕获段错误和总线错误。在调用用户 Map 或 Reduce 操作之前,MapReduce 库将参数的序列号存储在全局变量中。如果用户代码产生信号,信号处理程序向 MapReduce master 发送一个包含序列号的"最后遗言"UDP 数据包。当 master 看到特定记录上有多次故障时,它指示在相应 Map 或 Reduce 任务的下一次重新执行中应跳过该记录。

4.7 本地执行

调试 Map 或 Reduce 函数中的问题可能很棘手,因为实际计算发生在分布式系统中,通常在数千台机器上,工作分配决策由 master 动态做出。为了便于调试、性能分析和小规模测试,我们开发了 MapReduce 库的替代实现,在本地机器上顺序执行 MapReduce 操作的所有工作。向用户提供控制,以便将计算限制到特定的 map 任务。用户使用特殊标志调用程序,然后可以轻松地使用他们觉得有用的任何调试或测试工具(例如 gdb)。

4.8 状态信息

master 运行一个内部 HTTP 服务器并导出一组状态页面供人类查看。状态页面显示计算的进度,如已完成多少任务、多少正在进行中、输入字节数、中间数据字节数、输出字节数、处理速率等。页面还包含指向每个任务生成的标准错误和标准输出文件的链接。用户可以使用这些数据来预测计算需要多长时间,以及是否应该为计算添加更多资源。这些页面也可用于确定计算何时比预期慢得多。

此外,顶级状态页面显示哪些 worker 已失败,以及它们失败时正在处理哪些 map 和 reduce 任务。此信息在尝试诊断用户代码中的 bug 时很有用。

4.9 计数器

MapReduce 库提供计数器功能来计数各种事件的发生次数。例如,用户代码可能想要计数处理的单词总数或索引的德文文档数等。

要使用此功能,用户代码创建一个命名的计数器对象,然后在 Map 和/或 Reduce 函数中适当地递增计数器。例如:

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

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

来自各台 worker 机器的计数器值定期传播到 master(搭载在 ping 响应上)。master 聚合来自成功的 map 和 reduce 任务的计数器值,并在 MapReduce 操作完成时将它们返回给用户代码。当前的计数器值也显示在 master 状态页面上,以便人类可以观察实时计算的进度。在聚合计数器值时,master 消除相同 map 或 reduce 任务重复执行的影响以避免重复计数。(重复执行可能由我们使用备份任务以及因故障而重新执行任务引起。)

MapReduce 库自动维护一些计数器值,如处理的输入键/值对数和产生的输出键/值对数。

用户发现计数器功能对于检查 MapReduce 操作的行为很有用。例如,在某些 MapReduce 操作中,用户代码可能希望确保产生的输出对数恰好等于处理的输入对数,或者处理的德文文档比例在总处理文档数的某个可容忍比例范围内。


5 性能

在本节中,我们测量了 MapReduce 在大规模机器集群上运行两个计算的性能。一个计算在大约一 TB 的数据中搜索特定模式。另一个计算对大约一 TB 的数据进行排序。

这两个程序代表了 MapReduce 用户编写的大量实际程序的一个大子集——一类程序将数据从一种表示转换为另一种表示,另一类程序从大数据集中提取少量有趣的数据。

5.1 集群配置

所有程序在由大约 1800 台机器组成的集群上执行。每台机器有两个 2GHz Intel Xeon 处理器(启用超线程),4GB 内存,两个 160GB IDE 磁盘,以及一个千兆以太网链路。机器排列在两级树形交换网络中,根部大约有 100-200 Gbps 的聚合带宽可用。所有机器都在同一个托管设施中,因此任何一对机器之间的往返时间小于一毫秒。

在 4GB 内存中,大约 1-1.5GB 被集群上运行的其他任务预留。程序在周末下午执行,此时 CPU、磁盘和网络基本空闲。

5.2 Grep

grep 程序扫描 10^10 条 100 字节的记录,搜索一个相对罕见的三字符模式(该模式出现在 92,337 条记录中)。输入被分割为大约 64MB 的片段(M=15000),整个输出放在一个文件中(R=1)。

图 2 显示了计算随时间的进展。Y 轴显示输入数据的扫描速率。随着更多机器被分配到此 MapReduce 计算,速率逐渐上升,当分配了 1764 个 worker 时达到超过 30 GB/s 的峰值。随着 map 任务完成,速率开始下降,在计算开始约 80 秒时降至零。整个计算从开始到结束大约需要 150 秒。这包括大约一分钟的启动开销。开销是由于将程序传播到所有 worker 机器,以及与 GFS 交互以打开 1000 个输入文件集合并获取本地性优化所需信息的延迟。

5.3 排序

sort 程序对 10^10 条 100 字节的记录(大约一 TB 数据)进行排序。此程序以 TeraSort 基准测试为模型。

排序程序由不到 50 行用户代码组成。一个三行的 Map 函数从文本行中提取 10 字节的排序键,并将键和原始文本行作为中间键/值对输出。我们使用内置的 Identity 函数作为 Reduce 操作符。此函数将中间键/值对原样传递为输出键/值对。最终排序输出写入一组 2 副本 GFS 文件(即程序的输出写入 2 TB)。

与之前一样,输入数据被分割为 64MB 的片段(M=15000)。我们将排序输出分区为 4000 个文件(R=4000)。分区函数使用键的初始字节将其分离到 R 个部分之一。我们此基准测试的分区函数具有关于键分布的内置知识。在通用排序程序中,我们会添加一个预处理 MapReduce 操作来收集键的样本,并使用采样键的分布来计算最终排序传递的分割点。

图 3(a) 显示了排序程序的正常执行进展。左上方的图显示输入读取速率。速率峰值约为 13 GB/s,并且相当快地消失,因为所有 map 任务在 200 秒内完成。注意输入速率低于 grep。这是因为 sort 的 map 任务花费大约一半的时间和 I/O 带宽将中间输出写入本地磁盘。grep 的相应中间输出大小可以忽略不计。

中间左侧的图显示数据通过网络从 map 任务发送到 reduce 任务的速率。此 shuffle 在第一个 map 任务完成后立即开始。图中的第一个驼峰是针对第一批大约 1700 个 reduce 任务的(整个 MapReduce 分配了大约 1700 台机器,每台机器一次最多执行一个 reduce 任务)。大约在计算进行 300 秒时,这些第一批 reduce 任务中的一些完成,我们开始为剩余的 reduce 任务 shuffle 数据。所有 shuffle 在计算进行约 600 秒时完成。

底部左侧的图显示 reduce 任务将排序数据写入最终输出文件的速率。在第一个 shuffle 期结束和写入期开始之间有一个延迟,因为机器正忙于排序中间数据。写入以大约 2-4 GB/s 的速率持续一段时间。所有写入在计算进行约 850 秒时完成。包括启动开销,整个计算需要 891 秒。这与当前报告的 TeraSort 基准测试最佳结果 1057 秒相似。

需要注意的几点:由于我们的本地性优化,输入速率高于 shuffle 速率和输出速率——大多数数据从本地磁盘读取,绕过了我们相对带宽受限的网络。shuffle 速率高于输出速率,因为输出阶段写入排序数据的两个副本(出于可靠性和可用性原因,我们制作输出的两个副本)。我们写入两个副本是因为这是底层文件系统提供的可靠性和可用性机制。如果底层文件系统使用纠删码而非复制,写入数据的网络带宽需求将会降低。

5.4 备份任务的效果

在图 3(b) 中,我们展示了禁用备份任务的排序程序执行。执行流程与图 3(a) 所示类似,只是有一个非常长的尾部,几乎没有写入活动。960 秒后,除了 5 个 reduce 任务外所有任务都已完成。然而这最后几个落后者直到 300 秒后才完成。整个计算需要 1283 秒,经过时间增加了 44%。

5.5 机器故障

在图 3(c) 中,我们展示了在计算进行几分钟后故意杀死 1746 个 worker 进程中的 200 个的排序程序执行。底层集群调度器立即在这些机器上重启了新的 worker 进程(因为只有进程被杀死,机器仍然正常运行)。

worker 死亡表现为负输入速率,因为一些先前完成的 map 工作消失了(因为相应的 map worker 被杀死)需要重做。此 map 工作的重新执行发生得相对较快。包括启动开销,整个计算在 933 秒内完成(仅比正常执行时间增加 5%)。


6 经验

我们在 2003 年 2 月编写了 MapReduce 库的第一个版本,并在 2003 年 8 月对其进行了重大增强,包括本地性优化、worker 机器间任务执行的动态负载均衡等。从那时起,我们惊喜地发现 MapReduce 库对我们处理的问题种类具有广泛的适用性。

它已在 Google 内部的广泛领域中使用,包括:

  • 大规模机器学习问题,
  • Google News 和 Froogle 产品的聚类问题,
  • 提取用于生成热门查询报告的数据(例如 Google Zeitgeist),
  • 为新实验和产品提取网页属性(例如从大型网页语料库中提取地理位置用于本地化搜索),以及
  • 大规模图计算。

图 4 显示了随时间推移检入我们主要源代码管理系统的独立 MapReduce 程序数量的显著增长,从 2003 年初的 0 个增长到 2004 年 9 月底的近 900 个独立实例。MapReduce 之所以如此成功,是因为它使得编写一个简单程序并在半小时内在一千台机器上高效运行成为可能,极大地加速了开发和原型设计周期。此外,它允许没有分布式和/或并行系统经验的程序员轻松利用大量资源。

在每个作业结束时,MapReduce 库记录关于作业使用的计算资源的统计信息。在表 1 中,我们展示了 2004 年 8 月在 Google 运行的 MapReduce 作业子集的一些统计信息。

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

指标数值
作业数量29,423
平均作业完成时间634 秒
使用的机器天数79,186 天
读取的输入数据3,288 TB
产生的中间数据758 TB
写入的输出数据193 TB
每作业平均 worker 机器数157
每作业平均 worker 死亡数1.2
每作业平均 map 任务数3,351
每作业平均 reduce 任务数55
唯一 map 实现数395
唯一 reduce 实现数269
唯一 map/reduce 组合数426

6.1 大规模索引

迄今为止,我们使用 MapReduce 最重要的用途之一是对生成 Google 网页搜索服务使用的数据结构的生产索引系统进行完全重写。索引系统以我们的抓取系统检索的大型文档集合作为输入,存储为一组 GFS 文件。这些文档的原始内容超过 20 TB 的数据。索引过程作为五到十个 MapReduce 操作的序列运行。使用 MapReduce(而不是先前版本索引系统中的临时分布式传递)提供了几个好处:

  • 索引代码更简单、更小、更容易理解,因为处理容错、分布和并行化的代码隐藏在 MapReduce 库中。例如,计算的一个阶段从大约 3800 行 C++ 代码减少到使用 MapReduce 表达时的大约 700 行。

  • MapReduce 库的性能足够好,我们可以将概念上不相关的计算分开,而不是将它们混合在一起以避免对数据的额外遍历。这使得更改索引过程变得容易。例如,在旧索引系统中需要几个月才能完成的一个更改,在新系统中只需几天即可实现。

  • 索引过程变得更容易操作,因为由机器故障、慢机器和网络问题引起的大多数问题都由 MapReduce 库自动处理,无需操作员干预。此外,通过向索引集群添加新机器可以很容易地提高索引过程的性能。


7 相关工作

许多系统提供了受限的编程模型,并利用这些限制来自动并行化计算。例如,关联函数可以在 N 个处理器上以 log N 时间对 N 元素数组的所有前缀进行计算,使用并行前缀计算。MapReduce 可以被视为基于我们在大规模现实世界计算中的经验对这些模型的一些简化和提炼。更重要的是,我们提供了一个可扩展到数千个处理器的容错实现。

相比之下,大多数并行处理系统只在较小规模上实现,并将处理机器故障的细节留给程序员。

批量同步编程和一些 MPI 原语提供了更高级别的抽象,使程序员更容易编写并行程序。这些系统与 MapReduce 之间的关键区别是 MapReduce 利用受限的编程模型来自动并行化用户程序并提供透明的容错。

我们的本地性优化从活动磁盘等技术中获得灵感,其中计算被推送到靠近本地磁盘的处理元素中,以减少通过 I/O 子系统或网络发送的数据量。我们运行在直接连接少量磁盘的普通处理器上,而不是直接在磁盘控制器处理器上运行,但总体方法是相似的。

我们的备份任务机制类似于 Charlotte 系统采用的急切调度机制。简单急切调度的一个缺点是,如果给定任务导致重复故障,整个计算将无法完成。我们通过跳过坏记录的机制修复了此问题的一些实例。

MapReduce 实现依赖于一个内部集群管理系统,该系统负责在大量共享机器上分发和运行用户任务。虽然不是本文的重点,但集群管理系统在精神上类似于 Condor 等其他系统。

作为 MapReduce 库一部分的排序功能在操作上类似于 NOW-Sort。源机器(map worker)对待排序数据进行分区并将其发送到 R 个 reduce worker 之一。每个 reduce worker 在本地排序其数据(如果可能在内存中)。当然 NOW-Sort 没有使我们的库广泛适用的用户可定义的 Map 和 Reduce 函数。

River 提供了一个编程模型,其中进程通过通过分布式队列发送数据相互通信。与 MapReduce 一样,River 系统试图在由异构硬件或系统扰动引入的不均匀性存在的情况下提供良好的平均情况性能。River 通过仔细调度磁盘和网络传输来实现均衡的完成时间。MapReduce 有不同的方法。通过限制编程模型,MapReduce 框架能够将问题分区为大量细粒度任务。这些任务在可用的 worker 上动态调度,使更快的 worker 处理更多任务。受限的编程模型还允许我们在作业结束时调度任务的冗余执行,这大大减少了在不均匀性(如慢速或卡住的 worker)存在时的完成时间。

BAD-FS 与 MapReduce 有非常不同的编程模型,与 MapReduce 不同,它针对跨广域网执行作业。然而,有两个根本的相似之处。(1) 两个系统都使用冗余执行来从故障导致的数据丢失中恢复。(2) 两者都使用本地性感知调度来减少通过拥塞网络链路发送的数据量。

TACC 是一个旨在简化构建高可用性网络服务的系统。与 MapReduce 一样,它依赖重新执行作为实现容错的机制。


8 结论

MapReduce 编程模型已在 Google 成功用于许多不同目的。我们将此成功归因于几个原因。首先,该模型易于使用,即使对于没有并行和分布式系统经验的程序员也是如此,因为它隐藏了并行化、容错、本地性优化和负载均衡的细节。其次,大量问题可以轻松地表达为 MapReduce 计算。例如,MapReduce 用于生成 Google 生产网页搜索服务的数据、排序、数据挖掘、机器学习以及许多其他系统。第三,我们开发了一个可扩展到由数千台机器组成的大规模集群的 MapReduce 实现。该实现高效地利用这些机器资源,因此适合用于 Google 遇到的许多大型计算问题。

我们从这项工作中学到了几件事。首先,限制编程模型使得并行化和分布计算以及使这些计算具有容错性变得容易。其次,网络带宽是稀缺资源。因此,我们系统中的许多优化旨在减少通过网络发送的数据量:本地性优化允许我们从本地磁盘读取数据,将中间数据的单个副本写入本地磁盘可节省网络带宽。第三,冗余执行可用于减少慢速机器的影响,并处理机器故障和数据丢失。


致谢

Josh Levenberg 在修订和扩展用户级 MapReduce API 方面发挥了重要作用,基于他使用 MapReduce 的经验和其他人的增强建议,添加了许多新功能。MapReduce 从 Google 文件系统读取输入并将输出写入其中。我们感谢 Mohit Aron、Howard Gobioff、Markus Gutschke、David Kramer、Shun-Tak Leung 和 Josh Redstone 在开发 GFS 方面的工作。我们还感谢 Percy Liang 和 Olcan Sercinoglu 在开发 MapReduce 使用的集群管理系统方面的工作。Mike Burrows、Wilson Hsieh、Josh Levenberg、Sharon Perl、Rob Pike 和 Debby Wallach 对本文早期草稿提供了有益的评论。匿名 OSDI 审稿人和我们的 shepherd Eric Brewer 提供了许多关于论文可以改进的领域的有用建议。最后,我们感谢 Google 工程组织内所有 MapReduce 用户提供的有益反馈、建议和 bug 报告。


附录 A 单词频率

本节包含一个程序,用于计算命令行上指定的一组输入文件中每个唯一单词的出现次数。

#include "mapreduce/mapreduce.h"

// 用户的 map 函数
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; ) {
      // 跳过前导空白
      while ((i < n) && isspace(text[i]))
        i++;
      // 找到单词结尾
      int start = i;
      while ((i < n) && !isspace(text[i]))
        i++;
      if (start < i)
        Emit(text.substr(start, i-start), "1");
    }
  }
};
REGISTER_MAPPER(WordCounter);

// 用户的 reduce 函数
class Adder : public Reducer {
  virtual void Reduce(ReduceInput* input) {
    // 遍历具有相同键的所有条目并增加值
    int64 value = 0;
    while (!input->done()) {
      value += StringToInt(input->value());
      input->NextValue();
    }
    // 输出 input->key() 的总和
    Emit(IntToString(value));
  }
};
REGISTER_REDUCER(Adder);

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

  MapReduceSpecification spec;

  // 将输入文件列表存储到 "spec" 中
  for (int i = 1; i < argc; i++) {
    MapReduceInput* input = spec.add_input();
    input->set_format("text");
    input->set_filepattern(argv[i]);
    input->set_mapper_class("WordCounter");
  }

  // 指定输出文件:
  //   /gfs/test/freq-00000-of-00100
  //   /gfs/test/freq-00001-of-00100
  //   ...
  MapReduceOutput* out = spec.output();
  out->set_filebase("/gfs/test/freq");
  out->set_num_tasks(100);
  out->set_format("text");
  out->set_reducer_class("Adder");

  // 可选:在 map 任务内进行部分求和以节省网络带宽
  out->set_combiner_class("Adder");

  // 调优参数:最多使用 2000 台机器和每任务 100 MB 内存
  spec.set_machines(2000);
  spec.set_map_megabytes(100);
  spec.set_reduce_megabytes(100);

  // 现在运行
  MapReduceResult result;
  if (!MapReduce(spec, &result)) abort();

  // 完成:'result' 结构包含关于计数器、耗时、
  // 使用的机器数等信息。
  return 0;
}
Licensed under CC BY-NC-SA 4.0
Build by Oight
使用 Hugo 构建
主题 StackJimmy 设计