跳转到主要内容

6 复制

与可能出错的东西比,“不可能”出错的东西最显著的特点就是:一旦真的出错,通常就彻底玩完了。

—— 道格拉斯・亚当斯,《基本无害》(1992)

复制(replication)意味着在通过网络连接的多台机器上保留相同数据的副本。正如 “分布式与单节点系统” 中所讨论的,我们希望复制数据,可能出于以下原因:

  • 使数据在地理上更接近用户(从而降低访问延迟)
  • 即使系统的一部分发生故障,系统仍能继续工作(从而提高可用性)
  • 增加能够处理读查询的机器数量(从而提高读取吞吐量)

本章假设数据集足够小,每台机器都能保存整个数据集的副本。在 第 7 章 中,我们将放宽这一假设,讨论单台机器无法容纳的大型数据集如何进行 分片(sharding,也称 分区,partitioning)。再往后的章节会讨论复制数据系统中可能出现的各种故障,以及应对这些故障的方法。

如果要复制的数据不随时间变化,复制就很简单:只需把数据复制到每个节点一次,便大功告成。复制的全部难点都在于处理被复制数据的 变更,这正是本章的主题。我们将讨论三类在节点之间复制变更的算法:单主复制(single-leader replication)、多主复制(multi-leader replication)和 无主复制(leaderless replication)。几乎所有分布式数据库都采用其中一种。三者各有利弊,本章将逐一详述。

复制需要考虑许多权衡,例如采用 同步复制(synchronous replication)还是 异步复制(asynchronous replication),以及如何处理失效的副本。这些往往都是数据库的配置选项;具体细节因数据库而异,但不同实现背后的基本原理大体相通。本章将讨论这些选择带来的后果。

数据库复制算得上是老生常谈:自 20 世纪 70 年代有人开始研究以来,其基本原理并没有太大变化 1,因为网络的根本约束也一直未变。即便如此,最终一致性(eventual consistency)等概念仍常常引起困惑。在 “复制延迟的问题” 中,我们会更精确地说明最终一致性,并讨论 读己之写(read-your-writes)、单调读(monotonic reads)等保证。

备份与复制

你可能会问:有了复制,是否还需要备份?答案是肯定的,因为二者目的不同。副本会迅速把一个节点上的写入反映到其他节点,而备份保存的是数据在过去某一时刻的快照,以便恢复到先前状态。如果不慎删除了某些数据,复制帮不上忙,因为删除操作也会传播到所有副本;要恢复这些数据,仍然需要备份。

事实上,复制与备份往往相辅相成。正如 “设置新的副本” 中将要看到的,备份有时是建立复制的一环;反过来,归档复制日志也可以成为备份流程的一部分。

有些数据库会在内部维护过去状态的不可变快照,相当于一种内置备份。不过,这意味着旧版本和当前状态要保存在同一类存储介质上。数据量很大时,把旧数据的备份放在针对低频访问优化的对象存储中,往往比放在主存储中便宜;主存储只需保留数据库的当前状态。

单主复制

每个保存数据库拷贝的节点都称为一个 副本(replica)。存在多个副本时,一个问题不可避免:如何确保所有数据最终都出现在所有副本上?

数据库的每次写入都必须由每个副本处理,否则各副本就会包含不同的数据。最常见的解决方案称为 基于领导者的复制(leader-based replication),也称 主备复制(primary-backup)或 主动/被动复制(active/passive)。其工作原理如下(见 图 6-1):

  1. 其中一个副本被指定为 领导者(leader,也称 主库,primary,或 源,source 2)。客户端要写入数据库时,必须把请求发给领导者;领导者首先将新数据写入本地存储。
  2. 其他副本称为 追随者(follower,也称 只读副本,read replica,备库,secondary,或 热备,hot standby)。领导者每次把新数据写入本地存储,也会把数据变更作为 复制日志(replication log)或 变更流(change stream)发送给所有追随者。每个追随者取得领导者的日志,按照领导者处理写入的相同顺序应用所有写入,从而更新本地的数据库副本。
  3. 客户端读取数据库时,可以查询领导者,也可以查询任意追随者;但只有领导者接受写入(从客户端的角度看,追随者是只读的)。
单主复制把所有写入都发往指定的领导者,再由领导者将变更流发送给各追随者副本。
图 6-1 单主复制把所有写入都发往指定的领导者,再由领导者将变更流发送给各追随者副本。

如果数据库做了分片(见 第 7 章),每个分片都有一个领导者。不同分片的领导者可以位于不同节点,但每个分片仍必须有且只有一个领导者。在 “多主复制” 中,我们将讨论另一种模型:同一分片可以同时有多个领导者。

单主复制应用极为广泛。它是 PostgreSQL、MySQL、Oracle Data Guard 3 和 SQL Server Always On 可用性组 4 等许多关系数据库的内置功能;MongoDB、DynamoDB 5 等文档数据库,Kafka 等消息代理,DRBD 等复制块设备,以及一些网络文件系统也采用这种方式。Raft 等许多共识算法同样以单个领导者为基础;CockroachDB 6、TiDB 7、etcd、RabbitMQ 法定人数队列等系统用它来实现复制,并在原领导者失效时自动选举新领导者(第 10 章 将更详细地讨论共识)。

说明

较早的资料中可能会出现 主从复制(master–slave replication)一词。它与基于领导者的复制含义相同,但如今普遍认为这种说法具有冒犯性,应当避免使用 8。

同步复制与异步复制

复制系统的一个重要细节是,复制究竟 同步(synchronous)进行还是 异步(asynchronous)进行。(在关系数据库中,这通常是一个配置项;其他系统往往固定采用其中一种。)

设想 图 6-1 中的情形:某网站用户更新个人头像。客户端在某个时刻向领导者发出更新请求,不久后领导者收到请求;领导者随后在某个时刻把数据变更转发给追随者,并最终通知客户端更新成功。图 6-2 展示了其中一种可能的时序。

基于领导者的复制,其中一个追随者同步复制,另一个异步复制。
图 6-2 基于领导者的复制,其中一个追随者同步复制,另一个异步复制。

在 图 6-2 的例子中,发往追随者 1 的复制是 同步 的:领导者必须等到追随者 1 确认收到写入,才能向用户报告成功,也才能让其他客户端看到这次写入。发往追随者 2 的复制则是 异步 的:领导者发出消息,但不等待追随者响应。

图中追随者 2 处理消息前有一段明显的延迟。通常复制相当快:大多数数据库系统不到一秒便能把变更应用到追随者,但它们并不保证复制一定能在多久内完成。有时追随者可能落后领导者几分钟甚至更久,例如追随者正从失效中恢复、系统在容量极限附近运行,或节点间网络出现问题。

同步复制的优点是,追随者保证拥有与领导者一致的最新数据副本。领导者突然失效时,可以确信数据仍可从追随者取得。缺点是,如果同步追随者没有响应——无论因为崩溃、网络故障还是其他原因——写入就无法继续。领导者必须阻塞所有写入,直至同步副本重新可用。

因此,把所有追随者都设为同步并不现实:任意一个节点停机都会拖垮整个系统。实践中,数据库所谓的同步复制,通常是指 一个 追随者同步,其余追随者异步。如果同步追随者不可用或过慢,就把某个异步追随者切换为同步。这样可以保证至少有两个节点持有最新数据:领导者和一个同步追随者。这种配置有时也称为 半同步(semi-synchronous)。

有些系统会同步更新 多数 副本(例如含领导者在内的 5 个副本中更新 3 个),其余少数副本异步更新。这就是 法定人数 的一个例子,我们会在 “读写仲裁” 中进一步讨论。采用共识协议自动选举领导者的系统经常使用多数法定人数,第 10 章 将再次谈到这个问题。

有时,基于领导者的复制会配置为完全异步。如果领导者失效且无法恢复,所有尚未复制到追随者的写入都会丢失。也就是说,即使已经向客户端确认成功,写入仍不保证持久。不过,完全异步配置也有一个优点:即使所有追随者都已落后,领导者仍可继续处理写入。

削弱持久性听起来不像是划算的取舍,但异步复制依然应用广泛,尤其是在追随者很多或分布于不同地理位置时 9。我们会在 “复制延迟的问题” 中再谈这一点。

设置新的副本

有时需要设置新的追随者,也许是为了增加副本数量,也许是为了替换失效的节点。怎样才能确保新追随者拿到领导者数据的准确副本?

简单地把数据文件从一个节点复制到另一个节点通常并不够:客户端一直在向数据库写入,数据始终处于变化之中,普通的文件复制会在不同时间点读到数据库的不同部分,所得结果可能毫无意义。

可以锁住数据库,让磁盘上的文件保持一致,但这样数据库将无法接受写入,违背了高可用的目标。好在设置追随者通常不需要停机。从概念上讲,其过程如下:

  1. 取得领导者数据库在某个时刻的一致快照;如果可能,不要锁住整个数据库。大多数数据库都提供这一功能,因为备份同样需要它。有些情况下需要借助第三方工具,例如 MySQL 的 Percona XtraBackup。
  2. 把快照复制到新的追随者节点。
  3. 追随者连接领导者,请求从快照生成之后发生的所有数据变更。这要求快照与领导者复制日志中的准确位置相关联。不同系统对这个位置有不同称呼:PostgreSQL 称之为 日志序列号(log sequence number,LSN);MySQL 则有 binlog 位点(binlog coordinates)和 全局事务标识符(global transaction identifiers,GTIDs)两套机制。
  4. 追随者处理完快照之后积压的数据变更时,就称它已经 赶上进度。此后,它可以继续随时处理领导者产生的数据变更。

设置追随者的实际步骤因数据库而异。有些系统完全自动完成这一过程,另一些系统则需要管理员手工执行一套颇为晦涩的多步骤流程。

还可以把复制日志归档到对象存储;再定期把整个数据库的快照保存到对象存储,就形成了一套很好的数据库备份与灾难恢复方案。建立新追随者时,步骤 1 和 2 也可以通过从对象存储下载这些文件来完成。例如,WAL-G 为 PostgreSQL、MySQL 和 SQL Server 提供了这种功能,Litestream 则为 SQLite 提供了类似功能。

以对象存储为后端的数据库

对象存储不只能用于归档。许多数据库已经开始用 Amazon Web Services S3、Google Cloud Storage、Azure Blob Storage 等对象存储为在线查询提供数据。把数据库数据存入对象存储有很多好处:

  • 对象存储比其他云存储方案便宜。云数据库因而可以把查询频率较低的数据放到成本更低、延迟更高的存储中,同时用内存、SSD 和 NVMe 保存工作集。
  • 对象存储还提供多可用区、双地区或多地区复制,并给出很高的持久性保证;数据库也由此可以避开跨可用区网络费用。
  • 数据库可以利用对象存储的 条件写入 功能——本质上是一种 比较并设置(CAS)操作——来实现事务和领导者选举 10 11。
  • 把多个数据库的数据放在同一对象存储中,可以简化数据集成,尤其是在使用 Apache Parquet、Apache Iceberg 等开放格式时。

这些好处把事务、领导者选举和复制的责任转交给对象存储,从而大幅简化数据库架构。

以对象存储实现复制的系统也必须面对一些权衡。尤其是,对象存储的读写延迟远高于本地磁盘或 EBS 之类的虚拟块设备。许多云服务商还按 API 调用次数收费,迫使系统把读写合并成批次以降低成本,而批处理又会进一步增加延迟。此外,许多对象存储没有标准的文件系统接口,未集成对象存储的系统便无法直接利用它。用户空间文件系统(FUSE)等接口允许运维人员把对象存储桶挂载成文件系统,应用程序不必知道数据实际位于对象存储中。然而,许多面向对象存储的 FUSE 接口并不支持非顺序写入、符号链接等 POSIX 功能,而系统可能依赖这些功能。

不同系统处理这些权衡的方式各不相同。有些采用 分层存储 架构,把不常访问的数据放到对象存储,把新数据或常用数据放在 SSD、NVMe 乃至内存等更快的存储介质上。另一些系统以对象存储作为主存储层,但另用 Amazon EBS、Neon Safekeepers 12 等低延迟存储系统保存 WAL。近来,一些系统走得更远,采用 零磁盘架构(ZDA):所有数据都持久化到对象存储,磁盘和内存只用作缓存。节点因此无需保存持久状态,运维也大为简化。WarpStream、Confluent Freight、Buf 的 Bufstream 和 Redpanda Serverless 都是采用零磁盘架构构建的 Kafka 兼容系统。几乎所有现代云数据仓库也采用类似架构,向量搜索引擎 Turbopuffer 和云原生 LSM 存储引擎 SlateDB 亦是如此。

