# 一致性与共识

LLMS 索引： [llms.txt](/llms.txt)

---

<a id="ch_consistency"></a>

![](/map/ch09.png)

> *古谚有云：“出海切勿带两只航海钟；要么带一只，要么带三只。”*
>
> 弗雷德里克・P・布鲁克斯，《人月神话：软件工程随笔》（1995）

正如[第 9 章](/ch9#ch_distributed)所述，分布式系统里可能出错的事情很多。要让服务在这些故障发生时仍能正确运行，就必须设法容忍故障。

*复制*（*replication*）是实现容错最有力的工具之一。然而，正如[第 6 章](/ch6#ch_replication)所示，把同一份数据复制到多个副本，也带来了不一致的风险。读请求可能由尚未追上进度的副本处理，返回陈旧结果；如果多个副本都能接受写入，还必须解决不同副本上并发写入的值之间的冲突。从总体上看，处理这类问题有两种彼此竞争的思路：

最终一致性（eventual consistency）
: 这种思路把系统采用复制这一事实暴露给应用，由应用开发者处理随之而来的不一致与冲突。采用[“多主复制”](/ch6#sec_replication_multi_leader)和[“无主复制”](/ch6#sec_replication_leaderless)的系统常常使用这种方式。

强一致性（strong consistency）
: 这种思路认为，应用不应操心复制的内部细节，系统应当表现得仿佛只有一个节点。它让应用开发者的工作简单得多，代价则是更强的一致性会损害性能，而且有些故障在最终一致系统中尚可容忍，却会令强一致系统停摆。

一如既往，哪种方式更好取决于具体应用。如果应用允许用户离线修改数据，那么正如[“同步引擎与本地优先软件”](/ch6#sec_replication_offline_clients)所述，最终一致性不可避免。然而，应用要正确处理最终一致性也很困难。如果各副本位于通信快速而可靠的数据中心，强一致性的成本通常可以接受，因而往往更合适。

本章将深入讨论强一致性，重点考察三个方面：

1. “强一致性”这个说法相当含糊，因此我们先给出一个更精确的目标：*线性一致性*（*linearizability*）。
2. 接着讨论 ID 和时间戳的生成。这个问题看似与一致性无关，实际上二者关系密切。
3. 最后探讨分布式系统如何既实现线性一致性，又保持容错能力；答案在于 *共识*（*consensus*）算法。

在此过程中，我们会看到，分布式系统中什么可以做到、什么无法做到，受到一些根本限制。

本章讨论的内容素以难以正确实现而著称。一个系统在没有故障时运行良好并不难，难的是它可能在某种设计者未曾考虑的不利故障组合下彻底崩溃。为帮助我们推理这些边界情况，研究者发展出了大量理论；借助这些理论，我们才能构建真正稳健的容错系统。

本章只能浅尝辄止：我们会采用非形式化的直观解释，避开算法的繁琐细节、形式化模型和证明。如果你打算认真从事共识系统或类似基础设施的工作，就必须深入掌握相关理论，否则很难让系统真正可靠。与往常一样，本章参考文献可以作为进一步学习的起点。



## 线性一致性 {#sec_consistency_linearizability}

要让复制数据库尽可能简单易用，最好让它表现得仿佛根本没有复制。这样，用户就不必操心复制延迟、冲突和其他不一致问题：既能获得容错的好处，又不必承担思考多个副本所带来的复杂性。

这就是 *线性一致性*（*linearizability*）[^1]（也称为 *原子一致性*，*atomic consistency* [^2]，*强一致性*，*strong consistency*，*即时一致性*，*immediate consistency*，或 *外部一致性*，*external consistency* [^3]）背后的思想。线性一致性的精确定义相当微妙，本节余下部分会逐步展开；其基本思想，是让系统看起来仿佛只有一份数据，所有操作都原子地作用于这份数据。有了这种保证，即使实际上存在多个副本，应用也无需关心它们。

在线性一致的系统中，只要一个客户端成功完成写入，此后所有客户端从数据库读取时，都必须能看到刚写入的值。要维持“只有一份数据”的假象，就必须保证读到的是最近写入的最新值，而不是来自陈旧缓存或副本的旧值。换句话说，线性一致性是一种 *新鲜度保证*（*recency guarantee*）。下面用一个不满足线性一致性的系统来说明这一点。

**图 10-1.** 如果这个数据库满足线性一致性，那么 Alice 的读取应返回 1 而不是 0，或者 Bob 的读取应返回 0 而不是 1。

![如果这个数据库满足线性一致性，那么 Alice 的读取应返回 1 而不是 0，或者 Bob 的读取应返回 0 而不是 1。](/fig/ddia_1001.png)

[图 10-1](/ch10/#fig_consistency_linearizability_0)展示了一个不满足线性一致性的体育网站 [^4]。Aaliyah 和 Bryce 坐在同一个房间里，都在用手机关注自己喜爱球队的比赛结果。终场比分刚刚公布，Aaliyah 刷新页面，看到了获胜方，便兴奋地告诉 Bryce。Bryce 将信将疑地按下自己手机上的 *刷新*，但请求被路由到一个落后的数据库副本，于是页面仍显示比赛正在进行。

如果两人同时刷新，得到不同结果倒不那么令人意外，因为谁也不知道服务器究竟在什么时刻处理了各自的请求。然而 Bryce 知道，自己是在听见 Aaliyah 喊出终场比分 *之后* 才按下刷新按钮、发起查询的，因此他有理由期待查询结果至少不比 Aaliyah 看到的更旧。结果却返回了陈旧数据，这就违反了线性一致性。

### 什么使系统具有线性一致性？ {#sec_consistency_lin_definition}

为了更好地理解线性一致性，我们再看几个例子。[图 10-2](/ch10/#fig_consistency_linearizability_1)展示了三个客户端如何并发读写线性一致数据库中的同一个对象 *x*。在分布式系统理论中，*x* 称为 *寄存器*（*register*）；在实际系统里，它可以是键值存储中的一个键、关系数据库中的一行，或文档数据库中的一个文档。

**图 10-2.** 如果读请求与写请求并发，则可能返回旧值，也可能返回新值。

![如果读请求与写请求并发，则可能返回旧值，也可能返回新值。](/fig/ddia_1002.png)


为简单起见，[图 10-2](/ch10/#fig_consistency_linearizability_1)只展示客户端看到的请求，不涉及数据库内部。每根横条代表客户端发出的一次请求：左端是请求发出的时刻，右端是客户端收到响应的时刻。由于网络延迟变化不定，客户端不知道数据库究竟何时处理了请求，只知道处理一定发生在请求发出与响应到达之间。

在这个例子中，寄存器有两种类型的操作：

* *read*(*x*) ⇒ *v* 表示客户端请求读取寄存器 *x*，数据库返回值 *v*。
* *write*(*x*, *v*) ⇒ *r* 表示客户端请求把寄存器 *x* 设为 *v*，数据库返回响应 *r*（可以是 *ok* 或 *error*）。

在[图 10-2](/ch10/#fig_consistency_linearizability_1)中，*x* 的初始值为 0，客户端 C 发出写请求，要把它改为 1。在此期间，客户端 A 和 B 不断轮询数据库，读取最新值。它们可能得到哪些响应？

* 客户端 A 的第一次读取在写入开始前就已完成，因此必然返回旧值 0。
* 客户端 A 的最后一次读取在写入完成后才开始，因此在线性一致的数据库中必然返回新值 1，因为该读取一定是在写入之后处理的。
* 凡是时间上与写操作重叠的读取，都可能返回 0 或 1，因为我们不知道数据库处理读取时，写入究竟是否已经生效。这些读取与写入是 *并发* 的。

但这还不足以完整描述线性一致性。如果与写入并发的读取可以任意返回旧值或新值，那么在写入期间，读者可能看到值在新旧之间来回跳变。这不符合我们对“只有一份数据”的系统的预期。

要使系统满足线性一致性，还需要增加一条约束，如[图 10-3](/ch10/#fig_consistency_linearizability_2)所示。

**图 10-3.** 如果 Alice 和 Bob 拥有完美时钟，线性一致性要求读取返回 x \= 1，因为对 x 的读取开始于 x \= 1 写入完成之后。

![如果 Alice 和 Bob 拥有完美时钟，线性一致性要求读取返回 x \= 1，因为对 x 的读取开始于 x \= 1 写入完成之后。](/fig/ddia_1003.png)


在线性一致的系统中，可以设想在写操作起止之间存在某个时刻，*x* 的值在那一刻原子地从 0 变为 1。因此，只要某个客户端已经读到新值 1，此后所有读取也都必须返回 1，即使写操作本身尚未结束。

[图 10-3](/ch10/#fig_consistency_linearizability_2)用箭头标出了这种时序依赖。客户端 A 最先读到新值 1；A 的读取刚一返回，B 就开始了新的读取。由于 B 的读取严格晚于 A 的读取，它也必须返回 1，哪怕 C 的写入仍在进行。（这与[图 10-1](/ch10/#fig_consistency_linearizability_0)中 Aaliyah 和 Bryce 的情形相同：Aaliyah 已经读到新值之后，Bryce 也理应读到新值。）

还可以进一步细化时序图，把每个操作视为在某个时刻原子生效 [^5]，如[图 10-4](/ch10/#fig_consistency_linearizability_3)这个更复杂的例子所示。除了 *read* 和 *write*，图中又加入了第三种操作：

* *cas*(*x*, *v* old, *v* new) ⇒ *r* 表示客户端请求执行原子 *比较并设置*（*compare-and-set*）操作（参见[“条件写入（比较并设置）”](/ch8#sec_transactions_compare_and_set)）。如果寄存器 *x* 的当前值等于 *v* old，就原子地把它改为 *v* new；否则保持不变并返回错误。*r* 是数据库的响应（*ok* 或 *error*）。

[图 10-4](/ch10/#fig_consistency_linearizability_3)在每次操作的横条内画了一根竖线，表示我们认为该操作实际生效的时刻。把这些标记依次连起来，必须得到寄存器的一条合法读写序列——每次读取都应返回最近一次写入所设置的值。

线性一致性要求，连接这些操作标记的线只能沿时间向前移动（从左向右），绝不能倒退。这项要求保证了前面所说的新鲜度：一旦新值已经被写入或读到，此后的读取就都必须看到这个值，直至它再次被覆盖。

**图 10-4.** 对 x 的读取与 x \= 1 的写入并发。由于不知道操作的确切时序，读取可以返回 0 或 1。

![对 x 的读取与 x \= 1 的写入并发。由于不知道操作的确切时序，读取可以返回 0 或 1。](/fig/ddia_1004.png)


[图 10-4](/ch10/#fig_consistency_linearizability_3)中有几个细节值得注意：

* 客户端 B 先发出读取 *x* 的请求，随后 D 请求把 *x* 设为 0，A 又请求把 *x* 设为 1；但 B 最终读到了 1，也就是 A 写入的值。这没有问题：它说明数据库先处理 D 的写入，再处理 A 的写入，最后才处理 B 的读取。这个顺序虽然不同于请求发出的顺序，却仍然合法，因为三次请求彼此并发。也许 B 的读请求在网络中耽搁了一会儿，直到两次写入之后才抵达数据库。
* 客户端 B 在 A 收到数据库确认“写入 1 成功”的响应之前，就已经读到了 1。这也没有问题，只说明数据库发给 A 的 *ok* 响应在网络中有所延迟。
* 这个模型不作任何事务隔离假设，其他客户端随时都可能修改值。例如，C 先读到 1，随后又读到 2，是因为两次读取之间 B 修改了该值。原子的比较并设置（*cas*）操作可以检查某个值是否已被其他客户端并发修改：B 和 C 的 *cas* 请求成功，而 D 的 *cas* 请求失败，因为数据库处理该请求时，*x* 已经不再等于 0。
* 客户端 B 最后一次读取（阴影横条）不满足线性一致性。该读取与 C 的 *cas* 写入并发，后者把 *x* 从 2 改为 4。若没有其他请求，B 返回 2 原本是允许的；但在 B 开始读取之前，客户端 A 已经读到了新值 4，因此 B 不能再读到比 A 更旧的值。这仍然是[图 10-1](/ch10/#fig_consistency_linearizability_0)中 Aaliyah 与 Bryce 的同一种情形。

以上就是线性一致性的直观含义，形式化定义 [^1] 对此有更精确的描述。可以记录所有请求与响应的时序，再检查它们能否排成一条合法的顺序序列，以此检验系统行为是否满足线性一致性；只是这种检验的计算成本很高 [^6] [^7]。

正如事务除了可串行化之外还有各种[“弱隔离级别”](/ch8#sec_transactions_isolation_levels)，复制系统除了线性一致性之外，也有许多较弱的一致性模型 [^8]。我们在[“复制延迟的问题”](/ch6#sec_replication_lag)中见过的 *写后读*、*单调读* 和 *一致前缀读*，就是这类较弱保证。线性一致性不仅包含所有这些保证，还提供得更多。本章将集中讨论线性一致性——实际系统中常用的最强一致性模型。



<a id="sidebar_consistency_serializability"></a>

> [!TIP] 线性一致性与可串行化
>
> 线性一致性很容易与[“可串行化”](/ch8#sec_transactions_serializability)混淆，因为两个名称看起来都像是在说“可以排成某种顺序”。但二者是完全不同的保证，必须加以区分：
>
> 可串行化
> : 可串行化是事务的一项隔离属性；每个事务可以读写 *多个对象*（行、文档或记录）。它保证事务的行为等同于按 *某种* 串行顺序执行：先完整执行一个事务，再完整执行下一个事务，彼此不交错。这个串行顺序可以不同于事务实际运行的顺序 [^9]。
>
> 线性一致性
> : 线性一致性是对寄存器（即 *单个对象*）读写的保证。它不会把操作组合成事务，因此无法防止[“写偏差与幻读”](/ch8#sec_transactions_write_skew)这类涉及多个对象的问题。但线性一致性是一项 *新鲜度* 保证：如果一个操作在另一个操作开始前已经结束，那么后一个操作必须观察到至少与前一个操作同样新的状态。可串行化没有这项要求，例如它允许陈旧读取 [^10]。
>
> （*顺序一致性* 又是另外一回事 [^8]，但我们不会在这里讨论它。）
>
> 数据库可以同时提供可串行化与线性一致性，这种组合称为 *严格可串行化*（*strict serializability*）或 *强单副本可串行化*（*strong one-copy serializability*，*strong-1SR*）[^11] [^12]。单节点数据库通常满足线性一致性。对于采用[“可串行化快照隔离（SSI）”](/ch8#sec_transactions_ssi)等乐观方法的分布式数据库，情况要复杂一些。例如，CockroachDB 提供可串行化以及一定的读取新鲜度保证，却不提供严格可串行化 [^13]，因为后者要求事务之间进行代价高昂的协调 [^14]。
>
> 也可以把较弱的隔离级别与线性一致性组合，或把较弱的一致性模型与可串行化组合。事实上，一致性模型与隔离级别在很大程度上可以独立选择 [^15] [^16]。

### 依赖线性一致性 {#sec_consistency_linearizability_usage}

线性一致性在什么情况下有用？查看体育比赛的终场比分也许只是个无关紧要的例子：结果陈旧几秒，通常不会造成实际损失。然而在少数领域，线性一致性却是系统正确工作的必要条件。

#### 锁定与领导者选举 {#locking-and-leader-election}

采用单主复制的系统必须确保领导者确实只有一个，而不是同时出现多个领导者（脑裂）。一种选举办法是使用租约：每个启动的节点都尝试获取租约，成功者成为领导者 [^17]。无论底层机制如何实现，都必须满足线性一致性，绝不能让两个不同节点同时获得同一份租约。

Apache ZooKeeper [^18]、etcd 等协调服务经常用于实现分布式租约和领导者选举。它们借助共识算法，以容错方式提供线性一致的操作（本章稍后会讨论这些算法）。要正确实现租约和领导者选举，还有许多微妙细节，例如[“分布式锁和租约”](/ch9#sec_distributed_lock_fencing)所述的栅栏问题。Apache Curator 等库在 ZooKeeper 之上封装了更高层的惯用方案，可以减轻这项工作；而所有这些协调任务的根基，仍是线性一致的存储服务。


> [!NOTE]
> 严格来说，ZooKeeper 的写入满足线性一致性，但读取可能陈旧，因为它不保证读请求一定由当前领导者处理 [^18]。从版本 3 开始，etcd 默认提供线性一致的读取。



一些分布式数据库还会在细得多的粒度上使用分布式锁，例如 Oracle Real Application Clusters（RAC）[^19]。RAC 为每个磁盘页设置一把锁，多个节点共享同一套磁盘存储。由于这些线性一致的锁位于事务执行的关键路径上，RAC 部署通常会使用专用的集群互连网络，让数据库节点彼此通信。

#### 约束与唯一性保证 {#sec_consistency_uniqueness}

唯一性约束在数据库中十分常见。例如，用户名或电子邮件地址必须唯一标识一位用户；文件存储服务也不能同时存在路径和文件名完全相同的两个文件。如果要在写入时强制执行这种约束——也就是说，两个人并发创建同名用户或文件时，必须让其中一人收到错误——就需要线性一致性。

这种情形其实与锁很相似：用户注册服务时，可以看作在获取所选用户名的“锁”。这个操作也很像原子的比较并设置：只要用户名尚未被占用，就把它设为认领该名称的用户 ID。

类似的问题还有很多：确保银行账户余额永不为负，商品销量不超过仓库库存，或不让两个人同时订到同一航班、同一剧场的同一个座位。这些约束都要求存在一个由所有节点共同认可的最新值，例如账户余额、库存数量或座位占用状态。

实际应用有时可以宽松处理此类约束。例如航班超售后，可以把旅客转到另一趟航班，并为造成的不便提供补偿。在这种情况下，线性一致性未必必要；[“及时性与完整性”](/ch13#sec_future_integrity)会进一步讨论这种宽松解释的约束。

不过，关系数据库里常见的硬性唯一约束确实需要线性一致性。外键约束、属性约束等其他约束则可以不依赖线性一致性来实现 [^20]。

#### 跨通道时序依赖 {#cross-channel-timing-dependencies}

请注意[图 10-1](/ch10/#fig_consistency_linearizability_0)中的一个细节：如果 Aaliyah 没有喊出比分，Bryce 就不会知道自己的查询结果已经陈旧。他只会在几秒后再次刷新，并最终看到终场比分。之所以能察觉违反线性一致性的现象，只是因为系统里还有一条额外的通信通道——Aaliyah 的声音传到了 Bryce 耳中。

计算机系统中也会出现类似情形。假设某网站允许用户上传视频，后台进程会把视频转码为画质较低的版本，以便通过慢速网络流式播放。系统架构和数据流如[图 10-5](/ch10/#fig_consistency_transcoder)所示。

视频转码器必须收到明确指令才会执行转码作业，这条指令由 Web 服务器通过消息队列发送（参见[“消息传递系统”](/ch12#sec_stream_messaging)）。Web 服务器不会把整段视频塞进队列，因为大多数消息代理是为短消息设计的，而视频可能有数十兆字节甚至更大。它会先把视频写入文件存储服务，确认写入完成后，再将转码指令放入队列。

**图 10-5.** 一个不满足线性一致性的系统：Alice 和 Bob 在不同时刻看到上传的图像，因此 Bob 的请求建立在陈旧数据之上。

![一个不满足线性一致性的系统：Alice 和 Bob 在不同时刻看到上传的图像，因此 Bob 的请求建立在陈旧数据之上。](/fig/ddia_1005.png)


如果文件存储服务满足线性一致性，这套系统就能正常工作；否则便可能出现竞态条件：消息队列（[图 10-5](/ch10/#fig_consistency_transcoder)中的步骤 3 和 4）也许比存储服务内部的复制传播得更快。这样一来，转码器获取原始视频时（步骤 5），可能读到文件的旧版本，甚至什么也读不到。如果它转码了旧版本，文件存储中的原始视频与转码版本就会永久不一致。

问题的根源在于，Web 服务器与转码器之间存在两条不同的通信通道：文件存储和消息队列。没有线性一致性提供的新鲜度保证，两条通道之间就可能发生竞态。这与[图 10-1](/ch10/#fig_consistency_linearizability_0)完全类似：一条通道是数据库复制，另一条则是从 Aaliyah 嘴里到 Bryce 耳中的现实声音。

能接收推送通知的移动应用也可能遇到类似竞态：应用收到通知后会向服务器获取相关数据；如果读请求可能落到滞后的副本上，推送通知也许很快就到了，紧随其后的数据读取却看不到通知所指的更新。

线性一致性不是避免这类竞态的唯一办法，却是最容易理解的一种。如果额外的通信通道由你控制——消息队列属于这种情况，Aaliyah 和 Bryce 之间的交流则不属于——也可以采用与[“读己之写”](/ch6#sec_replication_ryw)类似的替代方案，只是系统会更加复杂。


### 实现线性一致性系统 {#sec_consistency_implementing_linearizable}

看过线性一致性的几个用途后，接下来思考如何实现一个提供线性一致语义的系统。

线性一致性本质上要求系统“表现得仿佛只有一份数据，而且所有操作都原子地作用于它”，因此最简单的实现就是真的只保存一份数据。可惜这种方式无法容错：保存这份数据的节点一旦失效，数据就会丢失，至少也会在节点恢复前无法访问。

让我们重新审视[第 6 章](/ch6#ch_replication)介绍的复制方法，看看它们能否实现线性一致性：

单主复制（可能线性一致）
: 在单主复制系统中，领导者保存用于写入的主副本，其他节点上的追随者则维护备份。只要所有读写都由领导者处理，通常就有可能满足线性一致性。但这依赖一个前提：你必须确切知道谁是领导者。正如[“分布式锁和租约”](/ch9#sec_distributed_lock_fencing)所述，节点完全可能误以为自己仍是领导者；如果这个自以为是的领导者继续处理请求，就很可能破坏线性一致性 [^21]。采用异步复制时，故障切换甚至可能丢失已经提交的写入，同时违反持久性与线性一致性。

 对单主数据库进行分片、让每个分片拥有各自的领导者，不会影响线性一致性，因为它只保证单个对象。跨分片事务则是另一个问题（参见[“分布式事务”](/ch8#sec_transactions_distributed)）。

共识算法（很可能线性一致）
: 有些共识算法本质上是增加了自动领导者选举和故障切换的单主复制。它们经过精心设计以避免脑裂，因而能够安全实现线性一致的存储。例如，ZooKeeper 使用 Zab 共识算法 [^22]，etcd 使用 Raft [^23]。不过，系统采用了共识，并不等于它的所有操作都满足线性一致性：如果某个节点处理读取前没有确认自己仍是领导者，那么在刚刚选出新领导者时，它就可能返回陈旧结果。

多主复制（非线性一致）
: 多主复制系统通常不满足线性一致性，因为多个节点会并发处理写入，再把结果异步复制到其他节点。因此，它们可能产生需要[“处理写入冲突”](/ch6#sec_replication_write_conflicts)的并发写入。

无主复制（很可能不满足线性一致性）
: 对于采用无主复制的系统（Dynamo 风格，参见[“无主复制”](/ch6#sec_replication_leaderless)），有人声称，只要要求仲裁读写满足 *w* + *r* > *n*，就能得到“强一致性”。这取决于具体算法以及“强一致性”的定义，但通常并不准确。

 Cassandra 和 ScyllaDB 等系统采用基于日历时钟的“最后写入者胜”来解决冲突，这几乎肯定不满足线性一致性，因为时钟偏差使时间戳无法保证与事件的实际顺序一致（参见[“对同步时钟的依赖”](/ch9#sec_distributed_clocks_relying)）。即使采用仲裁读写，仍然可能出现违反线性一致性的行为，下一节会给出例子。

#### 线性一致性与仲裁 {#sec_consistency_quorum_linearizable}

直觉上，Dynamo 风格模型中的仲裁读写似乎应当满足线性一致性。但当网络延迟变化不定时，仍然可能出现竞态，如[图 10-6](/ch10/#fig_consistency_leaderless)所示。

**图 10-6.** 当网络延迟变化不定时，仅靠法定人数不足以保证线性一致性。

![当网络延迟变化不定时，仅靠法定人数不足以保证线性一致性。](/fig/ddia_1006.png)


在[图 10-6](/ch10/#fig_consistency_leaderless)中，*x* 的初始值为 0。一个写入客户端把请求发往全部三个副本（*n* = 3，*w* = 3），要将 *x* 更新为 1。与此同时，客户端 A 从两个节点读取，达到读法定人数（*r* = 2），并在其中一个节点上看到了新值 1；同样与写入并发的客户端 B，则从另外两个节点读取，两个节点都返回旧值 0。

尽管满足法定人数条件 *w* + *r* > *n*，这次执行仍不满足线性一致性：B 的请求开始于 A 的请求完成之后，却返回了旧值，而 A 已经读到新值。（这又是[图 10-1](/ch10/#fig_consistency_linearizability_0)中 Aaliyah 和 Bryce 的情形。）

可以让 Dynamo 风格的法定人数读写满足线性一致性，但要牺牲性能：读取方必须先同步完成[“追赶错过的写入”](/ch6#sec_replication_read_repair)所述的读修复，再把结果返回应用 [^24]；写入方则必须在写入前先读取达到法定人数的节点的最新状态，取得此前所有写入中的最大时间戳，并确保新写入使用更大的时间戳 [^25] [^26]。Riak 因为性能代价而不执行同步读修复。Cassandra 的仲裁读确实会等待读修复完成 [^27]，但它使用日历时钟生成时间戳，因而仍然不满足线性一致性。

而且，这种方式只能实现线性一致的读写；无法实现线性一致的比较并设置，因为后者需要共识算法 [^28]。

总之，最稳妥的假设是：采用 Dynamo 风格复制的无主系统即使使用仲裁读写，也不提供线性一致性。

### 线性一致性的代价 {#sec_linearizability_cost}

既然有些复制方式能够提供线性一致性，有些不能，我们就有必要更仔细地考察它的利弊。

[第 6 章](/ch6#ch_replication)已经讨论过不同复制方式的适用场景。例如，对多地区复制而言，多主复制往往是不错的选择（参见[“跨地域运行”](/ch6#sec_replication_multi_dc)）。[图 10-7](/ch10/#fig_consistency_cap_availability)展示了这样一种部署。

**图 10-7.** 如果网络分区使客户端无法联系足够多的副本，它们就无法处理写入。

![如果网络分区使客户端无法联系足够多的副本，它们就无法处理写入。](/fig/ddia_1007.png)


考虑两个地区之间网络中断时会发生什么。假设各地区内部的网络仍然正常，客户端也能访问本地区域，但两个地区彼此无法通信。这种情况称为 *网络分区*（*network partition*）。

在多主数据库中，每个地区都可以继续正常工作：一地的写入原本就异步复制到另一地，因此网络中断期间只需暂存排队，连接恢复后再相互交换。

若采用单主复制，领导者必然位于其中一个地区。所有写入和线性一致读取都必须发给领导者；因此，连接到追随者所在地区的客户端，必须跨地区同步地把读写请求发送到领导者所在地区。

单主配置下，一旦地区间网络中断，追随者所在地区的客户端便无法联系领导者，因而既不能写入数据库，也不能执行线性一致读取。它们仍可从追随者读取，但结果可能陈旧，不满足线性一致性。如果应用要求线性一致的读写，那么所有无法联系领导者的地区都会在网络中断期间变得不可用。

能直接连接领导者所在地区的客户端不受影响，应用在那里仍可正常工作；但只能访问追随者所在地区的客户端会一直停摆，直到网络链路修复。

#### CAP 定理 {#the-cap-theorem}

这个问题并非单主复制与多主复制特有：任何线性一致的数据库，无论怎样实现，都会面对同样的困境。它也不限于多地区部署；任何不可靠的网络都可能发生这种情况，即使是在同一地区内。具体权衡如下：

* 如果应用 *要求* 线性一致性，而网络故障使一部分副本与其他副本失去联系，那么断开的副本就不能继续处理请求：它们只能等待网络恢复，或立即返回错误；无论哪种方式，都变得 *不可用*。这种选择有时称为 *CP*（网络分区时保持一致）。
* 如果应用 *不要求* 线性一致性，就可以让各副本在彼此断开时仍独立处理请求，例如采用多主复制。这样，应用在网络故障期间仍然 *可用*，但行为不满足线性一致性。这种选择称为 *AP*（网络分区时保持可用）。

因此，不要求线性一致性的应用可以更好地容忍网络问题。这一认识通常称为 *CAP 定理* [^29] [^30] [^31] [^32]，由 Eric Brewer 于 2000 年命名，不过早在 20 世纪 70 年代，分布式数据库设计者就已了解这种权衡 [^33] [^34] [^35]。

CAP 最初只是一个没有精确定义的经验法则，目的是引发对数据库权衡的讨论。当时许多分布式数据库都专注于在共享存储的机器集群上提供线性一致语义 [^19]；CAP 则鼓励数据库工程师探索更广阔的分布式无共享系统设计空间，而后者更适合承载大规模 Web 服务 [^36]。CAP 推动了这种观念转变，也帮助催生了 NoSQL 运动以及 21 世纪头十年中期涌现的大批新型数据库技术。

> [!TIP] 帮不上忙的 CAP 定理
>
> CAP 有时被概括成 *一致性、可用性、分区容错性，三者择二*。遗憾的是，这种说法会误导人 [^32]：网络分区是一类故障，并不是可以自由取舍的选项——不管你愿不愿意，它都会发生。
>
> 网络正常时，系统完全可以同时提供一致性（线性一致性）与全面可用性；发生网络故障时，才必须在线性一致性和全面可用性之间取舍。因此，更准确的说法是：*发生分区时，要么一致，要么可用* [^37]。网络越可靠，面临这种选择的次数就越少，但它终究无法彻底避免。
>
> CP/AP 分类还有几个缺陷 [^4]。其中的 *一致性* 被严格定义为线性一致性，定理对较弱的一致性模型只字未提；*可用性* 的形式化定义 [^30] 也不符合这个词的通常含义 [^38]。许多通常认为高度可用、具备容错能力的系统，其实并不满足 CAP 那套特殊的可用性定义。还有些系统设计者出于充分理由，既不提供线性一致性，也不提供 CAP 所假设的那种可用性，因此既不能归为 CP，也不能归为 AP [^39] [^40]。
>
> 总而言之，围绕 CAP 的误解与混淆太多，它无助于我们更好地理解系统，因此最好避免使用 CAP 这一框架。
>
> 形式化的 CAP 定理 [^30] 适用范围极窄：它只考察一种一致性模型（线性一致性）和一种故障（网络分区；Google 的数据表明，网络分区造成的事故不到 8% [^41]），对网络延迟、节点失效以及其他权衡都没有说明。因此，CAP 虽然在历史上影响深远，对实际系统设计却几乎没有指导价值 [^4] [^38]。
>
> 有人尝试把 CAP 推广到更一般的情形。例如，*PACELC 原则* 指出，即便网络正常，系统设计者也可能为了降低延迟而削弱一致性 [^39] [^40] [^42]：发生网络分区（P）时，需要在可用性（A）与一致性（C）之间选择；否则（E），没有分区时，则可能在低延迟（L）与一致性（C）之间选择。不过，这个定义继承了 CAP 的若干问题，例如对一致性和可用性的定义仍然有违直觉。
>
> 分布式系统领域还有许多更有意义的不可能性结论 [^43]，CAP 也早已被更精确的结果取代 [^44] [^45]。如今，它主要只剩下历史意义。

#### 线性一致性与网络延迟 {#linearizability-and-network-delays}

线性一致性虽然很有用，实际满足它的系统却少得出人意料。甚至现代多核 CPU 上的 RAM 也不满足线性一致性 [^46]：一个 CPU 核上的线程写入某个内存地址后，另一个核上的线程稍晚读取同一地址，也不保证能看见前一个线程写入的值，除非使用 *内存屏障* 或 *栅栏* [^47]。

原因在于，每个 CPU 核都有自己的缓存和存储缓冲区。默认情况下，内存访问会先经过缓存，修改再异步写回主存。访问缓存远快于访问主存 [^48]，因此这种机制对现代 CPU 的性能不可或缺。但系统里也因此出现了多份数据——主存中一份，各级缓存里可能还有几份——而且它们异步更新，于是失去了线性一致性。

为什么要作这种取舍？用 CAP 定理解释多核处理器的内存一致性模型毫无意义：在同一台计算机内，我们通常假定通信可靠，也不指望某个 CPU 核与计算机其他部分断开后还能正常工作。这里牺牲线性一致性是为了 *性能*，而不是容错 [^39]。

许多不提供线性一致性保证的分布式数据库也是如此：主要目的是提高性能，而非增强容错 [^42]。线性一致性很慢，而且始终如此，并非只在网络故障时才慢。

有没有更高效的线性一致存储实现？答案似乎是否定的。Attiya 和 Welch [^49] 证明：要获得线性一致性，读写请求的响应时间至少与网络延迟的不确定程度成正比。在大多数计算机网络这种延迟高度不稳定的环境里（参见[“超时和无界延迟”](/ch9#sec_distributed_queueing)），线性一致读写的响应时间必然很高。不存在更快的线性一致性算法，而较弱的一致性模型却可以快得多，因此这种取舍对延迟敏感型系统十分重要。[“及时性与完整性”](/ch13#sec_future_integrity)将讨论如何在不牺牲正确性的前提下绕开线性一致性。


## ID 生成器和逻辑时钟 {#sec_consistency_logical}

许多应用在创建数据库记录时，都要为其分配某种唯一 ID，作为以后引用该记录的主键。单节点数据库通常使用自增整数，它的优点是只需 64 位即可存储；如果能确定记录数永远不会超过 40 亿，甚至可以只用 32 位，不过这样做颇有风险。

自增 ID 还有一个好处：ID 顺序可以反映记录的创建顺序。例如，[图 10-8](/ch10/#fig_consistency_id_generator)中的聊天应用会在消息发出时为其分配自增 ID。按 ID 递增顺序展示消息，对话就能保持合理的先后关系：Aaliyah 的问题得到 ID 1，而 Bryce 随后的回答得到更大的 ID 3。

**图 10-8.** 两个不同节点可能生成相互冲突的 ID。

![两个不同节点可能生成相互冲突的 ID。](/fig/ddia_1008.png)


这种单节点 ID 生成器也是一个线性一致系统。每次取 ID 都会原子地递增计数器并返回递增前的值，这称为 *获取并增加*（fetch-and-add）操作。线性一致性保证：如果 Aaliyah 的消息在 Bryce 开始发消息之前已经发布完成，那么 Bryce 的 ID 必须更大。[图 10-8](/ch10/#fig_consistency_id_generator)中 Aaliyah 与 Caleb 的消息彼此并发，因此线性一致性不规定二者的 ID 顺序，只要求它们互不相同。

内存中的单节点 ID 生成器很容易实现：直接使用 CPU 提供的原子递增指令，就能让多个线程安全地更新同一个计数器。要把计数器做成持久的稍微麻烦一些，否则节点崩溃重启后计数器会复位，产生重复 ID。但更棘手的问题是：

* 单节点 ID 生成器不具备容错能力，因为这个节点本身就是单点故障。
* 如果要在另一个地区创建记录，仅仅为了取得 ID，可能就得跨越半个地球往返一次，速度很慢。
* 写入吞吐量很高时，这个节点可能成为瓶颈。

ID 生成器还有几种替代方案：

分片 ID 分配
: 可以让多个节点分别分配 ID，例如一个只生成偶数，另一个只生成奇数。更一般地，可以在 ID 中预留若干位来保存分片编号。这种 ID 仍然紧凑，却失去了顺序含义：看到 ID 为 16 和 17 的两条聊天消息，并不能断定消息 16 先发出，因为两个 ID 来自不同节点，而其中一个节点的进度可能领先于另一个。

预分配 ID 块
: 单节点生成器不必逐个发放 ID，也可以一次分配一整块。例如，节点 A 取得 1 到 1,000，节点 B 取得 1,001 到 2,000；此后各节点可在自己的区间内独立发号，快用完时再申请下一块。但这种方案同样无法保证正确顺序：一条消息可能先从 1,001 到 2,000 的区间取得 ID，而稍后另一节点发出的消息却从 1 到 1,000 的区间取得了更小的 ID。

随机 UUID
: 可以采用 *通用唯一标识符*（UUID），也称 *全局唯一标识符*（GUID）。它最大的好处是，任意节点都能在本地生成，无需通信；代价是占用更多空间（128 位）。UUID 有多个版本，最简单的第 4 版本质上是一个足够长的随机数，两个节点碰巧选中同一值的概率微乎其微。可惜这类 ID 的顺序也是随机的，比较两个 ID 无法判断哪个更新。

为日历时钟时间戳补充唯一性信息
: 如果各节点通过 NTP 让日历时钟大致准确，可以把时间戳放在 ID 的高位，再用额外信息填充其余位，保证即使时间戳相同，完整 ID 仍然唯一。例如，可以加入分片编号和分片内自增序列号，或一段足够长的随机值。第 7 版 UUID [^50]、Twitter Snowflake [^51]、ULID [^52]、Hazelcast Flake ID 生成器、MongoDB ObjectID 等许多方案都采用这种思路 [^50]。这类 ID 生成器既可以在应用代码中实现，也可以放在数据库内部 [^53]。

这些方案都能生成唯一 ID——至少碰撞概率低到几乎可以忽略——但它们提供的顺序保证远弱于单节点自增方案。

正如[“用于事件排序的时间戳”](/ch9#sec_distributed_lww)所述，日历时钟时间戳至多只能给出近似顺序：如果较早的写入读取了略快的时钟，较晚的写入读取了略慢的时钟，时间戳顺序就可能与事件实际顺序相反。非单调时钟还可能突然跳变，甚至让同一节点生成的时间戳顺序出错。因此，基于日历时钟的 ID 生成器通常不满足线性一致性。

使用原子钟或 GPS 接收器进行高精度时钟同步，可以减少这种顺序错乱。但如果无需特殊硬件，也能生成唯一且顺序正确的 ID，当然更好。这正是 *逻辑时钟*（*logical clock*）要解决的问题。

### 逻辑时钟 {#sec_consistency_timestamps}

在[“不可靠的时钟”](/ch9#sec_distributed_clocks)中，我们讨论了日历时钟和单调时钟。二者都属于 *物理时钟*，度量的是经过了多少秒（或毫秒、微秒等）。

分布式系统还经常使用另一类时钟，称为 *逻辑时钟*（*logical clock*）。物理时钟是计算流逝秒数的硬件设备；逻辑时钟则是一种计算已发生事件数量的算法。因此，逻辑时间戳不能告诉你现在是几点，却 *可以* 相互比较，判断哪个较早、哪个较晚。

逻辑时钟的要求通常是：

* 时间戳紧凑（只有几个字节）且唯一；
* 任意两个时间戳都可以比较，也就是构成 *全序*；
* 时间戳顺序与因果关系 *一致*：如果操作 A 先于 B 发生，那么 A 的时间戳小于 B 的时间戳。（我们在[“‘先发生’关系与并发”](/ch6#sec_replication_happens_before)中讨论过因果关系。）

单节点 ID 生成器满足这些要求，而前面几种分布式 ID 生成器不满足因果顺序要求。

#### Lamport 时间戳 {#lamport-timestamps}

幸运的是，有一种简单方法可以生成与因果关系 *一致* 的逻辑时间戳，并将其用作分布式 ID。这就是 Leslie Lamport 于 1978 年提出的 *Lamport 时钟* [^54]；介绍它的论文如今已成为分布式系统领域引用次数最多的论文之一。

[图 10-9](/ch10/#fig_consistency_lamport_ts)展示了 Lamport 时钟如何用于[图 10-8](/ch10/#fig_consistency_id_generator)中的聊天示例。每个节点都有唯一标识符；[图 10-9](/ch10/#fig_consistency_lamport_ts)使用“Aaliyah”“Bryce”和“Caleb”作为标识，实际系统则可以使用随机 UUID 等值。此外，每个节点维护一个计数器，记录自己处理过多少次操作。Lamport 时间戳就是一个二元组（*计数器*，*节点 ID*）。不同节点的计数器值有时相同，但加入节点 ID 后，每个时间戳仍然唯一。

**图 10-9.** Lamport 时间戳给出了与因果关系一致的全序。

![Lamport 时间戳给出了与因果关系一致的全序。](/fig/ddia_1009.png)


节点每生成一次时间戳，都会先递增本地计数器并使用新值。节点每次看到其他节点生成的时间戳时，如果其中的计数器值大于自己的本地值，就把本地计数器向前推进到同一个值。

在[图 10-9](/ch10/#fig_consistency_lamport_ts)中，Aaliyah 发出自己的消息时尚未看见 Caleb 的消息，Caleb 也一样。假设两人的计数器初始值都是 0，他们各自将其递增为 1，并把新值附在消息上。Bryce 收到这两条消息后，把自己的计数器推进到 1；随后他回复 Aaliyah 的消息，再把本地计数器递增为 2，并将 2 附在回复上。

比较两个 Lamport 时间戳时，先比较计数器值。例如，(2, “Bryce”) 大于 (1, “Aaliyah”)，也大于 (1, “Caleb”)。若计数器相同，再按通常的字符串字典序比较节点 ID。因此，本例中的时间戳顺序为 (1, “Aaliyah”) < (1, “Caleb”) < (2, “Bryce”)。

#### 混合逻辑时钟 {#hybrid-logical-clocks}

Lamport 时间戳很适合表示事件发生顺序，但也有一些局限：

* 它与物理时间没有直接关系，所以无法据此查找某个具体日期发布的所有消息；物理时间必须另行保存。
* 如果两个节点从不通信，一个节点的计数器增长就永远不会反映到另一个节点上。因此，不同节点在大致同一时刻生成的事件，计数器值可能相差悬殊。

*混合逻辑时钟*（*hybrid logical clock*，HLC）兼具物理日历时钟的优点与 Lamport 时钟的顺序保证 [^55]。它像物理时钟一样以秒或微秒计数；又像 Lamport 时钟一样，在看到其他节点更大的时间戳时，把自己的本地值向前推进到对方的时间戳。因此，如果某个节点的时钟偏快，其他节点与它通信后也会相应地把时钟向前推进。

混合逻辑时钟每次生成时间戳时也会递增，从而保证始终单调向前，即使底层物理时钟因 NTP 校正而向后跳变也不例外。因此，混合逻辑时钟可能略微领先于底层物理时钟；算法会尽量把这项误差控制在最小范围内。

因此，混合逻辑时间戳几乎可以像普通日历时间戳一样使用，同时还多了一项性质：其顺序与先发生关系一致。它不依赖特殊硬件，只要求各时钟大致同步。CockroachDB 就使用了混合逻辑时钟。

#### Lamport/混合逻辑时钟 vs. 向量时钟 {#lamporthybrid-logical-clocks-vs-vector-clocks}

在[“多版本并发控制（MVCC）”](/ch8#sec_transactions_snapshot_impl)中，我们讨论过快照隔离的一种常见实现：为每个事务分配事务 ID，让它看见 ID 较小的事务所做的写入，同时隐藏 ID 较大的事务所做的写入。Lamport 时钟和混合逻辑时钟很适合生成这些事务 ID，因为它们可以保证快照与因果关系一致 [^56]。

多个时间戳并发生成时，这些算法会任意规定它们之间的顺序。因此，只看两个时间戳，通常无法判断二者是并发生成，还是一个先于另一个发生。（在[图 10-9](/ch10/#fig_consistency_lamport_ts)中，因为 Aaliyah 与 Caleb 的消息具有相同计数器值，可以断定它们彼此并发；但计数器值不同时，就无法作出这种判断。）

如果需要判断记录是否并发创建，就得采用另一种算法，例如 *向量时钟*（*vector clock*）。它的缺点是时间戳大得多，可能需要为系统中的每个节点保存一个整数。有关并发检测的更多细节，参见[“检测并发写入”](/ch6#sec_replication_concurrent)。

### 线性一致的 ID 生成器 {#sec_consistency_linearizable_id}

Lamport 时钟与混合逻辑时钟虽然提供了实用的顺序保证，却仍弱于前面那种线性一致的单节点 ID 生成器。回想一下，线性一致性要求：只要请求 A 在请求 B 开始前已经完成，B 的 ID 就必须更大，即使二者从未彼此通信。Lamport 时钟只能保证节点新生成的时间戳大于它此前 *见过* 的所有时间戳，对未曾见过的时间戳则无法作出任何保证。

[图 10-10](/ch10/#fig_consistency_permissions)展示了不满足线性一致性的 ID 生成器会造成什么问题。假设某个社交网站的用户 A 想把一张难为情的照片只分享给朋友。A 的账户原本公开；A 先在笔记本电脑上把账户设为私密，随后又用手机上传照片。这些操作由 A 依次完成，因此 A 有理由认为，上传的照片会受新的私密账户权限保护。

**图 10-10.** 一个使用 Lamport 时间戳的权限系统。

![一个使用 Lamport 时间戳的权限系统。](/fig/ddia_1010.png)


账户权限与照片保存在两个独立数据库中（也可能是同一数据库的不同分片），并假设二者都用 Lamport 时钟或混合逻辑时钟为写入分配时间戳。照片数据库没有读取账户数据库，因此它的本地计数器可能稍微落后，导致照片上传得到的时间戳反而小于账户设置更新的时间戳。

接着，假设一位并非 A 好友的访客正在浏览 A 的个人资料，而这次读取由实现快照隔离的 MVCC 处理。访客读取的时间戳可能大于照片上传，却小于账户权限更新。系统于是认定，在这个快照中账户仍是公开的，并把本不该让访客看到的照片显示出来。

这个问题有几种可能的修复方式。也许照片数据库应在写入前读取用户账户状态，但开发者很容易漏掉这项检查。如果 A 的所有操作都来自同一台设备，设备上的应用或许能跟踪该用户最近一次写入的时间戳；可本例中用户同时使用笔记本电脑和手机，事情就没那么简单。

这里最简单的解决方案，是使用线性一致的 ID 生成器，保证照片上传取得的 ID 一定大于账户权限变更的 ID。

#### 实现线性一致的 ID 生成器 {#implementing-a-linearizable-id-generator}

确保 ID 分配满足线性一致性的最简单办法，确实是使用单个节点。这个节点只需在收到请求时原子递增计数器并返回结果，同时持久化计数器值，以免崩溃重启后产生重复 ID，再用单主复制实现容错。实际系统也采用这种方案：受 Google Percolator [^57] 启发，TiDB/TiKV 将其称为 *时间戳预言机*（timestamp oracle）。

可以通过批量预留来避免每次请求都写盘和复制。ID 生成器先写入一条记录，声明某一批 ID；这条记录持久化并完成复制后，节点便可按顺序把其中的 ID 发给客户端。在这一批快要用完前，再提前持久化并复制下一批。这样，节点崩溃重启或故障切换到追随者时，可能会跳过一些 ID，却绝不会发出重复或顺序错误的 ID。

ID 生成器很难分片：多个分片独立发号后，就无法再保证整体顺序满足线性一致性。它也很难分布到多个地区，因此地理分布式数据库中的所有 ID 请求，都得发往某个固定地区的节点。好在 ID 生成器的工作十分简单，单个节点就能承担很高的请求吞吐量。

如果不想采用单节点 ID 生成器，还有 Google Spanner 的做法，参见[“用于全局快照的同步时钟”](/ch9#sec_distributed_spanner)。它依赖一种物理时钟：返回的不是单个时间戳，而是一个时间戳区间，表示时钟读数的不确定性；系统会等待这个不确定区间完全过去之后再返回结果。

只要不确定区间确实可靠——真实物理时间始终落在区间内——这个过程同样能保证：若一个请求在另一个请求开始前完成，后一个请求就取得更大的时间戳。它无需节点间通信就能实现线性一致的 ID 分配；即使请求来自不同地区，也能正确排序，不必等待跨地区往返。代价是必须有相应的软硬件支持，既要让时钟高度同步，又要计算必要的不确定区间。

#### 使用逻辑时钟强制约束 {#enforcing-constraints-using-logical-clocks}

在[“约束与唯一性保证”](/ch10#sec_consistency_uniqueness)中，我们看到，线性一致的比较并设置可以用来实现分布式锁、唯一性约束等构造。于是自然会问：逻辑时钟或线性一致的 ID 生成器，是否也足以实现这些功能？

答案是：还不够。若多个节点都在争抢同一把锁或注册同一个用户名，可以用逻辑时钟为请求分配时间戳，并选出时间戳最小者。如果时钟满足线性一致性，那么此后所有请求生成的时间戳都更大，也就不可能再出现比当前胜者更小的时间戳。

可惜问题还有一半没有解决：节点怎样知道自己的时间戳就是最小的？要做到确定无疑，它必须收到 *每一个* 可能生成时间戳的其他节点的消息 [^54]。只要其中一个节点失效，或因网络问题无法联系，整个系统就会停滞，因为谁也无法排除那个节点拥有更小时间戳的可能。这显然不是我们想要的容错系统。

要以容错方式实现锁、租约等构造，需要比逻辑时钟或 ID 生成器更强的工具：我们需要共识。



## 共识 {#sec_consistency_consensus}

本章已经见过好几个例子：只用单个节点时，事情十分简单；一旦还要求容错，就会困难得多：

* 只设一个领导者，并让所有读写都由它处理，数据库便可以具有线性一致性。可是，如果这个领导者失效，怎样才能在避免脑裂的同时完成故障切换？怎样确保某个自认为仍是领导者的节点，其实没有在此期间被投票罢免？
* 单节点上的线性一致 ID 生成器，不过是一个带有原子“获取并增加”指令的计数器；但如果这个节点崩溃了呢？
* 原子比较并设置（CAS）操作用途广泛：例如，多个进程争抢锁或租约时决定谁能获得它，或者确保给定名称的文件或用户具有唯一性。在单个节点上，CAS 可能只需一条 CPU 指令；但怎样才能让它容错？

事实证明，这些都是同一个分布式系统基本问题的不同实例：*共识*（*consensus*）。共识是分布式计算中最重要、最基本的问题之一；同时，它也出了名地难以正确实现 [^58] [^59]，许多系统都曾在这里栽过跟头。至此，我们已经讨论过复制（[第 6 章](/ch6#ch_replication)）、事务（[第 8 章](/ch8#ch_transactions)）、系统模型（[第 9 章](/ch9#ch_distributed)）以及线性一致性（本章），终于可以正面处理共识问题了。

最著名的共识算法包括视图戳复制（Viewstamped Replication，VSR）[^60] [^61]、Paxos [^58] [^62] [^63] [^64]、Raft [^23] [^65] [^66] 和 Zab [^18] [^22] [^67]。这些算法有不少相似之处，但并不完全相同 [^68] [^69]。它们采用非拜占庭系统模型：网络通信可以被任意延迟或丢弃，节点可以崩溃、重启或断开连接；但除此之外，算法假定节点都会正确遵守协议，不会采取恶意行为。

另一些共识算法可以容忍一部分拜占庭节点，也就是不正确遵守协议的节点（例如，它们会向不同节点发送相互矛盾的消息）。这类算法通常假定，出现拜占庭故障的节点少于三分之一 [^26] [^70]。这样的 *拜占庭容错*（BFT）共识算法会用于区块链 [^71]。不过，正如 [“拜占庭故障”](/ch9#sec_distributed_byzantine) 所述，BFT 算法不在本书的讨论范围之内。

> [!TIP] 共识的不可能性
>
> 你可能听说过 FLP 结果 [^72]——这个名字取自三位作者 Fischer、Lynch 和 Paterson。它证明：只要存在节点崩溃的可能，就没有一种算法能保证 *始终* 达成共识。分布式系统必须假定节点可能崩溃，这岂不是意味着可靠的共识根本不可能？可我们现在又在讨论实现共识的算法，这究竟是怎么回事？
>
> 首先，FLP 并没有说共识永远无法达成，只是说我们无法保证共识算法 *每次都能* 终止。此外，FLP 的证明针对异步系统模型中的确定性算法（见 [“系统模型与现实”](/ch9#sec_distributed_system_model)），也就是说，算法不能使用任何时钟或超时。只要允许算法根据超时怀疑另一个节点可能已经崩溃——哪怕这种怀疑偶尔会出错——共识就变得可以解决 [^73]。甚至只需允许算法使用随机数，也足以绕过这个不可能性结果 [^74]。
>
> 因此，尽管 FLP 的共识不可能性结果在理论上极为重要，分布式系统在实践中通常仍能达成共识。

### 共识的多面性 {#sec_consistency_faces}

共识可以用几种不同的方式表达：

* *单值共识*（*single-value consensus*）与原子 *比较并设置* 操作非常相似，可以用来实现锁、租约和唯一性约束。
* 构建 *仅追加日志*（*append-only log*）同样需要共识；这个问题通常形式化为 *全序广播*（*total order broadcast*）。有了日志，就可以构建 *状态机复制*（*state machine replication*）、基于领导者的复制、事件溯源以及其他许多有用的机制。
* 多数据库或多分片事务的 *原子提交*（*atomic commitment*），要求所有参与者就是否提交或中止事务达成一致。

下面很快就会逐一探讨这些形式。事实上，这几个问题彼此等价：只要有一种算法能解决其中一个问题，就可以把它转换成其他任意一种问题的解法。这是一个相当深刻、甚至有些出人意料的洞见！也正因为如此，尽管它们表面上截然不同，我们仍可以把它们统统归入“共识”这一范畴。

#### 单值共识 {#single-value-consensus}

共识的标准表述涉及让多个节点就单个值达成一致。例如：

* 采用单主复制的数据库初次启动时，或现有领导者失效时，可能有多个节点同时试图成为领导者。同样，多个节点也可能竞相获取同一把锁或同一份租约。共识可以帮助它们决定由谁胜出。
* 如果几个人同时试图预订飞机上的最后一个座位、剧院里的同一个座位，或用相同的用户名注册账户，共识算法可以确定哪一个请求应当成功。

更一般地说，一个或多个节点可以 *提议* 某个值，而共识算法从中 *决定* 一个值。在上述例子中，每个节点都可以提议自己的 ID，算法则决定哪个节点 ID 应当成为新的领导者、租约持有者或机票／戏票的购买者。在这种形式化定义下，共识算法必须满足以下属性 [^26]：

一致同意
: 任意两个节点都不会作出不同的决定。

完整性
: 节点一旦决定了某个值，就不能再决定另一个值来改变主意。

有效性
: 如果某个节点决定了值 *v*，那么 *v* 必须由某个节点提议过。

终止
: 每个未崩溃的节点最终都能决定某个值。

如果需要决定多个值，可以为每个值分别运行一个共识算法实例。例如，可以为剧院中每个可预订的座位分别运行一次共识，从而为每个座位作出一项决定（确定一位买家）。

一致同意与完整性定义了共识的核心思想：所有节点都决定相同的结果，而且一旦作出决定，就不能再改变主意。有效性排除了平凡的解法：例如，无论节点提议什么，算法都一律决定 `null`；这种算法满足一致同意与完整性，却不满足有效性。

如果不关心容错，满足前三个属性很容易：只需把一个节点硬编码为“独裁者”，由它作出所有决定即可。但这个节点一旦失效，系统便再也无法作出任何决定——这与没有故障切换的单主复制如出一辙。所有困难都源于容错这一要求。

终止属性形式化了容错的含义。它实质上要求共识算法不能永远无所作为——换句话说，算法必须取得进展。即使一部分节点失效，其余节点仍须作出决定。（终止是一种活性属性，另外三种则是安全属性，见 [“安全性与活性”](/ch9#sec_distributed_safety_liveness)。）

如果崩溃的节点还有可能恢复，当然可以等它回来。然而，共识必须保证，即使某个节点突然消失而且永远不再回来，系统仍能作出决定。（不要只想象软件崩溃；想象一场地震引发山体滑坡，彻底摧毁了节点所在的数据中心。必须假定这个节点被埋在 30 英尺深的泥土下，再也不会重新上线。）

当然，如果 *所有* 节点都已崩溃，无一仍在运行，那么任何算法都不可能决定任何事情。算法能容忍的故障数量存在上限：事实上，可以证明，任何共识算法都要求至少多数节点正常运行，才能保证终止 [^73]。这组多数节点可以安全地构成法定人数（见 [“读写仲裁”](/ch6#sec_replication_quorum_condition)）。

因此，终止属性以少于一半节点崩溃或不可达为前提。不过，即使多数节点失效，或网络出现严重问题，大多数共识算法仍能保证安全属性——一致同意、完整性与有效性——始终成立 [^75]。所以，大规模中断可能令系统无法处理请求，却不能迫使共识系统作出彼此矛盾的决定，从而破坏其正确性。

#### 比较并设置作为共识 {#compare-and-set-as-consensus}

比较并设置（CAS）操作会检查某个对象的当前值是否等于预期值。如果相等，它就以原子方式把对象更新为新值；否则保持对象不变并返回错误。

有了容错且线性一致的 CAS 操作，就很容易解决共识问题：先把对象设为空值；每个想提议某个值的节点都调用 CAS，把预期值设为空，把新值设为自己要提议的值（假定该值非空）。对象最终被设成什么值，共识决定的就是什么值。

反过来，有了共识的解法，也可以实现 CAS：每当一个或多个节点想以相同的预期值执行 CAS 时，就通过共识协议提议各次 CAS 调用中的新值，再把对象设为共识所决定的值。新值未获选中的 CAS 调用返回错误。预期值不同的 CAS 调用，则分别运行共识协议。

由此可见，CAS 与共识彼此等价 [^28] [^73]。两者在单个节点上都很简单，难点都在于如何实现容错。我们曾在 [“以对象存储为后端的数据库”](/ch6#sec_replication_object_storage) 中见过分布式环境下的 CAS：对象存储的条件写入操作只有在当前客户端上次读取之后，同名对象没有被其他客户端创建或修改时，才允许写入发生。

不过，线性一致的读写寄存器不足以解决共识。FLP 结果表明，在异步崩溃停止模型中，确定性算法无法解决共识 [^72]；但我们在 [“线性一致性与仲裁”](/ch10#sec_consistency_quorum_linearizable) 中看到，线性一致的寄存器可以在这一模型下通过法定人数读写来实现 [^24] [^25] [^26]。由此可以推知，线性一致的寄存器无法解决共识。

#### 共享日志作为共识 {#sec_consistency_shared_logs}

我们已经见过多种日志，例如复制日志、事务日志和预写日志。日志存储一系列 *日志条目*（*log entry*），任何读取者都会以相同顺序看到相同的条目。有时，只有一个写入者有权向日志追加新条目；而在 *共享日志*（*shared log*）中，多个节点都可以请求追加条目。单主复制就是一个例子：任何客户端都可以请求领导者执行写入，领导者把写入追加到复制日志，随后所有追随者都按照与领导者相同的顺序应用这些写入。

更形式化地说，共享日志支持两种操作：请求把一个值加入日志，以及读取日志条目。它必须满足以下属性：

最终追加
: 如果某个节点请求把一个值加入日志，而且该节点没有崩溃，那么它最终必须能在某条日志条目中读到这个值。

可靠交付
: 日志条目不会丢失：如果某个节点读到了某条日志条目，那么每个未崩溃的节点最终也必须读到它。

仅追加
: 节点一旦读到某条日志条目，该条目便不可再改变；新条目只能追加在它之后，不能插入它之前。节点重新读取日志时，必须以初次读取时的相同顺序看到相同条目，即使它曾经崩溃并重启也不例外。

一致性
: 如果两个节点都读到了某条日志条目 *e*，那么在 *e* 之前，它们必须以相同顺序读到完全相同的日志条目序列。

有效性
: 如果某个节点读到了一条包含某值的日志条目，那么此前必有某个节点请求把这个值加入日志。


> [!NOTE]
> 共享日志在形式上称为 *全序广播*（*total order broadcast*）、*原子广播*（*atomic broadcast*）或 *全序组播*（*total order multicast*）协议 [^26] [^76] [^77]。这些术语只是用不同说法描述同一件事：请求把一个值加入日志称为“广播”这个值，读取日志条目则称为“交付”这条日志。


有了共享日志的实现，就很容易解决共识问题：每个想提议某个值的节点都请求把它加入日志，而第一条日志条目中读出的值就是决定值。由于所有节点都按相同顺序读取日志条目，它们必然会就哪个值最先交付达成一致 [^28]。

反过来，有了共识的解法，也可以实现共享日志。具体细节稍显复杂，但基本思路如下 [^73]：

1. 为日志中每个未来的条目预留一个槽位，并为每个槽位分别运行一个共识算法实例，以决定该条目应当包含什么值。
2. 节点想向日志加入某个值时，就为一个尚未决定的槽位提议这个值。
3. 共识算法为某个槽位作出决定，而且此前所有槽位也都已有决定后，就把决定值追加为新的日志条目；其后所有已经连续作出决定的槽位，其决定值也一并追加到日志。
4. 如果提议值没有被某个槽位选中，想加入这个值的节点就改为向后面的槽位重新提议。

这说明，共识等价于全序广播，也等价于共享日志。没有故障切换的单主复制不满足活性要求，因为领导者一旦崩溃，系统就会停止交付消息。一如既往，真正的挑战在于如何安全、自动地完成故障切换。

#### 获取并增加作为共识 {#fetch-and-add-as-consensus}

我们在 [“线性一致的 ID 生成器”](/ch10#sec_consistency_linearizable_id) 中看到，线性一致的 ID 生成器距离解决共识只差一步，却终究还差一点。这样的 ID 生成器可以用“获取并增加”操作实现：它以原子方式递增计数器，并返回计数器的旧值。

有了 CAS 操作，实现获取并增加很容易：先读取计数器的值，再执行一次 CAS，把预期值设为刚才读到的值，把新值设为旧值加一。如果 CAS 失败，就从头重试，直到成功为止。存在争用时，这种实现不如原生的获取并增加操作高效，但两者在功能上等价。既然共识可以实现 CAS，自然也可以实现获取并增加。

反过来，如果有了容错的获取并增加操作，能否解决共识？假设计数器初值为零，每个想提议某个值的节点都调用获取并增加来递增计数器。由于这项操作具有原子性，其中一个节点会读到初始值零，其他节点读到的值则都至少已经递增过一次。

现在规定，读到零的节点胜出，其提议值成为决定值。这对读到零的节点当然没问题，其他节点却陷入了困境：它们知道自己没有胜出，却不知道其他节点中究竟谁赢了。胜者可以发消息告知其他节点，但如果它还没来得及发送消息就崩溃了呢？其余节点将一直悬而未决，无法决定任何值，共识也就不能终止。它们也不能改选另一个节点，因为读到零的节点日后仍可能恢复，并有充分理由决定自己提议的值。

有一个例外：我们能够确定提议值的节点不超过两个。此时，两个节点可以先互相发送各自的提议值，再分别执行获取并增加。读到零的节点决定自己的值，读到一的节点则决定另一个节点的值。这样就解决了两个节点之间的共识问题，因此我们说获取并增加的 *共识数*（*consensus number*）为二 [^28]。相比之下，CAS 和共享日志可以在任意数量的节点提议值时解决共识，所以它们的共识数为 ∞（无穷大）。

#### 原子提交作为共识 {#atomic-commitment-as-consensus}

在 [“分布式事务”](/ch8#sec_transactions_distributed) 中，我们见过 *原子提交*（*atomic commitment*）问题：参与分布式事务的所有数据库或分片，要么全都提交事务，要么全都中止。我们还见过 *两阶段提交*（*two-phase commit*，2PC）算法，它依赖一个构成单点故障的协调者。

共识与原子提交是什么关系？乍看之下，两者十分相似——都要求节点达成某种一致。不过，它们之间存在一项重要区别：共识可以决定任意一个被提议的值；原子提交则要求，只要 *任何* 参与者投票中止，算法就 *必须* 中止。更准确地说，原子提交必须满足以下属性 [^78]：

一致同意
: 任意两个节点都不会决定不同的结果。

完整性
: 节点一旦决定了某个结果，就不能再决定另一个结果来改变主意。

有效性
: 如果某个节点决定提交，那么此前所有节点都必须投票提交；只要任何节点投票中止，所有节点都必须中止。

非平凡性
: 如果所有节点都投票提交，而且没有发生通信超时，那么所有节点都必须决定提交。

终止
: 每个未崩溃的节点最终都能决定提交或中止。

有效性确保事务只有在所有节点都同意时才能提交；非平凡性则确保算法不能简单地一律中止（但只要节点间发生任何通信超时，它就允许中止）。另外三项属性与共识基本相同。

有了共识的解法，可以用多种方式解决原子提交 [^78] [^79]。其中一种做法如下：准备提交事务时，每个节点把自己的提交或中止票发送给其他所有节点。某个节点如果收到自己以及所有其他节点的提交票，就通过共识算法提议“提交”；如果收到中止票或遇到超时，就通过共识算法提议“中止”。节点得知共识算法的决定后，再据此提交或中止事务。

在这种算法中，只有所有节点都投票提交，才可能有人提议“提交”。只要有任何节点投票中止，共识算法收到的所有提议都会是“中止”。如果所有节点都投票提交，但部分通信发生超时，就可能有些节点提议“中止”，另一些节点提议“提交”；这时最终提交还是中止并不重要，只要所有节点采取相同的决定即可。

反过来，有了容错的原子提交协议，也可以解决共识。每个想提议某个值的节点都在一组法定人数节点上发起事务，并在每个节点上执行一次单节点 CAS：如果寄存器尚未被另一个事务设值，就将其设为本节点的提议值。CAS 成功时节点投票提交，否则投票中止。如果原子提交协议决定提交这个事务，其值就成为共识的决定值；如果原子提交中止，提议节点就用一个新事务重试。

由此可见，原子提交与共识也彼此等价。

### 共识的实践 {#sec_consistency_total_order}

前面已经看到，单值共识、CAS、共享日志与原子提交彼此等价：其中任意一个问题的解法，都可以转换为其他问题的解法。这个理论洞见很有价值，却没有回答一个实际问题：共识有这么多种表述，实践中究竟哪一种最有用？

答案是：大多数共识系统都提供共享日志，也就是全序广播。Raft、视图戳复制和 Zab 直接提供共享日志；Paxos 提供的是单值共识，但实践中，采用 Paxos 的系统大多使用名为 Multi-Paxos 的扩展，它同样提供共享日志。

#### 使用共享日志 {#sec_consistency_smr}

共享日志与数据库复制十分契合：如果每条日志条目都代表一次数据库写入，而每个副本都以相同顺序、使用确定性逻辑处理相同的写入，那么所有副本最终都会处于一致状态。这个思想称为 *状态机复制*（*state machine replication*）[^80]，也是我们在 [“事件溯源与 CQRS”](/ch3#sec_datamodels_events) 中见过的事件溯源原理。共享日志对流处理也很有用，我们将在 [第 12 章](/ch12#ch_stream) 看到这一点。

类似地，共享日志还可以实现可串行化事务。正如 [“实际串行执行”](/ch8#sec_transactions_serial) 所述，如果每条日志条目都代表一个将以存储过程形式执行的确定性事务，而且每个节点都以相同顺序执行这些事务，那么事务就是可串行化的 [^81] [^82]。


> [!NOTE]
> 采用强一致性模型的分片数据库，通常会为每个分片分别维护一份日志。这样做提高了可伸缩性，却也限制了数据库能跨分片提供的一致性保证（例如一致快照与外键引用）。跨分片的可串行化事务并非不能实现，但需要额外协调 [^83]。


共享日志的另一个强大之处在于，它很容易改造成其他形式的共识：

* 前面已经说明怎样用它实现单值共识和 CAS：只需决定日志中最先出现的值。
* 如果需要许多个单值共识实例（例如，剧院里每个被多人争订的座位各一个实例），可以在日志条目中写入座位号，并以包含该座位号的第一条日志条目作为决定。
* 如果需要原子的获取并增加操作，可以把要加到计数器上的数字写进日志条目；计数器的当前值就是截至目前所有日志条目中数字的总和。还可以对日志条目简单计数，以此生成栅栏令牌（见 [“用栅栏机制隔离僵尸与延迟请求”](/ch9#sec_distributed_fencing_tokens)）。例如，在 ZooKeeper 中，这个序列号称为 `zxid` [^18]。

#### 从单主复制到共识 {#from-single-leader-replication-to-consensus}

前面已经看到，如果由一个“独裁者”节点作出决定，单值共识就很容易；类似地，如果只有一个领导者有权向共享日志追加条目，共享日志也很容易实现。问题在于：这个节点失效时，怎样才能实现容错？

传统的单主复制数据库并没有解决这个问题，而是把领导者故障切换留给管理员手工操作。人类的反应速度终究有限，这种方式难免造成相当长的停机时间，也不满足共识的终止属性。共识要求算法能够自动选出新的领导者。（并非所有共识算法都有领导者，但常用算法通常都有 [^84]。）

然而，这里有一个难题。前面讨论脑裂时说过，所有节点必须就谁是领导者达成一致，否则两个不同节点可能都自认为是领导者，进而作出彼此矛盾的决定。如此看来，选举领导者需要共识，解决共识又需要领导者。怎样才能跳出这个先有鸡还是先有蛋的困局？

事实上，共识算法并不要求任何时刻都只有一个领导者。它们提供的是稍弱一些的保证：协议定义一个 *纪元编号*（*epoch number*；Paxos 称为 *投票编号*，*ballot number*，视图戳复制称为 *视图编号*，*view number*，Raft 称为 *任期编号*，*term number*），并保证每个纪元内的领导者是唯一的。

如果一个节点在指定的超时时间内始终没有收到现任领导者的消息，因而认为它已经失效，这个节点就可能发起投票，选举新的领导者。这次选举会取得一个大于以往所有纪元的新纪元编号。如果两个不同纪元的领导者发生冲突——也许前任领导者其实并未失效——纪元编号较高的领导者说了算。

领导者要向共享日志追加下一条记录，必须先确认不存在纪元编号更高的其他领导者，否则后者可能会追加不同的条目。它可以向一组法定人数节点收集选票，这组节点通常是多数，但并非总是如此 [^85]。只有在不知道任何更高纪元领导者的情况下，节点才会投赞成票。

因此，这里需要两轮投票：第一轮选举领导者；第二轮对领导者提议追加的下一条日志条目进行表决。两轮投票的法定人数必须相互重叠：如果某项提议表决通过，投赞成票的节点中，至少要有一个参加过最近一次成功的领导者选举 [^85]。所以，如果提议表决通过，而且投票过程没有发现编号更高的纪元，现任领导者便可以断定，没有纪元编号更高的领导者当选，因而可以安全地把提议条目追加到日志 [^26] [^86]。

这两轮投票表面上很像两阶段提交，实际上却是两种截然不同的协议。在共识算法中，任何节点都可以发起选举，而且只需一组法定人数节点响应；在 2PC 中，只有协调者能请求投票，并且必须从 *每个* 参与者那里得到“同意”票，事务才能提交。

#### 共识的微妙之处 {#subtleties-of-consensus}

Raft、Multi-Paxos、Zab 和视图戳复制都采用这一基本结构：先由一组法定人数节点投票选举领导者；此后，领导者想追加的每一条日志条目，还要经过另一组法定人数节点投票 [^68] [^69]。每条新日志条目都要同步复制到一组法定人数节点，才会向发起写入的客户端确认成功。这样即使现任领导者失效，日志条目也不会丢失。

不过，魔鬼藏在细节里，而这些算法的差异也恰恰体现在细节中。例如，旧领导者失效并选出新领导者后，算法必须保证，新领导者会保留旧领导者在失效前已经追加的所有日志条目。Raft 的做法是：只有日志至少与多数追随者一样新鲜的节点，才有资格成为新领导者 [^69]。Paxos 则允许任何节点成为新领导者，但要求它先从其他节点补齐日志，才能开始追加自己的新条目。


> [!TIP] 领导者选举中的一致性与可用性
>
> 如果希望共识算法严格保证 [“共享日志作为共识”](/ch10#sec_consistency_shared_logs) 中列出的属性，那么新领导者在处理任何写入或线性一致读取之前，必须掌握所有已经确认的日志条目。这是保证上述属性的必要条件。如果一个数据陈旧的节点成为新领导者，它可能会改写旧领导者已经写入的日志条目，从而违反共享日志的仅追加属性。
>
> 有些系统会选择削弱共识属性，以便更快地从领导者失效中恢复。例如，Kafka 可以启用 *非同步副本选举*（unclean leader election），允许任何副本成为领导者，即使它没有追上最新进度。另外，在采用异步复制的数据库中，领导者失效时，根本无法保证任何追随者已经赶上最新进度。
>
> 放弃“新领导者必须掌握最新数据”这一要求，或许能提高性能与可用性，却也如同在薄冰上行走，因为共识理论已经不再适用。没有故障时，系统固然可以正常工作；但一旦遇到 [第 9 章](/ch9#ch_distributed) 讨论的种种问题，就很容易造成大量数据丢失或损坏。

另一个微妙之处是：旧领导者在失效前已经提议了某条日志条目，但追加该条目的投票尚未结束，算法应当如何处理。关于这些细节，可以参阅本章末尾的参考文献 [^23] [^69] [^86]。

对于用共识算法做复制的数据库，不仅写入要转成日志条目并复制到一组法定人数节点。如果还要保证线性一致读取，读取也必须像写入一样经过法定人数投票，以确认那个自认为是领导者的节点确实仍掌握最新数据。etcd 的线性一致读取就是这样实现的。

大多数共识算法的标准形式都假定节点集合固定不变：节点可以下线后重新上线，但哪些节点有权投票，在集群创建时便已确定。实践中，却经常需要在系统配置中添加新节点或移除旧节点。共识算法因此扩展出了 *重新配置* 功能。向系统增加新地区，或把系统从一个位置迁往另一个位置时，这项功能尤其有用：可以先加入新节点，再移除旧节点。

#### 共识的利弊 {#pros-and-cons-of-consensus}

共识算法虽然复杂而微妙，却是分布式系统领域的一项重大突破。共识本质上就是“正确实现的单主复制”：领导者失效时自动进行故障切换；即使面对 [第 9 章](/ch9#ch_distributed) 讨论的所有问题，也能保证已提交的数据不会丢失，系统绝不会发生脑裂。

既然带自动故障切换的单主复制本质上就是共识的一种定义，那么，任何提供自动故障切换、却没有采用经过验证的共识算法的系统，都很可能不安全 [^87]。当然，采用经过验证的共识算法，并不能保证整个系统一定正确——错误仍可能潜伏在许多其他角落——但至少是一个良好的开端。

尽管如此，共识并未无处不在，因为它的好处也有代价。共识系统总要有严格多数节点才能运行：要容忍一个节点故障，至少需要三个节点；要容忍两个节点故障，至少需要五个。每项操作都必须与一组法定人数节点通信，因此不能靠增加节点来提高吞吐量（实际上，每增加一个节点，算法反而会变慢）。如果网络分区把一部分节点同其余节点隔开，只有多数派所在的一侧能够取得进展，另一侧则会阻塞。

共识系统通常依靠超时来检测失效节点。在网络延迟变化很大的环境中，尤其是跨多个地理区域部署的系统，超时时间很难调好：设得太长，故障恢复会耗时很久；设得太短，又会触发大量不必要的领导者选举，导致性能极差——系统花在选举领导者上的时间，可能比花在有用工作上的还多。

有些共识算法对网络问题格外敏感。例如，Raft 已被发现存在一些棘手的边界情况 [^88] [^89]：即使整个网络都运行正常，只要有一条特定链路始终不可靠，领导权就可能在两个节点之间来回跳转，或者现任领导者不断被迫辞职，导致系统实际上永远无法取得进展。怎样设计对不可靠网络更稳健的算法，至今仍是一个开放的研究问题。

如果系统既希望高可用，又不愿承担共识的成本，真正可行的选择只有改用较弱的一致性模型，例如 [第 6 章](/ch6#ch_replication) 讨论的无主复制或多主复制。这些方法通常不提供线性一致性；不过，对于并不需要线性一致性的应用，这已经足够。


### 协调服务 {#sec_consistency_coordination}

任何想提供线性一致操作的分布式数据库，都能从共识算法中受益；许多现代分布式数据库确实也用共识算法做复制。不过，有一类系统尤其倚重共识：ZooKeeper、etcd、Consul 等 *协调服务*（*coordination service*）。它们表面上与普通键值存储相似，却不像大多数数据库那样以通用数据存储为目标。

协调服务的用途，是协调另一个分布式系统中的多个节点。例如，Kubernetes 依赖 etcd；Spark 和 Flink 在高可用模式下，则依赖后台运行的 ZooKeeper。协调服务只保存少量足以全部放入内存的数据（同时仍会写入磁盘以保证持久性），再用容错共识算法把这些数据复制到多个节点。

协调服务以 Google 的 Chubby 锁服务为蓝本 [^17] [^58]，把共识算法与几项对构建分布式系统格外有用的功能结合在一起：

锁与租约
: 前面已经看到，共识系统可以实现具备容错能力的原子比较并设置（CAS）操作。协调服务以此实现锁和租约：多个节点并发争抢同一份租约时，只有一个能够成功。

栅栏机制
: 正如 [“分布式锁和租约”](/ch9#sec_distributed_lock_fencing) 所述，以租约保护某项资源时，需要用 *栅栏机制* 防止客户端在进程暂停或网络严重延迟时相互干扰。共识系统可以为每条日志条目分配单调递增的 ID，以此生成栅栏令牌（ZooKeeper 使用 `zxid` 和 `cversion`，etcd 使用修订号）。

故障检测
: 客户端与协调服务维持一个长期会话，并定期交换心跳，确认对方是否仍然存活。即使连接暂时中断或某台服务器失效，客户端持有的租约仍然有效；但如果心跳中断的时间超过租约超时，协调服务就会认定客户端已经失效，并释放其租约。（ZooKeeper 把这种随会话到期自动消失的条目称为 *临时节点*，*ephemeral node*。）

变更通知
: 客户端可以要求协调服务在某些键发生变化时主动发送通知。这样一来，客户端便能得知另一个客户端何时加入集群（根据它写入协调服务的值），或何时失效（它的会话超时，临时节点随之消失），无须再频繁轮询服务来发现变化。

故障检测与变更通知本身并不需要共识；但把它们同确实需要共识的原子操作、栅栏机制组合起来，对分布式协调就格外有用。

> [!TIP] 用协调服务管理配置
>
> 应用与基础设施通常都有超时时间、线程池大小等配置参数。有时，人们会把这类配置数据以键值对形式存入协调服务。进程启动时加载最新配置，并订阅此后的变更通知。配置变化后，进程可以立即采用新设置，也可以通过重启来加载最新设置。
>
> 配置管理本身并不需要协调服务的共识能力；不过，如果系统本来就在运行协调服务，顺便利用它的通知功能会很方便。另一种办法是让进程定期从文件或 URL 轮询配置更新，从而不必依赖专门的协调服务。

#### 将工作分配给节点 {#allocating-work-to-nodes}

如果某个进程或服务有多个实例，需要从中选出一个领导者或主实例，协调服务就很有用。领导者失效时，应当由其他某个节点接管。这对单主数据库必不可少，也适用于作业调度器等有状态系统。

另一种场景是分配分片资源（数据库、消息流、文件存储、分布式 Actor 系统等）：系统需要决定每个分片交给哪个节点。新节点加入集群时，需要把一部分分片从现有节点移到新节点，以重新平衡负载；节点被移除或失效时，则要由其他节点接手它的工作。

审慎组合协调服务中的原子操作、临时节点与通知机制，就可以完成这类任务。实现得当时，应用能在无需人工干预的情况下自动从故障中恢复。即使已经有 Apache Curator 之类基于 ZooKeeper 客户端 API 提供高层工具的库，这仍然不是一件容易的事；但无论如何，都远胜于从头实现所需的共识算法——后者极易埋下错误。

专用协调服务还有一个优点：无论依赖它协调的分布式系统有多少节点，协调服务本身都可以只运行在一组固定节点上，通常是三个或五个。例如，一个存储系统拥有数千个分片；若让数千个节点共同运行共识算法，效率会低得惊人。把共识“外包”给少量运行协调服务的节点，要合理得多。

协调服务管理的数据通常变化缓慢，例如“IP 地址为 10.1.1.23 的节点是分片 7 的领导者”这类分配关系，往往几分钟或几小时才会变一次。协调服务并非用来存储每秒可能变化数千次的数据；这类数据更适合交给常规数据库。或者，也可以用 Apache BookKeeper [^90] [^91] 之类的工具，复制服务内部快速变化的状态。

#### 服务发现 {#service-discovery}

ZooKeeper、etcd 和 Consul 也经常用于 *服务发现*（*service discovery*），即找出应当连接哪个 IP 地址才能访问特定服务（见 [“负载均衡器、服务发现和服务网格”](/ch5#sec_encoding_service_discovery)）。云环境中的虚拟机不断创建和销毁，通常无法预先知道服务的 IP 地址。常见做法是让服务在启动时把自己的网络端点注册到服务注册表，供其他服务查找。

用协调服务做服务发现很方便：故障检测与变更通知功能，让客户端很容易跟踪服务实例的增减。如果系统已经用协调服务管理租约、锁或领导者选举，继续用它做服务发现也顺理成章，因为它本来就知道哪个节点应当接收服务请求。

不过，服务发现使用共识往往有些杀鸡用牛刀。这个场景通常不要求线性一致性；高可用与低延迟反而更加重要，因为一旦服务发现不可用，整个系统都会停摆。因此，更常见的做法是缓存服务发现信息，并接受结果可能略有陈旧。例如，基于 DNS 的服务发现就通过多层缓存获得良好的性能与可用性。

为支持这类用途，ZooKeeper 提供了 *观察者*（observer）节点。这些副本会接收日志、维护一份 ZooKeeper 数据副本，却不参与共识算法的投票。观察者上的读取可能陈旧，因而不具备线性一致性；但即使网络中断，读取仍可继续，而且缓存还能提高系统所能支持的读取吞吐量。

## 总结 {#summary}

本章探讨了容错系统中的强一致性：它是什么，又该怎样实现。我们深入研究了强一致性的一种常用形式化定义——线性一致性。它要求复制数据表现得仿佛只有一个副本，而且所有操作都以原子方式作用于这个副本。当某些数据在读取时必须是最新的，或需要解决竞态条件时（例如多个节点并发创建同名文件），线性一致性就很有用。

线性一致性很有吸引力，因为它易于理解：数据库的行为就像单线程程序中的一个变量。但它的缺点是速度慢，网络延迟较大时尤其如此。许多复制算法都无法保证线性一致性，哪怕乍看之下似乎能够提供强一致性。

接着，我们把线性一致性的概念应用到 ID 生成器。单节点自增计数器具有线性一致性，却不能容错。许多分布式 ID 生成方案也无法保证，ID 的顺序与事件实际发生的顺序一致。Lamport 时钟、混合逻辑时钟等逻辑时钟，可以给出与因果关系一致的顺序，却不具备线性一致性。

由此，我们引出了共识：达成共识，就是以一种所有节点都认同结果、而且事后不能改变主意的方式作出决定。许多问题实际上都可以归约为共识，而且彼此等价——换句话说，只要有一个问题的解法，就能把它转换成其他所有问题的解法。这些等价问题包括：

线性一致的比较并设置操作
: 寄存器必须根据当前值是否等于操作给出的参数，以原子方式 *决定* 是否设置新值。

锁和租约
: 多个客户端并发争抢锁或租约时，锁会 *决定* 由哪一个成功获取。

唯一性约束
: 多个事务并发创建键相同、彼此冲突的记录时，约束必须 *决定* 允许哪一个，又让哪一个因违反约束而失败。

共享日志
: 多个节点并发请求向日志追加条目时，日志会 *决定* 条目的追加顺序。全序广播也与之等价。

原子事务提交
: 参与分布式事务的所有数据库节点，必须作出相同的 *决定*：提交事务，或是中止事务。

线性一致的获取并增加操作
: 这项操作可以实现 ID 生成器。多个节点可以并发调用它，而它会 *决定* 各节点递增计数器的顺序。这种方法实际上只能解决两个节点之间的共识，前面几种则适用于任意数量的节点。

如果只有一个节点，或者愿意把决策权交给一个节点，所有这些问题都很简单。单主数据库正是如此：全部决策权都归领导者所有，这也是此类数据库能够提供线性一致操作、唯一性约束、复制日志等功能的原因。

可是，一旦这个领导者失效，或网络中断令它不可达，系统便无法继续取得进展，只能等待人工完成故障切换。Raft、Paxos 等广泛使用的共识算法，本质上就是内置了自动领导者选举与故障切换的单主复制。

共识算法经过精心设计，能保证故障切换期间不会丢失任何已经提交的写入，也不会让系统陷入多个节点同时接受写入的脑裂状态。为此，每次写入和每次线性一致读取都必须得到一组法定人数节点（通常是多数节点）的确认。这个过程代价不菲，跨地理区域时尤其如此；但若想获得共识所提供的强一致性与容错能力，这项代价无法避免。

ZooKeeper、etcd 等协调服务同样建立在共识算法之上。它们提供锁、租约、故障检测和变更通知等功能，有助于管理分布式应用的状态。如果要做的事情可以归约为共识，而且还必须容错，最好使用协调服务。它不能保证实现一定正确，却很可能帮上忙。

共识算法复杂而微妙，但背后有一套自 20 世纪 80 年代以来不断发展的丰富理论体系。这套理论让我们能够构建出这样的系统：既容忍 [第 9 章](/ch9#ch_distributed) 讨论的各种故障，又保证数据不被破坏。这是分布式系统工程的一项非凡成就；本章末尾的参考文献列出了其中若干重要工作。

不过，共识并不总是正确的工具。有些系统并不需要它所提供的强一致性，以较弱的一致性换取更高可用性和更好性能，反而更加合适。在这些场景中，人们通常采用 [第 6 章](/ch6#ch_replication) 讨论过的无主复制或多主复制；本章介绍的逻辑时钟，对这类系统也很有帮助。

### 参考文献

[^1]: Maurice P. Herlihy and Jeannette M. Wing. [Linearizability: A Correctness Condition for Concurrent Objects](https://cs.brown.edu/~mph/HerlihyW90/p463-herlihy.pdf). *ACM Transactions on Programming Languages and Systems* (TOPLAS), volume 12, issue 3, pages 463–492, July 1990. [doi:10.1145/78969.78972](https://doi.org/10.1145/78969.78972) 
[^2]: Leslie Lamport. [On interprocess communication](https://www.microsoft.com/en-us/research/publication/interprocess-communication-part-basic-formalism-part-ii-algorithms/). *Distributed Computing*, volume 1, issue 2, pages 77–101, June 1986. [doi:10.1007/BF01786228](https://doi.org/10.1007/BF01786228) 
[^3]: David K. Gifford. [Information Storage in a Decentralized Computer System](https://bitsavers.org/pdf/xerox/parc/techReports/CSL-81-8_Information_Storage_in_a_Decentralized_Computer_System.pdf). Xerox Palo Alto Research Centers, CSL-81-8, June 1981. Archived at [perma.cc/2XXP-3JPB](https://perma.cc/2XXP-3JPB) 
[^4]: Martin Kleppmann. [Please Stop Calling Databases CP or AP](https://martin.kleppmann.com/2015/05/11/please-stop-calling-databases-cp-or-ap.html). *martin.kleppmann.com*, May 2015. Archived at [perma.cc/MJ5G-75GL](https://perma.cc/MJ5G-75GL) 
[^5]: Kyle Kingsbury. [Call Me Maybe: MongoDB Stale Reads](https://aphyr.com/posts/322-call-me-maybe-mongodb-stale-reads). *aphyr.com*, April 2015. Archived at [perma.cc/DXB4-J4JC](https://perma.cc/DXB4-J4JC) 
[^6]: Kyle Kingsbury. [Computational Techniques in Knossos](https://aphyr.com/posts/314-computational-techniques-in-knossos). *aphyr.com*, May 2014. Archived at [perma.cc/2X5M-EHTU](https://perma.cc/2X5M-EHTU) 
[^7]: Kyle Kingsbury and Peter Alvaro. [Elle: Inferring Isolation Anomalies from Experimental Observations](https://www.vldb.org/pvldb/vol14/p268-alvaro.pdf). *Proceedings of the VLDB Endowment*, volume 14, issue 3, pages 268–280, November 2020. [doi:10.14778/3430915.3430918](https://doi.org/10.14778/3430915.3430918) 
[^8]: Paolo Viotti and Marko Vukolić. [Consistency in Non-Transactional Distributed Storage Systems](https://arxiv.org/abs/1512.00168). *ACM Computing Surveys* (CSUR), volume 49, issue 1, article no. 19, June 2016. [doi:10.1145/2926965](https://doi.org/10.1145/2926965) 
[^9]: Peter Bailis. [Linearizability Versus Serializability](http://www.bailis.org/blog/linearizability-versus-serializability/). *bailis.org*, September 2014. Archived at [perma.cc/386B-KAC3](https://perma.cc/386B-KAC3) 
[^10]: Daniel Abadi. [Correctness Anomalies Under Serializable Isolation](https://dbmsmusings.blogspot.com/2019/06/correctness-anomalies-under.html). *dbmsmusings.blogspot.com*, June 2019. Archived at [perma.cc/JGS7-BZFY](https://perma.cc/JGS7-BZFY) 
[^11]: Peter Bailis, Aaron Davidson, Alan Fekete, Ali Ghodsi, Joseph M. Hellerstein, and Ion Stoica. [Highly Available Transactions: Virtues and Limitations](https://www.vldb.org/pvldb/vol7/p181-bailis.pdf). *Proceedings of the VLDB Endowment*, volume 7, issue 3, pages 181–192, November 2013. [doi:10.14778/2732232.2732237](https://doi.org/10.14778/2732232.2732237), extended version published as [arXiv:1302.0309](https://arxiv.org/abs/1302.0309) 
[^12]: Philip A. Bernstein, Vassos Hadzilacos, and Nathan Goodman. [*Concurrency Control and Recovery in Database Systems*](https://www.microsoft.com/en-us/research/people/philbe/book/). Addison-Wesley, 1987. ISBN: 978-0-201-10715-9, available online at [*microsoft.com*](https://www.microsoft.com/en-us/research/people/philbe/book/). 
[^13]: Andrei Matei. [CockroachDB’s consistency model](https://www.cockroachlabs.com/blog/consistency-model/). *cockroachlabs.com*, February 2021. Archived at [perma.cc/MR38-883B](https://perma.cc/MR38-883B) 
[^14]: Murat Demirbas. [Strict-serializability, but at what cost, for what purpose?](https://muratbuffalo.blogspot.com/2022/08/strict-serializability-but-at-what-cost.html) *muratbuffalo.blogspot.com*, August 2022. Archived at [perma.cc/T8AY-N3U9](https://perma.cc/T8AY-N3U9) 
[^15]: Ben Darnell. [How to talk about consistency and isolation in distributed DBs](https://www.cockroachlabs.com/blog/db-consistency-isolation-terminology/). *cockroachlabs.com*, February 2022. Archived at [perma.cc/53SV-JBGK](https://perma.cc/53SV-JBGK) 
[^16]: Daniel Abadi. [An explanation of the difference between Isolation levels vs. Consistency levels](https://dbmsmusings.blogspot.com/2019/08/an-explanation-of-difference-between.html). *dbmsmusings.blogspot.com*, August 2019. Archived at [perma.cc/QSF2-CD4P](https://perma.cc/QSF2-CD4P) 
[^17]: Mike Burrows. [The Chubby Lock Service for Loosely-Coupled Distributed Systems](https://research.google/pubs/pub27897/). At *7th USENIX Symposium on Operating System Design and Implementation* (OSDI), November 2006. 
[^18]: Flavio P. Junqueira and Benjamin Reed. [*ZooKeeper: Distributed Process Coordination*](https://www.oreilly.com/library/view/zookeeper/9781449361297/). O’Reilly Media, 2013. ISBN: 978-1-449-36130-3 
[^19]: Murali Vallath. [*Oracle 10g RAC Grid, Services & Clustering*](https://www.oreilly.com/library/view/oracle-10g-rac/9781555583217/). Elsevier Digital Press, 2006. ISBN: 978-1-555-58321-7 
[^20]: Peter Bailis, Alan Fekete, Michael J. Franklin, Ali Ghodsi, Joseph M. Hellerstein, and Ion Stoica. [Coordination Avoidance in Database Systems](https://arxiv.org/abs/1402.2237). *Proceedings of the VLDB Endowment*, volume 8, issue 3, pages 185–196, November 2014. [doi:10.14778/2735508.2735509](https://doi.org/10.14778/2735508.2735509) 
[^21]: Kyle Kingsbury. [Call Me Maybe: etcd and Consul](https://aphyr.com/posts/316-call-me-maybe-etcd-and-consul). *aphyr.com*, June 2014. Archived at [perma.cc/XL7U-378K](https://perma.cc/XL7U-378K) 
[^22]: Flavio P. Junqueira, Benjamin C. Reed, and Marco Serafini. [Zab: High-Performance Broadcast for Primary-Backup Systems](https://marcoserafini.github.io/assets/pdf/zab.pdf). At *41st IEEE International Conference on Dependable Systems and Networks* (DSN), June 2011. [doi:10.1109/DSN.2011.5958223](https://doi.org/10.1109/DSN.2011.5958223) 
[^23]: Diego Ongaro and John K. Ousterhout. [In Search of an Understandable Consensus Algorithm](https://www.usenix.org/system/files/conference/atc14/atc14-paper-ongaro.pdf). At *USENIX Annual Technical Conference* (ATC), June 2014. 
[^24]: Hagit Attiya, Amotz Bar-Noy, and Danny Dolev. [Sharing Memory Robustly in Message-Passing Systems](https://www.cs.huji.ac.il/course/2004/dist/p124-attiya.pdf). *Journal of the ACM*, volume 42, issue 1, pages 124–142, January 1995. [doi:10.1145/200836.200869](https://doi.org/10.1145/200836.200869) 
[^25]: Nancy Lynch and Alex Shvartsman. [Robust Emulation of Shared Memory Using Dynamic Quorum-Acknowledged Broadcasts](https://groups.csail.mit.edu/tds/papers/Lynch/FTCS97.pdf). At *27th Annual International Symposium on Fault-Tolerant Computing* (FTCS), June 1997. [doi:10.1109/FTCS.1997.614100](https://doi.org/10.1109/FTCS.1997.614100) 
[^26]: Christian Cachin, Rachid Guerraoui, and Luís Rodrigues. [*Introduction to Reliable and Secure Distributed Programming*](https://www.distributedprogramming.net/), 2nd edition. Springer, 2011. ISBN: 978-3-642-15259-7, [doi:10.1007/978-3-642-15260-3](https://doi.org/10.1007/978-3-642-15260-3) 
[^27]: Niklas Ekström, Mikhail Panchenko, and Jonathan Ellis. [Possible Issue with Read Repair?](https://lists.apache.org/thread/wwsjnnc93mdlpw8nb0d5gn4q1bmpzbon) Email thread on *cassandra-dev* mailing list, October 2012. 
[^28]: Maurice P. Herlihy. [Wait-Free Synchronization](https://cs.brown.edu/~mph/Herlihy91/p124-herlihy.pdf). *ACM Transactions on Programming Languages and Systems* (TOPLAS), volume 13, issue 1, pages 124–149, January 1991. [doi:10.1145/114005.102808](https://doi.org/10.1145/114005.102808) 
[^29]: Armando Fox and Eric A. Brewer. [Harvest, Yield, and Scalable Tolerant Systems](https://radlab.cs.berkeley.edu/people/fox/static/pubs/pdf/c18.pdf). At *7th Workshop on Hot Topics in Operating Systems* (HotOS), March 1999. [doi:10.1109/HOTOS.1999.798396](https://doi.org/10.1109/HOTOS.1999.798396) 
[^30]: Seth Gilbert and Nancy Lynch. [Brewer’s Conjecture and the Feasibility of Consistent, Available, Partition-Tolerant Web Services](https://www.comp.nus.edu.sg/~gilbert/pubs/BrewersConjecture-SigAct.pdf). *ACM SIGACT News*, volume 33, issue 2, pages 51–59, June 2002. [doi:10.1145/564585.564601](https://doi.org/10.1145/564585.564601) 
[^31]: Seth Gilbert and Nancy Lynch. [Perspectives on the CAP Theorem](https://groups.csail.mit.edu/tds/papers/Gilbert/Brewer2.pdf). *IEEE Computer Magazine*, volume 45, issue 2, pages 30–36, February 2012. [doi:10.1109/MC.2011.389](https://doi.org/10.1109/MC.2011.389) 
[^32]: Eric A. Brewer. [CAP Twelve Years Later: How the ‘Rules’ Have Changed](https://sites.cs.ucsb.edu/~rich/class/cs293-cloud/papers/brewer-cap.pdf). *IEEE Computer Magazine*, volume 45, issue 2, pages 23–29, February 2012. [doi:10.1109/MC.2012.37](https://doi.org/10.1109/MC.2012.37) 
[^33]: Susan B. Davidson, Hector Garcia-Molina, and Dale Skeen. [Consistency in Partitioned Networks](https://www.cs.rice.edu/~alc/old/comp520/papers/DGS85.pdf). *ACM Computing Surveys*, volume 17, issue 3, pages 341–370, September 1985. [doi:10.1145/5505.5508](https://doi.org/10.1145/5505.5508) 
[^34]: Paul R. Johnson and Robert H. Thomas. [RFC 677: The Maintenance of Duplicate Databases](https://tools.ietf.org/html/rfc677). Network Working Group, January 1975. 
[^35]: Michael J. Fischer and Alan Michael. [Sacrificing Serializability to Attain High Availability of Data in an Unreliable Network](https://sites.cs.ucsb.edu/~agrawal/spring2011/ugrad/p70-fischer.pdf). At *1st ACM Symposium on Principles of Database Systems* (PODS), March 1982. [doi:10.1145/588111.588124](https://doi.org/10.1145/588111.588124) 
[^36]: Eric A. Brewer. [NoSQL: Past, Present, Future](https://www.infoq.com/presentations/NoSQL-History/). At *QCon San Francisco*, November 2012. 
[^37]: Adrian Cockcroft. [Migrating to Microservices](https://www.infoq.com/presentations/migration-cloud-native/). At *QCon London*, March 2014. 
[^38]: Martin Kleppmann. [A Critique of the CAP Theorem](https://arxiv.org/abs/1509.05393). arXiv:1509.05393, September 2015. 
[^39]: Daniel Abadi. [Problems with CAP, and Yahoo’s little known NoSQL system](https://dbmsmusings.blogspot.com/2010/04/problems-with-cap-and-yahoos-little.html). *dbmsmusings.blogspot.com*, April 2010. Archived at [perma.cc/4NTZ-CLM9](https://perma.cc/4NTZ-CLM9) 
[^40]: Daniel Abadi. [Hazelcast and the Mythical PA/EC System](https://dbmsmusings.blogspot.com/2017/10/hazelcast-and-mythical-paec-system.html). *dbmsmusings.blogspot.com*, October 2017. Archived at [perma.cc/J5XM-U5C2](https://perma.cc/J5XM-U5C2) 
[^41]: Eric Brewer. [Spanner, TrueTime & The CAP Theorem](https://research.google.com/pubs/archive/45855.pdf). *research.google.com*, February 2017. Archived at [perma.cc/59UW-RH7N](https://perma.cc/59UW-RH7N) 
[^42]: Daniel J. Abadi. [Consistency Tradeoffs in Modern Distributed Database System Design](https://www.cs.umd.edu/~abadi/papers/abadi-pacelc.pdf). *IEEE Computer Magazine*, volume 45, issue 2, pages 37–42, February 2012. [doi:10.1109/MC.2012.33](https://doi.org/10.1109/MC.2012.33) 
[^43]: Nancy A. Lynch. [A Hundred Impossibility Proofs for Distributed Computing](https://groups.csail.mit.edu/tds/papers/Lynch/podc89.pdf). At *8th ACM Symposium on Principles of Distributed Computing* (PODC), August 1989. [doi:10.1145/72981.72982](https://doi.org/10.1145/72981.72982) 
[^44]: Prince Mahajan, Lorenzo Alvisi, and Mike Dahlin. [Consistency, Availability, and Convergence](https://apps.cs.utexas.edu/tech_reports/reports/tr/TR-2036.pdf). University of Texas at Austin, Department of Computer Science, Tech Report UTCS TR-11-22, May 2011. Archived at [perma.cc/SAV8-9JAJ](https://perma.cc/SAV8-9JAJ) 
[^45]: Hagit Attiya, Faith Ellen, and Adam Morrison. [Limitations of Highly-Available Eventually-Consistent Data Stores](https://www.cs.tau.ac.il/~mad/publications/podc2015-replds.pdf). At *ACM Symposium on Principles of Distributed Computing* (PODC), July 2015. [doi:10.1145/2767386.2767419](https://doi.org/10.1145/2767386.2767419) 
[^46]: Peter Sewell, Susmit Sarkar, Scott Owens, Francesco Zappa Nardelli, and Magnus O. Myreen. [x86-TSO: A Rigorous and Usable Programmer’s Model for x86 Multiprocessors](https://www.cl.cam.ac.uk/~pes20/weakmemory/cacm.pdf). *Communications of the ACM*, volume 53, issue 7, pages 89–97, July 2010. [doi:10.1145/1785414.1785443](https://doi.org/10.1145/1785414.1785443) 
[^47]: Martin Thompson. [Memory Barriers/Fences](https://mechanical-sympathy.blogspot.com/2011/07/memory-barriersfences.html). *mechanical-sympathy.blogspot.co.uk*, July 2011. Archived at [perma.cc/7NXM-GC5U](https://perma.cc/7NXM-GC5U) 
[^48]: Ulrich Drepper. [What Every Programmer Should Know About Memory](https://www.akkadia.org/drepper/cpumemory.pdf). *akkadia.org*, November 2007. Archived at [perma.cc/NU6Q-DRXZ](https://perma.cc/NU6Q-DRXZ) 
[^49]: Hagit Attiya and Jennifer L. Welch. [Sequential Consistency Versus Linearizability](https://courses.csail.mit.edu/6.852/01/papers/p91-attiya.pdf). *ACM Transactions on Computer Systems* (TOCS), volume 12, issue 2, pages 91–122, May 1994. [doi:10.1145/176575.176576](https://doi.org/10.1145/176575.176576) 
[^50]: Kyzer R. Davis, Brad G. Peabody, and Paul J. Leach. [Universally Unique IDentifiers (UUIDs)](https://www.rfc-editor.org/rfc/rfc9562). RFC 9562, IETF, May 2024. 
[^51]: Ryan King. [Announcing Snowflake](https://blog.x.com/engineering/en_us/a/2010/announcing-snowflake). *blog.x.com*, June 2010. Archived at [archive.org](https://web.archive.org/web/20241128214604/https%3A//blog.x.com/engineering/en_us/a/2010/announcing-snowflake) 
[^52]: Alizain Feerasta. [Universally Unique Lexicographically Sortable Identifier](https://github.com/ulid/spec). *github.com*, 2016. Archived at [perma.cc/NV2Y-ZP8U](https://perma.cc/NV2Y-ZP8U) 
[^53]: Rob Conery. [A Better ID Generator for PostgreSQL](https://bigmachine.io/2014/05/29/a-better-id-generator-for-postgresql/). *bigmachine.io*, May 2014. Archived at [perma.cc/K7QV-3KFC](https://perma.cc/K7QV-3KFC) 
[^54]: Leslie Lamport. [Time, Clocks, and the Ordering of Events in a Distributed System](https://www.microsoft.com/en-us/research/publication/time-clocks-ordering-events-distributed-system/). *Communications of the ACM*, volume 21, issue 7, pages 558–565, July 1978. [doi:10.1145/359545.359563](https://doi.org/10.1145/359545.359563) 
[^55]: Sandeep S. Kulkarni, Murat Demirbas, Deepak Madeppa, Bharadwaj Avva, and Marcelo Leone. [Logical Physical Clocks](https://doi.org/10.1007/978-3-319-14472-6_2). *18th International Conference on Principles of Distributed Systems* (OPODIS), December 2014. [doi:10.1007/978-3-319-14472-6\_2](https://doi.org/10.1007/978-3-319-14472-6_2)
[^56]: Manuel Bravo, Nuno Diegues, Jingna Zeng, Paolo Romano, and Luís Rodrigues. [On the use of Clocks to Enforce Consistency in the Cloud](http://sites.computer.org/debull/A15mar/p18.pdf). *IEEE Data Engineering Bulletin*, volume 38, issue 1, pages 18–31, March 2015. Archived at [perma.cc/68ZU-45SH](https://perma.cc/68ZU-45SH) 
[^57]: Daniel Peng and Frank Dabek. [Large-Scale Incremental Processing Using Distributed Transactions and Notifications](https://www.usenix.org/legacy/event/osdi10/tech/full_papers/Peng.pdf). At *9th USENIX Conference on Operating Systems Design and Implementation* (OSDI), October 2010. 
[^58]: Tushar Deepak Chandra, Robert Griesemer, and Joshua Redstone. [Paxos Made Live – An Engineering Perspective](https://www.read.seas.harvard.edu/~kohler/class/08w-dsi/chandra07paxos.pdf). At *26th ACM Symposium on Principles of Distributed Computing* (PODC), June 2007. [doi:10.1145/1281100.1281103](https://doi.org/10.1145/1281100.1281103) 
[^59]: Will Portnoy. [Lessons Learned from Implementing Paxos](https://blog.willportnoy.com/2012/06/lessons-learned-from-paxos.html). *blog.willportnoy.com*, June 2012. Archived at [perma.cc/QHD9-FDD2](https://perma.cc/QHD9-FDD2) 
[^60]: Brian M. Oki and Barbara H. Liskov. [Viewstamped Replication: A New Primary Copy Method to Support Highly-Available Distributed Systems](https://pmg.csail.mit.edu/papers/vr.pdf). At *7th ACM Symposium on Principles of Distributed Computing* (PODC), August 1988. [doi:10.1145/62546.62549](https://doi.org/10.1145/62546.62549) 
[^61]: Barbara H. Liskov and James Cowling. [Viewstamped Replication Revisited](https://pmg.csail.mit.edu/papers/vr-revisited.pdf). Massachusetts Institute of Technology, Tech Report MIT-CSAIL-TR-2012-021, July 2012. Archived at [perma.cc/56SJ-WENQ](https://perma.cc/56SJ-WENQ) 
[^62]: Leslie Lamport. [The Part-Time Parliament](https://www.microsoft.com/en-us/research/publication/part-time-parliament/). *ACM Transactions on Computer Systems*, volume 16, issue 2, pages 133–169, May 1998. [doi:10.1145/279227.279229](https://doi.org/10.1145/279227.279229) 
[^63]: Leslie Lamport. [Paxos Made Simple](https://www.microsoft.com/en-us/research/publication/paxos-made-simple/). *ACM SIGACT News*, volume 32, issue 4, pages 51–58, December 2001. Archived at [perma.cc/82HP-MNKE](https://perma.cc/82HP-MNKE) 
[^64]: Robbert van Renesse and Deniz Altinbuken. [Paxos Made Moderately Complex](https://people.cs.umass.edu/~arun/590CC/papers/paxos-moderately-complex.pdf). *ACM Computing Surveys* (CSUR), volume 47, issue 3, article no. 42, February 2015. [doi:10.1145/2673577](https://doi.org/10.1145/2673577) 
[^65]: Diego Ongaro. [Consensus: Bridging Theory and Practice](https://github.com/ongardie/dissertation). PhD Thesis, Stanford University, August 2014. Archived at [perma.cc/5VTZ-2ADH](https://perma.cc/5VTZ-2ADH) 
[^66]: Heidi Howard, Malte Schwarzkopf, Anil Madhavapeddy, and Jon Crowcroft. [Raft Refloated: Do We Have Consensus?](https://www.cl.cam.ac.uk/research/srg/netos/papers/2015-raftrefloated-osr.pdf) *ACM SIGOPS Operating Systems Review*, volume 49, issue 1, pages 12–21, January 2015. [doi:10.1145/2723872.2723876](https://doi.org/10.1145/2723872.2723876) 
[^67]: André Medeiros. [ZooKeeper’s Atomic Broadcast Protocol: Theory and Practice](http://www.tcs.hut.fi/Studies/T-79.5001/reports/2012-deSouzaMedeiros.pdf). Aalto University School of Science, March 2012. Archived at [perma.cc/FVL4-JMVA](https://perma.cc/FVL4-JMVA) 
[^68]: Robbert van Renesse, Nicolas Schiper, and Fred B. Schneider. [Vive La Différence: Paxos vs. Viewstamped Replication vs. Zab](https://arxiv.org/abs/1309.5671). *IEEE Transactions on Dependable and Secure Computing*, volume 12, issue 4, pages 472–484, September 2014. [doi:10.1109/TDSC.2014.2355848](https://doi.org/10.1109/TDSC.2014.2355848) 
[^69]: Heidi Howard and Richard Mortier. [Paxos vs Raft: Have we reached consensus on distributed consensus?](https://arxiv.org/abs/2004.05074). At *7th Workshop on Principles and Practice of Consistency for Distributed Data* (PaPoC), April 2020. [doi:10.1145/3380787.3393681](https://doi.org/10.1145/3380787.3393681) 
[^70]: Miguel Castro and Barbara H. Liskov. [Practical Byzantine Fault Tolerance and Proactive Recovery](https://www.microsoft.com/en-us/research/wp-content/uploads/2017/01/p398-castro-bft-tocs.pdf). *ACM Transactions on Computer Systems*, volume 20, issue 4, pages 396–461, November 2002. [doi:10.1145/571637.571640](https://doi.org/10.1145/571637.571640) 
[^71]: Shehar Bano, Alberto Sonnino, Mustafa Al-Bassam, Sarah Azouvi, Patrick McCorry, Sarah Meiklejohn, and George Danezis. [SoK: Consensus in the Age of Blockchains](https://doi.org/10.1145/3318041.3355458). At *1st ACM Conference on Advances in Financial Technologies* (AFT), October 2019. [doi:10.1145/3318041.3355458](https://doi.org/10.1145/3318041.3355458)
[^72]: Michael J. Fischer, Nancy Lynch, and Michael S. Paterson. [Impossibility of Distributed Consensus with One Faulty Process](https://groups.csail.mit.edu/tds/papers/Lynch/jacm85.pdf). *Journal of the ACM*, volume 32, issue 2, pages 374–382, April 1985. [doi:10.1145/3149.214121](https://doi.org/10.1145/3149.214121) 
[^73]: Tushar Deepak Chandra and Sam Toueg. [Unreliable Failure Detectors for Reliable Distributed Systems](https://courses.csail.mit.edu/6.852/08/papers/CT96-JACM.pdf). *Journal of the ACM*, volume 43, issue 2, pages 225–267, March 1996. [doi:10.1145/226643.226647](https://doi.org/10.1145/226643.226647) 
[^74]: Michael Ben-Or. [Another Advantage of Free Choice: Completely Asynchronous Agreement Protocols](https://homepage.cs.uiowa.edu/~ghosh/BenOr.pdf). At *2nd ACM Symposium on Principles of Distributed Computing* (PODC), August 1983. [doi:10.1145/800221.806707](https://doi.org/10.1145/800221.806707) 
[^75]: Cynthia Dwork, Nancy Lynch, and Larry Stockmeyer. [Consensus in the Presence of Partial Synchrony](https://groups.csail.mit.edu/tds/papers/Lynch/jacm88.pdf). *Journal of the ACM*, volume 35, issue 2, pages 288–323, April 1988. [doi:10.1145/42282.42283](https://doi.org/10.1145/42282.42283) 
[^76]: Xavier Défago, André Schiper, and Péter Urbán. [Total Order Broadcast and Multicast Algorithms: Taxonomy and Survey](https://dspace.jaist.ac.jp/dspace/bitstream/10119/4883/1/defago_et_al.pdf). *ACM Computing Surveys*, volume 36, issue 4, pages 372–421, December 2004. [doi:10.1145/1041680.1041682](https://doi.org/10.1145/1041680.1041682) 
[^77]: Hagit Attiya and Jennifer Welch. *Distributed Computing: Fundamentals, Simulations and Advanced Topics*, 2nd edition. John Wiley & Sons, 2004. ISBN: 978-0-471-45324-6, [doi:10.1002/0471478210](https://doi.org/10.1002/0471478210) 
[^78]: Rachid Guerraoui. [Revisiting the Relationship Between Non-Blocking Atomic Commitment and Consensus](https://citeseerx.ist.psu.edu/pdf/5d06489503b6f791aa56d2d7942359c2592e44b0). At *9th International Workshop on Distributed Algorithms* (WDAG), September 1995. [doi:10.1007/BFb0022140](https://doi.org/10.1007/BFb0022140) 
[^79]: Jim N. Gray and Leslie Lamport. [Consensus on Transaction Commit](https://dsf.berkeley.edu/cs286/papers/paxoscommit-tods2006.pdf). *ACM Transactions on Database Systems* (TODS), volume 31, issue 1, pages 133–160, March 2006. [doi:10.1145/1132863.1132867](https://doi.org/10.1145/1132863.1132867) 
[^80]: Fred B. Schneider. [Implementing Fault-Tolerant Services Using the State Machine Approach: A Tutorial](https://www.cs.cornell.edu/fbs/publications/SMSurvey.pdf). *ACM Computing Surveys*, volume 22, issue 4, pages 299–319, December 1990. [doi:10.1145/98163.98167](https://doi.org/10.1145/98163.98167) 
[^81]: Alexander Thomson, Thaddeus Diamond, Shu-Chun Weng, Kun Ren, Philip Shao, and Daniel J. Abadi. [Calvin: Fast Distributed Transactions for Partitioned Database Systems](https://cs.yale.edu/homes/thomson/publications/calvin-sigmod12.pdf). At *ACM International Conference on Management of Data* (SIGMOD), May 2012. [doi:10.1145/2213836.2213838](https://doi.org/10.1145/2213836.2213838) 
[^82]: Mahesh Balakrishnan, Dahlia Malkhi, Ted Wobber, Ming Wu, Vijayan Prabhakaran, Michael Wei, John D. Davis, Sriram Rao, Tao Zou, and Aviad Zuck. [Tango: Distributed Data Structures over a Shared Log](https://www.microsoft.com/en-us/research/publication/tango-distributed-data-structures-over-a-shared-log/). At *24th ACM Symposium on Operating Systems Principles* (SOSP), November 2013. [doi:10.1145/2517349.2522732](https://doi.org/10.1145/2517349.2522732) 
[^83]: Mahesh Balakrishnan, Dahlia Malkhi, Vijayan Prabhakaran, Ted Wobber, Michael Wei, and John D. Davis. [CORFU: A Shared Log Design for Flash Clusters](https://www.usenix.org/system/files/conference/nsdi12/nsdi12-final30.pdf). At *9th USENIX Symposium on Networked Systems Design and Implementation* (NSDI), April 2012. 
[^84]: Vasilis Gavrielatos, Antonios Katsarakis, and Vijay Nagarajan. [Odyssey: the impact of modern hardware on strongly-consistent replication protocols](https://vasigavr1.github.io/files/Odyssey_Eurosys_2021.pdf). At *16th European Conference on Computer Systems* (EuroSys), April 2021. [doi:10.1145/3447786.3456240](https://doi.org/10.1145/3447786.3456240) 
[^85]: Heidi Howard, Dahlia Malkhi, and Alexander Spiegelman. [Flexible Paxos: Quorum Intersection Revisited](https://drops.dagstuhl.de/opus/volltexte/2017/7094/pdf/LIPIcs-OPODIS-2016-25.pdf). At *20th International Conference on Principles of Distributed Systems* (OPODIS), December 2016. [doi:10.4230/LIPIcs.OPODIS.2016.25](https://doi.org/10.4230/LIPIcs.OPODIS.2016.25) 
[^86]: Martin Kleppmann. [Distributed Systems lecture notes](https://www.cl.cam.ac.uk/teaching/2425/ConcDisSys/dist-sys-notes.pdf). *University of Cambridge*, October 2024. Archived at [perma.cc/SS3Q-FNS5](https://perma.cc/SS3Q-FNS5) 
[^87]: Kyle Kingsbury. [Call Me Maybe: Elasticsearch 1.5.0](https://aphyr.com/posts/323-call-me-maybe-elasticsearch-1-5-0). *aphyr.com*, April 2015. Archived at [perma.cc/37MZ-JT7H](https://perma.cc/37MZ-JT7H) 
[^88]: Heidi Howard and Jon Crowcroft. [Coracle: Evaluating Consensus at the Internet Edge](https://conferences.sigcomm.org/sigcomm/2015/pdf/papers/p85.pdf). At *Annual Conference of the ACM Special Interest Group on Data Communication* (SIGCOMM), August 2015. [doi:10.1145/2829988.2790010](https://doi.org/10.1145/2829988.2790010) 
[^89]: Tom Lianza and Chris Snook. [A Byzantine failure in the real world](https://blog.cloudflare.com/a-byzantine-failure-in-the-real-world/). *blog.cloudflare.com*, November 2020. Archived at [perma.cc/83EZ-ALCY](https://perma.cc/83EZ-ALCY) 
[^90]: Ivan Kelly. [BookKeeper Tutorial](https://github.com/ivankelly/bookkeeper-tutorial). *github.com*, October 2014. Archived at [perma.cc/37Y6-VZWU](https://perma.cc/37Y6-VZWU) 
[^91]: Jack Vanlightly. [Apache BookKeeper Insights Part 1 — External Consensus and Dynamic Membership](https://medium.com/splunk-maas/apache-bookkeeper-insights-part-1-external-consensus-and-dynamic-membership-c259f388da21). *medium.com*, November 2021. Archived at [perma.cc/3MDB-8GFB](https://perma.cc/3MDB-8GFB)
