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.com和anchor: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会话过期,主服务器会自行终止。不过如前所述,主服务器故障不会改变子表到子表服务器的分配关系。
当集群管理系统启动一个主服务器时,它需要先发现当前的子表分配情况,然后才能进行变更。主服务器启动时执行以下步骤:
- 主服务器在Chubby中获取一个唯一的主锁,防止同时存在多个主服务器实例。
- 主服务器扫描Chubby中的服务器目录,找出所有活跃的服务器。
- 主服务器与每个活跃子表服务器通信,获取每个服务器当前已分配的子表。
- 主服务器扫描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.