处理节点故障

系统中的任何节点都可能停机,既可能是故障意外所致,也可能是计划内维护,例如重启机器以安装内核安全补丁。能够在服务不中断的情况下逐个重启节点,对运维和维护大有裨益。因此,我们的目标是:即使个别节点失效,整个系统仍能继续运行,并尽可能减小节点停机的影响。

如何用基于领导者的复制实现高可用?

追随者失效:追赶恢复

每个追随者都会在本地磁盘上记录从领导者收到的数据变更。如果追随者崩溃后重启,或领导者与追随者之间的网络暂时中断,恢复起来相对容易:追随者可以从日志中得知故障发生前处理的最后一个事务。它随后连接领导者,请求自己断开期间发生的所有数据变更。应用完这些变更后,它便赶上领导者,可以像以前一样继续接收数据变更流。

追随者恢复在概念上很简单,性能上却可能很棘手。如果数据库写入吞吐量很高,或追随者离线很久,需要追赶的写入可能非常多。追赶期间,正在恢复的追随者和领导者都会承受很高负载——领导者还要把积压的写入发送给追随者。

所有追随者确认处理完某段写入日志后,领导者便可将其删除。但如果某个追随者长时间不可用,领导者必须作出选择:要么一直保留日志,等追随者恢复并赶上进度,但要冒领导者磁盘空间耗尽的风险;要么删除尚未得到该追随者确认的日志,这样追随者恢复上线后就无法靠日志追赶,只能从备份恢复。

领导者失效:故障切换

领导者失效处理起来更加棘手:必须把一个追随者提升为新领导者,重新配置客户端以便把写入发给新领导者,并让其他追随者开始消费新领导者的数据变更。这个过程称为 故障切换(failover)。

故障切换可以手工完成——通知管理员领导者已经失效,再由管理员执行必要步骤选出新领导者;也可以自动完成。自动故障切换通常包括以下步骤:

  1. 确认领导者失效。 可能出问题的地方很多:崩溃、断电、网络故障等等。没有万无一失的办法判断究竟出了什么问题,因此大多数系统只是使用超时:节点之间频繁往返传递消息,如果某个节点在一段时间内(例如 30 秒)没有响应,就认为它已经失效。(因计划维护而主动下线领导者则不适用这一规则。)
  2. 选择新的领导者。 可以通过选举完成(由剩余副本中的多数选出领导者),也可以由事先指定的 控制器节点 任命 13。最合适的候选者通常是从旧领导者收到数据变更最多、数据最新的副本,这样可以尽量减少数据损失。让所有节点就新领导者达成一致是一个共识问题,第 10 章 将详细讨论。
  3. 重新配置系统以使用新领导者。 客户端现在必须把写请求发给新领导者(参见 “请求路由”)。如果旧领导者恢复上线,它可能仍以为自己是领导者,没有意识到其他副本已经迫使它退位。系统必须保证旧领导者转为追随者,并承认新领导者。

故障切换过程中有很多地方可能出错:

  • 如果使用异步复制,新领导者可能没有收到旧领导者失效前的全部写入。选出新领导者后,原领导者若重新加入集群,那些未复制的写入该怎么办?与此同时,新领导者可能已经收到与之冲突的写入。最常见的办法是直接丢弃旧领导者尚未复制的写入,这意味着你原以为已经提交的写入其实并未持久保存。
  • 如果数据库内容还需要与数据库之外的存储系统协调,丢弃写入尤其危险。例如,GitHub 曾发生过一起事故 14:一个数据过时的 MySQL 追随者被提升为领导者。数据库用自增计数器为新行分配主键;新领导者的计数器落后于旧领导者,因而重复使用了旧领导者已经分配过的一些主键。这些主键同时用于 Redis 存储,主键重用造成 MySQL 与 Redis 数据不一致,最终使一些私有数据泄露给了错误的用户。
  • 在某些故障场景下(见 第 9 章),可能有两个节点都认为自己是领导者。这种情况称为 脑裂(split brain),非常危险:如果两个领导者都接受写入,而系统又没有冲突解决流程(参见 “多主复制”),数据很可能丢失或损坏。有些系统设有保险机制,一旦发现两个领导者便关闭其中一个节点;但机制设计不当时,也可能把两个节点都关闭 15。而且,等系统发现脑裂并关闭旧节点时,可能已经为时过晚,数据早已损坏。
  • 宣布领导者失效前,超时应该设为多长?超时越长,领导者确实失效时恢复所需的时间就越长;超时太短,又容易触发不必要的故障切换。例如,短暂的负载尖峰可能使节点响应时间超过超时值,网络抖动也可能延迟数据包。如果系统已经饱受高负载或网络问题困扰,不必要的故障切换只会让情况更糟。
说明

通过限制或关闭旧领导者来防止脑裂,称为 栅栏(fencing),更形象的说法是 爆彼之头(Shoot The Other Node In The Head,STONITH)。“分布式锁和租约” 将更详细地讨论栅栏机制。

这些问题没有简单的解决方案。因此,即使软件支持自动故障切换,一些运维团队还是更愿意手工执行。

故障切换最重要的是选出一个数据最新的追随者作为新领导者。采用同步或半同步复制时,应选择旧领导者确认写入前所等待的那个追随者;采用异步复制时,则可以选择日志序列号最大的追随者。这样能尽量减少故障切换时的数据损失:丢失几分之一秒内的写入或许尚可容忍,选中一个落后数天的追随者却可能是灾难性的。

节点失效、不可靠的网络,以及围绕副本一致性、持久性、可用性和延迟所作的权衡,都是分布式系统中的根本问题。第 9 章 和 第 10 章 将进一步深入讨论。

复制日志的实现

基于领导者的复制在底层究竟如何工作?实践中采用了几种不同的复制方式,下面逐一简要介绍。

基于语句的复制

最简单的情况下,领导者会记录自己执行的每个写入请求(即 语句),并把这份语句日志发送给追随者。对于关系数据库,这意味着每条 INSERT、UPDATE 或 DELETE 语句都会转发给追随者;每个追随者解析并执行这条 SQL 语句,就像它直接来自客户端一样。

虽然听起来很合理,但这种复制方式有许多可能出错的地方:

  • 任何调用非确定性函数的语句,都可能在各副本上产生不同的值,例如用 NOW() 取得当前日期和时间,或用 RAND() 取得随机数。
  • 如果语句使用自增列,或依赖数据库中的现有数据(例如 UPDATE …​ WHERE <some condition>),就必须在每个副本上按完全相同的顺序执行,否则可能产生不同结果。有多个事务并发执行时,这会成为限制。
  • 带有副作用的语句(例如触发器、存储过程、用户定义函数),可能在各副本上产生不同的副作用,除非这些副作用完全确定。

这些问题可以绕开。例如,领导者在记录语句时,可以用固定的返回值替换非确定性函数调用,从而让所有追随者得到相同的值。按固定顺序执行确定性语句,与 “事件溯源与 CQRS” 中介绍的事件溯源模型很相似。这种方法也称为 状态机复制(state machine replication);“使用共享日志” 将讨论其理论基础。

MySQL 5.1 以前使用基于语句的复制。由于日志相当紧凑,如今有时仍会采用这种方式;不过默认情况下,只要语句中存在任何非确定性,MySQL 就会切换到稍后介绍的基于行的复制。VoltDB 也采用基于语句的复制,并要求事务必须是确定性的,以保证安全 16。然而,实践中很难确保确定性,因此许多数据库更倾向于其他复制方式。

预写日志(WAL)传输

我们在 第 4 章 中看到,B 树存储引擎需要预写日志才能可靠工作:每次修改都先写入 WAL,以便崩溃后把树恢复到一致状态。WAL 包含把索引和堆恢复到一致状态所需的全部信息,因此同一份日志也能用来在另一个节点上建立副本:领导者除了把日志写入磁盘,还通过网络把它发送给追随者。追随者处理日志后,就会构建出与领导者完全相同的文件副本。

PostgreSQL、Oracle 等数据库采用这种复制方式 17 18。其主要缺点是,日志在非常低的层次描述数据:WAL 记录了哪个磁盘块中的哪些字节发生变化。因此,复制与存储引擎紧密耦合。数据库从一个版本升级到另一个版本、存储格式随之改变时,通常无法在领导者和追随者上运行不同版本的数据库软件。

这看起来只是一个微不足道的实现细节,却可能给运维带来巨大影响。如果复制协议允许追随者运行比领导者更新的软件版本,就可以先升级追随者,再执行故障切换,让一个升级后的节点成为新领导者,从而实现数据库软件的零停机升级。如果复制协议不允许版本不一致——WAL 传输往往如此——这类升级就需要停机。

逻辑(基于行)日志复制

另一种方法是让复制和存储引擎使用不同的日志格式,从而使复制日志与存储引擎的内部实现解耦。这种复制日志称为 逻辑日志(logical log),以区别于存储引擎的(物理,physical)数据表示。

关系数据库的逻辑日志通常由一系列记录组成,以行的粒度描述对数据库表的写入:

  • 插入一行时,日志包含所有列的新值。
  • 删除一行时,日志包含足以唯一标识被删行的信息,通常是主键;如果表没有主键,则需要记录所有列的旧值。
  • 更新一行时,日志包含足以唯一标识被更新行的信息,以及所有列的新值(或者至少包含所有发生变化的列的新值)。

修改多行的事务会生成多条这样的日志记录,随后再跟一条表示事务已经提交的记录。MySQL 配置为基于行的复制时,除了 WAL 之外,还会维护一份称为 binlog 的独立逻辑复制日志。PostgreSQL 则把物理 WAL 解码成行插入、更新和删除事件,以实现逻辑复制 19。

由于逻辑日志与存储引擎的内部实现解耦,更容易保持向后兼容,因而领导者与追随者可以运行不同版本的数据库软件,也就能以极少的停机时间升级到新版本 20。

逻辑日志格式也更容易由外部应用程序解析。如果要把数据库内容发送到外部系统,例如送入数据仓库做离线分析,或构建自定义索引和缓存 21,这一点会很有用。这种技术称为 变更数据捕获(change data capture,CDC),我们会在 “变更数据捕获” 中再次谈到它。

复制延迟的问题

容忍节点失效只是使用复制的一个原因。正如 “分布式与单节点系统” 中提到的,其他原因还包括可伸缩性(处理单台机器无法承受的请求量)和延迟(把副本放在地理上更靠近用户的位置)。

基于领导者的复制要求所有写入经过一个节点,但只读查询可以发往任意副本。对于读多写少的工作负载——在线服务往往如此——有一种很有吸引力的方案:建立许多追随者,把读请求分散到这些追随者上。这样既能减轻领导者的负载,也能由附近的副本处理读请求。

在这种 读扩展 架构中,只需增加追随者,便可提高只读请求的处理能力。不过,这种办法实际上只适用于异步复制。如果尝试同步复制到所有追随者,任意一个节点失效或网络中断都会让整个系统无法写入。节点越多,越可能有某个节点停机,因此完全同步的配置会极不可靠。

不幸的是,应用程序从 异步(asynchronous)追随者读取时,如果追随者落后,就可能看到过时的信息。数据库于是显得不一致:同时在领导者和追随者上执行同一查询,结果可能不同,因为追随者尚未反映所有写入。这种不一致只是暂时的——如果停止写入并等待一段时间,追随者最终会赶上领导者,恢复一致。因此,这种现象称为 最终一致性(eventual consistency)22。

说明

最终一致性 一词由 Douglas Terry 等人提出 23,经 Werner Vogels 推广 24,后来成了许多 NoSQL 项目的口号。不过,最终一致性并非 NoSQL 数据库独有:关系数据库采用异步复制时,其追随者也有相同特性。

“最终”一词有意保持模糊:一般而言,副本能落后多久并没有上限。正常运行时,从写入发生在领导者上,到变更反映在追随者上,两者之间的延迟——即 复制延迟——可能只有几分之一秒,在实践中难以察觉。但当系统在容量极限附近运行或网络出现问题时,延迟很容易增至几秒甚至几分钟。

