【转载】创新的扩散:生成式学习的最新进展

原文地址:A Diffusion of Innovations: Recent Developments in Generative Learning,by Sebastian Raschka, on 2022-11-05

创新的扩散:生成式学习的最新进展

文章与趋势

生成式学习将走向何方?我们将去往何处?

前些年,生成式领域的核心是用于文本生成的大型语言Transformer模型。尽管文本与视觉方向的Transformer模型仍是众多研究者和AI公司的关注重点,但扩散模型近来已然抢占了它们的风头。

随着扩散模型(完整论文:https://arxiv.org/abs/2006.11239 )的出现,生成式学习领域近来掀起了一场重大变革。大约一年半前,也就是2021年5月11日,我们初次见识到扩散模型的表现超越了生成对抗网络(GAN)。而如今,扩散模型已然无处不在。

什么是扩散模型?

从本质上讲,扩散模型属于生成式模型。在训练过程中,模型会分多个步骤向输入数据中逐步添加随机噪声,对数据进行扰动。随后,模型通过逆向过程学习如何反转噪声、对输出进行去噪,从而还原出原始数据。在推理阶段,扩散模型便可以基于噪声输入生成全新的数据。

figure01
扩散模型示意图

扩散模型的应用

扩散模型首批真正广为人知的代表性应用,是一系列用于生成艺术作品与图像的大型模型;下面这些名字你大概率已经有所耳闻:

需要注意的是,尽管这些模型能够生成逼真且富有艺术感的图像,但相关争议也在不断发酵——这些模型的训练数据均爬取自互联网,并未征得原创艺术家的许可。此外,艺术家群体也对AI这一竞争者表达了担忧。我对知识产权相关法律了解有限,无法提供更多深入见解,但我记得去年GitHub Copilot(https://github.com/features/copilot)基于GitHub代码训练时,也引发过类似的批评。

我不确定这会将艺术家群体置于何种境地。就目前我们所知的情况来看,AI并不会取代程序员;更准确地说,是不会使用AI的程序员会被会使用AI的程序员取代——这和90年代“深蓝”战胜卡斯帕罗夫后国际象棋领域发生的变化类似。我不确定这对艺术家而言意味着什么,但既然精灵已经逃出瓶子,再也收不回去,或许我们只能顺势而为,学会利用这些新技术发挥其最大价值。

从图像到视频

继上述第一波图像生成扩散模型之后,各家公司近期开始向更高阶迈进,陆续推出能够根据文本提示生成视频片段的模型:

我猜再过几周,我们就能看到Stability.ai推出首个开源的同类视频生成模型了。另外值得一提的是,OpenAI目前还没有推出自己的视频生成扩散模型。(顺带一提,结合他们近期发布的字幕生成工具Whisper来看,我很好奇他们是否正在研发视频片段重新对齐的API。)

另外需要说明的是,目前的扩散模型数量远比这里列举的多得多。想要感受其数量之庞大,可以在这个网站中搜索“diffusion”关键词:https://scorebasedgenerativemodeling.github.io/ 。或者也可以查看这个仓库,截至本文撰写时,里面已经收录了超过300篇相关论文:https://github.com/heejkoo/Awesome-Diffusion-Models

我猜想下一波标志性的应用方向会聚焦于3D环境(比如VR领域)与运动生成。目前这一方向已经出现了一些值得关注的成果:

生物学领域的应用

我曾指导过一名博士生,他的研究方向是用图神经网络生成分子,因此我对扩散模型在生物学领域的应用也深有感触。例如:

figure02
Foldingdiff示意图。来源:https://arxiv.org/abs/2209.15171v1

再结合基于Transformer的模型,比如https://www.biorxiv.org/content/10.1101/2021.05.24.445464v1https://www.nature.com/articles/s41586-021-03819-2 中的成果,我不禁好奇图神经网络的未来会走向何方。一方面,去年微软的这档播客(https://www.microsoft.com/en-us/research/podcast/machine-learning-molecular-simulation-and-the-opportunity-for-societal-good-with-chris-bishop-and-max-welling/ )也探讨过相关话题。但迄今为止,AI在分子数据领域最亮眼的应用,并非基于图神经网络实现的。其中一种解释是,我们尚未找到让图神经网络突破小分子限制、实现规模化扩展的方法,但未来如何仍有待观察。就像投资者会分散投资一样,我们也不该把所有希望都寄托在单一技术方向上。

拓展学习资源

想要深入了解扩散模型?下面是一份精选的实用资源清单:


数值计算拥有开放的未来吗?

作为Python Conda工具(https://docs.conda.io/en/latest/ )的长期爱好者与使用者,我近期很荣幸受邀参加了Anaconda十周年的圆桌论坛,论坛主题围绕开源与科学计算展开,同台参与的还有https://twitter.com/teoliphanthttps://twitter.com/DynamicWebPaige 以及https://datascience.columbia.edu/people/ryan-abernathey/

数值计算的未来会是开放的吗?我们如今正处于开源2.0阶段。我还记得十年前,我们还只是把代码项目上传到GitHub而已。一路走来,我们已经进步了太多!而现在的第二阶段,我们要学习的是如何长期维护项目,构建可持续发展的社区。

近期一个很好的例子,就是PyTorch加入Linux基金会这件事:https://www.linuxfoundation.org/blog/blog/welcoming-pytorch-to-the-linux-foundation

本次圆桌论坛的完整录像可在下方查看。


研究深度解读

在这个部分,我会精选一篇研究论文进行解读。为了贴合本期通讯的“扩散主题”,我选的当然是一篇扩散模型相关的论文!这是一篇用于生成表格数据的扩散模型:https://arxiv.org/abs/2209.15421

首先,为什么要研发生成表格数据的模型?主要有两方面的应用价值:

  1. 生成规模更大的数据集,以及在类别不平衡场景下对少数类进行过采样。
  2. 在原始数据无法共享的场景下(比如出于隐私保护考虑),生成逼真的合成数据集。

TabDDPM绝非首个尝试生成逼真表格数据的模型。在它之前,生成对抗网络(GAN)和变分自编码器(VAE)早已被广泛用于这一领域,代表性的研究包括https://arxiv.org/abs/2204.00401https://arxiv.org/abs/1907.00503

而TabDDPM这篇论文的核心观点是:扩散模型的表现能够超越同类的GAN与VAE模型(这和本文前面提到的图像、视频领域的发展规律相似)。

TabDDPM的工作原理是什么?

对于类别型(以及二值型)特征,TabDDPM采用多项扩散,添加均匀噪声;对于数值型特征,则使用通用的高斯扩散过程。随后,模型通过全连接网络(MLP)学习逆向扩散过程。

figure03
特征分布对比。来源:https://arxiv.org/abs/2209.15421

需要注意的是,TabDDPM并非只是死记硬背训练数据中的样本。作者通过分析“最近样本距离”(DCR)分布得出了这一结论——分布中存在大量非零的距离值。

后续针对合成数据集质量的分析,量化了在生成数据上训练的多个分类器的性能。例如,在15个分类与回归数据集上,随机森林、逻辑回归、决策树、CatBoost等模型在TabDDPM生成的合成数据上的表现,均优于在GAN或VAE生成的合成数据上的表现。

如果你想要尝试使用TabDDPM,可以在这里获取它的源代码:https://github.com/rotot0/tab-ddpm


开源亮点

用Whisper生成高质量字幕

OpenAI近期发布了开源语言模型Whisper(https://github.com/openai/whisper ),可用于为音频文件生成字幕(隐藏式字幕)与内容摘要。我亲自测试了一下,效果惊艳。和YouTube自带的字幕生成功能,以及我这些年试过的好几款付费服务相比,Whisper生成字幕的准确率要高得多。它对我的口音,以及深度学习领域的专业术语都能很好地识别,这点让我印象格外深刻。

希望这个工具对你也有帮助。我用它给我博客上所有的深度学习课程视频(https://sebastianraschka.com/blog/2021/dl-course.html )重新生成了字幕,脚本在这里:https://gist.github.com/rasbt/0d09932c861851f177bd8f13dc93b354

Jupyter笔记本中的PyTorch多GPU支持

了解我的人都知道,我是Jupyter笔记本的忠实用户——它是我做原型开发和教学时最喜欢的环境。但Jupyter笔记本有个小局限:由于多进程限制,它无法支持大多数现代深度学习多GPU训练方案。我很高兴地告诉大家,PyTorch Lightning 1.7版本新增了在Jupyter笔记本中对PyTorch的多GPU支持。你可以在这里了解更多详情:https://lightning.ai/pages/community/tutorial/multi-gpu-jupyter-notebooks/

特征选择中的特征分组功能

自荐一下:我的mlxtend库(http://rasbt.github.io/mlxtend/CHANGELOG/ )现在在特征选择和排列重要性算法中支持特征分组功能了。用户可以将多个相关特征(比如独热编码生成的一组特征)在选择过程中当作一个整体分组来处理。

figure04
独热编码特征的特征分组支持


教学小贴士

这里分享5个从计算层面优化深度神经网络推理性能的技巧。

(注:为了限定讨论范围,本文只介绍不改变模型架构的方法。会改变架构的技术,比如模型剪枝,会在后续的推文中讨论。)

(1) 并行化

并行化指的是将训练数据的(小)批次拆分成多个分块,目标是并行处理这些更小的数据分块。

figure05
并行化示意图

(2) 向量化

向量化用批量运算取代了开销高昂的for循环,一次性对多个元素执行相同操作。

figure06
用点积实现加权和的for循环向量化

(3) 循环分块

改变循环中数据的访问顺序,以充分利用硬件的内存布局与缓存机制。

figure07
算子融合的概念示意图

(4) 算子融合

如果有多个循环,尽量将它们合并成一个。

figure08
两次遍历与一次遍历计算均值和方差的对比

(5) 量化

量化通过降低数值精度(通常将浮点数转为整数)来加快计算速度、降低内存占用,同时尽量保证精度不受损失。

figure09
PyTorch Lightning中量化感知训练的示例


机器学习面试题

你是在巩固知识点,还是在准备面试?这个板块的问题或许能帮到你(也可能帮不上)。

  1. 自然语言处理中的分布假说是什么?它应用在哪些场景?适用范围有多广?
  2. 无状态训练和有状态训练的区别是什么?分别在什么情况下使用?
  3. 递归和动态规划的区别是什么?

(答案将在下期通讯中公布。)


精彩引语

“我认为五年之后,我们就不再需要做提示工程了。” ——https://greylock.com/greymatter/sam-altman-ai-for-the-next-era/


近期活动

全球顶尖的机器学习与AI顶会NeurIPS(https://neurips.cc/ )很快就要召开了!2022年11月28日至12月1日。

如果你要去参会的话,欢迎过来打个招呼!我会联合主持一场主题为“业界、学界,以及两者之间”(即从学术到产业的转型)的工作坊社交活动。更多细节后续公布!


学习与效率技巧

我最喜欢的学习新事物的工具之一是Anki(https://apps.ankiweb.net/ )。这是一款用于间隔重复学习的开源闪卡软件。间隔重复的意思是,系统会根据你对问题难度的判断,在特定间隔后才重复出现该问题(比如,第一次间隔3天,之后是2周、3个月、1年,以此类推)。

我从本科时期就开始用它了。它帮我最大化地利用阅读和学习时间,牢牢记住知识点。我坚信,它是帮我度过大学初期时光的核心工具之一。到现在已经用了十多年,我积累了成千上万张卡片。

我是怎么用的呢?读文章或者教材的时候,我会针对那些想记住的重点内容,给自己提一些“有意思”的问题。有时候我还会把截图作为答案加进去。

figure10
Anki中一张闪卡问题的示例

我每天会新增5到15张新卡片。另外,我每天会花大概10到15分钟复习,避免卡片堆积。

figure11
我Anki卡组里待复习的内容,有数千张卡片。

我通常在晚上放松休息前完成每日复习。不过有时候我也会用手机版,利用碎片时间复习卡片——比如剪头发早到了几分钟,或者在机场等行李的时候。


这份通讯是我的个人兴趣项目,没有直接的商业收入。不过如果你愿意支持我的话,可以考虑购买我的书籍:https://sebastianraschka.com/books 。如果你觉得这些书有深度、有帮助,也欢迎推荐给你的朋友和同事。

figure12

https://www.amazon.com/Machine-Learning-PyTorch-Scikit-Learn-scikit-learn-ebook-dp-B09NW48MR1/dp/B09NW48MR1/ , https://nostarch.com/machine-learning-and-ai-beyond-basics , and http://mng.bz/M96o

你的支持对我意义重大!非常感谢!

Bigtable:一种面向结构化数据的分布式存储系统(Bigtable: A Distributed Storage System for Structured Data)


OSDI 2006 会议论文
第205–218页

Bigtable: A Distributed Storage System for Structured Data

Bigtable:一种面向结构化数据的分布式存储系统

作者:Fay Chang, Jeffrey Dean, Sanjay Ghemawat, Wilson C. Hsieh, Deborah A. Wallach
Mike Burrows, Tushar Chandra, Andrew Fikes, Robert E. Gruber

邮箱:{fay,jeff,sanjay,wilsonh,kerr,m3b,tushar,fikes,gruber}@google.com

Google公司

摘要

Bigtable是一种用于管理结构化数据的分布式存储系统,其设计目标是扩展至超大规模:跨越数千台商用服务器,存储PB级别的数据。Google内部有诸多项目都使用Bigtable存储数据,包括网页索引、Google Earth以及Google Finance等。这些应用对Bigtable的需求差异极大,无论是数据规模(从URL、网页到卫星影像)还是延迟要求(从后端批量处理到实时数据服务)皆是如此。尽管需求千差万别,Bigtable仍成功为所有这些Google产品提供了灵活、高性能的解决方案。本文将介绍Bigtable提供的简洁数据模型(该模型赋予客户端对数据布局与格式的动态控制权),并阐述Bigtable的设计与实现细节。

1 引言

在过去两年半的时间里,我们在Google设计、实现并部署了一套用于管理结构化数据的分布式存储系统,将其命名为Bigtable。Bigtable的设计目标是能够可靠地扩展至PB级数据和数千台机器的规模,它实现了多项目标:广泛的适用性、可扩展性、高性能与高可用性。

目前已有超过60个Google产品和项目使用Bigtable,包括Google Analytics、Google Finance、Orkut、个性化搜索、Writely以及Google Earth。这些产品使用Bigtable承载各类高负载工作场景,从吞吐量导向的批量处理任务,到面向终端用户的低延迟数据服务皆有涉及。这些产品使用的Bigtable集群配置跨度极大,从寥寥数台到数千台服务器不等,存储的数据量最高可达数百TB。

在很多方面,Bigtable都与数据库类似:它与数据库共享了诸多实现策略。并行数据库[14]与内存数据库[13]都实现了可扩展性与高性能,但Bigtable提供的接口与这类系统不同。Bigtable不支持完整的关系型数据模型;相反,它为客户端提供了一套简洁的数据模型,支持对数据布局与格式的动态控制,同时允许客户端感知底层存储中数据的局部性特征。数据通过行名与列名进行索引,行名和列名可以是任意字符串。Bigtable将数据视为未解析的字符串,尽管客户端通常会将各种形式的结构化、半结构化数据序列化到这些字符串中。客户端可以通过精心设计Schema来控制数据的局部性。最后,Bigtable的Schema参数允许客户端动态控制数据是从内存提供服务,还是从磁盘提供服务。

第2节将更详细地描述数据模型,第3节概述客户端API。第4节简要介绍Bigtable所依赖的Google底层基础设施。第5节阐述Bigtable实现的核心原理,第6节介绍我们为提升Bigtable性能所做的一些优化。第7节给出Bigtable的性能测试数据。第8节举例说明Bigtable在Google的实际应用场景,第9节讨论我们在设计与运维Bigtable过程中总结的经验教训。最后,第10节介绍相关工作,第11节给出结论。

2 数据模型

Bigtable是一个稀疏的、分布式的、持久化的多维排序映射。该映射通过行键、列键和时间戳进行索引;映射中的每个值都是一段未解析的字节数组。

(行:字符串, 列:字符串, 时间:int64) → 字符串

图1:存储网页的示例表片段。行名为反转后的URL。contents列族包含页面内容,anchor列族包含引用该页面的所有锚点的文本。CNN主页同时被Sports Illustrated和MY-look主页引用,因此该行包含anchor:cnnsi.comanchor:my.look.ca两列。每个锚点单元格仅有一个版本;contents列有三个版本,时间戳分别为(t_3)、(t_5)和(t_6)。

在考察了类Bigtable系统的各类潜在应用场景后,我们确定了这一数据模型。作为推动我们部分设计决策的具体案例,假设我们需要保存大量网页副本及相关信息,供多个不同项目使用;我们将这个表命名为Webtable。在Webtable中,我们使用URL作为行键,将网页的各类属性作为列名,在contents:列下、以抓取时的时间戳存储网页内容,如图1所示。

表中的行键是任意字符串(目前最大支持64KB,不过大多数用户的典型长度为10-100字节)。对同一个行键下的所有数据执行的读写操作都是原子性的(无论该行中读写的列数量有多少),这一设计决策让客户端更容易理解在同一行存在并发更新时系统的行为。

Bigtable按照行键的字典序维护数据。表的行范围会被动态分区,每个行范围称为一个子表(tablet),它是分布式部署与负载均衡的基本单元。因此,短行范围的读取效率很高,通常只需要与少量机器通信。客户端可以通过合理设计行键,让数据访问具备良好的局部性,从而利用这一特性。例如在Webtable中,通过反转URL的主机名部分,将同一域名下的页面聚集为连续的行。例如,我们将maps.google.com/index.html的数据存储在行键com.google.maps/index.html下。将同一域名的页面相邻存储,让很多针对主机和域名的分析效率更高。

列族

列键被分组为称为**列族(column family)**的集合,列族是访问控制的基本单位。存储在同一个列族中的所有数据通常类型相同(我们会对同一个列族内的数据统一压缩)。在向某个列族下的任意列键存储数据之前,必须先创建该列族;列族创建完成后,族内的任意列键都可以使用。我们的设计初衷是,表中不同列族的数量应当较少(最多数百个),且列族在运行过程中很少变更。与之相对的是,表可以拥有无限数量的列。

列键的命名语法为:族名:限定符。列族名称必须是可打印字符,而限定符可以是任意字符串。Webtable的一个示例列族是language,用于存储网页的编写语言。我们在language族中只使用一个列键,存储每个网页的语言ID。该表另一个实用的列族是anchor;族中的每个列键代表一个单独的锚点,如图1所示。限定符是引用站点的名称,单元格内容是链接文本。

访问控制、磁盘与内存计费都在列族级别执行。在Webtable示例中,这些控制让我们可以管理多种不同类型的应用:有些应用添加新的基础数据,有些应用读取基础数据并生成衍生列族,还有些应用仅允许查看已有数据(甚至出于隐私原因,可能都不允许查看所有已有的列族)。

时间戳

Bigtable中的每个单元格都可以包含同一数据的多个版本;这些版本通过时间戳进行索引。Bigtable的时间戳是64位整数。它们可以由Bigtable自动分配,此时代表微秒级的“真实时间”;也可以由客户端应用显式赋值。需要避免冲突的应用必须自行生成唯一的时间戳。单元格的不同版本按时间戳降序存储,这样最新的版本会被优先读取。

为了降低版本化数据的管理成本,我们支持两种列族级别的配置,让Bigtable自动回收单元格版本。客户端可以指定只保留单元格的最近(n)个版本,或者只保留足够新的版本(例如,只保留最近七天内写入的值)。

在Webtable示例中,我们将contents:列下存储的爬取页面的时间戳,设置为这些页面版本实际被爬取的时间。通过上述自动回收机制,我们可以只保留每个页面的最近三个版本。

3 API

Bigtable API提供了创建和删除表、列族的函数,同时也提供了修改集群、表和列族元数据的功能,比如访问控制权限。

客户端应用可以在Bigtable中写入或删除数据、查询单行数据,或者遍历表中的部分数据。图2展示了使用RowMutation抽象执行一系列更新操作的C++代码。(为了保持示例简洁,省略了无关细节。)调用Apply会对Webtable执行一次原子更新:它为www.cnn.com添加一个锚点,同时删除另一个锚点。

// Open the table
Table *T = OpenOrDie("/bigtable/web/webtable");
// Write a new anchor and delete an old anchor
RowMutation r1(T, "com.cnn.www");
r1.Set("anchor:www.c-span.org", "CNN");
r1.Delete("anchor:www.abc.com");
Operation op;
Apply(&op, &r1);

图2:向Bigtable写入数据

图3展示了使用Scanner抽象遍历特定行中所有锚点的C++代码。客户端可以遍历多个列族,同时有多种机制可以限制扫描返回的行、列和时间戳。例如,我们可以将上述扫描限制为只返回列名匹配正则表达式anchor:*.cnn.com的锚点,或者只返回时间戳在当前时间十天以内的锚点。

Scanner scanner(T);
ScanStream *stream;
stream = scanner.FetchColumnFamily("anchor");
stream->SetReturnAllVersions();
scanner.Lookup("com.cnn.www");
for (; !stream->Done(); stream->Next()) {
    printf("%s %s %lld %s\n",
        scanner.RowName(),
        stream->ColumnName(),
        stream->MicroTimestamp(),
        stream->Value());
}

图3:从Bigtable读取数据

Bigtable还支持其他几项功能,允许用户以更复杂的方式操作数据。首先,Bigtable支持单行事务,可用于对单个行键下的数据执行原子的读-改-写序列。Bigtable目前不支持跨行键的通用事务,不过它提供了在客户端批量执行跨行写入的接口。其次,Bigtable允许将单元格用作整数计数器。最后,Bigtable支持在服务器的地址空间中执行客户端提供的脚本。这些脚本使用Google开发的一种数据处理语言Sawzall编写[28]。目前,我们基于Sawzall的API不允许客户端脚本写回Bigtable,但支持各种形式的数据转换、基于任意表达式的过滤,以及通过各类算子实现的聚合。

Bigtable可以与MapReduce[12]配合使用,MapReduce是Google开发的大规模并行计算框架。我们编写了一组包装器,让Bigtable既可以作为MapReduce作业的输入源,也可以作为输出目标。

4 基础构件

Bigtable构建于Google的多项其他基础设施之上。Bigtable使用分布式Google文件系统(GFS)[17]存储日志文件和数据文件。Bigtable集群通常运行在共享机器池中,这些机器上同时运行着各类其他分布式应用,Bigtable进程经常与其他应用的进程共享同一台机器。Bigtable依赖集群管理系统来调度作业、管理共享机器上的资源、处理机器故障以及监控机器状态。

Bigtable内部使用Google SSTable文件格式存储数据。SSTable提供了一种持久化、有序且不可变的键值映射,键和值都是任意字节串。SSTable支持两种操作:查询指定键对应的值,以及遍历指定键范围内的所有键值对。在内部,每个SSTable包含一系列数据块(通常每个块大小为64KB,该值可配置)。块索引存储在SSTable的末尾,用于定位数据块;打开SSTable时,索引会被加载到内存中。一次查询只需要一次磁盘寻道即可完成:我们先通过在内存索引中进行二分查找找到对应的数据块,然后从磁盘读取该块。可选地,SSTable可以完全映射到内存中,这样查询和扫描操作都无需访问磁盘。

Bigtable依赖一套高可用、持久化的分布式锁服务,名为Chubby[8]。一个Chubby服务包含5个活跃副本,其中一个被选举为主节点,主动处理请求。当大多数副本都在运行且可以互相通信时,服务处于可用状态。Chubby使用Paxos算法[9,23]保证副本在故障场景下保持一致。Chubby提供了一个由目录和小文件组成的命名空间,每个目录或文件都可以用作锁,对文件的读写都是原子性的。Chubby客户端库提供了Chubby文件的一致性缓存。每个Chubby客户端都会与Chubby服务维护一个会话。如果客户端无法在会话租约到期前续约,会话就会过期。客户端会话过期时,会丢失持有的所有锁和打开的文件句柄。Chubby客户端还可以在Chubby文件和目录上注册回调,以接收变更通知或会话过期通知。

Bigtable将Chubby用于多种任务:确保任意时刻最多只有一个活跃的主节点;存储Bigtable数据的引导位置(见5.1节);发现子表服务器并确认子表服务器失效(见5.2节);存储Bigtable的Schema信息(每个表的列族信息);以及存储访问控制列表。如果Chubby长时间不可用,Bigtable也会不可用。我们最近在横跨11个Chubby实例的14个Bigtable集群中测量了这一影响:由于Chubby不可用(由Chubby宕机或网络问题导致),导致Bigtable中部分数据无法访问的时间,占Bigtable服务器总运行时长的平均比例为0.0047%。受Chubby不可用影响最严重的单个集群,该比例为0.0326%。

5 实现

Bigtable的实现包含三个核心组件:嵌入每个客户端的库、一个主服务器,以及多个子表服务器。可以根据工作负载的变化,动态地向集群中添加或移除子表服务器。

主服务器负责将子表分配给子表服务器、检测子表服务器的加入与退出、平衡子表服务器的负载,以及回收GFS中的文件垃圾。此外,它还处理Schema变更,比如创建表和列族。

每个子表服务器管理一组子表(通常每个子表服务器管理10到1000个子表)。子表服务器处理其所加载子表的读写请求,同时在子表增长过大时进行分裂。

与许多单主分布式存储系统一样[17,21],客户端数据不经过主服务器转发:客户端直接与子表服务器通信进行读写。由于Bigtable客户端不依赖主服务器获取子表位置信息,大多数客户端从不与主服务器通信。因此,实际运行中主服务器的负载很低。

一个Bigtable集群存储多张表。每张表由一组子表组成,每个子表包含一个行范围内的所有数据。初始时,每张表只有一个子表。随着表不断增长,它会自动分裂为多个子表,默认每个子表大小约为100-200 MB。

5.1 子表定位

我们使用一种类似B+树[10]的三层层级结构来存储子表位置信息(图4)。

图4:子表定位层级结构

第一层是存储在Chubby中的一个文件,其中包含**根表(root tablet)**的位置。根表包含了一个特殊的METADATA表中所有子表的位置。每个METADATA子表包含一组用户表的位置。根表只是METADATA表的第一个子表,但它被特殊处理——永远不会分裂——以确保子表定位层级最多只有三层。

METADATA表中,每个子表的位置信息以行键存储,行键由子表所属的表标识和子表的结束行编码而成。每个METADATA行在内存中存储约1KB的数据。在METADATA子表大小上限为128MB的合理限制下,我们的三层定位方案足以寻址(2^{34})个子表(若每个子表128MB,则对应(2^{61})字节的数据)。

客户端库会缓存子表位置。如果客户端不知道某个子表的位置,或者发现缓存的位置信息有误,它就会递归地向上查询子表定位层级。如果客户端缓存为空,定位算法需要三次网络往返,包括一次从Chubby读取数据。如果客户端缓存失效,定位算法最多可能需要六次往返,因为只有在查询未命中时才会发现缓存条目过期(假设METADATA子表不会频繁移动)。尽管子表位置存储在内存中,无需访问GFS,我们仍通过客户端库的预取机制进一步降低了常见场景下的开销:每当读取METADATA表时,会同时读取多个子表的元数据。

我们还在METADATA表中存储辅助信息,包括每个子表的所有事件日志(比如服务器开始加载该子表的时间)。这些信息对调试和性能分析很有帮助。

5.2 子表分配

每个子表同一时间只会分配给一个子表服务器。主服务器维护着所有活跃子表服务器的集合,以及子表到子表服务器的当前分配关系,包括未分配的子表。当存在未分配的子表,且有子表服务器有足够空间承载该子表时,主服务器会向该子表服务器发送子表加载请求,完成分配。

Bigtable通过Chubby跟踪子表服务器的状态。子表服务器启动时,会在Chubby的指定目录下创建一个唯一命名的文件,并获取该文件的排他锁。主服务器监控这个目录(称为服务器目录)来发现子表服务器。如果子表服务器丢失了排他锁(比如由于网络分区导致服务器丢失了Chubby会话),它就会停止为其承载的子表提供服务。(Chubby提供了一种高效机制,让子表服务器无需产生网络流量即可检查自己是否仍持有锁。)只要文件仍然存在,子表服务器就会尝试重新获取该文件的排他锁。如果文件已不存在,子表服务器就永远无法再提供服务,因此会自行终止。每当子表服务器退出时(比如集群管理系统将该子表服务器的机器从集群中移除),它会尝试释放锁,以便主服务器可以更快地重新分配它的子表。

主服务器负责检测子表服务器何时停止服务,并尽快重新分配这些子表。为了检测子表服务器是否仍在服务,主服务器会定期询问每个子表服务器其锁的状态。如果子表服务器报告已丢失锁,或者主服务器连续多次尝试都无法联系到某个服务器,主服务器就会尝试获取该服务器文件的排他锁。如果主服务器成功获取锁,说明Chubby正常运行,而该子表服务器要么已宕机,要么无法连接到Chubby,此时主服务器会删除该服务器的文件,确保它永远无法再提供服务。服务器文件被删除后,主服务器就可以将之前分配给该服务器的所有子表移入未分配子表集合。为了避免Bigtable集群因主服务器与Chubby之间的网络问题而出现风险,若主服务器的Chubby会话过期,主服务器会自行终止。不过如前所述,主服务器故障不会改变子表到子表服务器的分配关系。

当集群管理系统启动一个主服务器时,它需要先发现当前的子表分配情况,然后才能进行变更。主服务器启动时执行以下步骤:

  1. 主服务器在Chubby中获取一个唯一的主锁,防止同时存在多个主服务器实例。
  2. 主服务器扫描Chubby中的服务器目录,找出所有活跃的服务器。
  3. 主服务器与每个活跃子表服务器通信,获取每个服务器当前已分配的子表。
  4. 主服务器扫描METADATA表,获取所有子表的集合。扫描过程中,每当发现一个尚未分配的子表,主服务器就将其加入未分配子表集合,使其可以被分配。

一个复杂之处在于,METADATA表的扫描必须在METADATA子表分配完成后才能进行。因此,在开始扫描(步骤4)之前,如果步骤3中没有发现根表的分配信息,主服务器会将根表加入未分配子表集合。这一操作保证了根表一定会被分配。由于根表包含所有METADATA子表的名称,主服务器扫描完根表后就能知晓所有METADATA子表。

只有在创建或删除表、将两个已有子表合并为一个更大的子表,或者将一个已有子表分裂为两个更小的子表时,现有子表的集合才会发生变化。主服务器能够跟踪这些变化,因为除了子表分裂之外,其他操作都由主服务器发起。子表分裂比较特殊,因为它是由子表服务器发起的。子表服务器通过在METADATA表中记录新子表的信息来提交分裂操作,分裂提交后会通知主服务器。如果分裂通知丢失(无论是子表服务器还是主服务器宕机导致),当主服务器要求子表服务器加载一个已经分裂的子表时,主服务器会发现新的子表。子表服务器会通知主服务器分裂事件,因为它在METADATA表中找到的子表条目,只包含主服务器要求加载的子表的一部分范围。

5.3 子表服务

子表的持久化状态存储在GFS中,如图5所示。更新操作会提交到存储重做记录的提交日志中。在这些更新里,最近提交的更新存储在内存中一个排序的缓冲区里,称为内存表(memtable);更早的更新则存储在一系列SSTable中。

图5:子表结构

要恢复一个子表,子表服务器会从METADATA表中读取该子表的元数据。元数据包含构成该子表的所有SSTable列表,以及一组重做点——重做点是指向提交日志的指针,日志中可能包含该子表的数据。服务器将所有SSTable的索引读入内存,并重做重做点之后提交的所有更新,以此重建内存表。

当写请求到达子表服务器时,服务器会先检查请求格式是否合法,以及发送者是否有权执行该更新操作。授权检查通过读取Chubby文件中的允许写入者列表完成(该列表几乎总能在Chubby客户端缓存中命中)。合法的更新会被写入提交日志。我们使用组提交机制来提升大量小更新操作的吞吐量[13,16]。写入提交完成后,更新内容会被插入到内存表中。

当读请求到达子表服务器时,同样会进行格式检查与权限验证。合法的读操作会在SSTable序列与内存表的合并视图上执行。由于SSTable和内存表都是字典序排序的数据结构,合并视图可以高效地构建。

在子表分裂和合并过程中,正常的读写操作可以继续进行。

5.4 合并

随着写操作不断执行,内存表的大小会持续增长。当内存表大小达到阈值时,就会被冻结,同时创建一个新的内存表,被冻结的内存表会转换为SSTable并写入GFS。这个过程称为小合并(minor compaction),它有两个目的:减少子表服务器的内存占用,以及减少服务器宕机时恢复过程中需要从提交日志读取的数据量。合并过程中,正常的读写操作可以继续进行。

每次小合并都会生成一个新的SSTable。如果放任这种情况持续下去,读操作可能需要合并来自任意数量SSTable的更新。为此,我们通过在后台定期执行**合并归并(merging compaction)**来限制SSTable的数量。合并归并会读取若干个SSTable和内存表的内容,写出一个新的SSTable。合并完成后,输入的SSTable和内存表就可以被丢弃。

将所有SSTable重写为单个SSTable的合并归并,称为大合并(major compaction)。非大合并生成的SSTable可以包含特殊的删除条目,用于屏蔽仍在使用的旧SSTable中的已删除数据。而大合并生成的SSTable不包含任何删除信息和已删除数据。Bigtable会遍历所有子表,定期对它们执行大合并。大合并让Bigtable可以回收已删除数据占用的资源,同时确保已删除数据及时从系统中消失,这对于存储敏感数据的服务至关重要。

6 性能优化

上一节描述的基础实现,还需要进行多项优化,才能达到用户要求的高性能、高可用与高可靠性。本节将更详细地介绍部分实现细节,以展示这些优化点。

局部组

客户端可以将多个列族组合为一个局部组(locality group)。每个子表中的每个局部组都会生成独立的SSTable。将通常不会被同时访问的列族拆分到不同的局部组,可以提升读取效率。例如,Webtable中的页面元数据(比如语言、校验和)可以放在一个局部组,而页面内容可以放在另一个组:只需要读取元数据的应用,就不必遍历所有页面内容。

此外,还可以在局部组级别指定一些实用的调优参数。例如,可以将某个局部组声明为内存型。内存型局部组的SSTable会被懒加载到子表服务器的内存中,加载完成后,访问该局部组内的列族就无需访问磁盘。该特性适用于访问频繁的小批量数据:我们在内部将其用于METADATA表的位置列族。

压缩

客户端可以控制局部组的SSTable是否启用压缩,以及使用哪种压缩格式。用户指定的压缩格式会应用于每个SSTable数据块(块大小可以通过局部组的调优参数控制)。虽然逐个块压缩会损失一些空间收益,但好处是只需解压部分数据就能读取SSTable的小片段,无需解压整个文件。

许多客户端使用一种自定义的两遍压缩方案。第一遍使用Bentley和McIlroy的算法[6],在大窗口内压缩长公共字符串。第二遍使用快速压缩算法,在16KB的小窗口内查找重复数据。两遍压缩的速度都非常快——在现代机器上,压缩速度可达100-200 MB/s,解压速度可达400-1000 MB/s。

尽管我们选择压缩算法时更看重速度而非空间缩减率,但这套两遍压缩方案的表现却出人意料。例如在Webtable中,我们使用该压缩方案存储网页内容。在一项实验中,我们将大量文档存储在一个压缩局部组中,实验中我们只存储每个文档的一个版本,而非所有可用版本。该方案实现了10:1的空间缩减率。这远好于Gzip对HTML页面通常3:1或4:1的压缩率,原因在于Webtable的行布局:同一主机的所有页面都相邻存储,这让Bentley-McIlroy算法可以识别出同一主机页面中大量的公共模板内容。不仅是Webtable,许多应用都会通过设计行名,让相似数据聚集在一起,从而获得极佳的压缩比。当在Bigtable中存储同一数据的多个版本时,压缩比会更高。

读性能缓存

为了提升读取性能,子表服务器使用两级缓存。**扫描缓存(Scan Cache)**是高层缓存,缓存SSTable接口返回给子表服务器代码的键值对。**块缓存(Block Cache)**是低层缓存,缓存从GFS读取的SSTable数据块。扫描缓存最适合重复读取相同数据的应用;块缓存适合读取的数据与近期读取的数据位置相近的场景(比如顺序读取,或者在热点行中读取同一局部组内的不同列)。

布隆过滤器

如5.3节所述,读操作需要从构成子表状态的所有SSTable中读取数据。如果这些SSTable不在内存中,就可能需要多次磁盘访问。我们允许客户端指定为特定局部组的SSTable创建布隆过滤器[7],以此减少访问次数。布隆过滤器可以快速判断一个SSTable中是否可能包含指定行/列对的数据。对于某些应用,用少量子表服务器内存存储布隆过滤器,可以大幅减少读操作所需的磁盘寻道次数。布隆过滤器的使用还意味着,大多数对不存在的行或列的查询都不需要访问磁盘。

提交日志实现

如果我们为每个子表单独维护一个提交日志文件,那么GFS中会同时写入大量文件。取决于每个GFS服务器底层文件系统的实现,这些写入操作可能需要大量磁盘寻道,才能写入不同的物理日志文件。此外,每个子表单独使用日志文件也会降低组提交优化的效果,因为每个提交组的规模会更小。为了解决这些问题,我们将所有更新追加到每个子表服务器的单个提交日志中,将不同子表的更新混合在同一个物理日志文件中[18,20]。

使用单一日志在正常运行时能带来显著的性能收益,但会让恢复过程变得复杂。当一台子表服务器宕机后,它承载的子表会被转移到大量其他子表服务器上:每个服务器通常只加载原服务器的少数子表。要恢复一个子表的状态,新的子表服务器需要从原服务器写入的提交日志中,重做该子表对应的更新。然而,这些子表的更新是混合在同一个物理日志文件中的。一种方案是,每个新的子表服务器都读取完整的提交日志文件,只应用自身需要恢复的子表对应的条目。但在这种方案下,如果100台机器各分配到宕机服务器的一个子表,那么日志文件会被读取100次(每台服务器读一次)。

我们通过先对提交日志条目按⟨表, 行名, 日志序列号⟩排序,避免了重复读取日志。排序后的输出中,同一个子表的所有更新都是连续的,因此只需一次磁盘寻道加上顺序读取,就能高效地获取这些数据。为了并行化排序过程,我们将日志文件分成64MB的段,在不同的子表服务器上并行排序每个段。排序过程由主服务器协调,当子表服务器表明需要从某个提交日志文件恢复更新时,就会启动该过程。

将提交日志写入GFS有时会出现性能波动,原因多种多样:比如写入涉及的某台GFS服务器宕机,或者访问某组GFS服务器的网络路径出现拥塞、负载过高。为了避免更新操作受GFS延迟尖峰的影响,每个子表服务器实际上有两个日志写入线程,各自写入自己的日志文件;同一时间只有一个线程处于活跃状态。如果活跃日志文件的写入性能变差,就切换到另一个线程,由新的活跃线程写入提交日志队列中的更新。日志条目包含序列号,让恢复过程可以跳过日志切换导致的重复条目。

加速子表恢复

如果主服务器将一个子表从一台子表服务器迁移到另一台,源子表服务器会先对该子表执行一次小合并。这次合并减少了子表服务器提交日志中未合并的状态量,从而缩短恢复时间。合并完成后,子表服务器停止为该子表提供服务。在真正卸载子表之前,子表服务器还会再执行一次(通常非常快的)小合并,以消除第一次小合并期间到达的、仍在日志中的未合并状态。第二次小合并完成后,子表就可以加载到另一台服务器上,且无需重做任何日志条目。

利用不可变性

除了SSTable缓存之外,由于所有生成的SSTable都是不可变的,Bigtable系统的其他很多部分也得到了简化。例如,读取SSTable时,我们不需要对文件系统的访问做任何同步。因此,行级的并发控制可以非常高效地实现。唯一同时被读写访问的可变数据结构只有内存表。为了减少读取内存表时的竞争,我们将每个内存表行实现为写时复制,允许读写操作并行进行。

由于SSTable是不可变的,永久删除已删除数据的问题就转化为回收过时SSTable的垃圾回收问题。每个子表的SSTable都注册在METADATA表中。主服务器通过对SSTable集合执行标记-清扫式垃圾回收[25]来移除过时的SSTable,其中METADATA表包含了所有活跃SSTable的根。

在小型Bigtable集群中,要删除过时文件,主服务器会定期扫描对应目录树,收集集群所有现有文件的名称,然后扫描METADATA表获取集群所有活跃文件的集合,最后删除前者有而后者没有的文件。(子表的METADATA行中包含了存储该子表数据的所有文件的名称。)这种方法的扩展性不佳,因此在大型Bigtable集群中,主服务器将该过程拆分为更小的块,每次只读取部分目录树和METADATA表,就能检测并删除该部分目录树中的失效文件。

最后,SSTable的不可变性让我们可以快速分裂子表。我们不需要为每个子子表生成一整套新的SSTable,而是让子子表共享父表的SSTable。

7 性能评估

我们搭建了一个包含N台子表服务器的Bigtable集群,通过改变N的取值来衡量Bigtable的性能与可扩展性。子表服务器配置为使用1GB内存,写入到由1786台机器组成的GFS集群,每台机器配备两块400GB IDE硬盘。N台客户端机器生成测试用的Bigtable负载。(我们使用与子表服务器数量相同的客户端,确保客户端不会成为瓶颈。)每台机器配备两颗双核Opteron 2GHz处理器,足够的物理内存以容纳所有运行进程的工作集,以及一条千兆以太网链路。机器组成两层树形交换网络,根节点的总带宽约为100-200 Gbps。所有机器都在同一个机房内,因此任意两台机器之间的往返时间都小于1毫秒。

子表服务器与主服务器、测试客户端、GFS服务器都运行在同一组机器上。每台机器都运行一个GFS服务器,部分机器还运行子表服务器、客户端进程,或者其他同时在使用该资源池的作业进程。

R是测试中涉及的不同Bigtable行键的数量。R的取值使得每个基准测试中,每台子表服务器读写约1GB的数据。

顺序写基准测试使用名称从0到(R-1)的行键。行键空间被划分为(10N)个大小相等的范围,由中央调度器分配给N个客户端:客户端完成上一个分配的范围后,就会被分配下一个可用范围。这种动态分配有助于缓解客户端机器上其他进程导致的性能波动。我们在每个行键下写入一个字符串,每个字符串都是随机生成的,因此无法压缩。此外,不同行键下的字符串互不相同,因此无法实现跨行压缩。

随机写基准测试与顺序写类似,区别在于写入前会先对行键做R取模的哈希运算,这样写入负载在整个测试期间会大致均匀地分布在整个行空间。

顺序读基准测试生成行键的方式与顺序写完全相同,但它不写入行键下的数据,而是读取该行键下存储的字符串(由之前的顺序写基准测试写入)。类似地,随机读基准测试对应随机写基准测试的操作。

扫描基准测试与顺序读基准测试类似,但使用Bigtable API提供的扫描功能,遍历行范围内的所有值。使用扫描可以减少基准测试执行的RPC次数,因为一次RPC就可以从子表服务器获取大量值。

内存随机读基准测试与随机读基准测试类似,但存储测试数据的局部组被标记为内存型,因此读取操作由子表服务器的内存提供服务,无需访问GFS。仅在该基准测试中,我们将每台子表服务器的数据量从1GB减少到100MB,使其可以舒适地容纳在子表服务器的可用内存中。

图6展示了读写1000字节值时,各项基准测试的性能表现,包含两个视角:表格展示了每台子表服务器的每秒操作数,折线图展示了总每秒操作数。

图6:每秒读写1000字节值的数量。表格展示单台服务器的速率;图表展示总速率。

实验类型 1台 50台 250台 500台
随机读 1212 593 479 241
内存随机读 10811 8511 8000 6250
随机写 8850 3745 3425 2000
顺序读 4425 2463 2625 2469
顺序写 8547 3623 2451 1905
扫描 15385 10526 9524 7843

单台子服务器性能

首先来看单台子表服务器的性能。随机读的速度比其他所有操作慢一个数量级以上。每次随机读都需要通过网络将一个64KB的SSTable块从GFS传输到子表服务器,而其中只有一个1000字节的值会被使用。子表服务器每秒执行约1200次读取,相当于从GFS读取约75MB/s的数据。该带宽足以让子表服务器的CPU达到饱和,因为网络协议栈、SSTable解析以及Bigtable代码都存在开销,同时也几乎足以让系统中的网络链路达到饱和。对于这种访问模式的Bigtable应用,大多数都会将块大小调小到更小的值,通常是8KB。

内存随机读的速度要快得多,因为每次1000字节的读取都由子表服务器的本地内存完成,无需从GFS获取64KB的大数据块。

随机写和顺序写的性能都优于随机读,因为每个子表服务器将所有写入请求都追加到单个提交日志中,并使用组提交将这些写入高效地传输到GFS。随机写和顺序写的性能没有显著差异;两种情况下,所有写入子表服务器的操作都记录在同一个提交日志中。

顺序读的性能优于随机读,因为每个从GFS获取的64KB SSTable块都会存储在块缓存中,用于服务接下来的64次读取请求。

扫描的速度更快,因为子表服务器可以在一次客户端RPC中返回大量值,因此RPC开销被分摊到大量值上。

可扩展性

当系统中的子表服务器数量从1台增加到500台时,总吞吐量大幅提升,增长超过100倍。例如,内存随机读的性能随着子表服务器数量增长500倍,提升了近300倍。出现这种表现的原因是,该基准测试的性能瓶颈是单个子表服务器的CPU。

然而,性能并非线性增长。对于大多数基准测试,从1台增加到50台子表服务器时,单台服务器的吞吐量会出现显著下降。这种下降是由多服务器配置下的负载不均衡导致的,通常是因为其他进程在争抢CPU和网络资源。我们的负载均衡算法试图处理这种不均衡,但无法做到完美,主要有两个原因:为了减少子表迁移次数,再平衡操作会被限流(子表迁移时会有短暂的不可用时间,通常不到1秒);并且基准测试产生的负载会随着测试进行而发生变化。

随机读基准测试的扩展性最差(服务器数量增长500倍,总吞吐量仅增长100倍)。出现这种表现的原因是(如前所述),每读取1000字节,我们就要通过网络传输一个64KB的大块。这种传输会使网络中各种共享的千兆链路达到饱和,因此随着机器数量增加,单台服务器的吞吐量会显著下降。

8 实际应用

截至2006年8月,Google的各个机器集群中运行着388个非测试用途的Bigtable集群,总计约有24500台子表服务器。表1展示了每个集群的子表服务器数量分布。其中许多集群用于开发目的,因此有大量空闲时间。一组包含14个繁忙集群、共8069台子表服务器的系统,每秒处理总量超过120万次请求,入站RPC流量约741 MB/s,出站RPC流量约16 GB/s。

表1:Bigtable集群的子表服务器数量分布

子表服务器数量 集群数量
0 – 19 259
20 – 49 47
50 – 99 20
100 – 499 50
> 500 12

表2提供了当前在用的部分表的相关数据。有些表存储的数据用于向用户提供服务,有些则存储批量处理的数据;这些表在总大小、平均单元格大小、内存服务数据占比以及表Schema复杂度方面差异很大。在本节余下部分,我们将简要介绍三个产品团队如何使用Bigtable。

表2:生产环境中部分表的特性。表大小(压缩前测量)与单元格数量为近似值。未启用压缩的表未给出压缩比。

项目名称 表大小(TB) 压缩比 单元格数量(十亿) 列族数量 局部组数量 内存占比 是否延迟敏感
Crawl 800 11% 1000 16 8 0%
Crawl 50 33% 200 2 2 0%
Google Analytics 20 29% 10 1 1 0%
Google Analytics 200 14% 80 1 1 0%
Google Base 2 31% 10 29 3 15%
Google Earth 0.5 64% 8 7 2 33%
Google Earth 70 9 8 3 0%
Orkut 9 0.9 8 5 1%
个性化搜索 4 47% 6 93 11 5%

8.1 Google Analytics

Google Analytics(analytics.google.com)是一项帮助网站站长分析网站流量模式的服务。它提供汇总统计数据,比如每日独立访客数、每日每个URL的页面浏览量,以及网站追踪报告,比如用户在浏览特定页面后完成购买的比例。

为了实现该服务,网站站长会在网页中嵌入一小段JavaScript程序。每当页面被访问时,该程序就会被调用,将请求的各类信息记录到Google Analytics中,比如用户标识符以及被访问页面的信息。Google Analytics对这些数据进行汇总,再提供给网站站长。

我们简要介绍Google Analytics使用的两张表。原始点击表(约200TB)为每个终端用户会话维护一行,行名是一个元组,包含网站名称和会话创建时间。这种Schema保证了访问同一网站的会话是连续的,且按时间顺序排列。该表压缩后为原始大小的14%。

汇总表(约20TB)包含每个网站的各类预定义汇总数据。该表由原始点击表通过定期调度的MapReduce作业生成。每个MapReduce作业从原始点击表中提取最近的会话数据。整个系统的吞吐量受限于GFS的吞吐量。该表压缩后为原始大小的29%。

8.2 Google Earth

Google运营着一系列服务,为用户提供全球高分辨率卫星影像访问,包括基于网页的Google Maps界面(maps.google.com)和Google Earth(earth.google.com)客户端软件。这些产品允许用户在地球表面浏览:可以平移、查看、标注不同分辨率的卫星影像。该系统使用一张表进行数据预处理,另一组表用于向客户端提供数据服务。

预处理流水线使用一张表存储原始影像。在预处理过程中,影像表中的每一行对应一个单独的地理分片。行名的设计确保了相邻的地理分片相邻存储。该表包含一个列族,用于跟踪每个分片的数据来源。这个列族拥有大量的列:基本上每个原始数据影像对应一列。由于每个分片仅由少量影像拼接而成,因此该列族非常稀疏。

预处理流水线重度依赖基于Bigtable的MapReduce来转换数据。在部分MapReduce作业中,整个系统每台子表服务器的数据处理速度超过1 MB/s。

服务系统使用一张表为存储在GFS中的数据建立索引。这张表相对较小(约500GB),但必须在每个数据中心每秒承载数万次查询,且延迟很低。因此,该表部署在数百台子表服务器上,并且包含内存型列族。

8.3 个性化搜索

个性化搜索(www.google.com/psearch)是一项可选服务,会记录用户在Google各类产品(如网页搜索、图片搜索、新闻搜索)中的查询和点击行为。用户可以浏览自己的搜索历史,回看之前的查询和点击记录,还可以基于自己的Google历史使用模式,获取个性化的搜索结果。

个性化搜索将每个用户的数据存储在Bigtable中。每个用户有唯一的用户ID,并以该ID作为行名。所有用户行为都存储在一张表中,每种行为类型都对应一个独立的列族(例如,有一个列族存储所有网页查询)。每个数据元素的Bigtable时间戳,就是对应用户行为发生的时间。个性化搜索通过基于Bigtable的MapReduce生成用户画像,这些画像用于个性化实时搜索结果。

个性化搜索的数据在多个Bigtable集群间复制,以提升可用性,并降低客户端距离带来的延迟。个性化搜索团队最初在Bigtable之上构建了一套客户端复制机制,保证所有副本的最终一致性。当前系统则使用内置于服务器的复制子系统。

个性化搜索存储系统的设计,允许其他团队在自己的列中添加新的用户维度信息,现在已有许多其他Google产品使用该系统存储用户级的配置选项与设置。多个团队共享一张表,导致列族数量异常多。为了帮助支持共享,我们在Bigtable中增加了简单的配额机制,限制共享表中任意特定客户端的存储消耗量;该机制为使用该系统存储用户信息的各个产品团队提供了一定的隔离性。

9 经验教训

在设计、实现、维护和支持Bigtable的过程中,我们获得了宝贵的经验,也总结出了一些有意思的教训。

我们学到的一个教训是:大型分布式系统容易受到多种类型的故障影响,而不仅仅是许多分布式协议中假设的标准网络分区和故障停止失效。例如,我们遇到过由以下所有原因导致的问题:内存与网络数据损坏、时钟大幅偏移、机器挂死、长时间且不对称的网络分区、我们依赖的其他系统(比如Chubby)的Bug、GFS配额溢出,以及计划内和计划外的硬件维护。随着我们对这些问题积累了更多经验,我们通过修改各类协议来解决它们。例如,我们为RPC机制增加了校验和。我们还通过移除系统中一个模块对另一个模块的假设,解决了一些问题。例如,我们不再假设Chubby操作只会返回固定集合中的错误。

我们学到的另一个教训是:推迟添加新功能,直到明确新功能的使用方式,这一点非常重要。例如,我们最初计划在API中支持通用事务。但由于当时没有立即的使用需求,我们并没有实现。现在,有了众多运行在Bigtable上的实际应用后,我们得以考察它们的真实需求,并且发现大多数应用只需要单行事务。在人们提出的分布式事务需求中,最重要的用途是维护二级索引,我们计划添加一套专门的机制来满足这一需求。新机制不会像分布式事务那样通用,但效率会更高(尤其是对于跨越数百行甚至更多行的更新),并且能更好地与我们的跨数据中心乐观复制方案配合。

从运维Bigtable的实践中,我们学到的一个实用教训是:完善的系统级监控至关重要(即既要监控Bigtable本身,也要监控使用Bigtable的客户端进程)。例如,我们扩展了RPC系统,对采样的RPC保留详细的追踪,记录为该RPC执行的所有关键操作。该特性让我们能够发现并解决很多问题,比如子表数据结构上的锁竞争、提交Bigtable更新时GFS写入缓慢、METADATA子表不可用时对METADATA表的访问卡住等。另一个有效监控的例子是,每个Bigtable集群都在Chubby中注册,这让我们可以追踪所有集群,了解它们的规模、运行的软件版本、承载的流量大小,以及是否存在异常高延迟等问题。

我们学到的最重要的教训是:简洁设计的价值。考虑到我们系统的规模(约10万行非测试代码),以及代码会随着时间以意想不到的方式演进,我们发现代码与设计的清晰性对于代码维护和调试帮助极大。

这方面的一个例子是我们的子表服务器成员协议。我们最初的协议很简单:主服务器定期向子表服务器发放租约,如果租约过期,子表服务器就自行终止。然而,在存在网络问题的情况下,该协议会显著降低可用性,因此我们多次重新设计协议,直到它能良好运行。但最终的协议过于复杂,并且依赖于Chubby中很少被其他应用使用的特性。我们发现自己花费了大量时间调试各种极端边界情况,不仅是Bigtable代码,还有Chubby代码。最终,我们废弃了该协议,改用了一套更新、更简洁的协议,它只依赖Chubby中被广泛使用的特性。

10 相关工作

Boxwood项目[24]的组件在某些方面与Chubby、GFS和Bigtable有重叠,因为它提供了分布式一致性、锁、分布式块存储和分布式B树存储。在存在重叠的每个领域,Boxwood的组件定位似乎都比对应的Google服务更低层。Boxwood项目的目标是为构建文件系统或数据库等更上层服务提供基础设施,而Bigtable的目标是直接为需要存储数据的客户端应用提供支持。

近年来,许多项目都致力于解决在广域网上提供分布式存储或更上层服务的问题,通常是“互联网规模”的。这包括分布式哈希表方面的工作,始于CAN[29]、Chord[32]、Tapestry[37]和Pastry[30]等项目。这些系统解决的问题是Bigtable不会遇到的,比如高度变化的带宽、不可信的参与者、频繁的配置变更;去中心化控制和拜占庭容错也不是Bigtable的目标。

就向应用开发者提供的分布式数据存储模型而言,我们认为分布式B树或分布式哈希表提供的键值对模型局限性太大。键值对是有用的构建块,但不应该是提供给开发者的唯一构建块。我们选择的模型比简单的键值对更丰富,支持稀疏的半结构化数据。尽管如此,它仍然足够简洁,可以实现非常高效的平面文件存储,并且(通过局部组)具备足够的透明度,让用户可以调优系统的关键行为。

多家数据库厂商都开发了可以存储海量数据的并行数据库。Oracle的实时应用集群数据库[27]使用共享磁盘存储数据(Bigtable使用GFS),并使用分布式锁管理器(Bigtable使用Chubby)。IBM的DB2并行版[4]基于与Bigtable类似的无共享[33]架构,每个DB2服务器负责表中的一部分行,存储在本地关系型数据库中。这两款产品都提供完整的带事务的关系模型。

Bigtable的局部组实现了与其他列式存储系统类似的压缩和磁盘读取性能收益,包括C-Store[1,34],以及Sybase IQ[15,36]、SenSage[31]、KDB+[22]等商业产品,还有MonetDB/X100中的ColumnBM存储层[38]。另一个将数据垂直和水平分区到平面文件中、实现了良好数据压缩比的系统是AT&T的Daytona数据库[19]。局部组不支持CPU缓存级别的优化,比如Ailamaki等人描述的优化[2]。

Bigtable使用内存表和SSTable存储子表更新的方式,与日志结构合并树(LSM-tree)[26]存储索引更新的方式类似。在这两个系统中,排序后的数据都会先缓存在内存中,再写入磁盘,读取时必须合并内存和磁盘中的数据。

C-Store与Bigtable有很多共同特点:两个系统都使用无共享架构,都有两种不同的数据结构,一种用于近期写入,一种用于存储长期数据,并且都有机制将数据从一种形式转换为另一种形式。两个系统在API上差异显著:C-Store的行为类似关系型数据库,而Bigtable提供更低层的读写接口,设计目标是每台服务器每秒支持数千次这类操作。C-Store还是一个“读优化的关系型数据库管理系统”,而Bigtable在读密集型和写密集型应用上都能提供良好的性能。

Bigtable的负载均衡器需要解决一些与无共享数据库相同的负载和内存均衡问题(比如[11,35])。我们的问题相对更简单:(1)我们不需要考虑同一数据的多个副本,也不需要考虑由于视图或索引导致的其他形式的数据副本;(2)我们让用户指定哪些数据放在内存,哪些数据留在磁盘,而不是动态地进行判断;(3)我们不需要执行或优化复杂查询。

11 结论

我们介绍了Bigtable——Google内部一套用于存储结构化数据的分布式系统。Bigtable集群自2005年4月起投入生产使用,在此之前我们花费了约7人年的时间进行设计与实现。截至2006年8月,已有超过60个项目在使用Bigtable。用户认可Bigtable实现提供的性能与高可用性,也认可当资源需求随时间变化时,只需向系统中添加更多机器,就能扩展集群容量。

考虑到Bigtable与众不同的接口,一个有意思的问题是:用户适应Bigtable的使用难度有多大。新用户有时会不确定如何最佳地使用Bigtable接口,尤其是如果他们习惯了支持通用事务的关系型数据库。尽管如此,众多Google产品成功使用Bigtable的事实证明,我们的设计在实践中表现良好。

我们正在实现Bigtable的几项额外功能,比如支持二级索引,以及构建跨数据中心、多主副本复制的Bigtable基础设施。我们也开始将Bigtable作为服务提供给各个产品团队,这样各个团队就不需要维护自己的集群。随着服务集群规模扩大,我们需要在Bigtable内部处理更多的资源共享问题[3,5]。

最后,我们发现在Google构建自己的存储解决方案有显著的优势。通过为Bigtable设计自己的数据模型,我们获得了极大的灵活性。此外,我们对Bigtable实现、以及Bigtable所依赖的其他Google基础设施的掌控力,意味着我们可以在瓶颈和低效问题出现时就将其消除。

致谢

感谢匿名审稿人、David Nagle,以及我们的指导Brad Calder,感谢他们对本文的反馈。Bigtable系统从Google内部众多用户的反馈中获益良多。此外,感谢以下人员对Bigtable做出的贡献:Dan Aguayo、Sameer Ajmani、Zhifeng Chen、Bill Coughran、Mike Epstein、Healfdene Goguen、Robert Griesemer、Jeremy Hylton、Josh Hyman、Alex Khesin、Joanna Kulik、Alberto Lerner、Sherry Listgarten、Mike Maloney、Eduardo Pinheiro、Kathy Polizzi、Frank Yellin以及Arthur Zwiegencew。

参考文献

[1] ABADI, D. J., MADDEN, S. R., AND FERREIRA, M. C. Integrating compression and execution in column-oriented database systems. Proc. of SIGMOD (2006).
[2] AILAMAKI, A., DEWITT, D. J., HILL, M. D., AND SKOUNAKIS, M. Weaving relations for cache performance. In The VLDB Journal (2001), pp. 169–180.
[3] BANGA, G., DRUSCHEL, P., AND MOGUL, J. C. Resource containers: A new facility for resource management in server systems. In Proc. of the 3rd OSDI (Feb. 1999), pp. 45–58.
[4] BARU, C. K., FECTEAU, G., GOYAL, A., HSIAO, H., JHINGRAN, A., PADMANABHAN, S., COPELAND, G. P., AND WILSON, W. G. DB2 parallel edition. IBM Systems Journal 34, 2 (1995), 292–322.
[5] BAVIER, A., BOWMAN, M., CHUN, B., CULLER, D., KARLIN, S., PETERSON, L., ROSCOE, T., SPALINK, T., AND WAWRZONIAK, M. Operating system support for planetary-scale network services. In Proc. of the 1st NSDI (Mar. 2004), pp. 253–266.
[6] BENTLEY, J. L., AND MCILROY, M. D. Data compression using long common strings. In Data Compression Conference (1999), pp. 287–295.
[7] BLOOM, B. H. Space/time trade-offs in hash coding with allowable errors. CACM 13, 7 (1970), 422–426.
[8] BURROWS, M. The Chubby lock service for loosely-coupled distributed systems. In Proc. of the 7th OSDI (Nov. 2006).
[9] CHANDRA, T., GRIESEMER, R., AND REDSTONE, J. Paxos made live — An engineering perspective. In Proc. of PODC (2007).
[10] COMER, D. Ubiquitous B-tree. Computing Surveys 11, 2 (June 1979), 121–137.
[11] COPELAND, G. P., ALEXANDER, W., BOUGHTER, E. E., AND KELLER, T. W. Data placement in Bubba. In Proc. of SIGMOD (1988), pp. 99–108.
[12] DEAN, J., AND GHEMAWAT, S. MapReduce: Simplified data processing on large clusters. In Proc. of the 6th OSDI (Dec. 2004), pp. 137–150.
[13] DEWITT, D., KATZ, R., OLKEN, F., SHAPIRO, L., STONEBRAKER, M., AND WOOD, D. Implementation techniques for main memory database systems. In Proc. of SIGMOD (June 1984), pp. 1–8.
[14] DEWITT, D. J., AND GRAY, J. Parallel database systems: The future of high performance database systems. CACM 35, 6 (June 1992), 85–98.
[15] FRENCH, C. D. One size fits all database architectures do not work for DSS. In Proc. of SIGMOD (May 1995), pp. 449–450.
[16] GAWLICK, D., AND KINKADE, D. Varieties of concurrency control in IMS/VS fast path. Database Engineering Bulletin 8, 2 (1985), 3–10.
[17] GHEMAWAT, S., GOBIOFF, H., AND LEUNG, S.-T. The Google file system. In Proc. of the 19th ACM SOSP (Dec. 2003), pp. 29–43.
[18] GRAY, J. Notes on database operating systems. In Operating Systems — An Advanced Course, vol. 60 of Lecture Notes in Computer Science. Springer-Verlag, 1978.
[19] GREER, R. Daytona and the fourth-generation language Cymbal. In Proc. of SIGMOD (1999), pp. 525–526.
[20] HAGMANN, R. Reimplementing the Cedar file system using logging and group commit. In Proc. of the 11th SOSP (Dec. 1987), pp. 155–162.
[21] HARTMAN, J. H., AND OUSTERHOUT, J. K. The Zebra striped network file system. In Proc. of the 14th SOSP (Asheville, NC, 1993), pp. 29–43.
[22] KX.COM. kx.com/products/database.php. Product page.
[23] LAMPORT, L. The part-time parliament. ACM TOCS 16, 2 (1998), 133–169.
[24] MACCORMICK, J., MURPHY, N., NAJORK, M., THEKKATH, C. A., AND ZHOU, L. Boxwood: Abstractions as the foundation for storage infrastructure. In Proc. of the 6th OSDI (Dec. 2004), pp. 105–120.
[25] MCCARTHY, J. Recursive functions of symbolic expressions and their computation by machine. CACM 3, 4 (Apr. 1960), 184–195.
[26] O’NEIL, P., CHENG, E., GAWLICK, D., AND O’NEIL, E. The log-structured merge-tree (LSM-tree). Acta Inf. 33, 4 (1996), 351–385.
[27] ORACLE.COM. www.oracle.com/technology/products/database/clustering/index.html. Product page.
[28] PIKE, R., DORWARD, S., GRIESEMER, R., AND QUINLAN, S. Interpreting the data: Parallel analysis with Sawzall. Scientific Programming Journal 13, 4 (2005), 227–298.
[29] RATNASAMY, S., FRANCIS, P., HANDLEY, M., KARP, R., AND SHENKER, S. A scalable content-addressable network. In Proc. of SIGCOMM (Aug. 2001), pp. 161– 172.
[30] ROWSTRON, A., AND DRUSCHEL, P. Pastry: Scalable, distributed object location and routing for large-scale peer-to-peer systems. In Proc. of Middleware 2001 (Nov. 2001), pp. 329–350.
[31] SENSAGE.COM. sensage.com/products-sensage.htm. Product page.
[32] STOICA, I., MORRIS, R., KARGER, D., KAASHOEK, M. F., AND BALAKRISHNAN, H. Chord: A scalable peer-to-peer lookup service for Internet applications. In Proc. of SIGCOMM (Aug. 2001), pp. 149–160.
[33] STONEBRAKER, M. The case for shared nothing. Database Engineering Bulletin 9, 1 (Mar. 1986), 4–9.
[34] STONEBRAKER, M., ABADI, D. J., BATKIN, A., CHEN, X., CHERNIACK, M., FERREIRA, M., LAU, E., LIN, A., MADDEN, S., O’NEIL, E., O’NEIL, P., RASIN, A., TRAN, N., AND ZDONIK, S. C-Store: A column-oriented DBMS. In Proc. of VLDB (Aug. 2005), pp. 553– 564.
[35] STONEBRAKER, M., AOKI, P. M., DEVINE, R., LITWIN, W., AND OLSON, M. A. Mariposa: A new architecture for distributed data. In Proc. of the Tenth ICDE (1994), IEEE Computer Society, pp. 54–65.
[36] SYBASE.COM. www.sybase.com/products/databaseservers/sybaseiq. Product page.
[37] ZHAO, B. Y., KUBIATOWICZ, J., AND JOSEPH, A. D. Tapestry: An infrastructure for fault-tolerant wide-area location and routing. Tech. Rep. UCB/CSD-01-1141, CS Division, UC Berkeley, Apr. 2001.
[38] ZUKOWSKI, M., BONCZ, P. A., NES, N., AND HEMAN, S. MonetDB/X100 — A DBMS in the CPU cache. IEEE Data Eng. Bull. 28, 2 (2005), 17–22.

【转载】谷歌文件系统(The Google File System)

The Google File System

谷歌文件系统

Sanjay Ghemawat, Howard Gobioff, and Shun-Tak Leung
Google $^(^∗)$


摘要

我们设计并实现了谷歌文件系统(Google File System, GFS)——这是一个面向大规模分布式数据密集型应用的可扩展分布式文件系统。它能够在廉价的商用硬件上运行,同时提供容错能力,并向大量客户端提供高聚合性能。

尽管与此前的分布式文件系统有着诸多相同目标,但我们的设计源于对应用负载和技术环境的观察(包括当前现状与未来预期),这些观察反映出我们与一些早期文件系统假设存在显著分歧。这促使我们重新审视传统选择,并探索截然不同的设计切入点。

该文件系统已成功满足了我们的存储需求。它作为存储平台在谷歌内部被广泛部署,用于生成和处理各项服务所使用的数据,以及需要大规模数据集的研发工作。截至目前,最大的集群在超过一千台机器的数千块磁盘上提供了数百TB的存储容量,并被数百个客户端并发访问。

在本文中,我们介绍了为支持分布式应用而设计的文件系统接口扩展,讨论了设计中的诸多方面,并给出了微观基准测试与实际应用场景下的性能测量结果。

分类与主题描述
D [4]: 3-分布式文件系统

通用术语
设计、可靠性、性能、测量

关键词
容错、可扩展性、数据存储、集群存储

*作者联系方式:{sanjay, hgobioff, shuntak}@google.com

未经许可,不得以营利或商业优势为目的制作或分发本文的数字或纸质副本;副本需标注本声明及首页完整引用信息。如需其他形式的复制、再发布、在服务器上发布或向列表重新分发,需要事先获得特定许可和/或支付费用。


1. 引言

我们设计并实现了谷歌文件系统(GFS),以满足谷歌数据处理需求的快速增长。GFS与此前的分布式文件系统有着许多相同目标,例如性能、可扩展性、可靠性和可用性。然而,其设计的驱动力源于对我们应用负载和技术环境的关键观察(包括当前现状与未来预期),这些观察反映出与一些早期文件系统设计假设的显著差异。我们重新审视了传统选择,并在设计空间中探索了截然不同的切入点。

首先,组件失效是常态而非异常。文件系统由数百甚至数千台存储机构成,这些机器均由廉价的商用部件组装而成,并且被数量相当的客户端机器访问。从部件的数量和质量来看,几乎可以保证其中一些在任何给定时间都无法正常工作,还有一些将无法从当前故障中恢复。我们见过由应用程序bug、操作系统bug、人为错误,以及磁盘、内存、连接器、网络和电源故障导致的各种问题。因此,持续监控、错误检测、容错和自动恢复必须成为系统不可或缺的组成部分。

其次,文件规模远超传统标准。数GB的文件十分常见。每个文件通常包含大量应用对象,例如网页文档。当我们常规处理数十亿个对象、数TB级的快速增长数据集时,即使文件系统能够支持,管理数十亿个KB级大小的文件也难以操作。因此,I/O操作和块大小等设计假设与参数都需要重新考量。

第三,绝大多数文件的修改方式是追加新数据,而非覆盖已有数据。文件内的随机写入在实际中几乎不存在。文件一旦写入,就只会被读取,而且通常是顺序读取。各类数据都具有这些特征:有些是数据处理程序扫描的大型存储库;有些是运行中的应用持续生成的数据流;有些是归档数据;还有些是在一台机器上生成、在另一台机器上处理的中间结果(可以是同时处理,也可以是后续处理)。鉴于大文件的这种访问模式,追加操作成为性能优化和原子性保证的重点,而在客户端缓存数据块则失去了吸引力。

第四,应用程序与文件系统API的协同设计通过提升灵活性让整个系统受益。例如,我们放宽了GFS的一致性模型,极大地简化了文件系统,同时不会给应用程序带来沉重负担。我们还引入了原子追加操作,使得多个客户端可以并发地向同一个文件追加数据,而无需在它们之间进行额外同步。这些将在本文后续部分详细讨论。

目前已有多个GFS集群部署用于不同用途。最大的集群拥有超过1000个存储节点、超过300TB的磁盘存储,并被数百台不同机器上的客户端持续密集访问。


2. 设计概述

2.1 设计假设

在设计满足我们需求的文件系统时,我们以一系列假设为指导,这些假设既带来挑战,也带来机遇。我们在前文已经提及一些关键观察,现在将更详细地阐述这些假设。

  • 系统由大量廉价商用组件构建,这些组件经常发生故障。系统必须持续监控自身,并例行地、及时地检测、容忍和从组件故障中恢复。
  • 系统存储数量适中的大文件。我们预期有数百万个文件,每个文件通常为100MB或更大。数GB的文件是常见情况,应当被高效管理。小文件必须得到支持,但无需针对它们进行优化。
  • 工作负载主要由两种读取组成:大规模流式读取和小规模随机读取。在大规模流式读取中,单次操作通常读取数百KB,更常见的是1MB或更多。来自同一客户端的连续操作通常读取文件的连续区域。小规模随机读取通常在任意偏移处读取几KB。注重性能的应用通常会将小规模读取进行批量和排序,以在文件中稳步推进,而非来回跳转。
  • 工作负载中也有大量向文件追加数据的大规模顺序写入。典型的操作大小与读取类似。文件一旦写入,就很少再被修改。文件中任意位置的小规模写入是被支持的,但无需保证高效。
  • 系统必须为并发向同一文件追加的多个客户端高效实现定义明确的语义。我们的文件常被用作生产者-消费者队列或多路归并。数百个生产者(每台机器运行一个)会并发地向一个文件追加数据。具有最小同步开销的原子性至关重要。文件可能在之后被读取,或者消费者可能同时读取文件。
  • 高持续带宽比低延迟更重要。我们的大多数目标应用都优先考虑以高速率批量处理数据,而很少有应用对单次读取或写入有严格的响应时间要求。

2.2 接口

GFS提供了熟悉的文件系统接口,尽管它并未实现POSIX之类的标准API。文件以层次化目录结构组织,并通过路径名标识。我们支持创建、删除、打开、关闭、读取和写入文件等常规操作。

此外,GFS还拥有快照(snapshot)和记录追加(record append)操作。快照以低成本创建文件或目录树的副本。记录追加允许多个客户端并发地向同一文件追加数据,同时保证每个客户端追加操作的原子性。它可用于实现多路归并结果和生产者-消费者队列,多个客户端可以同时向其中追加数据而无需额外加锁。我们发现这些类型的文件在构建大型分布式应用时极具价值。快照和记录追加将分别在第3.4节和第3.3节进一步讨论。

2.3 架构

一个GFS集群由一个主服务器(master)和多个块服务器(chunkserver)组成,并被多个客户端访问,如图1所示。其中每个角色通常都是运行用户级服务器进程的商用Linux机器。只要机器资源允许,并且运行可能不稳定的应用代码所导致的可靠性下降是可接受的,在同一台机器上同时运行块服务器和客户端是很容易实现的。

文件被划分为固定大小的块(chunk)。每个块在创建时由主服务器分配一个不变且全局唯一的64位块句柄(chunk handle)来标识。块服务器将块以Linux文件的形式存储在本地磁盘上,并根据块句柄和字节范围来读写块数据。为了可靠性,每个块在多个块服务器上进行复制。默认情况下,我们存储三个副本,不过用户可以为文件命名空间的不同区域指定不同的复制级别。

主服务器维护所有文件系统元数据。这包括命名空间、访问控制信息、文件到块的映射,以及块的当前位置。它还控制系统范围的活动,例如块租约管理、孤立块的垃圾回收,以及块服务器之间的块迁移。主服务器定期通过心跳消息与每个块服务器通信,向其下达指令并收集其状态。

链接到每个应用程序中的GFS客户端代码实现了文件系统API,并代表应用程序与主服务器和块服务器通信以读写数据。客户端与主服务器交互以进行元数据操作,但所有承载数据的通信都直接发往块服务器。我们不提供POSIX API,因此无需接入Linux的vnode层。

客户端和块服务器都不缓存文件数据。客户端缓存几乎没有益处,因为大多数应用要么流式读取超大文件,要么工作集太大而无法缓存。不做缓存简化了客户端和整个系统,因为消除了缓存一致性问题。(不过客户端确实会缓存元数据。)块服务器无需缓存文件数据,因为块以本地文件形式存储,Linux的缓冲区缓存已经将频繁访问的数据保存在内存中。

2.4 单主服务器

采用单主服务器极大地简化了我们的设计,并且使主服务器能够利用全局信息做出复杂的块放置和复制决策。然而,我们必须尽量减少主服务器在读写操作中的参与,以避免它成为瓶颈。客户端永远不通过主服务器读写文件数据。相反,客户端向主服务器询问它应该联系哪些块服务器。它将该信息缓存一段时间,并在后续多次操作中直接与块服务器交互。

让我们结合图1说明一次简单读取的交互过程。首先,利用固定的块大小,客户端将应用程序指定的文件名和字节偏移转换为文件内的块索引。然后,它向主服务器发送包含文件名和块索引的请求。主服务器回复相应的块句柄和副本位置。客户端以文件名和块索引为键缓存该信息。

随后,客户端向其中一个副本(最可能是最近的那个)发送请求。请求指定了块句柄和块内的字节范围。对同一块的后续读取无需再与主服务器交互,直到缓存信息过期或文件被重新打开。事实上,客户端通常在同一次请求中请求多个块,而主服务器也可以包含紧随请求块之后的那些块的信息。这些额外信息几乎无需额外成本,就避开了后续多次客户端-主服务器交互。

图1:GFS架构

2.5 块大小

块大小是关键设计参数之一。我们选择了64MB,这比典型的文件系统块大小要大得多。每个块副本以普通Linux文件的形式存储在块服务器上,并且仅在需要时才进行扩展。惰性空间分配避免了内部碎片造成的空间浪费——这或许是反对使用如此大块大小的最主要理由。

大块大小具有几个重要优势:

  1. 减少了客户端与主服务器交互的需求,因为同一块上的读写只需向主服务器请求一次块位置信息。对于我们的工作负载而言,这种减少尤为显著,因为应用大多顺序读写大文件。即使对于小规模随机读取,客户端也可以轻松缓存数TB工作集的所有块位置信息。
  2. 由于在大块上,客户端更可能对给定块执行多次操作,因此可以通过在较长时间内保持与块服务器的持久TCP连接来减少网络开销。
  3. 减少了主服务器上存储的元数据大小。这使我们能够将元数据保存在内存中,进而带来其他优势,我们将在第2.6.1节讨论。

另一方面,即使采用惰性空间分配,大块大小也有其缺点。小文件由少量块组成,可能只有一个。如果许多客户端访问同一个文件,存储这些块的块服务器可能会成为热点。在实践中,热点并不是主要问题,因为我们的应用大多顺序读取多块的大文件。

然而,当GFS最初被批处理队列系统使用时,热点确实出现了:一个可执行文件作为单块文件写入GFS,然后同时在数百台机器上启动。存储该可执行文件的少数块服务器因数百个并发请求而过载。我们通过提高此类可执行文件的复制因子,并让批处理队列系统错开应用启动时间,解决了这个问题。一个潜在的长期解决方案是允许客户端在这种情况下从其他客户端读取数据。

2.6 元数据

主服务器存储三类主要元数据:文件和块命名空间、文件到块的映射,以及每个块副本的位置。所有元数据都保存在主服务器的内存中。前两类(命名空间和文件到块的映射)还通过将变更记录到操作日志中进行持久化存储,该日志存储在主服务器本地磁盘并复制到远程机器。使用日志使我们能够简单、可靠地更新主服务器状态,并且在主服务器崩溃时不会有不一致的风险。主服务器不持久化存储块位置信息。相反,它在主服务器启动时以及块服务器加入集群时,向每个块服务器询问其块的情况。

2.6.1 内存数据结构

由于元数据存储在内存中,主服务器操作速度很快。此外,主服务器可以轻松高效地在后台定期扫描其全部状态。这种定期扫描用于实现块垃圾回收、块服务器故障时的重新复制,以及为负载均衡和磁盘空间均衡而进行的块迁移。

元数据全部存储在内存中也有潜在的缺点,即整个系统的容量(以及每个块服务器的块数量)受限于主服务器的内存大小。这在实践中并不是严重的限制。对于64MB的块,64MB的元数据可以支持数PB的数据。即使块大小更小,数百万个文件的元数据也只占用主服务器少量内存。因此,内存限制不会成为我们工作负载下的瓶颈,而为了获得内存元数据带来的简洁性、性能和灵活性,付出这个代价是值得的。

2.6.2 块位置

主服务器不持久化记录块副本的位置。它在启动时简单地从块服务器获取这些信息。之后,由于主服务器控制所有块的放置并通过定期心跳消息监控块服务器状态,因此它始终掌握最新信息。

这种设计避免了在块服务器加入、离开、重命名或故障时,主服务器与块服务器之间同步块位置信息的复杂问题。在一个拥有数百台块服务器的集群中,这些事件会持续发生。

换个角度看,块服务器是其自身块列表的最终权威。没有必要让主服务器维护块位置的一致视图,因为块服务器上的块集合可能随时因本地操作而改变。

2.6.3 操作日志

操作日志包含关键元数据变更的历史记录。它是GFS的核心。它不仅是元数据的唯一持久化记录,还作为定义并发操作顺序的逻辑时间线。文件和块,以及它们的版本(见第4.5节),都通过它们被创建时的逻辑时间唯一且永久地标识。

由于操作日志至关重要,我们必须可靠地存储它,并且在元数据变更持久化之前,不能让客户端看到变更。否则,即使块本身得以保留,我们实际上也会丢失整个文件系统或最近的客户端操作。因此,我们将其复制到多台远程机器上,并且只有在将相应的日志记录刷新到本地和远程磁盘之后,才响应客户端操作。主服务器在刷新前会将多条日志记录批量处理,从而减少刷新和复制对整体系统吞吐量的影响。

主服务器通过重放操作日志来恢复其文件系统状态。为了最小化启动时间,我们必须保持日志较小。当日志增长超过一定大小时,主服务器会对其状态设置检查点,这样它就可以通过从本地磁盘加载最新检查点,然后仅重放检查点之后的有限数量的日志记录来进行恢复。检查点采用紧凑的B树形式,可以直接映射到内存中,无需额外解析即可用于命名空间查找。这进一步加快了恢复速度并提高了可用性。

由于构建检查点可能需要一段时间,主服务器的内部状态采用了这样的结构:创建新检查点时不会延迟传入的变更操作。主服务器切换到新的日志文件,并在单独的线程中创建新检查点。新检查点包含切换之前的所有变更。对于拥有数百万个文件的集群,检查点可以在一分钟左右创建完成。完成后,它被写入本地和远程磁盘。

恢复只需要最新的完整检查点和后续的日志文件。旧的检查点和日志文件可以自由删除,不过我们会保留几份以防灾难性故障。检查点过程中的故障不会影响正确性,因为恢复代码会检测并跳过不完整的检查点。

2.7 一致性模型

GFS采用宽松的一致性模型,很好地支持了我们的高度分布式应用,同时实现起来相对简单高效。现在我们讨论GFS提供的保证,以及它们对应用程序的意义。我们还将强调GFS如何维护这些保证,但具体细节留待本文其他部分介绍。

2.7.1 GFS的保证

文件命名空间变更(例如文件创建)是原子的。它们完全由主服务器处理:命名空间锁保证了原子性和正确性(第4.1节);主服务器的操作日志定义了这些操作的全局全序(第2.6.3节)。

数据变更后文件区域的状态取决于变更类型、变更是否成功,以及是否存在并发变更。表1总结了结果。如果所有客户端无论从哪个副本读取,看到的数据始终相同,则该文件区域是一致的。如果文件数据变更后区域是一致的,并且客户端将看到变更写入的全部内容,则该区域是已定义的

当变更成功且没有并发写入者干扰时,受影响的区域是已定义的(并且隐含是一致的):所有客户端将始终看到变更写入的内容。并发的成功变更会使区域处于未定义但一致的状态:所有客户端看到相同的数据,但该数据可能不反映任何一个变更写入的内容。通常,它由来自多个变更的混合片段组成。失败的变更会使区域不一致(因此也是未定义的):不同客户端在不同时间可能看到不同的数据。

我们将在下文介绍应用程序如何区分已定义区域和未定义区域。应用程序无需进一步区分不同类型的未定义区域。

写入 记录追加
串行成功 已定义 已定义,中间夹杂不一致
并发成功 一致但未定义 已定义,中间夹杂不一致
失败 不一致

表1:变更后的文件区域状态

2.7.2 对应用程序的影响

GFS应用程序可以通过一些简单的技术来适应这种宽松的一致性模型,这些技术本来也适用于其他分布式系统。这些技术包括:依赖追加而非覆盖、写入校验和、写入唯一标识符,以及写入时写入预期的序列号。应用程序可以使用校验和来验证记录的完整性,并识别并丢弃填充或重复的记录。大多数应用程序只需要在恢复时重放少量记录。

对于大多数应用而言,追加写入语义远比覆盖写入更有用。追加操作天然高效,并且不存在并发写入者相互覆盖的问题。例如,在生产者-消费者队列中,生产者可以并发地向文件追加记录,而消费者可以通过检查点机制记录已处理的位置。

对于读取操作,只要应用程序能够容忍偶尔的不一致区域,就不需要额外的同步。例如,在数据处理作业中,读取数据的程序通常只扫描整个文件,并且能够处理偶尔出现的损坏或重复记录。

如果应用需要更强的一致性,可以使用以下机制:

  • 写入后校验:写入者可以在数据中包含校验和,读取者验证校验和。
  • 序列号:每条记录包含唯一的序列号,读取者可以据此检测重复和缺失。
  • 命名空间原子操作:文件创建等命名空间操作是原子的,可以用于同步。

3. 系统交互

我们现在详细描述客户端、块服务器和主服务器之间的交互,以实现数据变更、租约管理、数据流和原子追加。

3.1 租约与变更顺序

GFS使用租约(lease)机制来保持多个副本之间的变更顺序。主服务器向其中一个副本授予块租约,该副本称为主副本(primary)。主副本为块的所有变更选择一个串行顺序。所有副本都遵循这个顺序执行变更,从而保证全局一致性。

租约机制的设计目标是最小化主服务器的管理开销。租约初始超时时间为60秒。然而,只要块正在被写入,主副本就可以请求并通常获得主服务器的续期。这些续期请求和响应通过块服务器与主服务器之间的常规心跳消息捎带传输。

以下是写入操作的完整步骤(参见图2):

  1. 客户端向主服务器询问哪个块服务器持有该块的租约,以及其他副本的位置。如果没有任何副本持有租约,主服务器选择一个副本授予租约。
  2. 主服务器回复主副本和次要副本的位置。客户端缓存这些数据以便后续写入。只有当主副本不可达或者主副本回复称其不再持有租约时,客户端才需要重新联系主服务器。
  3. 客户端将数据推送到所有副本。客户端可以以任意顺序推送数据。每个块服务器在内部LRU缓存中保存数据,直到数据被使用或过期。通过将数据流与控制流分离,我们可以独立于数据推送顺序来调度昂贵的网络流量。我们将在第3.2节进一步讨论数据流。
  4. 一旦所有副本都确认收到数据,客户端就向主副本发送写入请求。该请求标识了之前推送到所有副本的数据。主副本为收到的所有变更分配连续的序列号,这提供了必要的串行化。它按序列号顺序将变更应用到自己的本地状态。
  5. 主副本将写入请求转发给所有次要副本。每个次要副本按照主副本分配的相同序列号顺序应用变更。
  6. 所有次要副本都回复主副本,表示它们已完成操作。
  7. 主副本回复客户端。任何副本上遇到的任何错误都将返回给客户端。如果发生错误,写入可能在主副本和部分次要副本上成功(如果在次要副本上失败,则步骤6可能不会全部完成)。客户端可以通过重试来处理错误。它在重试之前会重新执行步骤3到7。

对于失败的写入,区域可能处于不一致状态。客户端代码通过重试写入来处理这种情况。在向应用层报告失败之前,会在几个不同的块副本上重试几次。

图2:写入控制与数据流

3.2 数据流

我们将数据流与控制流分离,以高效利用网络。控制流从客户端流向主副本,再流向次要副本。而数据则以流水线方式沿着精心选择的服务器链线性推送,以充分利用每个机器的网络带宽。

目标是在避免网络瓶颈和高延迟链路的同时,最大化网络吞吐量。我们不通过树形结构广播数据,而是采用线性的、沿服务器链的流水线传输。在我们的网络环境中,每台机器的入站和出站带宽大致相当,而跨机架链路通常比机架内链路更慢、更拥塞。

具体来说,数据沿着块服务器链逐跳推送。每台机器收到数据后,立即将其转发给链中的下一台机器。这样可以充分利用每个机器的出站带宽,因为接收和转发可以并行进行。

例如,假设有三台块服务器A、B、C在不同机架上。客户端将数据发送给A,A转发给B,B转发给C。如果每跳需要时间t,那么整个传输大约需要3t时间。相比之下,如果客户端分别发送给三台机器,则需要3t的出站带宽,并且可能在客户端网络接口处形成瓶颈。

选择哪条链路由网络拓扑决定。理想情况下,数据在从客户端到最终副本的路径上,每台机器只经过一次。我们通过让每台机器将数据转发给尚未收到数据的最近副本,来近似这个目标。

流水线的另一个好处是,接收方可以在收到完整数据块之前就开始转发。由于我们使用TCP,一旦数据到达就可以立即转发,无需等待整个块。这大大减少了传输延迟。

3.3 原子记录追加

GFS提供了一种称为记录追加的原子追加操作。在传统的写入中,客户端指定写入偏移量。对同一区域的并发写入不具有串行性——该区域最终可能包含多个写入者的数据片段。而在记录追加中,客户端仅指定数据,GFS自动将数据追加到文件的至少一个偏移位置,并将该偏移位置返回给客户端。这保证了原子性,即使有多个客户端并发追加也是如此。

记录追加对于实现生产者-消费者队列至关重要。多个生产者可以并发地向同一个文件追加记录,而无需任何显式同步。每个消费者读取文件时,会看到完整的记录。

记录追加的实现与常规写入类似,但主副本需要处理一些额外的逻辑。客户端将数据推送到文件最后一个块的所有副本,然后向主副本发送请求。主副本检查追加该记录是否会导致块超过64MB的限制。如果不会,主副本将数据追加到自己的副本中,然后将相同的偏移量和长度发送给所有次要副本,并要求它们在完全相同的偏移处写入数据。

如果追加会使块超过最大大小,主副本会将当前块填充到最大大小,然后通知次要副本也这样做,并回复客户端,指示操作应该在下一个块重试。(记录追加的数据大小严格限制为最大块大小的四分之一,以保证最坏情况下的碎片率仍然可接受。)

当记录追加在某些副本上失败时,客户端会重试。因此,不同副本可能包含不同的数据——同一记录可能出现重复,或者在不同副本中出现在不同位置。GFS不保证所有副本在字节级别完全相同。它只保证数据作为一个整体被至少写入一次——这是原子性的基本含义。

应用程序必须能够处理重复和可能的填充记录。可以通过在每条记录中写入唯一标识符,然后扫描并跳过重复项来处理。应用程序还可以使用校验和来检测填充数据。

3.4 快照

快照操作几乎可以瞬间创建文件或目录树的副本,而不会中断正在进行的变更。我们使用标准的**写时复制(copy-on-write)**技术来实现快照。

当主服务器收到快照请求时:

  1. 它首先撤销对即将被快照的文件中所有块的租约。这确保任何后续写入都必须与主服务器交互以查找新的租约持有者。这给了主服务器先创建块副本的机会。
  2. 租约撤销或过期后,主服务器将操作记录到磁盘。然后,它通过复制源文件或目录树的元数据,将该日志记录应用到其内存状态。新创建的快照文件指向与源文件相同的块。

当客户端随后想要写入这些块之一时,它首先向主服务器请求当前租约持有者。主服务器注意到该块的引用计数大于1。它不会直接回复租约信息,而是要求每个块服务器创建该块的新副本。每个块服务器在本地创建新块后,主服务器可以将租约授予其中一个新副本,并回复客户端。客户端然后可以正常写入该块,而不会影响现有的快照。


4. 主服务器操作

主服务器执行所有命名空间操作。此外,它管理整个系统中的块副本:创建块、重新复制、重新均衡,以及垃圾回收。

4.1 命名空间管理与加锁

主服务器在执行任何命名空间操作之前,都会获取适当的锁以保证串行化。与传统文件系统不同,GFS没有每个目录的inode或目录条目数据结构。相反,它使用表示路径名的字符串的前缀树(trie)来表示命名空间。树中的每个节点都有一个关联的读写锁。

每个命名空间操作在执行前都会获取一组锁。通常,涉及路径/d1/d2/.../dn/leaf的操作会获取路径/d1/d1/d2、…、/d1/d2/.../dn上的读锁,以及完整路径/d1/d2/.../dn/leaf上的读锁或写锁。

例如,快照操作获取/home/save上的读锁,以及/home/user/save/user上的写锁。文件创建操作获取/home/home/user上的读锁,以及/home/user/foo上的写锁。这两个操作会被正确串行化,因为它们尝试获取/home/user上的冲突锁。

文件创建不需要在父目录上获取写锁,因为不存在”目录”或类似inode的数据结构需要防止修改。名称上的读锁足以防止父目录被删除、重命名或快照。

这种加锁方案的一个优良特性是允许同一目录中的并发变更。例如,多个文件创建可以在同一目录中并发执行:每个都获取目录名上的读锁和文件名上的写锁。目录名上的读锁足以防止目录被删除、重命名或快照。文件名上的写锁则串行化了两次创建同名文件的尝试。

由于命名空间可能有很多节点,读写锁对象采用惰性分配方式,并且在不再使用时删除。此外,锁按照一致的全序获取以防止死锁:首先按命名空间树的层级排序,同一层级内按字典序排序。

4.2 副本放置

GFS集群高度分布在多个机架上。块副本放置策略有两个目标:最大化数据可靠性和可用性,以及最大化网络带宽利用率。仅在机器之间分散副本是不够的——这只能防止机器故障,但不能防止整个机架故障(例如网络交换机、电源故障)。

因此,我们将块副本分布在不同的机架上。这确保即使整个机架损坏或离线,每个块也有副本幸存。这也意味着读取(尤其是大规模读取)可以利用多个机架的聚合带宽。另一方面,写入必须遍历多个机架,这是我们愿意接受的权衡。

副本放置还考虑磁盘空间利用率和负载均衡。主服务器在放置新块时,会选择磁盘空间利用率低于平均水平的块服务器。此外,它会限制每个块服务器上最近的创建数量,以防止写入风暴。当然,最终副本必须分布在不同机架上。

4.3 创建、重新复制、重新均衡

主服务器在三种情况下创建块副本:块创建、重新复制和重新均衡。

当创建一个新块时,主服务器选择在哪里放置初始副本。它考虑几个因素:

  1. 我们希望在磁盘空间利用率低于平均水平的块服务器上放置新副本。
  2. 我们希望限制每个块服务器上”最近”创建的数量,以防止写入流量激增。
  3. 如上所述,我们希望将副本分布在不同机架上。

当块的可用副本数量低于目标时,主服务器会重新复制该块。这可能由多种原因触发:块服务器不可用、块服务器报告其副本可能已损坏、磁盘出现故障,或者复制目标被提高。

每个需要重新复制的块都有优先级。优先级主要取决于缺失了多少副本以及缺失了多久。例如,丢失了两个副本的块比丢失了一个副本的块优先级更高。此外,属于活跃文件的块比属于最近删除文件的块优先级更高。最后,为了最小化对运行中工作负载的影响,我们提高了阻塞客户端的块的优先级。

主服务器按优先级顺序重新复制块。它选择一个块,然后指示某个块服务器直接从现有的有效副本”克隆”该块。目标副本的放置标准与新块类似:均衡磁盘空间、限制单个服务器上的并发克隆数量、跨机架分布。

为了防止克隆操作占用过多带宽,主服务器限制每个块服务器上的并发克隆数量。此外,每个块服务器通过限制向源块服务器的读取请求速率,来限制每个克隆操作消耗的带宽。

最后,主服务器定期重新均衡副本:它检查当前的副本分布,并移动副本以实现更好的磁盘空间和负载均衡。通过这个过程,主服务器逐渐填满新的块服务器,而不是立即用新块和随之而来的大量写入流量淹没它。新副本的放置标准与上述类似。此外,主服务器还必须选择移除哪个现有副本。通常,它优先移除磁盘空闲空间低于平均水平的块服务器上的副本,以均衡磁盘空间使用率。

4.4 垃圾回收

文件删除后,GFS不会立即回收可用的物理存储空间。它仅在常规垃圾回收期间惰性地回收,包括文件级和块级的回收。我们发现这种方法使系统更简单、更可靠。

4.4.1 机制

当应用程序删除文件时,主服务器会像其他变更一样立即记录删除操作。然而,它不会立即回收资源,而是将文件重命名为一个包含删除时间戳的隐藏名称。在主服务器定期扫描文件系统命名空间时,如果这些隐藏文件已存在超过三天(该间隔可配置),就会将其移除。在此之前,文件仍然可以通过新的特殊名称读取,并且可以通过重命名回正常名称来撤销删除。当隐藏文件从命名空间中移除时,其内存中的元数据被擦除。这实际上切断了它与所有块的链接。

在类似的块命名空间定期扫描中,主服务器识别孤立块(即无法从任何文件到达的块),并擦除这些块的元数据。在与主服务器定期交换的心跳消息中,每个块服务器报告其拥有的块的子集,主服务器回复所有已不存在于主服务器元数据中的块的标识。块服务器可以自由删除这些块的副本。

4.4.2 讨论

尽管在编程语言上下文中,分布式垃圾回收是一个需要复杂解决方案的难题,但在我们的场景中却相当简单。我们可以轻松识别所有对块的引用:它们都在主服务器唯一维护的文件到块映射中。我们也可以轻松识别所有块副本:它们是每个块服务器上指定目录下的Linux文件。任何主服务器不知道的副本都是”垃圾”。

与立即删除相比,垃圾回收方式回收存储有几个优势:

  1. 简单可靠:在组件故障频发的大规模分布式系统中,这一点至关重要。块创建可能在部分块服务器上成功而在其他服务器上失败,留下主服务器不知道存在的副本。副本删除消息可能丢失,主服务器必须记住在故障(自身故障和块服务器故障)后重发。垃圾回收提供了一种统一且可靠的方式来清理所有已知无用的副本。
  2. 批量处理,成本摊销:它将存储回收合并到主服务器的常规后台活动中,例如命名空间的定期扫描和与块服务器的心跳握手。因此,它以批量方式执行,成本被摊销。此外,它只在主服务器相对空闲时执行。主服务器可以更及时地响应需要及时处理的客户端请求。
  3. 安全网:存储回收的延迟为意外的、不可逆的删除提供了安全保障。

根据我们的经验,主要缺点是当存储紧张时,这种延迟有时会妨碍用户调整使用量。反复创建和删除临时文件的应用可能无法立即重用存储空间。我们通过以下方式解决这些问题:如果已删除文件被再次显式删除,则加速存储回收。我们还允许用户对命名空间的不同部分应用不同的复制和回收策略。例如,用户可以指定某个目录树中文件的所有块都不进行复制存储,并且任何删除的文件都会立即且不可撤销地从文件系统状态中移除。

4.5 过期副本检测

如果块服务器发生故障并且在停机期间错过了块的变更,块副本可能会过期。对于每个块,主服务器维护一个块版本号,以区分最新副本和过期副本。

每当主服务器授予块的新租约时,它都会增加块版本号,并通知所有最新副本。主服务器和这些副本都在它们的持久化状态中记录新的版本号。这发生在任何客户端收到通知之前,因此也发生在客户端可以开始写入块之前。如果另一个副本当前不可用,其块版本号不会被提升。当块服务器重启并报告其块集合及相应的版本号时,主服务器会检测到该块服务器拥有过期副本。如果主服务器看到的版本号大于其记录中的版本号,主服务器会认为自己在授予租约时发生了故障,因此采用更高的版本号作为最新版本。

主服务器在常规垃圾回收中移除过期副本。在此之前,当回复客户端的块信息请求时,它实际上认为过期副本根本不存在。作为另一项保障,当主服务器通知客户端哪个块服务器持有块的租约时,或者当它指示块服务器在克隆操作中从另一个块服务器读取块时,都会包含块版本号。客户端或块服务器在执行操作时验证版本号,以确保始终访问最新数据。


5. 容错与诊断

设计系统的最大挑战之一是处理频繁的组件故障。组件的质量和数量共同使得这些问题成为常态而非异常:我们不能完全信任机器,也不能完全信任磁盘。组件故障可能导致系统不可用,或者更糟的是,数据损坏。我们将讨论如何应对这些挑战,以及系统中内置的诊断工具。

5.1 高可用性

在GFS集群的数百台服务器中,总有一些在任何给定时间不可用。我们通过两种简单而有效的策略保持整个系统的高可用性:快速恢复和复制。

5.1.1 快速恢复

主服务器和块服务器都被设计为无论如何终止,都能在数秒内恢复状态并启动。事实上,我们不区分正常终止和异常终止;服务器通常通过直接杀死进程来关闭。客户端和其他服务器会经历轻微的中断,因为它们的未完成请求超时,然后重新连接到重启后的服务器并重试请求。

5.1.2 复制

如前所述,每个块都被复制到多个机架上的多个块服务器。用户可以为文件命名空间的不同区域指定不同的复制因子。默认值为三。当块服务器离线或检测到损坏的副本时,主服务器会使用现有副本重新创建块。

对于主服务器,其状态也被复制以实现高可用性。操作日志和检查点被复制到多台机器上。只有在日志记录被刷新到本地和所有远程副本之后,状态变更才被视为已提交。主服务器进程在任何时候只在一台机器上运行。当主服务器发生故障时,可以使用其在其他机器上的持久化状态在其他地方启动新的主服务器。

“影子”主服务器在主主服务器出现故障时提供只读访问。它们是影子,而非镜像,因为它们的状态可能略微滞后于主服务器,通常不到一秒。它们增强了文件系统在主服务器更新期间或主服务器故障切换期间的读取可用性。

5.2 数据完整性

每个块服务器都使用校验和来检测存储数据的损坏。考虑到每个块服务器上有许多磁盘,磁盘和IDE子系统级别的数据损坏并不罕见。

每个64MB的块被划分为64KB的块。每个块对应一个32位的校验和。与其他元数据一样,校验和存储在内存中,并持久化记录到磁盘,与用户数据分开。

对于读取,块服务器在返回数据之前验证相应块的校验和。这防止了损坏的数据被传播到其他客户端或块服务器。如果某个块损坏,块服务器向请求者返回错误,并通知主服务器。然后客户端可以从其他副本读取数据,主服务器可以从其他副本重新创建该块。

对于追加写入,校验和计算经过优化以匹配追加为主的工作负载。我们增量地更新部分填充的最后一个块的校验和,并为新追加的完整块计算新的校验和。

对于覆盖写入(即写入现有块的中间位置),我们必须读取被覆盖区域的第一个和最后一个块,验证它们的校验和,然后执行写入,最后重新计算校验和。这就是为什么随机写入效率较低的原因之一。

在空闲时间,块服务器可以扫描和验证不活跃块的校验和。这使我们能够检测未被读取的块的损坏。一旦检测到损坏,主服务器可以创建新的完好副本并删除损坏的副本。

5.3 诊断工具

广泛的诊断日志记录极大地帮助了问题隔离、调试和性能分析,而开销却微乎其微。GFS服务器生成各种事件日志和大量的RPC日志。

事件日志记录服务器生命周期中的重要事件,例如服务器启动和关闭、故障检测等。RPC日志记录进出的每个RPC请求和响应的详细信息,包括时间戳、消息大小等。通过关联不同服务器上的RPC日志,我们可以重建完整的交互历史来诊断问题。

日志还用于性能分析和工作负载研究。由于日志是追加写入的,对正在运行的系统影响很小。


6. 性能测量

我们现在展示GFS在微观基准和真实生产集群中的性能表现。

6.1 微观基准

我们在一个包含1台主服务器、2台主服务器副本、16台块服务器和16台客户端的集群上测量了性能。所有机器都配备双1.4GHz PIII处理器、2GB内存、一个80GB 5400rpm IDE磁盘和100Mbps全双工以太网。它们连接到一个HP ProCurve 2524交换机。整个GFS集群位于同一个机架上。

6.1.1 读取

图3(a)显示了聚合读取吞吐量。我们启动了N个客户端,每个客户端同时从整个文件系统中读取数据。每个客户端读取4GB的数据,总共读取N×4GB。读取被分为64KB的大小,随机偏移。

随着客户端数量从1增加到16,聚合读取吞吐量从10MB/s增长到94MB/s。这接近网络的理论最大值(16台客户端×100Mbps = 200MB/s全双工,但交换机背板只有约120MB/s的容量)。

单客户端的读取性能约为10MB/s,这也接近单个100Mbps链路的限制。

6.1.2 写入

图3(b)显示了聚合写入吞吐量。N个客户端各自向N个不同的文件写入1GB数据。写入同样以64KB为单位。

单客户端写入速度约为6MB/s。随着客户端数量增加到16,聚合写入达到约35MB/s。这低于读取吞吐量,部分原因是写入需要写入三个副本,并且每个副本都消耗网络带宽。

6.1.3 记录追加

图3(c)显示了记录追加的性能。N个客户端并发地向同一个文件追加记录。每条记录大小为64KB。

单客户端速度约为5MB/s。随着客户端数量增加,聚合吞吐量增长到约20MB/s。由于所有客户端都写入同一个文件,主副本成为瓶颈。这是预期的行为——记录追加并非设计用于数百个客户端同时向同一个文件写入的场景。对于大多数应用,多个生产者写入不同的文件更为常见。

图3:聚合吞吐量。上方曲线表示由网络拓扑限制的理论上限;下方曲线表示实测值。(a) 读取 (b) 写入 (c) 记录追加

6.2 真实世界集群

表2展示了谷歌两个生产集群的特征。集群X用于研究和开发,集群Y用于生产数据处理。

集群X 集群Y
块服务器数量 342 227
可用磁盘空间 72 TB 180 TB
文件数量 335万 735万
块数量 800万 1500万
元数据大小(主服务器) 53 MB 115 MB
每台块服务器平均读取速率 9 MB/s 15 MB/s
每台块服务器平均写入速率 3 MB/s 5 MB/s
主服务器操作速率 ~700 ops/s ~500 ops/s

表2:生产集群特征

这些数字表明,主服务器的元数据内存占用很小——只有几十MB。这验证了我们的说法,即主服务器内存不是瓶颈。

集群Y的磁盘空间更大,但文件和块数量更少,因为它的文件平均更大。两个集群的读写速率都显示出健康的流量水平。

6.3 恢复时间

我们测量了块服务器故障后的恢复时间。我们杀死了一个持有约15000个块(总计约1TB数据)的块服务器。为了限制对运行中工作负载的影响,重新复制被限速,并且集群还有其他写入活动。

所有块在23.2分钟内恢复到完整的复制因子。平均恢复速率约为440Mbps,这是合理的,因为它允许多个源块服务器并行发送数据。

在另一个测试中,我们杀死了两个块服务器(每个约16000个块)。这导致266个块减少到只有一个副本。这些高优先级块在2分钟内全部被重新复制。

6.4 工作负载分解

表3显示了两个集群上操作类型的细分。操作按涉及的字节数和操作计数来衡量。

集群X 集群Y
字节 操作 字节 操作
读取 83% 61% 95% 72%
写入 17% 16% 5% 11%
记录追加 0% 23% 0% 17%

表3:工作负载分解(字节和操作计数)

读取在两个集群中都占主导地位,这符合预期。记录追加操作在操作数量上占比很大,但字节数很少,因为记录通常很小。

值得注意的是,元数据操作(打开、创建、删除等)在操作计数中占相当大的比例。这是因为许多应用程序创建和删除大量临时文件。集群Y中元数据操作的比例较低,因为其自动化数据处理任务倾向于检查文件系统的部分内容以了解全局应用状态。相比之下,集群X的应用受到更明确的用户控制,并且通常预先知道所有需要的文件名。


7. 经验与教训

在构建和部署GFS的过程中,我们经历了各种问题,有些是运维方面的,有些是技术方面的。

最初,GFS被构想为我们生产系统的后端文件系统。随着时间推移,其用途扩展到包括研发任务。它最初对权限、配额等支持很少,但现在已经包含了这些功能的基本形式。生产系统纪律严明、受控良好,但用户有时并非如此。需要更多基础设施来防止用户相互干扰。

我们遇到的一些最大问题与磁盘和Linux相关。许多磁盘向Linux驱动程序声称它们支持一系列IDE协议版本,但实际上只对较新的版本能可靠响应。由于协议版本非常相似,这些磁盘大多能正常工作,但偶尔的不匹配会导致驱动器和内核对驱动器状态产生分歧。这会由于内核中的问题而静默地损坏数据。这个问题促使我们使用校验和来检测数据损坏,同时我们修改了内核以处理这些协议不匹配。

早些时候,我们在Linux 2.2内核上遇到了一些与fsync()成本相关的问题。它的成本与文件大小成正比,而不是与修改部分的大小成正比。这对于我们的大型操作日志来说是个问题,尤其是在我们实现检查点之前。我们曾一度使用同步写入来解决这个问题,并最终迁移到Linux 2.4。

另一个Linux问题是单个读写锁——地址空间中的任何线程在从磁盘分页(读锁)或在mmap()调用中修改地址空间(写锁)时都必须持有该锁。我们在轻负载下看到系统出现瞬时超时,并努力寻找资源瓶颈或偶发硬件故障。最终我们发现,当磁盘线程正在分页调入先前映射的数据时,这个锁会阻塞主网络线程将新数据映射到内存。由于我们主要受网络接口而非内存拷贝带宽的限制,我们通过用pread()替代mmap()来解决这个问题,代价是多了一次拷贝。

尽管偶尔会遇到问题,Linux代码的可用性一次又一次地帮助我们探索和理解系统行为。在适当的时候,我们改进内核,并与开源社区分享这些改动。


8. 相关工作

与AFS等其他大型分布式文件系统一样,GFS提供了位置无关的命名空间,使数据可以透明地移动以实现负载均衡或容错。与AFS不同的是,GFS将文件数据分散在存储服务器上,这种方式更类似于xFS和Swift,以提供聚合性能和更高的容错能力。

由于磁盘相对便宜,并且复制比更复杂的RAID方案更简单,GFS目前仅使用复制来实现冗余,因此比xFS或Swift消耗更多的原始存储空间。

与AFS、xFS、Frangipani和Intermezzo等系统不同,GFS不在文件系统接口之下提供任何缓存。我们的目标工作负载在单次应用运行中几乎没有数据重用,因为它们要么流式处理大型数据集,要么在其中随机查找并每次读取少量数据。

一些分布式文件系统(如Frangipani、xFS、明尼苏达大学的GFS和GPFS)移除了集中式服务器,依赖分布式算法来实现一致性和管理。我们选择集中式方法是为了简化设计、提高可靠性并获得灵活性。特别是,集中式主服务器使得实现复杂的块放置和复制策略变得容易得多,因为主服务器已经拥有大部分相关信息并控制其变化。我们通过保持主服务器状态较小并在其他机器上完整复制来解决容错问题。目前,我们的影子主服务器机制提供了可扩展性和高读取可用性。主服务器状态的更新通过追加到预写日志来持久化。因此,我们可以采用类似Harp中的主副本方案,以提供比当前方案更强一致性保证的高可用性。

我们正在解决与Lustre类似的问题,即向大量客户端提供聚合性能。然而,我们通过关注应用需求而非构建POSIX兼容的文件系统,大大简化了问题。此外,GFS假设有大量不可靠组件,因此容错是我们设计的核心。

GFS与NASD架构最为相似。NASD架构基于网络附加磁盘驱动器,而GFS使用商用机器作为块服务器,这与NASD原型中的做法相同。与NASD工作不同,我们的块服务器使用惰性分配的固定大小块,而不是可变长度对象。此外,GFS实现了生产环境所需的重新均衡、复制和恢复等功能。

与明尼苏达大学的GFS和NASD不同,我们不寻求改变存储设备的模型。我们专注于利用现有商用组件,解决复杂分布式系统的日常数据处理需求。

原子记录追加支持的生产者-消费者队列解决了与River中分布式队列类似的问题。River使用分布在机器上的基于内存的队列和精细的数据流控制,而GFS使用可以被多个生产者并发追加的持久化文件。River模型支持m到n的分布式队列,但缺乏持久化存储带来的容错能力,而GFS仅高效支持m到1的队列。多个消费者可以读取同一个文件,但它们必须协调以划分传入负载。


9. 结论

谷歌文件系统展示了在商用硬件上支持大规模数据处理工作负载所必需的品质。虽然一些设计决策是针对我们独特环境的,但许多决策可以应用于类似规模和成本意识的数据处理任务。

我们从根据当前和预期的应用负载与技术环境重新审视传统文件系统假设开始。我们的观察导向了设计空间中截然不同的切入点。我们将组件故障视为常态而非异常,针对主要是追加写入(可能是并发的)然后读取(通常是顺序的)的大文件进行优化,并且扩展和放宽了标准文件系统接口以改进整个系统。

我们的系统通过持续监控、复制关键数据以及快速自动恢复来提供容错。块复制使我们能够容忍块服务器故障。这些故障的频繁发生催生了一种新颖的在线修复机制,该机制定期且透明地修复损坏,并尽快补偿丢失的副本。此外,我们使用校验和来检测磁盘或IDE子系统级别的数据损坏——考虑到系统中磁盘的数量,这种情况变得非常普遍。

我们的设计为执行各种任务的许多并发读取者和写入者提供了高聚合吞吐量。我们通过将经过主服务器的文件系统控制与直接在块服务器和客户端之间传输的数据传输分离来实现这一点。大块大小和块租约使得常见操作中主服务器的参与度最小化,并将权限委托给变更数据的主副本。这使得简单的集中式主服务器成为可能,并且不会成为瓶颈。我们相信,网络栈的改进将解除当前的限制,使我们能够进一步提高聚合吞吐量。

我们已经展示了这样一个系统如何支持大规模生产环境,同时为研发提供灵活的平台。从这项工作中得出的关键教训是,对于大规模分布式系统,组件故障是常态,必须作为设计的一等公民来处理。通过接受故障、针对预期的访问模式进行优化,并放宽传统的一致性和接口约束,我们可以构建一个强大、高性能且具有成本效益的系统。


参考文献

[1] T. Anderson, M. Dahlin, J. Neefe, D. Patterson, D. Roselli, and R. Wang. Serverless network file systems. In Proceedings of the 15th ACM Symposium on Operating Systems Principles, pages 109-126, December 1995.

[2] E. D. Berger, K. S. McKinley, R. D. Blumofe, and P. R. Wilson. Hoard: A scalable memory allocator for multithreaded applications. In Proceedings of the Ninth International Conference on Architectural Support for Programming Languages and Operating Systems, pages 117-128, November 2000.

[3] D. A. D. R. J. M. D. S. H. D. C. B. J. Z. T. S. E. G. W. V. S. M. F. K. A. Swift: Using redundant disks to improve file system performance. In Proceedings of the 14th ACM Symposium on Operating Systems Principles, pages 94-105, December 1993.

[4] G. A. Gibson, D. F. Nagle, K. Amiri, J. Butler, F. W. Chang, H. Gobioff, C. Hardin, E. Riedel, D. Rochberg, and J. Zelenka. A cost-effective, high-bandwidth storage architecture. In Proceedings of the 8th International Conference on Architectural Support for Programming Languages and Operating Systems, pages 92-103, October 1998.

[5] J. H. Howard, M. L. Kazar, S. G. Menees, D. A. Nichols, M. Satyanarayanan, R. N. Sidebotham, and M. J. West. Scale and performance in a distributed file system. ACM Transactions on Computer Systems, 6(1):51-81, February 1988.

[6] J. J. Kistler and M. Satyanarayanan. Disconnected operation in the Coda file system. ACM Transactions on Computer Systems, 10(1):3-25, February 1992.

[7] B. Liskov, S. Ghemawat, R. Gruber, P. Johnson, L. Shrira, and M. Williams. Replication in the Harp file system. In Proceedings of the 13th ACM Symposium on Operating Systems Principles, pages 226-238, October 1991.

[8] P. J. Braam. The Lustre storage architecture. Cluster File Systems, Inc., November 2002.

[9] 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, June 1988.

[10] F. Schmuck and R. Haskin. GPFS: A shared-disk file system for large computing clusters. In Proceedings of the First USENIX Conference on File and Storage Technologies, pages 231-244, January 2002.

[11] S. A. Brandt, E. L. Miller, D. D. E. Long, and L. Xue. Efficient metadata management in large distributed file systems. In Proceedings of the 20th IEEE/11th NASA Goddard Conference on Mass Storage Systems and Technologies, pages 290-298, April 2003.

[12] C. A. Thekkath, T. Mann, and E. K. Lee. Frangipani: A scalable distributed file system. In Proceedings of the 16th ACM Symposium on Operating Systems Principles, pages 224-237, October 1997.

【转载】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;
}