Spanner:Google的全球分布式数据库
Spanner: Google’s Globally-Distributed Database
贡献:分布式数据库领域的划时代之作,被称为 “下一代分布式数据库的标杆”,重新定义了分布式关系数据库的技术路线。
James C. Corbett、Jeffrey Dean、Michael Epstein、Andrew Fikes、Christopher Frost、JJ Furman、Sanjay Ghemawat、Andrey Gubarev、Christopher Heiser、Peter Hochschild、Wilson Hsieh、Sebastian Kanthak、Eugene Kogan、李鸿毅、Alexander Lloyd、Sergey Melnik、David Mwaura、David Nagle、Sean Quinlan、Rajesh Rao、Lindsay Rolig、Yasushi Saito、Michal Szymaniak、Christopher Taylor、Ruth Wang、Dale Woodford
Google,Inc.
摘要
Spanner是谷歌推出的一款可扩展、多版本、全球分布且同步复制的数据库。它是首个实现全球范围数据分布、同时支持外部一致性分布式事务的系统。本文阐述了Spanner的架构设计、功能特性、各项设计决策背后的考量,以及一种能够暴露时钟不确定性的新型时间API。该API及其实现对于支持外部一致性以及诸多强大特性至关重要:包括全局范围内的非阻塞历史读取、无锁只读事务,以及原子模式变更。
1 引言
Spanner是由谷歌设计、研发并部署的一款可扩展全球分布式数据库。从最高抽象层面来看,它将数据分片到分布在全球各地数据中心的多组Paxos状态机上[21]。复制机制用于保障全球可用性与数据地理局部性,客户端可在副本之间自动故障转移。当数据量或服务器数量变化时,Spanner会自动在机器间重新分片;为了负载均衡与故障响应,它还会自动在机器间(甚至跨数据中心)迁移数据。Spanner的设计目标是可扩展至数百个数据中心、数百万台服务器,承载万亿级数据库行。
通过在同一大洲内甚至跨大洲复制数据,应用可以借助Spanner实现高可用性,即便遭遇广域自然灾害也不受影响。我们的首个客户是F1[35]——谷歌广告后端的重构版本。F1在美国全境部署了5个副本。大多数其他应用通常会在同一个地理区域的3到5个数据中心间复制数据,但这些数据中心的故障模式相对独立。也就是说,大多数应用会优先选择更低延迟而非更高可用性,只要能承受1到2个数据中心故障即可。
Spanner的核心重点是管理跨数据中心的复制数据,但我们也投入了大量时间在分布式系统基础设施之上设计并实现重要的数据库特性。尽管很多项目都在顺利使用Bigtable[9],但我们也不断收到用户反馈,称Bigtable在某些应用场景下使用困难:比如针对拥有复杂、持续演进的模式的应用,或是需要在广域复制下保证强一致性的场景。其他研究者也提出过类似观点[37]。谷歌内部很多应用选择使用Megastore[5],因为它支持半关系数据模型与同步复制,尽管它的写入吞吐量相对较低。因此,Spanner从类Bigtable的版本化键值存储演进为了一个时态多版本数据库。数据存储在具备模式定义的半关系表中;数据支持多版本,每个版本都会自动打上提交时间戳;旧版本数据可通过配置的垃圾回收策略进行清理;应用可以读取旧时间戳对应的数据。Spanner支持通用事务,并提供基于SQL的查询语言。
作为一款全球分布式数据库,Spanner具备多项引人关注的特性。首先,应用可以细粒度地动态控制数据的复制配置。应用可以指定约束条件,控制哪些数据中心存储哪些数据、数据与用户的距离(以控制读取延迟)、副本之间的距离(以控制写入延迟),以及维护的副本数量(以控制持久性、可用性与读取性能)。系统也可以动态、透明地在数据中心间迁移数据,以平衡各数据中心的资源使用率。其次,Spanner拥有两项在分布式数据库中难以实现的特性:它提供外部一致性[16]的读写,以及在同一时间戳下跨数据库的全局一致性读取。这些特性让Spanner能够在全球范围支持一致性备份、一致性MapReduce执行[12]以及原子模式更新,即便在有事务正在执行的情况下也不受影响。
这些特性的核心基础是:Spanner会为事务分配具有全局意义的提交时间戳,即便事务是分布式的也不例外。时间戳反映了事务的序列化顺序。此外,序列化顺序满足外部一致性(等价于线性一致性[20]):如果事务\(T_1\)提交发生在事务\(T_2\)开始之前,那么\(T_1\)的提交时间戳小于\(T_2\)的提交时间戳。Spanner是首个在全球范围提供这种保证的系统。
实现这些特性的关键是全新的TrueTime API及其实现。该API直接暴露时钟的不确定性,Spanner对时间戳的保证依赖于该实现提供的误差边界。如果不确定性较大,Spanner会降低速度以等待不确定性消除。谷歌的集群管理软件提供了TrueTime API的实现,通过使用多种现代时钟参考源(GPS与原子钟),该实现将不确定性控制在很小的范围(通常小于10毫秒)。
第2节阐述Spanner的实现架构、功能集,以及设计过程中的工程决策。第3节介绍全新的TrueTime API并概述其实现原理。第4节说明Spanner如何利用TrueTime实现外部一致性分布式事务、无锁只读事务以及原子模式更新。第5节给出Spanner性能与TrueTime表现的基准测试结果,并讨论F1的使用经验。第6、7、8节分别介绍相关工作、未来工作与总结结论。
2 实现架构
本节阐述Spanner实现的架构与设计原理,随后介绍用于管理复制与数据局部性的目录抽象——它也是数据移动的基本单元。最后介绍我们的数据模型,解释Spanner为何呈现为关系数据库而非键值存储的形态,以及应用如何控制数据局部性。
一个Spanner部署实例称为一个全域(universe)。由于Spanner管理全球数据,因此运行中的全域数量很少。目前我们运行着测试/沙箱全域、开发/生产全域,以及纯生产全域。
Spanner由多个区域(zone)组成,每个区域大致对应一个Bigtable服务器部署[9]。区域是部署管理的单元,所有区域构成了数据可复制的位置集合。随着新数据中心投入使用、旧数据中心下线,可以向运行中的系统添加或移除区域。区域也是物理隔离的单元:一个数据中心内可以有一个或多个区域,例如当不同应用的数据必须在同一数据中心的不同服务器组上分区部署时。
图1展示了一个Spanner全域中的服务器组成。每个区域包含一个区域主控(zonemaster)与一百到数千个Spanner服务器(spanserver)。区域主控负责将数据分配给Spanner服务器,Spanner服务器负责向客户端提供数据服务。每个区域的位置代理(location proxy)用于帮助客户端定位负责其数据的Spanner服务器。全域主控(universe master)与放置驱动器(placement driver)目前是单例节点。全域主控主要是一个控制台,展示所有区域的状态信息,用于交互式调试。放置驱动器以分钟为粒度,负责跨区域自动数据迁移。它会定期与Spanner服务器通信,找出需要迁移的数据——要么是为了满足更新后的复制约束,要么是为了负载均衡。出于篇幅考虑,我们仅详细介绍Spanner服务器。