延迟一旦大到这种程度,由此造成的不一致就不再只是理论问题,而会成为应用程序面临的真实问题。本节将重点介绍复制延迟容易引发的三类问题,并概述一些解决办法。

读己之写

许多应用程序允许用户提交数据,随后查看自己提交的内容。它可能是客户数据库中的一条记录、讨论主题下的一条评论,或其他类似内容。新数据必须写入领导者,但用户查看时可以从追随者读取。如果数据经常读取、很少写入,这种做法尤其合适。

异步复制在这里会产生问题,如 图 6-3 所示:用户刚写入数据不久便查看它时,新数据可能尚未到达该副本。在用户看来,自己刚刚提交的数据仿佛丢失了,当然会感到不满。

用户写入后,又从陈旧副本读取。要防止这种异常,需要写后读一致性。
图 6-3 用户写入后,又从陈旧副本读取。要防止这种异常,需要写后读一致性。

这种情况下需要 写后读一致性(read-after-write consistency),也称 读己之写一致性(read-your-writes consistency)23。它保证用户重新加载页面时,总能看到自己提交的更新。至于其他用户则不作保证:他们的更新可能过一段时间才会出现。但至少用户可以确信,自己的输入已经正确保存。

如何在基于领导者的复制系统中实现写后读一致性?有多种办法,例如:

  • 读取用户可能修改过的内容时,从领导者或同步更新的追随者读取;其他内容则从异步更新的追随者读取。这要求系统无需实际查询,就能判断某项内容是否可能被修改。例如,社交网络中的个人资料通常只能由本人编辑。因此可以制定一条简单规则:用户自己的资料总是从领导者读取,其他用户的资料则从追随者读取。
  • 如果应用程序中的大部分内容都可能由用户编辑,上述办法就没有效果,因为几乎所有内容都得从领导者读取,读扩展也就失去了意义。这时可以采用其他标准来决定是否从领导者读取。例如,记录最近一次更新时间,此后一分钟内的所有读取都发往领导者 25。也可以监控追随者的复制延迟,不向落后领导者超过一分钟的追随者发送查询。
  • 客户端可以记住自己最近一次写入的时间戳,系统再确保为该用户处理读取的副本至少已经反映到这个时间戳。如果副本不够新,就换一个副本处理读取,或让查询等待该副本赶上进度 26。这个时间戳既可以是 逻辑时间戳(例如表示写入顺序的日志序列号),也可以来自实际的系统时钟;后一种情况下,时钟必须准确同步,参见 “不可靠的时钟”。
  • 如果副本分布于多个地区(为了靠近用户或提高可用性),还会增加一层复杂性:凡是必须由领导者处理的请求,都要路由到领导者所在的地区。

同一用户通过多个设备访问服务时,例如同时使用桌面浏览器和移动应用,还会出现另一种复杂情况。这时可能需要提供 跨设备 写后读一致性:用户在一台设备上输入信息,随后在另一台设备上查看时,应当看到刚才输入的内容。

还需要考虑以下问题:

  • 依赖“记住用户最近一次更新时间戳”的方法会变得更难,因为一台设备上运行的代码不知道另一台设备做了哪些更新。这些元数据需要集中保存。
  • 如果副本分布在不同地区,来自不同设备的连接不保证会路由到同一地区。例如,用户的台式机使用家庭宽带,手机使用蜂窝网络,两台设备的网络路径可能完全不同。如果方案要求从领导者读取,可能首先要把该用户所有设备发出的请求都路由到同一地区。
地区与可用区

本书用 地区(region)表示同一地理位置上的一个或多个数据中心。云服务商通常会在同一地理区域设置多个数据中心,每个数据中心称为一个 可用区(availability zone),简称 区(zone)。因此,一个云地区由多个可用区组成。每个可用区都是位于独立物理设施中的数据中心,有自己的供电、制冷等基础设施。

同一地区内的可用区之间由高速网络连接,延迟足够低,因此大多数分布式系统可以把节点分散在多个可用区运行,就像它们位于同一个区一样。多可用区配置可以抵御某个可用区整体下线,却无法抵御一个地区内所有可用区均不可用的地区级中断。要在地区级中断时继续运行,分布式系统必须跨多个地区部署,而这可能带来更高延迟、更低吞吐量和更高的云网络费用。我们将在 “多主复制拓扑” 中进一步讨论这些权衡。这里暂且记住:本书所说的地区,是同一地理位置上的一组可用区或数据中心。

单调读

从异步追随者读取时可能发生的第二种异常,是用户可能会看到 时光倒流。

用户连续几次从不同副本读取时,就可能发生这种情况。例如,图 6-4 中,用户 2345 连续执行两次相同查询:第一次查询复制延迟较小的追随者,第二次查询延迟更大的追随者。(用户刷新网页、而每个请求被随机路由到不同服务器时,这种场景很容易出现。)第一次查询返回了用户 1234 最近添加的评论,第二次却什么也没返回,因为落后的追随者还没有收到这次写入。实际上,第二次查询观察到的系统状态比第一次更早。如果第一次查询本就没有结果,倒也问题不大,因为用户 2345 多半不知道用户 1234 刚刚添加过评论;但如果评论先出现又消失,就非常令人困惑。

用户先从较新的副本读取,随后从陈旧副本读取,时间仿佛倒退了。要防止这种异常,需要单调读。
图 6-4 用户先从较新的副本读取,随后从陈旧副本读取,时间仿佛倒退了。要防止这种异常,需要单调读。

单调读 22 保证这种异常不会发生。它弱于强一致性,却强于最终一致性。读取数据时仍可能看到旧值;单调读只保证同一用户顺序执行多次读取时,不会看到时间倒退——一旦读到较新的数据,以后就不会再读到更旧的数据。

实现单调读的一种办法,是确保每个用户始终从同一个副本读取(不同用户可以选择不同副本)。例如,可以根据用户 ID 的哈希选择副本,而不是随机选择。如果该副本失效,则需要把用户的查询重新路由到其他副本。

一致前缀读

第三种复制延迟异常违反了因果关系。设想 Poons 先生和 Cake 夫人有下面这段简短对话:

Poons 先生
Cake 夫人,你能看到多远的未来?
Cake 夫人
通常大约十秒钟,Poons 先生。

这两句话之间存在因果依赖:Cake 夫人听到 Poons 先生的问题,然后作出回答。

现在设想第三个人通过追随者旁听这段对话。Cake 夫人的话经过复制延迟较小的追随者,Poons 先生的话经过延迟更大的追随者(见 图 6-5)。于是,这位旁听者会听到:

Cake 夫人
通常大约十秒钟,Poons 先生。
Poons 先生
Cake 夫人,你能看到多远的未来?

在旁听者看来,Cake 夫人还没听到 Poons 先生提问,就已经回答了问题。这种通灵能力令人印象深刻,却也十分费解 27。

如果某些分片的复制速度慢于其他分片,观察者可能先看到答案,后看到问题。
图 6-5 如果某些分片的复制速度慢于其他分片,观察者可能先看到答案,后看到问题。

防止这种异常需要另一种保证:一致前缀读(consistent prefix reads)22。它保证,如果一系列写入按某个顺序发生,那么任何人读取这些写入时,也会看到它们以相同顺序出现。

这在分片(分区)数据库中尤其容易成为问题,我们将在 第 7 章 讨论这类数据库。如果数据库始终按相同顺序应用写入,读取就总会看到一致前缀,这种异常也不会发生。然而,在许多分布式数据库中,不同分片彼此独立运行,没有全局写入顺序。用户读取数据库时,可能看到一部分处于较旧状态,另一部分却处于较新状态。

一种解决办法是确保存在因果关系的写入落在同一个分片,但有些应用程序无法高效做到这一点。还有一些算法会显式跟踪因果依赖,我们会在 “先发生”关系与并发 中再次讨论。

复制延迟的解决方案

使用最终一致的系统时,值得认真考虑:如果复制延迟增加到几分钟甚至几小时,应用程序会怎样表现?如果答案是“没有影响”,那当然很好;如果会给用户带来糟糕体验,就必须把系统设计成能够提供写后读之类的更强保证。复制本来是异步的,却假装它是同步的,迟早会酿成问题。

如前所述,应用程序可以提供比底层数据库更强的保证,例如把某些读取发往领导者或同步更新的追随者。不过,在应用程序代码中处理这些问题既复杂又容易出错。

对应用程序开发者而言,最简单的编程模型,是选择一个能为副本提供强一致性保证(例如线性一致性,见 第 10 章)和 ACID 事务(见 第 8 章)的数据库。这样就可以基本忽略复制带来的挑战,把数据库看作只有一个节点。2010 年代初兴起的 NoSQL 运动曾宣扬一种观点:这些特性会限制可伸缩性,大规模系统不得不接受最终一致性。

此后,许多数据库开始在提供强一致性和事务的同时,保留分布式数据库在容错、高可用和可伸缩性方面的优势。正如 “关系模型与文档模型” 中提到的,为与 NoSQL 区分,这一趋势称为 NewSQL(尽管重点并不在 SQL 本身,而在可伸缩事务管理的新方法)。

如今虽已有可伸缩的强一致分布式数据库,一些应用程序仍有充分理由选择一致性保证较弱的其他复制方式:面对网络中断时,它们的韧性可能更强,开销也低于事务型系统。本章余下部分将继续探讨这些方法。

多主复制

到目前为止,本章只讨论了使用单个领导者的复制架构。虽然这是常见做法,但还有一些值得关注的选择。

单主复制有一个主要缺点:所有写入都必须经过唯一的领导者。无论出于什么原因,只要连接不上领导者——例如客户端与领导者之间的网络中断——就无法写入数据库。

单主复制模型可以自然地扩展为允许多个节点接受写入。复制仍以同样方式进行:每个处理写入的节点都必须把数据变更转发给其他所有节点。我们把这种配置称为 多主复制(multi-leader replication),也称 主动/主动复制(active/active)或 双向复制(bidirectional replication)。在这种配置中,每个领导者同时也是其他领导者的追随者。

与单主复制一样,多主复制也可以选择同步或异步。假设有两个领导者 A 和 B,现在要向 A 写入。如果写入必须从 A 同步复制到 B,那么两者之间的网络一旦中断,在网络恢复前就无法向 A 写入。这种同步多主复制提供的模型与单主复制极为相似:也就是说,这等同于把 B 设为领导者,由 A 把所有写请求转发给 B 执行。

因此,本节不再深入讨论同步多主复制,而把它视为等同于单主复制。以下讨论集中于异步多主复制:即使某个领导者与其他领导者之间的连接中断,它仍可处理写入。

跨地域运行

在单个地区内使用多主配置通常没有多少意义,因为所得好处很少能抵消额外的复杂性。不过,在某些场景中,这种配置确实合理。

设想一个数据库在多个地区都有副本,也许是为了在整个地区失效时仍能运行,也许是为了在地理上更接近用户。这种部署称为 地理分布式(geographically distributed)、跨地域分布式(geo-distributed)或 跨地域复制(geo-replicated)。采用单主复制时,领导者必须位于其中 一个 地区,所有写入都要经过该地区。

在多主配置中,每个 地区都可以有一个领导者。图 6-6 展示了这种架构:每个地区内部使用常规的领导者—追随者复制(追随者可以位于与领导者不同的可用区);地区之间,则由各地区的领导者把变更复制给其他地区的领导者。

跨多个地区的多主复制。
图 6-6 跨多个地区的多主复制。

下面比较单主与多主配置在多地区部署中的表现:

性能
在单主配置中,每次写入都必须通过互联网发往领导者所在的地区。这可能显著增加写入延迟,甚至违背多地区部署的初衷。在多主配置中,每次写入都能在本地地区处理,再异步复制到其他地区。地区间的网络延迟因而对用户不可见,感知到的性能可能更好。
容忍地区停机
在单主配置中,如果领导者所在的地区不可用,可以通过故障切换把另一个地区的追随者提升为领导者。在多主配置中,各地区可以彼此独立地继续运行;离线地区恢复上线后,复制会赶上进度。
容忍网络问题
即使地区之间使用专用连接,其流量也可能不如同一地区内不同可用区之间、或同一可用区内部的流量可靠。单主配置对地区间链路的问题极为敏感:一个地区的客户端要向另一个地区的领导者写入,就必须通过这条链路发出请求,并等待响应后才能完成操作。

异步复制的多主配置更能容忍网络问题:网络暂时中断期间,每个地区的领导者仍可独立处理写入。

一致性
单主系统可以提供可串行化事务等强一致性保证,我们会在 第 8 章 讨论。多主系统最大的缺点,是它所能提供的一致性要弱得多。例如,无法保证银行账户余额不会变成负数,也无法保证用户名唯一:不同领导者完全可能分别处理单独看来合法的写入(从账户支付一笔钱、注册某个用户名),但把它们与另一个领导者上的写入合在一起,就违反了约束。

这是分布式系统的一项根本限制 28。如果必须强制执行这类约束,最好采用单主系统。不过,正如 “处理写入冲突” 中将要看到的,对于不需要这类约束的大量应用程序,多主系统仍能提供有用的一致性属性。

多主复制不如单主复制常见,但 MySQL、Oracle、SQL Server、YugabyteDB 等许多数据库仍提供支持。有时它以外部附加组件的形式出现,例如 Redis Enterprise、EDB Postgres Distributed 和 pglogical 29。

许多数据库的多主复制都是后来加装的功能,因而常有细微的配置陷阱,也会与其他数据库功能产生出人意料的相互作用。例如,自增键、触发器和完整性约束都可能带来问题。因此,多主复制往往被视为应当尽量避开的危险领域 30。

多主复制拓扑

复制拓扑(replication topology)描述写入从一个节点传播到另一个节点时所经过的通信路径。如果只有两个领导者,如 图 6-9 所示,那么只有一种合理的拓扑:领导者 1 必须把所有写入发送给领导者 2,反之亦然。领导者超过两个时,则有多种拓扑可供选择。图 6-7 给出了三个例子。

多主复制可以采用的三种拓扑示例。
图 6-7 多主复制可以采用的三种拓扑示例。

最通用的是 图 6-7(c) 所示的 全对全(all-to-all)拓扑,每个领导者都会把写入发送给其他所有领导者。不过,实际系统也会采用限制更多的拓扑。例如在 环形拓扑(circular topology)中,每个节点从一个节点接收写入,再把这些写入连同自己的写入转发给另一个节点。另一种常见拓扑呈 星形(star topology):指定一个根节点,由它把写入转发给所有其他节点。星形拓扑还可以推广成树形。

说明

不要把星形网络拓扑与 星型模式 混为一谈;后者描述的是数据模型结构,参见 “星型与雪花型:分析模式”。

在环形和星形拓扑中,一次写入可能要经过多个节点才能到达所有副本。因此,节点必须转发从其他节点收到的数据变更。为避免无限复制循环,每个节点都有唯一标识符;复制日志中的每次写入都会标记自己经过的全部节点 31。节点收到带有自身标识符的数据变更时,会直接忽略它,因为这说明自己已经处理过该变更。

不同拓扑的问题

环形和星形拓扑有一个问题:只要一个节点失效,就可能中断其他节点之间的复制消息流;在该节点修复前,其他节点无法相互通信。虽然可以重配拓扑以绕过失效节点,但大多数部署都需要手工完成这一操作。连接更密集的拓扑(例如全对全)容错性更好,因为消息可以沿不同路径传播,避开单点故障。

另一方面,全对全拓扑也有问题。尤其是,不同网络链路的速度可能不同(例如受到网络拥塞影响),导致某些复制消息“超越”另一些消息,如 图 6-8 所示。

在多主复制中,写入抵达某些副本的顺序可能有误。
图 6-8 在多主复制中,写入抵达某些副本的顺序可能有误。

在 图 6-8 中,客户端 A 在领导者 1 上向表中插入一行,客户端 B 随后在领导者 3 上更新该行。但领导者 2 可能以相反顺序收到这两次写入:先收到更新(在它看来,这是要更新数据库中并不存在的行),稍后才收到本应先发生的插入。

这也是一个因果关系问题,与 “一致前缀读” 中看到的情况相似。更新依赖先前的插入,因此必须保证所有节点先处理插入,再处理更新。仅仅给每次写入附加时间戳并不够,因为不能相信各节点的时钟同步得足以让领导者 2 正确排列这些事件(见 第 9 章)。

要正确排列这些事件,可以采用本章稍后介绍的 版本向量(version vector;见 “检测并发写入”)。不过,许多多主复制系统并未使用可靠的更新排序技术,因而容易遇到 图 6-8 所示的问题。如果使用多主复制,应当了解这些风险,仔细阅读文档,并充分测试数据库,确认它确实提供了你以为它会提供的保证。

同步引擎与本地优先软件

应用程序需要在断网时继续工作,是另一个适合多主复制的场景。

以手机、笔记本电脑和其他设备上的日历应用为例。无论设备有没有联网,你都需要随时查看会议(发出读请求)和添加会议(发出写请求)。离线期间所作的变更,应在设备下次上线时与服务器及其他设备同步。

这种情况下,每台设备都有一个充当领导者的本地数据库副本,可以接受写入;各设备上的日历副本之间则通过异步多主复制过程进行同步。复制延迟可能长达数小时甚至数天,取决于设备何时重新联网。

从架构上看,这种配置相当于把地区间多主复制推到极致:每台设备都是一个“地区”,它们之间的网络连接极不可靠。

实时协作、离线优先和本地优先应用

此外,许多现代 Web 应用还提供 实时协作 功能,例如用于文档和电子表格的 Google Docs 与 Sheets、用于图形设计的 Figma,以及用于项目管理的 Linear。这些应用之所以响应迅速,是因为用户输入会立即反映在界面上,无需等待与服务器的一次网络往返;一位用户所作的编辑也会以很低的延迟呈现给协作者 32 33 34。

这同样形成了多主架构:每个打开共享文件的浏览器标签页都是一个副本,对文件所作的更新会异步复制到其他打开该文件的用户设备上。即使应用程序不支持离线编辑,只要多个用户可以不等服务器响应便各自编辑,它就已经是多主系统。

离线编辑与实时协作需要相似的复制基础设施:应用程序必须捕获用户对文件作出的所有变更,在线时立即发给协作者,离线时则先保存在本地,稍后再发送。同时,应用程序还要接收协作者的变更,将其合并到用户的本地文件副本,并更新界面以显示最新版本。多个用户并发修改文件时,还可能需要用冲突解决逻辑合并这些变更。

支持这一过程的软件库称为 同步引擎(sync engine)。这个想法由来已久,但“同步引擎”一词近来才受到关注 35 36 37。允许用户离线时继续编辑文件的应用程序称为 离线优先(offline-first)应用 38,它可以用同步引擎来实现。本地优先软件(local-first software)则不仅要支持离线优先,还要保证即使软件开发者关闭所有在线服务,协作应用仍能继续工作 39。一种实现方式是采用开放标准的同步协议,并让多个服务提供商都能支持这一协议 40。例如,Git 就是本地优先的协作系统(虽然它不支持实时协作),因为可以通过 GitHub、GitLab 或其他任意代码仓库托管服务进行同步。

同步引擎的利弊

如今构建 Web 应用的主流方式,是让客户端只保留极少的持久状态;每当需要显示新数据或更新数据时,就向服务器发出请求。使用同步引擎时则相反:客户端持有持久状态,与服务器的通信移到后台进行。这种方式有多项优点:

  • 数据在本地,用户界面的响应速度可以远快于等待服务调用返回数据。有些应用追求在图形系统的 下一帧 响应用户输入:对于刷新率为 60 Hz 的显示器,这意味着要在 16 毫秒内完成渲染。
  • 允许用户离线工作很有价值,尤其是在连接时断时续的移动设备上。使用同步引擎后,应用程序无需另设离线模式:离线不过是网络延迟变得非常大。
  • 与在应用代码中显式调用服务相比,同步引擎简化了前端应用的编程模型。正如 “远程过程调用(RPC)的问题” 中所述,每次服务调用都要处理错误。例如,更新服务器数据的请求失败后,用户界面必须以某种方式反映错误。同步引擎让应用直接读写几乎不会失败的本地数据,从而形成更具声明性的编程风格 41。
  • 要实时显示其他用户所作的编辑,需要接收变更通知,并据此高效更新用户界面。同步引擎与 响应式编程(reactive programming)模型结合,是实现这一功能的好办法 42。

如果能事先下载用户可能需要的全部数据,并持久保存在客户端,同步引擎的效果最好。这样一来,需要时就能离线访问;但也意味着,如果用户可以访问的数据量非常大,同步引擎便不适用。例如,下载用户自己创建的全部文件通常没问题(单个用户一般不会产生那么多数据),下载一个电子商务网站的全部商品目录则多半不合理。

Lotus Notes 在 20 世纪 80 年代率先采用了同步引擎的思想 43,尽管当时并未使用这个名称;日历等特定应用的同步功能也已存在多年。如今有不少通用同步引擎,其中一些依赖专有后端服务,例如 Google Firestore、Realm 或 Ditto;另一些提供开源后端,适合构建本地优先软件,例如 PouchDB/CouchDB、Automerge 或 Yjs。

多人视频游戏也需要立即响应玩家的本地操作,再与通过网络异步收到的其他玩家操作协调。在游戏开发术语中,与同步引擎对应的部分称为 网络代码(netcode)。网络代码所用的技术针对游戏需求高度定制 44,无法直接移植到其他软件,因此本书不再展开。

处理写入冲突

多主复制最大的难题——无论是地理分布式的服务端数据库,还是终端用户设备上的本地优先同步引擎——都是不同领导者上的并发写入可能彼此冲突,需要解决。

例如,图 6-9 展示了两个用户同时编辑一个维基页面。用户 1 把页面标题从 A 改为 B,用户 2 则独立地把标题从 A 改为 C。两位用户的变更都成功应用到各自的本地领导者,但异步复制这些变更时,系统发现了冲突。单主数据库不会遇到这个问题。

两个领导者并发更新同一条记录,造成写入冲突。
图 6-9 两个领导者并发更新同一条记录,造成写入冲突。
说明

我们称 图 6-9 中的两次写入为 并发写入,因为最初执行写入时,两者都不知道对方。它们在物理时间上是否真的同时发生并不重要;如果写入发生在离线期间,两者甚至可能相隔很久。真正重要的是,一次写入发生时,另一次写入是否已经生效。

我们会在 “检测并发写入” 中讨论数据库如何判断两次写入是否并发。现在先假定冲突已经能够检测,接下来考虑怎样解决才最合适。

冲突避免

一种策略是从一开始就避免冲突。例如,如果应用程序能确保某条记录的所有写入都经过同一个领导者,那么即使整个数据库采用多主复制,也不会产生冲突。同步引擎客户端离线更新时无法使用这种办法,但在跨地域复制的服务端系统中有时可行 30。

例如,在用户只能编辑自己数据的应用程序中,可以保证某位用户的请求总是路由到同一个地区,并使用该地区的领导者读写。不同用户可以有不同的“主”地区(也许根据与用户的地理距离来选择),但从任一用户的角度看,本质上仍是单主配置。

不过,有时需要改变某条记录的指定领导者:可能是一个地区不可用,必须把流量改送另一个地区;也可能是用户搬到了别处,现在离另一个地区更近。如果用户恰好在指定领导者切换期间执行写入,就可能发生冲突,必须用下面某种方法解决。因此,只要允许改变领导者,冲突避免就可能失效。

再举一个冲突避免的例子。假设要插入新记录,并用自增计数器生成唯一 ID。系统有两个领导者时,可以让一个只生成奇数,另一个只生成偶数。这样两个领导者便不会并发地把同一个 ID 分配给不同记录。“ID 生成器和逻辑时钟” 将讨论其他 ID 分配方案。

最后写入者胜(丢弃并发写入)

如果无法避免冲突,最简单的解决办法是给每次写入附加时间戳,并始终采用时间戳最大的值。例如在 图 6-9 中,假设用户 1 写入的时间戳大于用户 2。两个领导者都会判定页面的新标题应为 B,并丢弃把标题设为 C 的写入。如果两次写入碰巧拥有相同时间戳,还可以比较值来决定胜者(例如字符串可以选择字母顺序在前的值)。

这种方法称为 最后写入者胜(last write wins,LWW),因为时间戳最大的写入被视为“最后”一次写入。不过,这个名称容易误导:当两次写入像 图 6-9 中那样并发时,根本无从定义哪次较早、哪次较晚,因此并发写入的时间戳顺序实质上是随机的。