图1:Spanner服务器架构
2.1 Spanner服务器软件栈
本节聚焦Spanner服务器的实现,说明复制与分布式事务是如何分层构建在基于Bigtable的实现之上的。软件栈如图2所示。在最底层,每个Spanner服务器负责100到1000个称为**表块(tablet)**的数据结构实例。表块与Bigtable的表块抽象类似,它实现了一组如下映射:
(key:string, timestamp:int64) → string
与Bigtable不同,Spanner会为数据分配时间戳——这是Spanner更像多版本数据库而非键值存储的重要特征。表块的状态存储在一组类B树文件与预写日志中,全部存储在名为Colossus的分布式文件系统上(谷歌文件系统的后继者[15])。
为了支持复制,每个Spanner服务器在每个表块之上实现了一个Paxos状态机。(Spanner早期版本支持每个表块对应多个Paxos状态机,以实现更灵活的复制配置,但该设计的复杂性让我们放弃了这种方案。)每个状态机将其元数据与日志存储在对应的表块中。我们的Paxos实现支持基于时间的领导者租约,实现长生命周期领导者,租约时长默认为10秒。当前Spanner实现会将每次Paxos写入记录两次:一次记入表块日志,一次记入Paxos日志。这是权宜之计,我们未来会解决这个问题。我们的Paxos实现采用流水线化设计,以提升广域延迟下Spanner的吞吐量;但Paxos会按顺序应用写入操作(这是第4节依赖的前提)。
Paxos状态机用于实现一致复制的映射集合。每个副本的键值映射状态存储在对应的表块中。写入操作必须在领导者节点发起Paxos协议;读取操作可以在任何数据足够新的副本上,直接从底层表块访问状态。所有副本共同构成一个Paxos组。
在每个作为领导者的副本上,Spanner服务器都实现了一个锁表,用于实现并发控制。锁表维护两阶段锁的状态:将键范围映射到锁状态。(注意,长生命周期的Paxos领导者对于高效管理锁表至关重要。)在Bigtable与Spanner中,我们都针对长事务进行了设计(例如报表生成,可能需要数分钟),这类事务在冲突场景下使用乐观并发控制的表现很差。需要同步的操作(如事务读取)会在锁表中获取锁,其他操作则绕过锁表。
在每个作为领导者的副本上,Spanner服务器还实现了事务管理器,以支持分布式事务。事务管理器用于实现参与者领导者,组内的其他副本称为参与者从属节点。如果事务只涉及一个Paxos组(大多数事务都是如此),则可以绕过事务管理器,因为锁表与Paxos共同提供了事务保证。如果事务涉及多个Paxos组,这些组的领导者会协作执行两阶段提交。其中一个参与者组会被选为协调者:该组的参与者领导者称为协调者领导者,该组的从属节点称为协调者从属节点。每个事务管理器的状态存储在底层的Paxos组中(因此具备复制保障)。

图2:Spanner服务器软件栈
2.2 目录与放置
在键值映射集合之上,Spanner实现支持一种称为**目录(directory)**的分桶抽象,它是一组共享公共前缀的连续键。(使用“目录”这个术语是历史原因,更准确的术语应该是“桶”。)我们将在2.3节解释该前缀的来源。目录的支持让应用可以通过合理设计键来控制数据的局部性。
目录是数据放置的基本单元。同一个目录中的所有数据都拥有相同的复制配置。在Paxos组之间迁移数据时,是以目录为单位进行的,如图3所示。Spanner迁移目录的原因包括:为Paxos组减负、将频繁共同访问的目录放到同一个组、或是将目录迁移到离访问者更近的组。目录迁移可以在客户端操作持续进行的同时执行。通常一个50MB的目录可以在几秒内完成迁移。
一个Paxos组可以包含多个目录,这意味着Spanner的表块与Bigtable的表块不同:前者不一定是行空间上单个字典序连续的分区。相反,Spanner表块是一个容器,可以封装行空间的多个分区。我们做此设计是为了能够将频繁共同访问的多个目录并置在同一位置。
Movedir是用于在Paxos组之间迁移目录的后台任务[14]。Movedir也用于为Paxos组添加或移除副本[25],因为Spanner目前还不支持Paxos组内的配置变更。Movedir没有实现为单个事务,以避免大规模数据迁移阻塞正在进行的读写操作。相反,Movedir先登记开始迁移数据的事实,然后在后台迁移数据。当只剩极少量数据未迁移时,它会通过一个事务原子性地迁移这部分数据,并更新两个Paxos组的元数据。
目录也是应用可指定地理复制属性(简称放置策略)的最小单元。我们的放置规范语言在设计上分离了复制配置的管理职责。管理员控制两个维度:副本的数量与类型,以及副本的地理位置。他们会在这两个维度上创建一组命名选项(例如“北美,5副本加1个见证节点”)。应用通过为每个数据库和/或单个目录标记这些选项的组合,来控制数据的复制方式。例如,应用可以将每个终端用户的数据存储在各自的目录中,这样用户A的数据可以在欧洲有3个副本,用户B的数据可以在北美有5个副本。
为了便于说明,我们做了简化。实际上,当目录变得过大时,Spanner会将其拆分为多个分片(fragment)。不同分片可以由不同的Paxos组(也就是不同的服务器)提供服务。Movedir实际在组间迁移的是分片,而非整个目录。

图3:目录是Paxos组之间的数据移动单元
2.3 数据模型
Spanner向应用提供以下数据特性:基于模式化半关系表的数据模型、查询语言,以及通用事务。支持这些特性的演进由诸多因素推动。Megastore的流行证明了对模式化半关系表与同步复制的需求[5]。谷歌内部至少有300个应用在使用Megastore(尽管其性能相对较低),因为它的数据模型比Bigtable更易于管理,并且支持跨数据中心的同步复制(Bigtable仅支持跨数据中心的最终一致性复制)。谷歌使用Megastore的知名应用包括Gmail、Picasa、日历、安卓市场与AppEngine。同时,鉴于Dremel作为交互式数据分析工具的普及,Spanner显然需要支持类SQL的查询语言[28]。最后,Bigtable缺乏跨行事务的问题也屡遭诟病,Percolator[32]的部分设计初衷就是解决这一缺陷。一些研究者认为通用两阶段提交的支持成本过高,因为会带来性能或可用性问题[9,10,19]。我们认为,当事务过度使用成为瓶颈时,让应用程序员去处理性能问题,总比永远因为没有事务而绕路编码要好。而基于Paxos运行两阶段提交,可以缓解可用性问题。
应用数据模型分层构建在实现层支持的目录分桶键值映射之上。应用可以在一个全域中创建一个或多个数据库。每个数据库可以包含任意数量的模式化表。表的形态类似关系数据库表,包含行、列与版本化的值。我们不会详细介绍Spanner的查询语言,它大体是SQL,加上一些支持协议缓冲区(Protocol Buffer)类型字段的扩展。
Spanner的数据模型并非纯关系型,因为行必须有名称。更准确地说,每个表都要求包含一个或多个主键列组成的有序集合。这一要求正是Spanner仍保留键值存储特征的体现:主键构成了行的标识,每个表定义了从主键列到非主键列的映射。只有当行的键对应了某个值(即使是NULL)时,该行才存在。采用这种结构的好处是让应用可以通过键的设计来控制数据局部性。
图4给出了一个Spanner模式示例,用于按用户、按相册存储照片元数据。该模式语言与Megastore的类似,额外要求是每个Spanner数据库必须由客户端划分为一个或多个表层次结构。客户端应用通过INTERLEAVE IN声明在数据库模式中定义层次结构。层次结构最顶层的表称为目录表。目录表中键为K的每一行,加上其所有子孙表中字典序以K开头的行,共同构成一个目录。ON DELETE CASCADE表示删除目录表中的一行时,会级联删除所有关联的子行。图中还展示了示例数据库的交错布局:例如Albums(2,1)代表Albums表中user_id为2、album_id为1的行。这种通过表交错形成目录的设计意义重大,因为它让客户端可以描述多个表之间的局部性关系,这对于分片分布式数据库的性能至关重要。没有它,Spanner将无法知晓最重要的局部性关系。
CREATE TABLE Users {
uid INT64 NOT NULL, email STRING
} PRIMARY KEY(uid), DIRECTORY;
CREATE TABLE Albums {
uid INT64 NOT NULL, aid INT64 NOT NULL,
name STRING
} PRIMARY KEY(uid, aid),
INTERLEAVE IN PARENT Users ON DELETE CASCADE;