所以,LWW 的真正含义是:同一条记录在不同领导者上并发写入时,随机挑选其中一次作为胜者,其他写入则静默丢弃,即使它们都已在各自的领导者上成功处理。这样固然能让所有副本最终收敛到一致状态,代价却是数据丢失。

如果能够避免冲突——例如只插入以 UUID 等唯一键标识的记录,而且从不更新——LWW 就没有问题。但如果要更新现有记录,或不同领导者可能插入键相同的记录,就必须判断丢失更新对应用程序是否可以接受。不能接受时,应采用下面介绍的其他冲突解决方式。

LWW 还有一个问题:如果写入时间戳来自实时时钟(例如 Unix 时间戳),系统会对时钟同步极为敏感。假如一个节点的时钟快于其他节点,再尝试覆盖该节点写入的值时,新写入的时间戳可能反而更小,因而被忽略,尽管它显然发生得更晚。使用 逻辑时钟 可以解决这个问题,参见 “ID 生成器和逻辑时钟”。

手动冲突解决

如果不愿随机丢弃某些写入,下一个选择是手工解决冲突。你可能熟悉 Git 等版本控制系统中的做法:两个分支上的提交修改了同一文件的同一行,合并分支时就会产生合并冲突,必须先解决冲突才能完成合并。

在数据库中,让一次冲突阻塞整个复制过程,直到有人解决,显然并不现实。数据库通常会保存一条记录的所有并发写入值——例如 图 6-9 中的 B 和 C。这些值有时称为 兄弟值(siblings)。下次查询该记录时,数据库返回 全部 值,而不只是最新的一个。随后可以任意选择解决办法:在应用代码中自动处理(例如把 B 与 C 拼成“B/C”),或询问用户;最后再向数据库写回一个新值,消解冲突。

CouchDB 等系统采用这种冲突解决方式,但它也有不少问题:

  • 数据库 API 会发生变化。例如,维基页面标题原本只是字符串,现在却变成一组字符串;通常只有一个元素,发生冲突时却可能有多个。应用代码处理这种数据会相当别扭。
  • 让用户手工合并兄弟值,无论对应用开发者还是用户都是沉重负担:开发者必须制作冲突解决界面,用户则可能不明白自己为什么要做这件事、又该做什么。很多情况下,自动合并比打扰用户更合适。
  • 自动合并兄弟值如果不够谨慎,也会产生意外结果。例如,亚马逊购物车过去允许并发更新,再把任一兄弟值中出现的所有商品都保留下来,也就是取购物车的并集。如果顾客在一个兄弟值中删除商品,而另一个兄弟值仍含有该商品,已经删除的商品便会意外重现 45。图 6-10 展示了这种情况:设备 1 删除 Book,设备 2 同时删除 DVD,合并冲突后两件商品却都回来了。
  • 如果多个节点同时观察到冲突并各自解决,解决过程本身还可能引入新冲突,而且各节点给出的结果可能不一致。例如,如果没有固定排列顺序,一个节点可能把 B 和 C 合并成“B/C”,另一个却合并成“C/B”;再合并“B/C”与“C/B”时,结果可能变成“B/C/C/B”之类的怪东西。
亚马逊购物车异常示例:以并集方式合并购物车冲突时,已删除的商品可能重新出现。
图 6-10 亚马逊购物车异常示例:以并集方式合并购物车冲突时,已删除的商品可能重新出现。

自动冲突解决

对许多应用程序而言,处理冲突的最佳方式,是用算法自动把并发写入合并成一致状态。自动冲突解决可以保证所有副本 收敛 到同一状态:只要处理过相同的一组写入,各副本的状态就相同,与写入到达的顺序无关。

LWW 是冲突解决算法的一个简单例子。针对不同数据类型,人们还开发了更复杂的合并算法,目标是尽可能保留所有更新的预期效果,从而避免数据丢失:

  • 如果数据是文本(例如维基页面的标题或正文),可以检测相邻版本之间插入或删除了哪些字符。合并结果会保留任一兄弟值中的所有插入和删除。如果用户在同一位置并发插入文本,可以按确定性的顺序排列,确保所有节点得到相同结果。
  • 如果数据是元素集合,无论像待办事项列表一样有序,还是像购物车一样无序,都可以像合并文本那样跟踪插入和删除。为避免 图 6-10 中的购物车异常,算法会记住 Book 和 DVD 已被删除,因此合并结果为 Cart = {Soap}。
  • 如果数据是可以递增或递减的整数计数器(例如社交媒体帖子的点赞数),合并算法可以计算各兄弟值上分别发生了多少次递增和递减,再正确相加,既不重复计数,也不丢失更新。
  • 如果数据是键值映射,可以对同一个键下的值采用其他某种冲突解决算法;不同键上的更新则可以彼此独立地处理。

冲突解决并非无所不能。例如,如果规定一个列表最多包含五个元素,而多位用户并发添加元素,使总数超过五个,那么唯一的选择就是丢掉其中一些。即便如此,自动冲突解决仍足以构建许多实用应用。一旦决定构建可协作的离线优先或本地优先应用,冲突解决就不可避免,而自动化往往是最佳选择。

CRDT 与操作变换

实现自动冲突解决时,通常使用两类算法:无冲突复制数据类型(CRDT)46 和 操作变换(OT)47。二者的设计理念与性能特征不同,但都能自动合并前面提到的各种数据。

图 6-11 展示了 OT 和 CRDT 分别如何合并文本的并发更新。假设两个副本起初都保存文本“ice”。一个副本在开头插入字母“n”,得到“nice”;与此同时,另一个副本在末尾插入感叹号,得到“ice!”。

OT 与 CRDT 分别如何合并字符串中的两次并发插入。
图 6-11 OT 与 CRDT 分别如何合并字符串中的两次并发插入。

两类算法用不同方式得到合并结果“nice!”:

OT
记录字符插入或删除位置的索引:“n”插入索引 0,“!”插入索引 3。然后两个副本交换操作。在索引 0 插入“n”可以原样应用;但如果直接在状态“nice”的索引 3 插入“!”,结果会变成错误的“nic!e”。因此,必须根据已经应用的并发操作变换每个操作的索引。这里,为了计入较小索引处插入的“n”,需要把“!”的插入位置变换为索引 4。
CRDT
大多数 CRDT 不使用索引,而是给每个字符分配唯一且不可变的 ID,再据此确定插入和删除位置。例如在 图 6-11 中,“i”的 ID 是 1A,“c”的 ID 是 2A,依此类推。插入感叹号时,生成的操作既包含新字符的 ID(4B),也包含插入位置之前那个现有字符的 ID(3A)。要插在字符串开头,就把前驱字符 ID 设为“nil”。同一位置的并发插入按字符 ID 排列。这样无需变换操作,也能保证各副本收敛。

许多算法都建立在这些思路的不同变体上。列表和数组可以采用类似方法,把字符换成列表元素;键值映射等其他数据类型也很容易加入。OT 与 CRDT 在性能和功能上各有取舍,但也可以把二者的优点结合到同一种算法中 48。

OT 最常用于文本的实时协作编辑,例如 Google Docs 32;CRDT 则用于 Redis Enterprise、Riak、Azure Cosmos DB 等分布式数据库 49。面向 JSON 数据的同步引擎既可以用 CRDT 实现(如 Automerge、Yjs),也可以用 OT 实现(如 ShareDB)。

什么是冲突?

有些冲突显而易见。在 图 6-9 的例子中,两次写入并发修改同一条记录的同一个字段,把它设成两个不同的值。毫无疑问,这就是冲突。

另一些冲突则更为微妙,不易发现。以会议室预订系统为例,它记录哪个房间在什么时间由哪组人预订。应用程序必须确保同一时刻每个房间只分配给一组人,也就是说,同一房间的预订不能重叠。如果两项不同预订在同一时间占用同一房间,就会产生冲突。即使应用程序在允许预订前检查空闲情况,只要两次预订分别在不同领导者上进行,仍可能发生冲突。

这个问题没有现成的简短答案,不过在接下来的章节中,我们会逐步加深理解。第 8 章 将给出更多冲突示例;“排序事件以捕获因果关系” 则会讨论在复制系统中可伸缩地检测与解决冲突的方法。

无主复制

本章此前讨论的单主复制与多主复制,都基于同一个思路:客户端把写请求发给某个节点(领导者),再由数据库系统负责把写入复制到其他副本。领导者决定处理写入的顺序,追随者则按相同顺序应用领导者的写入。

另一些数据存储系统采取了不同办法:放弃领导者概念,允许任何副本直接接受客户端写入。最早的一些复制数据系统采用的就是无主模型 1 50,但在关系数据库占据主导地位的年代,这个思路几乎被遗忘。2007 年,亚马逊把它用于内部的 Dynamo 系统 45,无主架构由此再度流行。Riak、Cassandra 和 ScyllaDB 都是受 Dynamo 启发、采用无主复制模型的开源数据存储,因此这类数据库也称为 Dynamo 风格(Dynamo-style)数据库。

说明

最初的 Dynamo 系统只在论文中有所描述 45,从未在亚马逊之外发布。名称相近的 DynamoDB 是 AWS 后来推出的云数据库,但架构完全不同:它采用基于 Multi-Paxos 共识算法的单主复制 5。

在一些无主实现中,客户端直接把写入发给多个副本;另一些实现则由协调者节点代客户端完成这件事。不过,与有领导者的数据库不同,协调者不会强制规定写入顺序。我们将看到,这项设计差异会深刻影响数据库的使用方式。

当节点故障时写入数据库

假设一个数据库有三个副本,其中一个暂时不可用,也许正在重启以安装系统更新。在单主配置中,要继续处理写入,可能需要执行故障切换(参见 “处理节点故障”)。

无主配置则根本没有故障切换。图 6-12 展示了此时的情形:客户端(用户 1234)把写入并行发给三个副本;两个可用副本接受写入,不可用副本则错过了它。假设三个副本中有两个确认就足以判定写入成功:用户 1234 收到两个 ok 响应后,系统便认为写入成功,客户端直接忽略有一个副本漏掉写入这一事实。

节点停机后的仲裁写、仲裁读和读修复。
图 6-12 节点停机后的仲裁写、仲裁读和读修复。

现在设想不可用的节点恢复上线,客户端开始从它读取。该节点停机期间发生的写入都没有保存在这里,因此从它读取时,响应中可能包含 陈旧(过时)的值。

为解决这个问题,客户端读取数据库时,不能只把请求发给一个副本:读请求也要并行发往多个节点。不同节点可能给出不同响应,例如一个返回最新值,另一个返回陈旧值。

要分辨哪些响应是最新的、哪些已经过时,每个写入的值都必须带有版本号或时间戳,类似 “最后写入者胜(丢弃并发写入)” 中介绍的做法。客户端收到多个读取结果时,采用时间戳最大的值——即使只有一个副本返回该值,其他几个副本都返回旧值。更多细节参见 “检测并发写入”。

追赶错过的写入

复制系统应当保证所有数据最终都会复制到每个副本。不可用节点恢复上线后,怎样补上停机期间错过的写入?Dynamo 风格的数据存储会使用以下几种机制:

读修复(read repair)
客户端并行读取多个节点时,可以发现陈旧响应。例如在 图 6-12 中,用户 2345 从副本 3 得到版本 6 的值,从副本 1 和副本 2 得到版本 7 的值。客户端发现副本 3 的值已经过时,于是把较新的值写回这个副本。对于经常读取的值,这种方法很有效。
提示移交(hinted handoff)
某个副本不可用时,另一个副本可以替它保存写入,并把这些写入记录为 提示。原本应接收这些写入的副本恢复后,保存提示的副本会将它们发送过去,然后删除提示。即使某些值从未被读取、无法通过读修复更新,这个 移交 过程也能让副本赶上进度。
反熵(anti-entropy)
此外,还有一个后台进程定期查找副本之间的数据差异,把缺失的数据从一个副本复制到另一个。与基于领导者的复制日志不同,这个 反熵过程(anti-entropy process)并不按特定顺序复制写入,数据得到复制之前可能有很长延迟。

读写仲裁

在 图 6-12 的例子中,写入只在三个副本中的两个上完成,我们仍判定它成功。如果只有一个副本接受写入呢?这个下限究竟能压到多低?

如果能保证每次成功写入至少保存在三个副本中的两个上,那么最多只有一个副本是陈旧的。因此,只要读取至少两个副本,就可以确信其中至少一个是最新的。即使第三个副本停机或响应缓慢,读取仍能返回最新值。