图4:照片元数据的Spanner模式示例,以及INTERLEAVE IN声明对应的交错结构
3 TrueTime
表1:TrueTime API。参数t的类型为TStamp。
| 方法 | 返回值 |
|---|---|
TT.now() |
TTinterval类型:[earliest, latest] |
TT.after(t) |
如果时间t已确定过去,则返回true |
TT.before(t) |
如果时间t确定还未到来,则返回true |
本节介绍TrueTime API并概述其实现原理。大部分细节我们将在另一篇论文中阐述,本文的目标是展示这种API的价值。表1列出了该API的方法。TrueTime将时间显式表示为TTinterval——一个带边界时间不确定性的区间(这与标准时间接口不同,后者完全不会向客户端暴露不确定性)。TTinterval的端点类型为TTstamp。TT.now()方法返回一个TTinterval区间,该区间保证包含TT.now()调用发生时的绝对时间。时间纪元类似UNIX时间,但采用了闰秒平滑处理。定义瞬时误差边界为\(\epsilon\),即区间宽度的一半;平均误差边界为\(\bar{\epsilon}\)。TT.after()与TT.before()方法是基于TT.now()的便捷封装。
用函数\(t_{abs}(e)\)表示事件\(e\)的绝对时间。更正式地说,TrueTime保证:对于调用tt = TT.now(),有tt.earliest \(\leq t_{abs}(e_{now}) \leq\) tt.latest,其中\(e_{now}\)是该调用事件。
TrueTime使用的底层时间参考源是GPS与原子钟。采用两种时间参考源是因为它们的故障模式不同。GPS参考源的风险包括天线与接收器故障、本地无线电干扰、关联故障(例如闰秒处理错误、欺骗攻击等设计缺陷)以及GPS系统中断。原子钟的故障模式与GPS不相关,彼此之间也不相关,但长期来看会因频率误差产生显著漂移。
TrueTime的实现包括每个数据中心的一组时间主节点(time master),以及每台机器上的时间从守护进程(timeslave daemon)。大多数主节点配备带专用天线的GPS接收器,这些主节点在物理上分散部署,以降低天线故障、无线电干扰与欺骗攻击的影响。其余的主节点(我们称为“末日主节点”,Armageddon master)配备原子钟。原子钟并没有那么昂贵:一个末日主节点的成本与GPS主节点处于同一量级。所有主节点的时间参考会定期相互比对。每个主节点还会将其参考源的时间推进速率与本地时钟交叉校验,如果偏差过大就会自动下线。在两次同步之间,末日主节点会公布一个缓慢增长的时间不确定性,该值基于保守估算的最坏情况时钟漂移得出。GPS主节点公布的不确定性通常接近零。
每个守护进程会轮询多个主节点[29],以降低单个主节点出错带来的风险。其中一部分是邻近数据中心的GPS主节点,其余是较远数据中心的GPS主节点,以及一些末日主节点。守护进程使用Marzullo算法的变体[27]来检测并排除异常节点,然后将本地机器时钟与正常节点同步。为了防止本地时钟故障,频率偏移超出组件规格与运行环境得出的最坏边界的机器会被移出集群。
在两次同步之间,守护进程会公布一个缓慢增长的时间不确定性。\(\epsilon\)的取值基于保守估算的最坏情况本地时钟漂移得出,同时也受时间主节点的不确定性以及与主节点通信延迟的影响。在我们的生产环境中,\(\epsilon\)通常是时间的锯齿函数,每个轮询周期内在1到7毫秒之间波动。因此\(\bar{\epsilon}\)大部分时间为4毫秒。守护进程的轮询间隔目前为30秒,当前采用的漂移速率为200微秒/秒,这两者共同造成了0到6毫秒的锯齿边界。剩下的1毫秒来自与时间主节点的通信延迟。在故障场景下,锯齿形态可能会出现偏离。例如,偶尔的时间主节点不可用会导致整个数据中心的\(\epsilon\)升高。同样,机器过载与链路拥塞也可能导致局部的\(\epsilon\)尖峰。
表2:Spanner中的读写操作类型对比
| 操作类型 | 时间戳相关章节 | 并发控制方式 | 所需副本 |
|---|---|---|---|
| 读写事务 | 第4.1.2节 | 悲观并发控制 | 领导者节点 |
| 只读事务 | 第4.1.4节 | 无锁 | 时间戳分配需领导者节点;读取可在任意副本,需满足第4.1.3节条件 |
| 快照读取(客户端指定时间戳) | — | 无锁 | 任意副本,需满足第4.1.3节条件 |
| 快照读取(客户端指定陈旧度上限) | 第4.1.3节 | 无锁 | 任意副本,需满足第4.1.3节条件 |
4 并发控制
本节阐述如何利用TrueTime保证并发控制的正确性,以及如何基于这些特性实现外部一致性事务、无锁只读事务、非阻塞历史读取等功能。例如,基于这些特性可以保证:在时间戳t执行的全数据库审计读取,能且只能看到所有在t之前提交的事务的结果。
后续需要区分Paxos视角的写入(除非上下文明确,否则称为Paxos写入)与Spanner客户端写入。例如,两阶段提交的准备阶段会产生一次Paxos写入,但并没有对应的Spanner客户端写入。
4.1 时间戳管理
表2列出了Spanner支持的操作类型。Spanner实现支持读写事务、只读事务(预先声明的快照隔离事务)以及快照读取。单条写入以读写事务实现,非快照单条读取以只读事务实现。两者都具备内部重试机制(客户端无需自行实现重试循环)。
只读事务是一类具备快照隔离[6]性能优势的事务。只读事务必须预先声明不包含任何写入操作,它不是简单的不带写入的读写事务。只读事务中的读取在系统选定的时间戳执行,无需加锁,因此不会阻塞后续写入。只读事务的读取可以在任何数据足够新的副本上执行(见第4.1.3节)。
快照读取是对历史数据的读取,执行时无需加锁。客户端可以为快照读取指定一个时间戳,或者指定可接受的时间戳陈旧度上限,由Spanner选择时间戳。无论哪种方式,快照读取都可以在任何数据足够新的副本上执行。
对于只读事务与快照读取,一旦选定时间戳,提交就是必然的,除非该时间戳对应的数据已被垃圾回收。因此,客户端无需在重试循环中缓存结果。当服务器故障时,客户端可以通过复用时间戳与当前读取位置,在另一台服务器上内部继续执行查询。
4.1.1 Paxos领导者租约
Spanner的Paxos实现使用定时租约来实现长生命周期领导者(默认10秒)。候选领导者发送定时租约投票请求,当收到法定数量的租约投票时,该节点就获得了租约。副本在写入成功时会隐式延长其租约投票,领导者在租约即将到期时会请求续期。定义领导者的租约区间:从它确认获得法定数量租约投票开始,到它不再拥有法定数量租约投票(部分投票过期)时结束。Spanner依赖如下不相交不变量:对于每个Paxos组,任意两个Paxos领导者的租约区间互不相交。附录A阐述了该不变量的保障机制。
Spanner实现允许Paxos领导者通过释放从属节点的租约投票来主动退位。为了维护不相交不变量,Spanner对允许退位的时机做了约束。定义\(s_{max}\)为领导者使用过的最大时间戳,后续章节会说明\(s_{max}\)的推进时机。在退位之前,领导者必须等待,直到TT.after(s_max)为真。
4.1.2 为读写事务分配时间戳
事务型读写使用两阶段锁。因此,可以在获取所有锁之后、释放任何锁之前的任意时刻为事务分配时间戳。对于给定事务,Spanner将Paxos分配给该事务提交对应的Paxos写入的时间戳,作为事务的提交时间戳。
Spanner依赖如下单调性不变量:在每个Paxos组内,Spanner为Paxos写入分配的时间戳严格单调递增,即便在不同领导者之间也是如此。单个领导者副本自然可以按单调递增顺序分配时间戳。跨领导者的单调性通过不相交不变量保障:领导者只能在自身租约区间内分配时间戳。注意,每当分配了时间戳\(s\),就会将\(s_{max}\)推进到\(s\),以维护不相交性。
Spanner还维护如下外部一致性不变量:如果事务\(T_2\)的开始发生在事务\(T_1\)提交之后,那么\(T_2\)的提交时间戳必须大于\(T_1\)的提交时间戳。定义事务\(T_i\)的开始事件为\(e_i^{start}\),提交事件为\(e_i^{commit}\);事务\(T_i\)的提交时间戳为\(s_i\)。该不变量可表示为:\(t_{abs}(e_1^{commit}) < t_{abs}(e_2^{start}) \Rightarrow s_1 < s_2\)。事务执行与时间戳分配协议遵循两条规则,共同保证了该不变量,如下所述。定义写事务\(T_i\)的提交请求到达协调者领导者的事件为\(e_i^{server}\)。
起始规则:写事务\(T_i\)的协调者领导者分配提交时间戳\(s_i\),该时间戳不小于\(e_i^{server}\)之后计算的TT.now().latest值。注意参与者领导者在此不影响结果,第4.2.1节会说明它们如何参与下一条规则的实现。
提交等待规则:协调者领导者保证,在TT.after(s_i)为真之前,客户端无法看到\(T_i\)提交的任何数据。提交等待确保\(s_i\)小于\(T_i\)的绝对提交时间,即\(s_i < t_{abs}(e_i^{commit})\)。提交等待的实现见第4.2.1节。证明:
\(
\begin{aligned}
s_{1}&<t_{abs}(e_{1}^{commit})\quad&(\text{提交等待规则}) \\
t_{abs}(e_{1}^{commit})&<t_{abs}(e_{2}^{start})\quad&(\text{前提假设}) \\
t_{abs}(e_{2}^{start})&\leq t_{abs}(e_{2}^{server})\quad&(\text{因果关系}) \\
t_{abs}(e_{2}^{server})&\leq s_{2}\quad&(\text{起始规则}) \\
s_{1}&<s_{2}\quad&(\text{传递性})
\end{aligned}
\)
4.1.3 按时间戳提供读取服务
第4.1.2节所述的单调性不变量,让Spanner可以正确判断副本的状态是否足够新,能否满足读取需求。每个副本维护一个称为安全时间(\(t_{safe}\))的值,代表副本数据已更新到的最大时间戳。如果\(t \leq t_{safe}\),则该副本可以满足时间戳为t的读取请求。
定义\(t_{safe} = \min(t_{safe}^{Paxos}, t_{safe}^{TM})\),其中每个Paxos状态机有一个安全时间\(t_{safe}^{Paxos}\),每个事务管理器有一个安全时间\(t_{safe}^{TM}\)。\(t_{safe}^{Paxos}\)的定义更简单:它是已应用的最高Paxos写入的时间戳。由于时间戳单调递增且写入按顺序应用,对于Paxos层来说,不会再有时间戳小于等于\(t_{safe}^{Paxos}\)的写入。
如果没有处于准备阶段(已准备但未提交)的事务——也就是处于两阶段提交两个阶段之间的事务,那么副本的\(t_{safe}^{TM}\)为\(\infty\)。(对于参与者从属节点来说,\(t_{safe}^{TM}\)实际对应的是副本领导者的事务管理器状态,从属节点可以通过Paxos写入传递的元数据推断出该状态。)如果存在这类事务,那么受这些事务影响的状态就是不确定的:参与者副本还不知道这些事务是否会提交。正如第4.2.1节所述,提交协议保证每个参与者都知道已准备事务的时间戳下界。事务\(T_i\)的每个参与者领导者(属于组g)会为其准备记录分配一个准备时间戳\(s_{i,g}^{prepare}\)。协调者领导者保证事务的提交时间戳\(s_i\)大于等于所有参与者组g的\(s_{i,g}^{prepare}\)。因此,对于组g中的任意副本,所有在g上处于准备阶段的事务\(T_i\)满足:\(t_{safe}^{TM} = \min_i (s_{i,g}^{prepare}) – 1\)。
4.1.4 为只读事务分配时间戳
只读事务分两个阶段执行:分配时间戳\(s_{read}\)[8],然后以该时间戳作为快照读取,执行事务中的所有读操作。快照读取可以在任何数据足够新的副本上执行。
在事务开始后的任意时刻,简单地赋值\(s_{read} = TT.now().latest\),可以通过与第4.1.2节写事务类似的推导保证外部一致性。然而,如果\(t_{safe}\)还未推进到足够大,在\(s_{read}\)时刻执行数据读取就可能阻塞。(此外,选择\(s_{read}\)的值可能也会推进\(s_{max}\),以维护不相交性。)为了降低阻塞概率,Spanner应该选择能保证外部一致性的最旧时间戳。第4.2.2节会说明如何选择这样的时间戳。
4.2 实现细节
本节阐述之前省略的读写事务与只读事务的部分实践细节,以及用于实现原子模式变更的特殊事务类型的实现。随后介绍对基础方案的一些优化改进。
4.2.1 读写事务
与Bigtable类似,事务中的写入操作在客户端缓存,直到提交时才发送。因此,事务内的读取不会看到本事务写入的结果。这种设计在Spanner中非常适用,因为读取会返回所读数据的时间戳,而未提交的写入还未分配时间戳。
读写事务内的读取使用**伤害等待(wound-wait)**算法避免死锁[33]。客户端向对应组的领导者副本发起读取请求,领导者获取读锁后读取最新数据。在客户端事务保持打开期间,会发送保活消息,防止参与者领导者将事务超时。当客户端完成所有读取并缓存所有写入后,开始两阶段提交。客户端选择一个协调者组,向每个参与者的领导者发送提交消息,包含协调者标识与所有缓存的写入。由客户端驱动两阶段提交,可以避免数据在广域链路上传输两次。
非协调者的参与者领导者首先获取写锁,然后选择一个准备时间戳(必须大于之前分配给所有事务的时间戳,以保证单调性),并通过Paxos记录准备日志。随后每个参与者将其准备时间戳通知协调者。
协调者领导者同样先获取写锁,但跳过准备阶段。在收到所有其他参与者领导者的消息后,它为整个事务选择时间戳。提交时间戳\(s\)必须满足:大于等于所有准备时间戳(以满足第4.1.3节的约束)、大于协调者收到提交请求时的TT.now().latest值、大于该领导者之前分配给所有事务的时间戳(同样为了保证单调性)。随后协调者领导者通过Paxos记录提交日志(如果等待其他参与者超时则记录中止日志)。
在允许任何协调者副本应用提交记录之前,协调者领导者需要等待,直到TT.after(s)为真,以遵守第4.1.2节的提交等待规则。由于协调者领导者是基于TT.now().latest选择的\(s\),现在要等待到该时间戳确定已成为过去,因此预期等待时间至少为\(2*\bar{\epsilon}\)。该等待通常可以与Paxos通信重叠。提交等待结束后,协调者将提交时间戳发送给客户端与所有其他参与者领导者。每个参与者领导者通过Paxos记录事务结果。所有参与者在同一时间戳应用事务,然后释放锁。
4.2.2 只读事务
分配时间戳需要所有涉及读取的Paxos组之间进行协商。因此,Spanner要求每个只读事务都有一个范围表达式,用于概括整个事务将读取的键范围。对于单条查询,Spanner会自动推断范围。
如果范围中的数据仅由一个Paxos组提供服务,那么客户端将只读事务发送到该组的领导者。(当前Spanner实现仅在Paxos领导者处为只读事务分配时间戳。)领导者分配\(s_{read}\)并执行读取。对于单站点读取,Spanner的表现通常优于直接使用TT.now().latest。定义\((LastTS())\)为Paxos组最后一次提交写入的时间戳。如果没有处于准备阶段的事务,赋值\(s_{read} = (LastTS())\)显然满足外部一致性:事务会看到最后一次写入的结果,因此顺序在其之后。
如果范围中的数据由多个Paxos组提供服务,则有多种方案。最复杂的方案是与所有组的领导者进行一轮通信,基于\((LastTS())\)协商出\(s_{read}\)。Spanner目前实现了一种更简单的方案:客户端跳过协商轮次,直接以\(s_{read} = TT.now().latest\)执行读取(这可能需要等待安全时间推进)。事务中的所有读取都可以发送到数据足够新的副本。
4.2.3 模式变更事务
TrueTime让Spanner能够支持原子模式变更。如果使用标准事务是不现实的,因为参与者数量(数据库中的组数量)可能达到数百万。Bigtable在单个数据中心内支持原子模式变更,但其模式变更会阻塞所有操作。
Spanner的模式变更事务是标准事务的一种非阻塞变体。首先,它会被显式分配一个未来的时间戳,并在准备阶段登记。因此,跨数千台服务器的模式变更可以在对其他并发操作影响极小的情况下完成。其次,隐式依赖模式的读写操作,会与所有已登记的模式变更时间戳t进行同步:如果读写操作的时间戳早于t,则可以继续执行;如果晚于t,则必须阻塞等待模式变更事务完成。如果没有TrueTime,定义“模式变更在时刻t发生”是没有意义的。
4.2.4 优化改进
上述定义的\(t_{safe}^{TM}\)存在一个缺陷:单个处于准备阶段的事务就会阻塞\(t_{safe}\)的推进。结果,即便读取操作与该事务不冲突,也无法在更晚的时间戳执行。可以通过为\(t_{safe}^{TM}\)增加键范围到准备事务时间戳的细粒度映射,来消除这种伪冲突。该信息可以存储在锁表中——锁表本身已经维护了键范围到锁元数据的映射。当读取请求到达时,只需检查与之冲突的键范围对应的细粒度安全时间即可。
上述定义的\((LastTS())\)也有类似缺陷:如果一个事务刚提交,不冲突的只读事务仍然必须将\(s_{read}\)赋值为该事务之后的时间戳,因此读取执行可能被延迟。同样,可以通过在锁表中增加键范围到提交时间戳的细粒度映射来解决这个问题。(我们尚未实现该优化。)当只读事务到达时,可以取事务冲突的键范围对应的(LastTS())最大值作为时间戳,除非存在冲突的准备中事务(可通过细粒度安全时间判断)。
上述定义的\(t_{safe}^{Paxos}\)也存在缺陷:如果没有Paxos写入,它就无法推进。也就是说,对于时间戳t的快照读取,如果Paxos组的最后一次写入发生在t之前,就无法执行。Spanner利用领导者租约区间的不相交性解决了这个问题。每个Paxos领导者维护一个阈值,未来写入的时间戳都将大于该阈值,以此推进\(t_{safe}^{Paxos}\):它维护映射\((MinNextTS(n))\),表示Paxos序号n对应的下一个序号n+1的Paxos写入可能分配的最小时间戳。当副本应用完序号n的写入后,就可以将\(t_{safe}^{Paxos}\)推进到\((MinNextTS(n)) – 1\)。
单个领导者可以轻松保证其MinNextTS()承诺的有效性。由于MinNextTS()承诺的时间戳都在领导者租约区间内,不相交不变量保证了跨领导者的MinNextTS()承诺互不冲突。如果领导者希望将MinNextTS()推进到超出自身租约结束时间的位置,必须先续期租约。注意,\(s_{max}\)总是会推进到MinNextTS()中的最大值,以维护不相交性。
领导者默认每8秒推进一次MinNextTS()的值。因此,在没有准备中事务的情况下,空闲Paxos组中正常的从属节点,最坏情况下可以为8秒之前时间戳的读取提供服务。领导者也可以根据从属节点的需求按需推进MinNextTS()。
5 评估
我们首先测试Spanner在复制、事务与可用性方面的性能,然后给出TrueTime表现的数据,以及首个客户F1的案例研究。
5.1 微基准测试
表3给出了Spanner的部分微基准测试结果。测试在分时共享机器上进行:每个Spanner服务器运行在4GB内存、4核(AMD Barcelona 2200MHz)的调度单元上。客户端运行在独立机器上。每个区域包含一个Spanner服务器。客户端与区域都部署在网络延迟小于1毫秒的数据中心内。(这种部署非常常见:大多数应用不需要将所有数据分布到全球。)测试数据库包含50个Paxos组、2500个目录。操作为4KB大小的单条读写。所有读取在压缩后都从内存提供服务,因此我们只测量Spanner调用栈的开销。此外,测试前先执行一轮不计入统计的读取,以预热所有位置缓存。
表3:操作微基准测试。10次运行的均值与标准差。1D表示单副本且禁用提交等待。
| 副本数 | 延迟(毫秒) | 吞吐量(千次操作/秒) | ||||
|---|---|---|---|---|---|---|
| 写入 | 只读事务 | 快照读取 | 写入 | 只读事务 | 快照读取 | |
| 1D | 9.4 ± 0.6 | — | — | 4.0 ± 0.3 | — | — |
| 1 | 14.4 ± 1.0 | 1.4 ± 0.1 | 1.3 ± 0.1 | 4.1 ± 0.05 | 10.9 ± 0.4 | 13.5 ± 0.1 |
| 3 | 13.9 ± 0.6 | 1.3 ± 0.1 | 1.2 ± 0.1 | 2.2 ± 0.5 | 13.8 ± 0.32 | 38.5 ± 0.3 |
| 5 | 14.4 ± 0.4 | 1.4 ± 0.05 | 1.3 ± 0.04 | 2.8 ± 0.3 | 25.3 ± 0.52 | 50.0 ± 1.1 |
在延迟测试中,客户端发送的操作数量足够少,避免服务器端排队。从单副本测试结果可以看出,提交等待约为5毫秒,Paxos延迟约为9毫秒。随着副本数量增加,延迟基本保持不变,标准差有所降低,因为Paxos在组内所有副本上并行执行。副本越多,达成法定人数的延迟对单个从属节点变慢的敏感度就越低。
在吞吐量测试中,客户端发送足够多的操作,使服务器CPU达到饱和。快照读取可以在任何数据足够新的副本上执行,因此其吞吐量几乎随副本数量线性增长。单条读取的只读事务只能在领导者执行,因为时间戳分配必须在领导者进行。只读事务的吞吐量随副本数量增长,因为有效Spanner服务器数量增加了:在实验设置中,Spanner服务器数量等于副本数量,且领导者随机分布在各个区域。写入吞吐量也受益于同样的实验配置(这解释了从3副本到5副本吞吐量上升的原因),但随着副本数量增加,每次写入的工作量也线性增长,这一负面影响超过了前者的收益。
表4表明两阶段提交可以扩展到合理的参与者数量:该表总结了在3个区域(每个区域25个Spanner服务器)上运行的一组实验结果。扩展到50个参与者时,平均延迟与99分位延迟都处于合理范围;到100个参与者时,延迟开始显著上升。
表4:两阶段提交可扩展性。10次运行的均值与标准差。
| 参与者数量 | 延迟(毫秒) | |
|---|---|---|
| 均值 | 99分位 | |
| 1 | 17.0 ± 1.4 | 75.0 ± 34.9 |
| 2 | 24.5 ± 2.5 | 87.6 ± 35.9 |
| 5 | 31.5 ± 6.2 | 104.5 ± 52.2 |
| 10 | 30.0 ± 3.7 | 95.6 ± 25.4 |
| 25 | 35.5 ± 5.6 | 100.4 ± 42.7 |
| 50 | 42.7 ± 4.1 | 93.7 ± 22.9 |
| 100 | 71.4 ± 7.6 | 131.2 ± 17.6 |
| 200 | 150.5 ± 11.0 | 320.3 ± 35.1 |
5.2 可用性
图5展示了在多数据中心部署Spanner的可用性优势。图中给出了三组数据中心故障场景下的吞吐量实验结果,所有结果都对齐到同一时间轴。测试全域包含5个区域\(Z_i\),每个区域有25个Spanner服务器。测试数据库被分片为1250个Paxos组,100个测试客户端以总计5万次/秒的速率持续发送非快照读取请求。所有领导者都显式部署在\(Z_1\)中。在测试开始5秒时,杀死一个区域的所有服务器:非领导者组杀死\(Z_2\);硬切领导者组杀死\(Z_1\);软切领导者组杀死\(Z_1\),但会先通知所有服务器进行领导者交接。
杀死\(Z_2\)对读取吞吐量没有影响。杀死\(Z_1\)但给领导者时间将权限交接给其他区域,影响很小:吞吐量下降在图中不可见,大约为3-4%。另一方面,无预警杀死\(Z_1\)则会造成严重影响:完成速率几乎降到0。但随着领导者重新选举,系统吞吐量会上升到约10万次/秒,这是由于实验的两个效应:系统存在额外容量,且领导者不可用时操作会排队。因此吞吐量会先上升,然后回落并稳定在稳态速率。

图5:服务器宕机对吞吐量的影响
我们还可以看到Paxos领导者租约设为10秒的影响。杀死区域后,各个组的领导者租约到期时间会均匀分布在接下来的10秒内。当失效领导者的租约到期后,很快就会选举出新的领导者。大约在杀死区域10秒后,所有组都重新选出领导者,吞吐量完全恢复。更短的租约时间可以降低服务器故障对可用性的影响,但会增加租约续期的网络流量。我们正在设计并实现一种机制,让从属节点在领导者故障时主动释放Paxos领导者租约。
5.3 TrueTime
关于TrueTime需要回答两个问题:\(\epsilon\)是否真的是时钟不确定性的边界?\(\epsilon\)最坏情况下会有多大?对于第一个问题,最严重的风险是本地时钟漂移超过200微秒/秒:这会打破TrueTime的前提假设。我们的机器统计数据显示,CPU故障的概率是时钟故障的6倍。也就是说,相对于其他更严重的硬件问题,时钟问题极其罕见。因此我们认为,TrueTime的实现与Spanner依赖的其他任何软件一样可靠。
图6展示了在相距最远2200公里的多个数据中心中,数千台Spanner服务器上的TrueTime数据。图中绘制了\(\epsilon\)的90分位、99分位与99.9分位值,采样时机为时间从守护进程刚轮询完时间主节点时。该采样忽略了本地时钟不确定性带来的\(\epsilon\)锯齿波动,因此测量的是时间主节点的不确定性(通常为0)加上与主节点的通信延迟。