更一般地说,假设有 n 个副本,每次写入必须得到 w 个节点确认才算成功,每次读取则至少查询 r 个节点。(上述例子中,n = 3、w = 2、r = 2。)只要 w + r > n,读取时就有望得到最新值,因为所查询的 r 个节点中,至少有一个必然是最新的。遵守这些 r、w 取值的操作称为 仲裁读(quorum read)和 仲裁写(quorum write)50。可以把 r 和 w 看作一次读或写要成立所需的最低票数。

在 Dynamo 风格的数据库中,参数 n、w、r 通常都可以配置。常见做法是让 n 取奇数(通常为 3 或 5),并令 w = r = (n + 1) / 2(向上取整);不过也可以根据需要调整。例如,写少读多的工作负载可能适合设为 w = n、r = 1。这样读取更快,缺点是只要一个节点失效,所有数据库写入都会失败。

说明

集群中的节点数可以多于 n,但任意给定值只存储在 n 个节点上。这样就能对数据集进行分片,支持单个节点无法容纳的数据集。我们会在 第 7 章 继续讨论分片。

仲裁条件 w + r > n 使系统能够按以下方式容忍节点不可用:

  • 如果 w < n,有一个节点不可用时仍可处理写入。
  • 如果 r < n,有一个节点不可用时仍可处理读取。
  • 当 n = 3、w = 2、r = 2 时,可以像 图 6-12 那样容忍一个节点不可用。
  • 当 n = 5、w = 3、r = 3 时,可以容忍两个节点不可用,如 图 6-13 所示。

通常,读写请求总会并行发送到全部 n 个副本。参数 w 和 r 决定要等待多少个节点,也就是在判定读或写成功之前,n 个节点中必须有多少个报告成功。

如果 w + r > n,读取的 r 个副本中至少有一个必然见过最近一次成功写入。
图 6-13 如果 w + r > n,读取的 r 个副本中至少有一个必然见过最近一次成功写入。

如果可用节点少于所需的 w 或 r,写入或读取就会返回错误。节点不可用可能有许多原因:节点停机(崩溃或断电)、执行操作时出错(磁盘已满,无法写入)、客户端与节点之间网络中断,等等。我们只关心节点是否返回成功响应,无需区分故障的具体类型。

仲裁一致性的局限

如果有 n 个副本,并选择满足 w + r > n 的 w 和 r,通常可以期望每次读取都返回某个键最近写入的值。这是因为写入所涉及的节点集合与读取所涉及的节点集合必然有交集;也就是说,读取的节点中至少有一个保存着最新值,如 图 6-13 所示。

r 和 w 通常取节点的多数(多于 n/2),因为这样既能保证 w + r > n,又能容忍最多 n/2(向下取整)个节点失效。不过,法定人数并不一定非得是多数;真正重要的是,读操作与写操作所用的节点集合至少有一个共同节点。法定人数还可以有其他安排,为分布式算法的设计提供一定灵活性 51。

也可以把 w 和 r 设得更小,使 w + r ≤ n,即不满足仲裁条件。此时读写请求仍会发往 n 个节点,只是操作成功所需的成功响应更少。

w 和 r 越小,越容易读到陈旧值,因为读操作更可能没有覆盖保存最新值的节点。好处则是延迟更低、可用性更高:网络中断导致许多副本不可达时,系统仍有更大机会继续处理读写。只有可达副本数低于 w 或 r 时,数据库才会分别变得不可写或不可读。

然而,即使 w + r > n,仍有一些边缘情况会让一致性属性变得难以理解,例如:

  • 如果保存新值的节点失效,又从保存旧值的副本恢复数据,保存新值的副本数可能降到 w 以下,破坏仲裁条件。
  • 再平衡期间,一部分数据会从一个节点迁移到另一个节点(见 第 7 章),各节点对“某个值的 n 个副本应由哪些节点保存”可能看法不一,使读仲裁与写仲裁不再相交。
  • 如果读操作与写操作并发,读取可能看见并发写入的值,也可能看不见。尤其是,某次读取可能看到新值,后续读取却看到旧值,参见 “线性一致性与仲裁”。
  • 如果一次写入在部分副本上成功、在其余副本上失败(例如某些节点磁盘已满),总成功数少于 w,那么整体写入会判定失败,但成功副本上的写入不会回滚。这意味着即使系统报告写入失败,后续读取仍可能返回这次写入的值 52。
  • 如果数据库用实时时钟的时间戳判断写入的新旧(例如 Cassandra 和 ScyllaDB),另一个时钟较快的节点只要写过同一个键,后续写入就可能被静默丢弃。我们在 “最后写入者胜(丢弃并发写入)” 中已经见过这个问题,还会在 “对同步时钟的依赖” 中进一步讨论。
  • 如果两次写入并发发生,一个副本可能先处理其中一次,另一个副本则先处理另一次,由此产生冲突,与多主复制中的情况相似(参见 “处理写入冲突”)。我们会在 “检测并发写入” 中再谈这个问题。

因此,仲裁看似保证读取返回最近写入的值,实际却没有那么简单。Dynamo 风格数据库通常针对能够容忍最终一致性的场景优化。参数 w 和 r 可以调节读到陈旧值的概率 53,却不宜被当作绝对保证。

监控陈旧性

从运维角度看,监控数据库返回的结果是否最新非常重要。即使应用程序能够容忍陈旧读取,也必须了解复制是否健康。如果复制大幅落后,系统应发出告警,以便调查网络故障、节点过载等原因。

采用基于领导者的复制时,数据库通常会暴露复制延迟指标,供监控系统采集。这是因为写入在领导者和追随者上按相同顺序应用,每个节点都有自己在复制日志中的位置,也就是已经在本地应用了多少次写入。用领导者当前位置减去追随者当前位置,便可度量复制延迟。

无主复制系统没有固定的写入应用顺序,监控起来更加困难。副本为了移交而保存的提示数量可以作为一项健康指标,却很难作出有意义的解释 54。最终一致性有意给出了一项模糊保证,但为了可运维性,必须能够量化“最终”究竟有多远。

单主与无主复制的性能

基于单个领导者的复制系统能够提供强一致性保证,而无主系统很难甚至不可能做到。不过,正如 “复制延迟的问题” 中所见,在基于领导者的复制系统里,如果从异步更新的追随者读取,同样可能得到陈旧值。

从领导者读取可以保证响应最新,却存在性能问题:

  • 读取吞吐量受领导者处理能力限制;相比之下,读扩展可以把读取分散到异步更新的副本,但这些副本可能返回陈旧值。
  • 领导者失效后,必须等系统检测到故障并完成故障切换,才能继续处理请求。即使故障切换很快,响应时间的短暂上升也会被用户察觉;如果耗时很长,系统就会在此期间不可用。
  • 系统对领导者的性能问题极其敏感。如果领导者因过载或资源争用而响应缓慢,用户的响应时间也会立即增加。

无主架构的一大优点,是面对这些问题时韧性更强。系统无需故障切换,而且请求原本就会并行发往多个副本,因此一个副本变慢或不可用,对响应时间的影响很小:客户端只需采用响应较快的其他副本所返回的结果。采用最快响应的做法称为 请求对冲(request hedging),可以显著降低尾延迟 55。

无主系统之所以有这种韧性,根本原因是它不区分正常情况和故障情况。这对于处理所谓的 灰色失效(gray failure)尤其有利:节点并未彻底停机,却处于降级状态,处理请求异常缓慢 56;节点单纯过载时也是如此(例如节点离线一段时间后,靠提示移交恢复可能产生大量额外负载)。基于领导者的系统必须判断情况是否严重到需要故障切换,而故障切换本身又可能带来进一步中断;无主系统根本不需要作出这项判断。

当然,无主系统也可能遇到性能问题:

  • 即使不需要执行故障切换,也必须由一个副本发现另一个副本不可用,才能替它保存错过写入的提示。不可用副本恢复后,移交过程还要把这些提示发给它。在系统本已承压时,这会给副本增加额外负载 54。
  • 副本越多,法定人数越大,请求完成前必须等待的响应也越多。即使只等待最快的 r 或 w 个副本,即使所有请求并行发出,更大的 r 或 w 仍会提高遇到慢副本的概率,从而增加总体响应时间(参见 “响应时间指标的应用”)。
  • 大范围网络中断使客户端与大量副本断开时,可能根本无法组成法定人数。有些无主数据库允许任何可达副本接受写入,即使它不属于该键通常所在的副本集合(Riak 和 Dynamo 称之为 宽松仲裁,sloppy quorum 45;Cassandra 和 ScyllaDB 称之为 一致性级别 ANY)。后续读取不保证能看到这次写入,但对某些应用而言,这仍好过写入直接失败。

多主复制抵御网络中断的能力甚至可能强于无主复制,因为读写只需与一个领导者通信,而领导者可以与客户端位于同一地区。不过,一个领导者上的写入会异步传播给其他领导者,读取结果因而可能任意陈旧。仲裁读写提供了一种折中:既有良好的容错能力,也有很高概率读到最新数据。

多地区操作

我们此前把跨地区复制作为多主复制的一个用例(见 “多主复制”)。无主复制同样适合多地区运行,因为它本来就是为了容忍相互冲突的并发写入、网络中断和延迟尖峰而设计的。

Cassandra 和 ScyllaDB 在常规无主模型中实现多地区支持:客户端把写入直接发往所有地区的副本,并可选择多种一致性级别,规定请求至少得到多少响应才算成功。例如,可以要求所有地区的全部副本共同组成一个法定人数,也可以要求每个地区各自组成法定人数,或只要求客户端所在地区达到法定人数。本地法定人数无需等待其他地区的慢请求,但也更容易返回陈旧结果。

Riak 则把客户端与数据库节点之间的所有通信限制在本地地区,因此 n 表示一个地区内的副本数。数据库集群之间的跨地区复制在后台异步进行,方式与多主复制相似。

检测并发写入

与多主复制一样,无主数据库允许对同一个键并发写入,由此产生需要解决的冲突。冲突可能在写入发生时出现,但并非总是如此;它也可能到读修复、提示移交或反熵阶段才被发现。

问题在于,网络延迟会变化,系统还可能部分失效,所以事件抵达不同节点的顺序可能不同。例如,图 6-14 展示了客户端 A 和 B 同时写入三节点数据存储中的键 X:

  • 节点 1 收到 A 的写入,但由于短暂中断,一直没有收到 B 的写入。
  • 节点 2 先收到 A 的写入,再收到 B 的写入。
  • 节点 3 先收到 B 的写入,再收到 A 的写入。
Dynamo 风格数据存储中的并发写入没有明确定义的顺序。
图 6-14 Dynamo 风格数据存储中的并发写入没有明确定义的顺序。

如果每个节点一收到客户端写请求就直接覆盖键的值,各节点将永久不一致,如 图 6-14 最后的 get 请求所示:节点 2 认为 X 的最终值是 B,其他节点却认为是 A。

为了达到最终一致,各副本必须收敛到同一个值。可以采用 “处理写入冲突” 中讨论过的任意冲突解决机制,例如 Cassandra 和 ScyllaDB 使用的最后写入者胜、手工解决,或 “CRDT 与操作变换” 中介绍且 Riak 使用的 CRDT。

最后写入者胜很容易实现:给每次写入附加时间戳,时间戳较大的值总是覆盖较小的值。但时间戳无法告诉你两个值究竟是否冲突:它们可能是并发写入的,也可能先后写入。如果要显式解决冲突,系统必须更仔细地检测并发写入。

“先发生”关系与并发

怎样判断两个操作是否并发?先看几个例子来建立直觉:

  • 在 图 6-8 中,两次写入并不并发:A 的插入 先发生于 B 的递增,因为 B 所递增的值正是 A 插入的值。换句话说,B 的操作建立在 A 的操作之上,所以 B 必然发生得更晚。也可以说,B 因果依赖 于 A。
  • 图 6-14 中的两次写入则是并发的:每个客户端开始操作时,都不知道另一个客户端也在操作同一个键。因此,两次操作之间没有因果依赖。

如果操作 B 知道 A、依赖 A,或以某种方式建立在 A 之上,就称操作 A 先发生于(happens before)操作 B。一项操作是否先发生于另一项操作,是定义并发的关键。事实上,只要两个操作谁也不先发生于另一个——也就是说,谁都不知道对方——就可以称它们 并发(concurrent)57。