图6:TrueTime \(\epsilon\)值分布,采样时机为时间从守护进程刚轮询完时间主节点时。图中绘制了90分位、99分位与99.9分位值。
数据表明,决定\(\epsilon\)基准值的这两个因素通常都不是问题。然而,严重的尾延迟问题可能导致\(\epsilon\)升高。3月30日之后尾延迟的下降,得益于网络优化减少了瞬时链路拥塞。4月13日出现的\(\epsilon\)升高持续了约1小时,原因是一个数据中心的2台时间主节点因例行维护下线。我们会持续排查并消除TrueTime尖峰的诱因。
5.4 F1案例
2011年初,作为谷歌广告后端重构项目F1的一部分[35],Spanner开始在生产工作负载下进行实验性评估。该后端最初基于MySQL数据库,通过多种方式手动分片。未压缩数据集大小为数十TB,这相比很多NoSQL实例不算大,但已经给分片MySQL带来了诸多困难。MySQL分片方案将每个客户及其所有关联数据分配到固定分片。这种布局支持按客户使用索引与复杂查询处理,但应用业务逻辑需要感知分片逻辑。随着客户数量与数据量增长,对这个关乎营收的核心数据库进行重新分片成本极高。上一次重新分片耗费了两年多的高强度工作,涉及数十个团队的协调与测试,以最大限度降低风险。这种操作太过复杂,无法定期执行,因此团队不得不将部分数据存储在外部Bigtable中,以限制MySQL数据库的增长,这牺牲了事务特性与全数据查询能力。
F1团队选择使用Spanner有多个原因。首先,Spanner消除了手动重新分片的需求。其次,Spanner提供同步复制与自动故障转移。MySQL的主从复制模式下,故障转移难度大,存在数据丢失与停机风险。第三,F1需要强事务语义,这使得其他NoSQL系统都不适用。应用语义要求跨任意数据的事务与一致性读取。F1团队还需要数据的二级索引(由于Spanner尚未自动支持二级索引),他们利用Spanner事务自行实现了一致性全局索引。
现在所有应用写入默认都通过F1发送到Spanner,而非基于MySQL的应用栈。F1在美国西海岸有2个副本,东海岸有3个副本。选择这些副本站点是为了应对潜在重大自然灾害导致的中断,同时也匹配其前端站点的部署。实际使用中,Spanner的自动故障转移对他们几乎是透明的。尽管过去几个月发生过计划外集群故障,但F1团队最多只需要更新数据库模式,告诉Spanner优先将Paxos领导者部署在哪里,使其靠近迁移后的前端即可。
Spanner的时间戳语义让F1可以高效维护基于数据库状态计算的内存数据结构。F1维护所有变更的逻辑历史日志,该日志作为每个事务的一部分写入Spanner。F1通过某一时间戳的全量数据快照初始化内存结构,然后通过读取增量变更进行更新。
表5展示了F1中每个目录的分片数量分布。每个目录通常对应F1上层应用栈中的一个客户。绝大多数目录(也就是绝大多数客户)仅包含1个分片,这意味着对这些客户数据的读写肯定都在单台服务器上执行。分片数超过100的目录都是包含F1二级索引的表:对这类表的跨多个分片写入极其罕见。F1团队仅在未优化的批量数据加载事务中遇到过这种情况。
表5:F1中目录分片数量分布
| 分片数量 | 目录数量 |
|---|---|
| 1 | >1亿 |
| 2-4 | 341 |
| 5-9 | 5336 |
| 10-14 | 232 |
| 15-99 | 34 |
| 100-500 | 7 |
表6给出了从F1服务器端测量的Spanner操作延迟。在选择Paxos领导者时,东海岸数据中心的副本优先级更高。表中数据是从这些数据中心的F1服务器测得的。写入延迟的标准差较大,是因为锁冲突导致了明显的长尾。读取延迟的标准差更大,部分原因是Paxos领导者分布在两个数据中心,而只有其中一个数据中心的机器配备了SSD。此外,测量包含了两个数据中心系统内的所有读取:读取字节数的均值与标准差分别约为1.6KB和119KB。
表6:24小时内测得的F1视角操作延迟
| 操作类型 | 延迟(毫秒) | 调用次数 | |
|---|---|---|---|
| 均值 | 标准差 | ||
| 所有读取 | 8.7 | 376.4 | 215亿次 |
| 单站点提交 | 72.3 | 112.8 | 3120万次 |
| 多站点提交 | 103.0 | 52.2 | 3210万次 |
6 相关工作
Megastore[5]与DynamoDB[3]都提供了跨数据中心的一致性复制存储服务。DynamoDB提供键值接口,且仅在单个区域内复制。Spanner沿用了Megastore的半关系数据模型,甚至模式语言也很相似。Megastore的性能不高,它分层构建在Bigtable之上,带来了很高的通信成本。它也不支持长生命周期领导者,多个副本都可以发起写入。不同副本的所有写入在Paxos协议中必然冲突,即便它们逻辑上并不冲突:Paxos组的吞吐量会在每秒几次写入时就崩溃。Spanner提供了更高的性能、通用事务与外部一致性。
Pavlo等人[31]对比了数据库与MapReduce的性能[12]。他们指出,多项研究都在探索构建于分布式键值存储之上的数据库功能[1,4,7,41],这证明两个领域正在融合。我们同意这一结论,但我们的工作表明多层整合有其优势:例如,将并发控制与复制整合,降低了Spanner中提交等待的开销。
在复制存储之上分层事务的想法,最早至少可以追溯到Gifford的博士论文[16]。Scatter[17]是近年基于DHT的键值存储,在一致性复制之上分层实现了事务。Spanner专注于提供比Scatter更高层级的接口。Gray与Lamport[18]提出了一种基于Paxos的非阻塞提交协议,他们的协议比两阶段提交带来更多的消息开销,会加剧广域分布组的提交成本。Walter[36]提出了一种快照隔离变体,但仅支持单数据中心内,不支持跨数据中心。相比之下,我们的只读事务提供了更自然的语义,因为我们支持所有操作的外部一致性。
近年来有大量研究致力于降低或消除锁开销。Calvin[40]消除了并发控制:它预先分配时间戳,然后按时间戳顺序执行事务。H-Store[39]与Granola[11]各自支持自己的事务类型分类,其中部分类型可以避免加锁。这些系统都不提供外部一致性。Spanner通过支持快照隔离来解决竞争问题。
VoltDB[42]是一种分片内存数据库,支持广域主从复制用于灾难恢复,但不支持更通用的复制配置。它是所谓NewSQL的代表,NewSQL是业界对可扩展SQL的推动方向[38]。很多商业数据库都实现了历史读取功能,例如MarkLogic[26]与Oracle的Total Recall[30]。Lomet与Li[24]阐述了这种时态数据库的实现策略。
Farsite推导了相对于可信时钟参考的时钟不确定性边界(比TrueTime宽松得多)[13],Farsite中的服务器租约维护方式与Spanner维护Paxos租约的方式类似。松同步时钟在之前的研究中已被用于并发控制[2,23]。我们的工作表明,TrueTime让我们可以跨多组Paxos状态机进行全局时间推理。
7 未来工作
过去一年,我们大部分时间都在与F1团队合作,将谷歌的广告后端从MySQL迁移到Spanner。我们正在持续改进监控与支持工具,并优化性能。此外,我们还在提升备份/恢复系统的功能与性能。目前我们正在实现Spanner模式语言、二级索引自动维护,以及基于负载的自动重新分片。长期来看,我们计划研究几个特性:乐观并行读取可能是一个有价值的方向,但初步实验表明其正确实现并不简单。此外,我们计划最终支持Paxos配置的直接变更[22,34]。
考虑到很多应用会在距离较近的数据中心间复制数据,TrueTime的\(\epsilon\)可能会对性能产生显著影响。我们认为将\(\epsilon\)降低到1毫秒以下没有不可逾越的障碍。可以缩短时间主节点的查询间隔,更好的时钟晶体成本也相对较低。通过改进网络技术可以降低时间主节点的查询延迟,甚至可以通过其他时间分发技术完全避免该延迟。
最后,还有一些显而易见的改进方向。尽管Spanner在节点数量上具备可扩展性,但节点本地数据结构在复杂SQL查询上的性能相对较差,因为它们最初是为简单键值访问设计的。数据库领域的算法与数据结构可以大幅提升单节点性能。其次,根据客户端负载变化自动在数据中心间迁移数据是我们的长期目标,但要让这个目标真正发挥作用,我们还需要具备自动、协调地在数据中心间迁移应用进程的能力。进程迁移带来了更难的问题:如何管理数据中心间的资源获取与分配。
8 结论
总而言之,Spanner结合并拓展了两个研究领域的成果:数据库领域的易用半关系接口、事务与基于SQL的查询语言;系统领域的可扩展性、自动分片、容错、一致性复制、外部一致性与广域分布。从Spanner立项至今,我们用了5年多时间迭代出当前的设计与实现。迭代周期如此之长,部分原因是我们逐渐意识到,Spanner不应该只解决全局复制命名空间的问题,还应该聚焦于Bigtable所缺失的数据库特性。
我们的设计中有一点尤为突出:TrueTime是Spanner所有特性的核心支柱。我们的工作表明,将时钟不确定性具象化到时间API中,可以构建具备更强时间语义的分布式系统。此外,随着底层系统对时钟不确定性的约束越来越严格,强语义带来的开销也会越来越低。作为一个领域,我们在设计分布式算法时,不应再依赖松同步时钟与弱时间API。
致谢
很多人帮助改进了这篇论文:我们的论文指导Jon Howell,他付出了远超职责的努力;各位匿名审稿人;还有众多谷歌同事:Atul Adya、Fay Chang、Frank Dabek、Sean Dorward、Bob Gruber、David Held、Nick Kline、Alex Thomson与Joel Wein。我们的管理层对我们的工作与论文发表给予了大力支持:Aristotle Balogh、Bill Coughran、Urs Hölzle、Doron Meyer、Cos Nicolaou、Kathy Polizzi、Sridhar Ramaswany与Shivakumar Venkataraman。
我们的工作建立在Bigtable与Megastore团队的成果之上。F1团队(尤其是Jeff Shute)与我们紧密合作开发数据模型,并在排查性能与正确性问题上提供了极大帮助。平台团队(尤其是Luiz Barroso与Bob Felderman)助力TrueTime成为现实。最后,还有很多曾在我们团队工作过的谷歌同事:Ken Ashcraft、Paul Cychosz、Krzysztof Ostrowski、Amir Voskoobynik、Matthew Weaver、Theo Vassilakis与Eric Veach;以及近期加入团队的成员:Nathan Bales、Adam Beberg、Vadim Borisov、Ken Chen、Brian Cooper、Cian Cullinan、Robert-Jan Huijsman、Milind Joshi、Andrey Khorlin、Dawid Kuroczko、Laramie Leavitt、Eric Li、Mike Mammarella、Sunil Mushran、Simon Nielsen、Ovidiu Platon、Ananth Shrinivas、Vadim Suvorov与Marcel van der Holst。
参考文献
[1] Azza Abouzeid et al. “HadoopDB: an architectural hybrid of MapReduce and DBMS technologies for analytical workloads”. Proc. of VLDB. 2009, pp. 922–933.
[2] A. Adya et al. “Efficient optimistic concurrency control using loosely synchronized clocks”. Proc. of SIGMOD. 1995, pp. 23–34.
[3] Amazon. Amazon DynamoDB. 2012.
[4] Michael Armbrust et al. “PIQL: Success-Tolerant Query Processing in the Cloud”. Proc. of VLDB. 2011, pp. 181–192.
[5] Jason Baker et al. “Megastore: Providing Scalable, Highly Available Storage for Interactive Services”. Proc. of CIDR. 2011, pp. 223–234.
[6] Hal Berenson et al. “A critique of ANSI SQL isolation levels”. Proc. of SIGMOD. 1995, pp. 1–10.
[7] Matthias Brantner et al. “Building a database on S3”. Proc. of SIGMOD. 2008, pp. 251–264.
[8] A. Chan and R. Gray. “Implementing Distributed Read-Only Transactions”. IEEE TOSE SE-11.2 (Feb. 1985), pp. 205–212.
[9] Fay Chang et al. “Bigtable: A Distributed Storage System for Structured Data”. ACM TOCS 26.2 (June 2008), 4:1–4:26.
[10] Brian F. Cooper et al. “PNUTS: Yahoo!’s hosted data serving platform”. Proc. of VLDB. 2008, pp. 1277–1288.
[11] James Cowling and Barbara Liskov. “Granola: Low-Overhead Distributed Transaction Coordination”. Proc. of USENIX ATC. 2012, pp. 223–236.
[12] Jeffrey Dean and Sanjay Ghemawat. “MapReduce: a flexible data processing tool”. CACM 53.1 (Jan. 2010), pp. 72–77.
[13] John Douceur and Jon Howell. Scalable Byzantine-Fault-Quantifying Clock Synchronization. Tech. rep. MSR-TR-2003-67. MS Research, 2003.
[14] John R. Douceur and Jon Howell. “Distributed directory service in the Farsite file system”. Proc. of OSDI. 2006, pp. 321–334.
[15] Sanjay Ghemawat, Howard Gobioff, and Shun-Tak Leung. “The Google file system”. Proc. of SOSP. Dec. 2003, pp. 29–43.
[16] David K. Gifford. Information Storage in a Decentralized Computer System. Tech. rep. CSL-81-8. PhD dissertation. Xerox PARC, July 1982.
[17] Lisa Glendenning et al. “Scalable consistency in Scatter”. Proc. of SOSP. 2011.
[18] Jim Gray and Leslie Lamport. “Consensus on transaction commit”. ACM TODS 31.1 (Mar. 2006), pp. 133–160.
[19] Pat Helland. “Life beyond Distributed Transactions: an Apostate’s Opinion”. Proc. of CIDR. 2007, pp. 132–141.
[20] Maurice P. Herlihy and Jeannette M. Wing. “Linearizability: a correctness condition for concurrent objects”. ACM TOPLAS 12.3 (July 1990), pp. 463–492.
[21] Leslie Lamport. “The part-time parliament”. ACM TOCS 16.2 (May 1998), pp. 133–169.
[22] Leslie Lamport, Dahlia Malkhi, and Lidong Zhou. “Reconfiguring a state machine”. SIGACT News 41.1 (Mar. 2010), pp. 63–73.
[23] Barbara Liskov. “Practical uses of synchronized clocks in distributed systems”. Distrib. Comput. 6.4 (July 1993), pp. 211–219.
[24] David B. Lomet and Feifei Li. “Improving Transaction-Time DBMS Performance and Functionality”. Proc. of ICDE (2009), pp. 581–591.
[25] Jacob R. Lorch et al. “The SMART way to migrate replicated stateful services”. Proc. of EuroSys. 2006, pp. 103–115.
[26] MarkLogic. MarkLogic 5 Product Documentation. 2012.
[27] Keith Marzullo and Susan Owicki. “Maintaining the time in a distributed system”. Proc. of PODC. 1983, pp. 295–305.
[28] Sergey Melnik et al. “Dremel: Interactive Analysis of Web-Scale Datasets”. Proc. of VLDB. 2010, pp. 330–339.
[29] D.L. Mills. Time synchronization in DCNET hosts. Internet Project Report IEN–173. COMSAT Laboratories, Feb. 1981.
[30] Oracle. Oracle Total Recall. 2012.
[31] Andrew Pavlo et al. “A comparison of approaches to large-scale data analysis”. Proc. of SIGMOD. 2009, pp. 165–178.
[32] Daniel Peng and Frank Dabek. “Large-scale incremental processing using distributed transactions and notifications”. Proc. of OSDI. 2010, pp. 1–15.
[33] Daniel J. Rosenkrantz, Richard E. Stearns, and Philip M. Lewis II. “System level concurrency control for distributed database systems”. ACM TODS 3.2 (June 1978), pp. 178–198.
[34] Alexander Shraer et al. “Dynamic Reconfiguration of Primary/Backup Clusters”. Proc. of USENIX ATC. 2012, pp. 425–438.
[35] Jeff Shute et al. “F1 — The Fault-Tolerant Distributed RDBMS Supporting Google’s Ad Business”. Proc. of SIGMOD. May 2012, pp. 777–778.
[36] Yair Sovran et al. “Transactional storage for geo-replicated systems”. Proc. of SOSP. 2011, pp. 385–400.
[37] Michael Stonebraker. Why Enterprises Are Uninterested in NoSQL. 2010.
[38] Michael Stonebraker. Six SQL Urban Myths. 2010.
[39] Michael Stonebraker et al. “The end of an architectural era: (it’s time for a complete rewrite)”. Proc. of VLDB. 2007, pp. 1150–1160.
[40] Alexander Thomson et al. “Calvin: Fast Distributed Transactions for Partitioned Database Systems”. Proc. of SIGMOD. 2012, pp. 1–12.
[41] Ashish Thusoo et al. “Hive — A Petabyte Scale Data Warehouse Using Hadoop”. Proc. of ICDE. 2010, pp. 996–1005.
[42] VoltDB. VoltDB Resources. 2012.
附录A Paxos领导者租约管理
确保Paxos领导者租约区间不相交的最简单方法,是领导者在每次续期时,同步执行一次Paxos写入来记录租约区间。后续领导者读取该区间后,等待该区间结束即可。
利用TrueTime可以无需这些额外的日志写入,就能保证不相交性。候选的第i个领导者维护副本r的租约投票起始下界:\(v_{i,r}^{leader} = TT.now().earliest\),该值在\(e_{i,r}^{send}\)(领导者发送租约请求的事件)之前计算。每个副本r在\(e_{i,r}^{grant}\)时刻授予租约,该事件发生在\(e_{i,r}^{receive}\)(副本收到租约请求的事件)之后;租约结束时间为\(t_{i,r}^{end} = TT.now().latest + 10\),该值在\(e_{i,r}^{receive}\)之后计算。副本r遵守单投票规则:在TT.after(t_{i,r}^{end})为真之前,不会授予另一个租约投票。为了在副本r的不同生命周期中都强制执行该规则,Spanner在授予租约前,会在授予副本上记录租约投票日志;该日志写入可以搭载在现有的Paxos协议日志写入上。
当第i个领导者收到法定数量的投票时(事件\(e_i^{quorum}\)),它计算自身租约区间为:\(lease_i = [TT.now().latest, \min_r(v_{i,r}^{leader}) + 10]\)。当TT.before( min_r(v_{i,r}^{leader}) + 10 )为假时,领导者认为租约已过期。为了证明不相交性,我们利用如下事实:第i个与第i+1个领导者的法定投票集合中,必然存在一个共同的副本,记该副本为r0。证明:
\(
\begin{aligned}
lease_{i}.end=\min_{r}(v_{i,r}^{leader})+10\quad(\text{根据定义})
\end{aligned}
\)
\(
\min_{r}(v_{i,r}^{leader})+10\leq v_{i,r0}^{leader}+10
\)
v_{i,r0}^{leader}+10\leq t_{abs}(e_{i,r0}^{send})+10 \tag*{(最小值性质)}
\) \(
t_{abs}(e_{i,r0}^{send})+10\leq t_{abs}(e_{i,r0}^{receive})+10 \tag*{(因果关系)}
\) \(
t_{abs}(e_{i,r0}^{receive})+10\leq t_{i,r0}^{end} \tag*{(根据定义)}
\) \(
t_{i,r0}^{end}<t_{abs}(e_{i+1,r0}^{grant}) \tag*{(单投票规则)}
\) \(
t_{abs}(e_{i+1,r0}^{grant})\leq t_{abs}(e_{i+1}^{quorum}) \tag*{(因果关系)}
\) \(
t_{abs}(e_{i+1}^{quorum})\leq lease_{i+1}.start \tag*{(根据定义)}
\)