因此,对于任意两个操作 A 与 B,只有三种可能:A 先发生于 B;B 先发生于 A;或者 A 与 B 并发。我们需要一种算法判断两次操作是否并发。如果一项操作先发生于另一项,后发生的操作就应覆盖先前操作;如果两者并发,则出现了需要解决的冲突。

并发、时间与相对论

两项操作似乎只有在“同一时刻”发生时才应称为并发,实际上它们在物理时间上是否重叠并不重要。由于分布式系统中的时钟问题,判断两件事是否恰好同时发生相当困难,我们会在 第 9 章 进一步讨论。

定义并发时,精确时间并不重要:只要两项操作彼此都不知道对方,就称它们并发,无论它们实际发生在什么物理时刻。人们有时把这个原理与物理学中的狭义相对论联系起来 57。狭义相对论提出,信息传播不可能超过光速。因此,如果相隔一定距离的两个事件之间,时间差短于光传播这段距离所需的时间,它们就不可能相互影响。

在计算机系统中,即使按光速计算,一项操作原则上来得及影响另一项,两者仍可能并发。例如,当时网络很慢或已经中断,两项操作即使相隔一段时间,仍会因网络问题而彼此无法知晓。

捕获先发生关系

下面看一种算法,它可以判断两项操作是并发的,还是一项先发生于另一项。为简单起见,先从只有一个副本的数据库开始。弄清单副本的做法后,再推广到拥有多个副本的无主数据库。

图 6-15 展示了两个客户端并发地向同一个购物车添加商品。(如果这个例子太无聊,也可以设想两名空中交通管制员并发地把飞机加入各自正在监视的空域。)购物车最初为空,两个客户端先后共向数据库发出五次写入:

  1. 客户端 1 把 milk 加入购物车。这是该键的第一次写入,服务器成功保存它并分配版本 1;然后把值和版本号一起返回给客户端。
  2. 客户端 2 把 eggs 加入购物车,却不知道客户端 1 同时加入了 milk(它以为 eggs 是购物车中唯一的商品)。服务器为这次写入分配版本 2,把 eggs 和 milk 保存为两个独立的值(兄弟值),再把 两个 值连同版本号 2 一起返回给客户端。
  3. 客户端 1 不知道客户端 2 的写入,又想加入 flour,因此它认为购物车内容应为 [milk, flour]。它把这个值连同服务器此前给出的版本号 1 一起发送。服务器可以根据版本号判断:[milk, flour] 取代了先前的 [milk],却与 [eggs] 并发。因此,服务器为 [milk, flour] 分配版本 3,覆盖版本 1 的 [milk],保留版本 2 的 [eggs],并把剩下的两个值都返回给客户端。
  4. 与此同时,客户端 2 想加入 ham,并不知道客户端 1 刚刚加入 flour。客户端 2 在上一次响应中收到了 [milk] 和 [eggs],于是将两者合并,再加入 ham,形成新值 [eggs, milk, ham]。它把这个值连同先前的版本号 2 一起发给服务器。服务器判断版本 2 可以覆盖 [eggs],但与 [milk, flour] 并发;剩下的两个值便是版本 3 的 [milk, flour] 和版本 4 的 [eggs, milk, ham]。
  5. 最后,客户端 1 想加入 bacon。它此前在版本 3 的响应中收到 [milk, flour] 和 [eggs],于是合并二者,加入 bacon,把最终值 [milk, flour, eggs, bacon] 连同版本号 3 发给服务器。这个值覆盖 [milk, flour]([eggs] 已在上一步被覆盖),却与 [eggs, milk, ham] 并发,因此服务器保留这两个并发值。
捕获两个客户端并发编辑购物车时的因果依赖。
图 6-15 捕获两个客户端并发编辑购物车时的因果依赖。

图 6-15 中各操作之间的数据流,在 图 6-16 中以图形表示。箭头指出哪项操作 先发生于 另一项,也就是说,后发生的操作 知道 或 依赖 先发生的操作。在这个例子里,客户端从未完全掌握服务器上的最新数据,因为始终有另一项操作并发进行。但值的旧版本最终会被覆盖,而且不会丢失任何写入。

图 6-15 中因果依赖关系的图示。
图 6-16 图 6-15 中因果依赖关系的图示。

请注意,服务器仅凭版本号就能判断两项操作是否并发,无需解释值本身,因此值可以是任意数据结构。算法如下:

  • 服务器为每个键维护一个版本号;每次写入该键时递增版本号,并把新版本号与写入值一同保存。
  • 客户端读取一个键时,服务器返回所有兄弟值(即尚未被覆盖的全部值)以及最新版本号。客户端写入前必须先读取。
  • 客户端写入一个键时,必须带上前一次读取所得的版本号,还必须把上次读取收到的所有值合并起来,例如使用 CRDT,或询问用户。写请求的响应与读取相似,也会返回所有兄弟值,因此可以像购物车例子那样连续执行多次写入。
  • 服务器收到带有特定版本号的写入时,可以覆盖版本号不高于它的所有值,因为这些值已被合并进新值;版本号更高的值则必须保留,因为它们与传入写入并发。

写入带上前一次读取所得的版本号,就说明这次写入基于哪个先前状态。如果写入不含版本号,它就与其他所有写入并发,因而不会覆盖任何内容,只会作为后续读取返回的值之一。

版本向量

图 6-15 的例子只有一个副本。如果没有领导者,而且多个副本都能接受写入,算法需要怎样改变?

图 6-15 用一个版本号捕获操作之间的依赖关系,但多个副本并发接受写入时,一个版本号就不够了。此时必须针对每个键,给 每个副本 分别维护版本号。副本处理写入时递增自己的版本号,同时记录自己见过的其他副本版本号。这些信息表明哪些值应当覆盖,哪些值应作为兄弟值保留。

所有副本的版本号集合称为 版本向量 58。这种思路有若干变体,其中最值得关注的也许是 点化版本向量(dotted version vector)59 60,Riak 2.0 采用了这种变体 61 62。这里不展开细节;它的工作方式与购物车例子非常相似。

与 图 6-15 中的版本号一样,读取时数据库副本会把版本向量发给客户端,随后写入时客户端必须再把它带回数据库。(Riak 把版本向量编码成一个字符串,称为 因果上下文,causal context。)版本向量让数据库能够区分覆盖写入和并发写入。

版本向量还保证:先从一个副本读取,再把写入发给另一个副本,是安全的。这样做可能产生兄弟值,但只要正确合并兄弟值,就不会丢失数据。

版本向量与向量时钟

版本向量 有时也称为 向量时钟(vector clock),尽管两者并不完全相同。区别十分微妙,细节请参阅相关文献 60 63 64。简而言之,比较副本状态时,应当使用版本向量。

总结

本章考察了复制问题。复制有多种用途:

高可用性
即使一台或多台机器、一个可用区乃至整个地区停机,系统仍能继续运行
离线运行
网络中断时,应用程序仍能继续工作
延迟
把数据放在地理上靠近用户的位置,让用户能够更快地与之交互
可伸缩性
把读取分散到多个副本,处理超出单台机器能力的读取量

复制的目标看似简单——在多台机器上保留相同数据的副本——实际却极其棘手。它要求我们仔细考虑并发、所有可能出错的环节,以及怎样应对故障造成的后果。至少要处理节点不可用和网络中断,而且这还没有算上软件缺陷或硬件错误引起的静默数据损坏等更隐蔽的故障。

我们讨论了三种主要的复制方式:

单主复制
客户端把所有写入发给一个节点(领导者),领导者再把数据变更事件流发送给其他副本(追随者)。读取可以在任意副本上执行,但追随者返回的结果可能陈旧。
多主复制
客户端把每次写入发给多个领导者中的一个,任意领导者都能接受写入。各领导者相互发送数据变更事件流,也会将其发送给追随者。
无主复制
客户端把每次写入发给多个节点,并行读取多个节点,从而发现并修复持有陈旧数据的节点。

每种方式都有优缺点。单主复制很流行,因为它相对容易理解,又能提供强一致性。多主复制和无主复制面对节点故障、网络中断和延迟尖峰时韧性更强,代价是必须解决冲突,而且只能提供较弱的一致性保证。

复制可以同步,也可以异步;发生故障时,这项选择会深刻影响系统行为。系统平稳运行时,异步复制可能很快,但仍必须弄清复制延迟增大或服务器失效时会发生什么。如果领导者失效,而你把一个异步更新的追随者提升为新领导者,最近提交的数据可能丢失。

我们考察了复制延迟可能造成的几种反常现象,并讨论了几种一致性模型,以便判断应用程序在复制延迟下应有怎样的行为:

写后读一致性
用户应当总能看到自己提交的数据。
单调读
用户看到某一时刻的数据后,不应在稍后又看到更早时刻的数据。
一致前缀读
用户看到的数据状态应符合因果关系,例如以正确顺序看到问题及其回答。

最后,我们讨论了多主复制和无主复制如何让所有副本最终收敛到一致状态:用版本向量或类似算法检测哪些写入并发,再用 CRDT 等冲突解决算法合并并发写入的值。最后写入者胜和手工解决冲突也是可选方案。

本章一直假设每个副本都保存整个数据库的完整拷贝,但对大型数据集而言,这并不现实。下一章将介绍 分片,使每台机器只需存储一部分数据。

参考文献


  1. B. G. Lindsay, P. G. Selinger, C. Galtieri, J. N. Gray, R. A. Lorie, T. G. Price, F. Putzolu, I. L. Traiger, and B. W. Wade. Notes on Distributed Databases. IBM Research, Research Report RJ2571(33471), July 1979. Archived at perma.cc/EPZ3-MHDD ↩︎ ↩︎

  2. Kenny Gryp. MySQL Terminology Updates. dev.mysql.com, July 2020. Archived at perma.cc/S62G-6RJ2 ↩︎

  3. Oracle Corporation. Oracle (Active) Data Guard 19c: Real-Time Data Protection and Availability. White Paper, oracle.com, March 2019. Archived at perma.cc/P5ST-RPKE ↩︎

  4. Microsoft. What is an Always On availability group? learn.microsoft.com, September 2024. Archived at perma.cc/ABH6-3MXF ↩︎

  5. Mostafa Elhemali, Niall Gallagher, Nicholas Gordon, Joseph Idziorek, Richard Krog, Colin Lazier, Erben Mo, Akhilesh Mritunjai, Somu Perianayagam, Tim Rath, Swami Sivasubramanian, James Christopher Sorenson III, Sroaj Sosothikul, Doug Terry, and Akshat Vig. Amazon DynamoDB: A Scalable, Predictably Performant, and Fully Managed NoSQL Database Service. At USENIX Annual Technical Conference (ATC), July 2022. ↩︎ ↩︎

  6. Rebecca Taft, Irfan Sharif, Andrei Matei, Nathan VanBenschoten, Jordan Lewis, Tobias Grieger, Kai Niemi, Andy Woods, Anne Birzin, Raphael Poss, Paul Bardea, Amruta Ranade, Ben Darnell, Bram Gruneir, Justin Jaffray, Lucy Zhang, and Peter Mattis. CockroachDB: The Resilient Geo-Distributed SQL Database. At ACM SIGMOD International Conference on Management of Data (SIGMOD), pages 1493–1509, June 2020. doi:10.1145/3318464.3386134 ↩︎

  7. Dongxu Huang, Qi Liu, Qiu Cui, Zhuhe Fang, Xiaoyu Ma, Fei Xu, Li Shen, Liu Tang, Yuxing Zhou, Menglong Huang, Wan Wei, Cong Liu, Jian Zhang, Jianjun Li, Xuelian Wu, Lingyu Song, Ruoxi Sun, Shuaipeng Yu, Lei Zhao, Nicholas Cameron, Liquan Pei, and Xin Tang. TiDB: a Raft-based HTAP database. Proceedings of the VLDB Endowment, volume 13, issue 12, pages 3072–3084. doi:10.14778/3415478.3415535 ↩︎

  8. Mallory Knodel and Niels ten Oever. Terminology, Power, and Inclusive Language in Internet-Drafts and RFCs. IETF Internet-Draft, August 2023. Archived at perma.cc/5ZY9-725E ↩︎

  9. Buck Hodges. Postmortem: VSTS 4 September 2018. devblogs.microsoft.com, September 2018. Archived at perma.cc/ZF5R-DYZS ↩︎

  10. Gunnar Morling. Leader Election With S3 Conditional Writes. www.morling.dev, August 2024. Archived at perma.cc/7V2N-J78Y ↩︎

  11. Vignesh Chandramohan, Rohan Desai, and Chris Riccomini. SlateDB Manifest Design. github.com, May 2024. Archived at perma.cc/8EUY-P32Z ↩︎

  12. Stas Kelvich. Why does Neon use Paxos instead of Raft, and what’s the difference? neon.tech, August 2022. Archived at perma.cc/SEZ4-2GXU ↩︎

  13. Dimitri Fontaine. An introduction to the pg_auto_failover project. tapoueh.org, November 2021. Archived at perma.cc/3WH5-6BAF ↩︎

  14. Jesse Newland. GitHub availability this week. github.blog, September 2012. Archived at perma.cc/3YRF-FTFJ ↩︎

  15. Mark Imbriaco. Downtime last Saturday. github.blog, December 2012. Archived at perma.cc/M7X5-E8SQ ↩︎

  16. John Hugg. ‘All In’ with Determinism for Performance and Testing in Distributed Systems. At Strange Loop, September 2015. ↩︎

  17. Hironobu Suzuki. The Internals of PostgreSQL. interdb.jp, 2017. ↩︎

  18. Amit Kapila. WAL Internals of PostgreSQL. At PostgreSQL Conference (PGCon), May 2012. Archived at perma.cc/6225-3SUX ↩︎

  19. Amit Kapila. Evolution of Logical Replication. amitkapila16.blogspot.com, September 2023. Archived at perma.cc/F9VX-JLER ↩︎

  20. Aru Petchimuthu. Upgrade your Amazon RDS for PostgreSQL or Amazon Aurora PostgreSQL database, Part 2: Using the pglogical extension. aws.amazon.com, August 2021. Archived at perma.cc/RXT8-FS2T ↩︎

  21. Yogeshwer Sharma, Philippe Ajoux, Petchean Ang, David Callies, Abhishek Choudhary, Laurent Demailly, Thomas Fersch, Liat Atsmon Guz, Andrzej Kotulski, Sachin Kulkarni, Sanjeev Kumar, Harry Li, Jun Li, Evgeniy Makeev, Kowshik Prakasam, Robbert van Renesse, Sabyasachi Roy, Pratyush Seth, Yee Jiun Song, Benjamin Wester, Kaushik Veeraraghavan, and Peter Xie. Wormhole: Reliable Pub-Sub to Support Geo-Replicated Internet Services. At 12th USENIX Symposium on Networked Systems Design and Implementation (NSDI), May 2015. ↩︎

  22. Douglas B. Terry. Replicated Data Consistency Explained Through Baseball. Microsoft Research, Technical Report MSR-TR-2011-137, October 2011. Archived at perma.cc/F4KZ-AR38 ↩︎ ↩︎ ↩︎

  23. Douglas B. Terry, Alan J. Demers, Karin Petersen, Mike J. Spreitzer, Marvin M. Theher, and Brent B. Welch. Session Guarantees for Weakly Consistent Replicated Data. At 3rd International Conference on Parallel and Distributed Information Systems (PDIS), September 1994. doi:10.1109/PDIS.1994.331722 ↩︎ ↩︎

  24. Werner Vogels. Eventually Consistent. ACM Queue, volume 6, issue 6, pages 14–19, October 2008. doi:10.1145/1466443.1466448 ↩︎

  25. Simon Willison. Reply to: “My thoughts about Fly.io (so far) and other newish technology I’m getting into”. news.ycombinator.com, May 2022. Archived at perma.cc/ZRV4-WWV8 ↩︎

  26. Nithin Tharakan. Scaling Bitbucket’s Database. atlassian.com, October 2020. Archived at perma.cc/JAB7-9FGX ↩︎

  27. Terry Pratchett. Reaper Man: A Discworld Novel. Victor Gollancz, 1991. ISBN: 978-0-575-04979-6 ↩︎

  28. Peter Bailis, Alan Fekete, Michael J. Franklin, Ali Ghodsi, Joseph M. Hellerstein, and Ion Stoica. Coordination Avoidance in Database Systems. Proceedings of the VLDB Endowment, volume 8, issue 3, pages 185–196, November 2014. doi:10.14778/2735508.2735509 ↩︎

  29. Yaser Raja and Peter Celentano. PostgreSQL bi-directional replication using pglogical. aws.amazon.com, January 2022. Archived at https://perma.cc/BUQ2-5QWN ↩︎

  30. Robert Hodges. If You *Must* Deploy Multi-Master Replication, Read This First. scale-out-blog.blogspot.com, April 2012. Archived at perma.cc/C2JN-F6Y8 ↩︎ ↩︎

  31. Lars Hofhansl. HBASE-7709: Infinite Loop Possible in Master/Master Replication. issues.apache.org, January 2013. Archived at perma.cc/24G2-8NLC ↩︎

  32. John Day-Richter. What’s Different About the New Google Docs: Making Collaboration Fast. drive.googleblog.com, September 2010. Archived at perma.cc/5TL8-TSJ2 ↩︎ ↩︎

  33. Evan Wallace. How Figma’s multiplayer technology works. figma.com, October 2019. Archived at perma.cc/L49H-LY4D ↩︎

  34. Tuomas Artman. Scaling the Linear Sync Engine. linear.app, June 2023. ↩︎

  35. Amr Saafan. Why Sync Engines Might Be the Future of Web Applications. nilebits.com, September 2024. Archived at perma.cc/5N73-5M3V ↩︎

  36. Isaac Hagoel. Are Sync Engines The Future of Web Applications? dev.to, July 2024. Archived at perma.cc/R9HF-BKKL ↩︎

  37. Sujay Jayakar. A Map of Sync. stack.convex.dev, October 2024. Archived at perma.cc/82R3-H42A ↩︎

  38. Alex Feyerke. Designing Offline-First Web Apps. alistapart.com, December 2013. Archived at perma.cc/WH7R-S2DS ↩︎

  39. Martin Kleppmann, Adam Wiggins, Peter van Hardenberg, and Mark McGranaghan. Local-first software: You own your data, in spite of the cloud. At ACM SIGPLAN International Symposium on New Ideas, New Paradigms, and Reflections on Programming and Software (Onward!), October 2019, pages 154–178. doi:10.1145/3359591.3359737 ↩︎

  40. Martin Kleppmann. The past, present, and future of local-first. At Local-First Conference, May 2024. ↩︎

  41. Conrad Hofmeyr. API Calling is to Sync Engines as jQuery is to React. powersync.com, November 2024. Archived at perma.cc/2FP9-7WJJ ↩︎

  42. Peter van Hardenberg and Martin Kleppmann. PushPin: Towards Production-Quality Peer-to-Peer Collaboration. At 7th Workshop on Principles and Practice of Consistency for Distributed Data (PaPoC), April 2020. doi:10.1145/3380787.3393683 ↩︎

  43. Leonard Kawell, Jr., Steven Beckhardt, Timothy Halvorsen, Raymond Ozzie, and Irene Greif. Replicated document management in a group communication system. At ACM Conference on Computer-Supported Cooperative Work (CSCW), September 1988. doi:10.1145/62266.1024798 ↩︎

  44. Ricky Pusch. Explaining how fighting games use delay-based and rollback netcode. words.infil.net and arstechnica.com, October 2019. Archived at perma.cc/DE7W-RDJ8 ↩︎

  45. Giuseppe DeCandia, Deniz Hastorun, Madan Jampani, Gunavardhan Kakulapati, Avinash Lakshman, Alex Pilchin, Swaminathan Sivasubramanian, Peter Vosshall, and Werner Vogels. Dynamo: Amazon’s Highly Available Key-Value Store. At 21st ACM Symposium on Operating Systems Principles (SOSP), October 2007. doi:10.1145/1323293.1294281 ↩︎ ↩︎ ↩︎ ↩︎

  46. Marc Shapiro, Nuno Preguiça, Carlos Baquero, and Marek Zawirski. A Comprehensive Study of Convergent and Commutative Replicated Data Types. INRIA Research Report no. 7506, January 2011. ↩︎

  47. Chengzheng Sun and Clarence Ellis. Operational Transformation in Real-Time Group Editors: Issues, Algorithms, and Achievements. At ACM Conference on Computer Supported Cooperative Work (CSCW), November 1998. doi:10.1145/289444.289469 ↩︎

  48. Joseph Gentle and Martin Kleppmann. Collaborative Text Editing with Eg-walker: Better, Faster, Smaller. At 20th European Conference on Computer Systems (EuroSys), March 2025. doi:10.1145/3689031.3696076 ↩︎

  49. Dharma Shukla. Azure Cosmos DB: Pushing the frontier of globally distributed databases. azure.microsoft.com, September 2018. Archived at perma.cc/UT3B-HH6R ↩︎

  50. David K. Gifford. Weighted Voting for Replicated Data. At 7th ACM Symposium on Operating Systems Principles (SOSP), December 1979. doi:10.1145/800215.806583 ↩︎ ↩︎

  51. Heidi Howard, Dahlia Malkhi, and Alexander Spiegelman. Flexible Paxos: Quorum Intersection Revisited. At 20th International Conference on Principles of Distributed Systems (OPODIS), December 2016. doi:10.4230/LIPIcs.OPODIS.2016.25 ↩︎

  52. Joseph Blomstedt. Bringing Consistency to Riak. At RICON West, October 2012. ↩︎

  53. Peter Bailis, Shivaram Venkataraman, Michael J. Franklin, Joseph M. Hellerstein, and Ion Stoica. Quantifying eventual consistency with PBS. The VLDB Journal, volume 23, pages 279–302, April 2014. doi:10.1007/s00778-013-0330-1 ↩︎

  54. Colin Breck. Shared-Nothing Architectures for Server Replication and Synchronization. blog.colinbreck.com, December 2019. Archived at perma.cc/48P3-J6CJ ↩︎ ↩︎

  55. Jeffrey Dean and Luiz André Barroso. The Tail at Scale. Communications of the ACM, volume 56, issue 2, pages 74–80, February 2013. doi:10.1145/2408776.2408794 ↩︎

  56. Peng Huang, Chuanxiong Guo, Lidong Zhou, Jacob R. Lorch, Yingnong Dang, Murali Chintalapati, and Randolph Yao. Gray Failure: The Achilles’ Heel of Cloud-Scale Systems. At 16th Workshop on Hot Topics in Operating Systems (HotOS), May 2017. doi:10.1145/3102980.3103005 ↩︎

  57. Leslie Lamport. Time, Clocks, and the Ordering of Events in a Distributed System. Communications of the ACM, volume 21, issue 7, pages 558–565, July 1978. doi:10.1145/359545.359563 ↩︎ ↩︎

  58. D. Stott Parker Jr., Gerald J. Popek, Gerard Rudisin, Allen Stoughton, Bruce J. Walker, Evelyn Walton, Johanna M. Chow, David Edwards, Stephen Kiser, and Charles Kline. Detection of Mutual Inconsistency in Distributed Systems. IEEE Transactions on Software Engineering, volume SE-9, issue 3, pages 240–247, May 1983. doi:10.1109/TSE.1983.236733 ↩︎

  59. Nuno Preguiça, Carlos Baquero, Paulo Sérgio Almeida, Victor Fonte, and Ricardo Gonçalves. Dotted Version Vectors: Logical Clocks for Optimistic Replication. arXiv:1011.5808, November 2010. ↩︎

  60. Giridhar Manepalli. Clocks and Causality - Ordering Events in Distributed Systems. exhypothesi.com, November 2022. Archived at perma.cc/8REU-KVLQ ↩︎ ↩︎

  61. Sean Cribbs. A Brief History of Time in Riak. At RICON, October 2014. Archived at perma.cc/7U9P-6JFX ↩︎

  62. Russell Brown. Vector Clocks Revisited Part 2: Dotted Version Vectors. riak.com, November 2015. Archived at perma.cc/96QP-W98R ↩︎

  63. Carlos Baquero. Version Vectors Are Not Vector Clocks. haslab.wordpress.com, July 2011. Archived at perma.cc/7PNU-4AMG ↩︎

  64. Reinhard Schwarz and Friedemann Mattern. Detecting Causal Relationships in Distributed Computations: In Search of the Holy Grail. Distributed Computing, volume 7, issue 3, pages 149–174, March 1994. doi:10.1007/BF02277859 ↩︎