9 分布式系统的麻烦

意外这东西挺有意思:你没碰上之前,它就从来不会发生。
A.A. 米尔恩,《小熊维尼和老灰驴的家》(1928)
正如 “可靠性与容错” 中所讨论的,要使一个系统可靠,就要确保即使出了问题(即发生故障),整个系统仍能继续工作。然而,要预见并处理所有可能的故障并非易事。开发者很容易把注意力主要放在正常路径上(毕竟大多数时候一切都运行良好!),而忽略会带来大量边界情况的故障。
如果希望系统在发生故障时依然可靠,就必须从根本上转变思维方式,把注意力放在各种可能出错的地方,即使出错的概率很低。某件事只有百万分之一的概率出错并不意味着可以不管:系统足够大时,百万分之一的事件每天都会发生。经验丰富的系统运维人员会告诉你,任何 可能 出错的事情,终究都会 出错。
而且,使用 分布式系统(distributed system)与在单台计算机上编写软件有着根本区别——最主要的区别,就是事情有了许多新奇而刺激的出错方式 1 2。本章将带你领略实践中会遇到的问题,并帮助你理解哪些东西可以依赖,哪些不可以。
为了理解我们面对的挑战,接下来让我们把悲观主义发挥到极致,考察分布式系统里可能出错的种种事情。我们将讨论网络问题(“不可靠的网络”),以及时钟和时序问题(“不可靠的时钟”)。这些问题造成的后果往往令人迷失方向,因此我们还要探讨如何认识分布式系统的状态,以及如何推断已经发生过的事情(“知识、真相和谎言”)。随后在 第 10 章 中,我们将通过一些例子看看,面对这些故障时如何实现容错。
故障与部分失效
当你在一台计算机上编写程序时,它通常会以相当可预测的方式运行:要么正常工作,要么彻底罢工。有缺陷的软件可能会让人觉得计算机偶尔也会“状态不好”(而重启往往能解决问题),但这通常只是软件写得糟糕所造成的表象。
从根本上说,单台计算机上的软件不应该时灵时不灵:只要硬件正常,同样的操作总会产生同样的结果(也就是 确定性的,deterministic)。如果硬件出了问题(例如内存损坏或连接器松动),后果通常是整个系统失效(例如内核恐慌、“蓝屏死机”或无法启动)。一台运行着良好软件的计算机,通常要么功能完好,要么完全失效,而不会停留在两者之间。
这是计算机设计中的有意选择:发生内部故障时,我们宁愿让计算机彻底崩溃,也不愿让它返回错误结果,因为后者既难处理又容易造成混乱。于是,计算机把其底层模糊而混乱的物理现实隐藏起来,呈现出一个以数学般的完美方式运行的理想化系统模型。CPU 指令每次都会做同样的事情;写入内存或磁盘的数据会原样保留,不会随机损坏。正如 “硬件与软件故障” 中所讨论的,事实并非真的如此——现实中,数据的确可能在没有任何警告的情况下损坏,CPU 偶尔也会悄无声息地给出错误结果——只不过这些情况足够罕见,通常可以忽略。
当软件运行在通过网络连接的多台计算机上时,情况就截然不同了。分布式系统中的故障要频繁得多,我们再也无法视而不见,只能直面物理世界混乱的现实。而在物理世界中,可能出错的事情多得惊人,下面这段轶事便是一个写照 3:
在我有限的从业经历中,我处理过单个数据中心(DC)内长期存在的网络分区、PDU(配电单元)故障、交换机故障、整个机架意外断电重启、整个数据中心骨干网故障、整个数据中心停电,以及一名低血糖司机开着福特皮卡撞进数据中心的 HVAC(供暖、通风与空调)系统。而我甚至还不是运维人员。
—— 柯达黑尔
在分布式系统中,即使其他部分工作正常,系统的某些部分也很可能以不可预知的方式发生故障。这叫作 部分失效(partial failure)。棘手之处在于,部分失效是 非确定性的(nondeterministic):任何涉及多个节点及其网络的操作,有时能够成功,有时却会莫名其妙地失败。正如我们稍后会看到的,你甚至可能根本 不知道 某件事究竟成功了没有!
正是这种非确定性和部分失效的可能性,让分布式系统如此难以驾驭 4。不过,如果分布式系统能够容忍部分失效,也会由此获得强大的能力。例如,你可以执行滚动升级:每次重启一个节点来安装软件更新,同时让整个系统始终不间断地工作。因此,借助容错,我们可以用不可靠的组件构建出比单节点系统更可靠的分布式系统。
不过,在实现容错之前,我们需要进一步了解究竟要容忍哪些故障。应当考虑尽可能广泛的故障——包括那些看来不太可能发生的故障——并在测试环境中人为制造这些情形,观察系统会怎样反应。在分布式系统中,多一些怀疑、悲观和偏执总会有所回报。
不可靠的网络
正如 “共享内存、共享磁盘与无共享架构” 中所讨论的,本书关注的分布式系统大多是 无共享系统,也就是一组通过网络连接的机器。网络是这些机器彼此通信的唯一途径——我们假定每台机器都有自己的内存和磁盘,一台机器无法直接访问另一台机器的内存或磁盘,只能通过网络向服务发送请求。即使存储本身是共享的(例如 Amazon S3),机器也仍然要通过网络与共享存储服务通信。
互联网和数据中心里的大多数内部网络(通常是以太网)都是 异步分组网络(asynchronous packet network)。在这种网络中,一个节点可以向另一个节点发送消息(即数据包),但网络既不保证消息何时到达,也不保证它一定能够到达。如果你发出请求并等待响应,可能会发生许多问题(其中一些如 图 9-1 所示):
- 请求可能已经丢失(或许有人拔掉了网线)。
- 请求可能还在队列中等待,稍后才会送达(或许网络或接收方过载了)。
- 远程节点可能已经失效(或许它崩溃了,或是被关闭了)。
- 远程节点可能只是暂时停止响应(或许它正经历一次漫长的垃圾回收暂停;参见 “进程暂停”),稍后又会恢复响应。
- 远程节点可能已经处理了请求,但响应在网络中丢失了(或许某台网络交换机配置有误)。
- 远程节点可能已经处理了请求,但响应被延迟,稍后才会送达(或许网络或你自己的机器过载了)。

发送方甚至无法判断数据包是否送达:唯一的办法是由接收方发回响应消息,而这个响应同样可能丢失或延迟。在异步网络中,这些情形无法区分:你掌握的唯一信息只是“尚未收到响应”。向另一个节点发送请求却没有收到响应时,不可能 判断原因究竟是什么。
处理这个问题的惯常办法是设置 超时(timeout):等待一段时间后便放弃,并假定响应不会再来。然而,即使发生超时,你仍然不知道远程节点是否收到了请求(如果请求仍在某处排队,那么即便发送方已经放弃,它仍可能在稍后送达接收方)。
TCP 的局限性
网络数据包有大小上限(通常只有几千字节),但许多应用程序需要发送无法装进单个数据包的消息,例如请求和响应。这类应用程序通常使用 TCP(传输控制协议)建立一条 连接(connection),把较大的数据流拆成一个个数据包,再在接收端重新组装起来。
下面关于 TCP 的大部分讨论,也适用于较新的替代方案 QUIC、WebRTC 使用的流控制传输协议(SCTP)、BitTorrent 的 uTP 协议,以及其他传输协议。关于它与 UDP 的比较,参见 “TCP 与 UDP”。
TCP 常被说成能提供“可靠”的传输,这里的“可靠”是指:它能检测并重传丢失的数据包,发现顺序错乱的数据包并将其恢复为正确顺序,还能用简单的校验和检测数据包损坏。TCP 也会判断应当以多快的速度发送数据,既尽可能快速传输,又不至于压垮网络或接收节点;这叫作 拥塞控制(congestion control)、流量控制(flow control)或 背压(backpressure)5。
当你把数据写入套接字来“发送”时,数据其实不会立即发出,只会先进入操作系统管理的缓冲区。拥塞控制算法判断目前有能力发送数据包后,才会从缓冲区取出一个数据包大小的数据,交给网络接口。数据包会经过若干交换机和路由器,最终由接收节点的操作系统把数据放进接收缓冲区,并向发送方发回确认包。直到这时,接收端操作系统才会通知应用程序又有数据到达 6。
既然 TCP 提供了“可靠性”,是不是就不必再操心网络不可靠了?遗憾的是,并非如此。如果在一定的超时时间内没有收到确认,TCP 会认定数据包必定已经丢失,但它同样无法判断丢掉的究竟是发出的数据包,还是返回的确认包。TCP 虽然可以重发,却不能保证重发的数据包一定能够到达。如果网线被拔了,TCP 可没法替你把它插回去。最终,经过可配置的超时时间后,TCP 会放弃重试,并向应用程序报告错误。
如果 TCP 连接因错误而关闭——或许是远程节点崩溃了,也可能是网络中断了——你无法知道远程节点究竟处理了多少数据 6。即使 TCP 已经确认某个数据包送达,也只说明远程节点的操作系统内核收到了它;应用程序仍可能在处理这份数据之前就崩溃。如果要确认请求成功,必须由应用程序本身明确返回成功响应 7。
尽管如此,TCP 依然非常有用,因为它让我们能够方便地收发无法装进单个数据包的消息。建立 TCP 连接后,还可以通过同一条连接发送多个请求和响应。通常的做法是先发送一个头部,注明紧随其后的消息有多少字节,再发送消息本身。HTTP 和许多 RPC 协议(参见 “流经服务的数据流:REST 与 RPC”)都是这样工作的。
实践中的网络故障
我们建设计算机网络已有几十年,按说早该找到让网络可靠的办法了。遗憾的是,我们至今仍未成功。系统性研究与大量轶事证据都表明,即使在由一家公司运营的数据中心这类受控环境中,网络问题也可能频繁得出人意料 8:
- 一项针对中型数据中心的研究发现,每月大约会发生 12 次网络故障,其中一半会断开一台机器,另一半则会断开整个机架 9。
- 另一项研究测量了架顶交换机、汇聚交换机和负载均衡器等组件的失效率 10。研究发现,增加冗余网络设备并不能像预想的那样大幅减少故障,因为它防不住人为错误(例如交换机配置错误),而人为错误正是停机的一大主因。
- 广域光纤链路的中断曾被归咎于奶牛 11、海狸 12 和鲨鱼 13(不过随着海底电缆的防护改善,鲨鱼咬坏电缆的事件已经越来越少 14)。当然,人类也难辞其咎:误配置 15、盗割电缆 16 和蓄意破坏 17 都曾造成事故。
- 在不同云区域之间的通信中,高百分位数的往返时间最长可达数 分钟 18。即使在同一个数据中心内,交换机软件升级出现问题并触发网络拓扑重配置时,数据包也可能延迟一分钟以上 19。因此,我们必须假定消息可能遭到任意长时间的延迟。
- 有时通信只会部分中断,能否连通取决于通信双方是谁。例如,A 与 B 可以通信,B 与 C 也可以通信,但 A 与 C 却无法通信 20 21。还有些故障更出人意料,比如某个网络接口有时会丢弃所有入站数据包,却仍能正常发出数据包 22:网络链路在一个方向上可用,并不保证反方向也可用。
- 即使网络中断只持续了很短时间,其后续影响也可能远远长于最初的问题 8 20 23。
当网络的一部分因网络故障而与其余部分隔绝时,这种情况有时称为 网络分区(network partition)或 网络分裂(netsplit),但它与其他类型的网络中断并没有本质区别。网络分区与存储系统的分片无关,后者有时也称为 分区(参见 第 7 章)。
即使你的环境很少遇到网络故障,故障 有可能 发生这一事实也意味着软件必须能够处理它。只要通过网络通信,就有可能失败——这一点无可回避。
如果没有明确定义并测试网络故障的处理方式,后果可能糟到没有下限。例如,即使网络已经恢复,集群仍可能陷入死锁,从此无法再处理请求 24;甚至还可能删掉你的全部数据 25。一旦软件落入设计者未曾预料的情形,它就可能做出任何出人意料的事情。
处理网络故障不一定意味着要 容忍 它。如果网络通常相当可靠,那么在发生网络问题时直接向用户显示错误消息,也可以是一种合理策略。不过,你必须知道软件会如何应对网络问题,并确保系统能够从中恢复。可以考虑有意触发网络问题,测试系统会如何反应(这叫作 故障注入,fault injection;参见 “故障注入”)。
故障检测
许多系统需要自动检测发生故障的节点。例如:
- 负载均衡器需要停止向已经宕机的节点发送请求(也就是把它 移出轮询池)。
- 在采用单主复制的分布式数据库中,如果领导者失效,就需要把某个追随者提升为新的领导者(参见 “处理节点故障”)。
遗憾的是,网络的不确定性让人很难判断节点是否仍在工作。在某些特定情况下,你也许能收到明确表明某处出了问题的反馈:
- 如果运行节点的机器可以访问,但目标端口没有进程监听(例如进程已经崩溃),操作系统会返回
RST或FIN数据包,从而关闭或拒绝 TCP 连接。 - 如果节点进程崩溃(或被管理员终止),但节点的操作系统仍在运行,可以由脚本把崩溃消息通知其他节点,让另一个节点不必等待超时便能迅速接管。例如,HBase 就采用这种方式 26。
- 如果你能访问数据中心里网络交换机的管理接口,就可以查询交换机,在硬件层面检测链路故障(例如远程机器是否已经断电)。但如果你通过互联网连接、身处无法访问交换机本身的共享数据中心,或是网络问题使管理接口也无法访问,这种办法就行不通了。
- 如果路由器确信你要连接的 IP 地址不可达,它可能会返回一个 ICMP“目标不可达”数据包。不过,路由器也没有什么神奇的故障检测能力——它同样受到网络中其他参与者所面对的那些限制。
能够迅速得知远程节点宕机固然有用,但不能指望总有这样的反馈。发生问题时,你或许会从协议栈的某一层收到错误响应,但通常必须假定自己什么响应也收不到。你可以重试几次,等待超时;如果超时时间内一直没有回复,最终便宣告该节点已经死亡。
超时和无界延迟
如果超时是检测故障的唯一可靠方法,那么超时时间应该设为多长?遗憾的是,这个问题没有简单答案。
较长的超时意味着要等很久才能宣告节点死亡(在此期间,用户可能只能等待,或不断看到错误消息)。较短的超时可以更快发现故障,却更容易把只是暂时变慢的节点(例如节点或网络出现负载峰值)误判为死亡。
过早宣告节点死亡会带来麻烦:如果节点其实仍然活着,并且正在执行某个操作(例如发送电子邮件),此时另一个节点又接管了它的工作,同一个操作最终可能执行两次。我们将在 “知识、真相和谎言”、第 10 章 以及 “数据库的端到端原则” 中更详细地讨论这个问题。
宣告节点死亡后,它承担的职责就要转移给其他节点,这会额外增加其他节点和网络的负担。如果系统本来就在高负载下勉力支撑,过早宣告节点死亡可能让局面进一步恶化。尤其可能出现这样的情况:节点其实没有死,只是过载导致响应缓慢;把它的负载转移给其他节点,又可能引发级联失效(极端情况下,所有节点互相宣告对方死亡,整个系统彻底停止工作——参见 “当过载系统无法恢复时”)。
设想一个虚构的系统,其网络保证数据包的最大延迟:每个数据包要么在时间 d 以内送达,要么丢失,绝不会在超过 d 之后才送达。再假定我们能够保证,任何尚未失效的节点总会在时间 r 以内处理完请求。在这种情况下,每个成功的请求都能保证在 2 d + r 以内收到响应;如果过了这么久仍未收到响应,就可以断定网络或远程节点没有正常工作。假如这些保证真的成立,那么 2 d + r 就是合理的超时时间。
遗憾的是,我们使用的大多数系统都不具备上述任何一项保证:异步网络具有 无界延迟(unbounded delay;也就是说,它会尽快尝试送达数据包,但数据包所需的传输时间没有上限),大多数服务器实现也不能保证一定在某个最长时间内处理完请求(参见 “响应时间保证”)。对故障检测而言,系统仅仅在大多数时候运行得很快还不够:如果超时时间设得很短,一次短暂的往返时间尖峰就足以打乱整个系统。
网络拥塞与排队
开车出行时,交通拥堵往往是造成行程时间波动的最大因素。同样,在计算机网络上,数据包延迟的变化通常也是排队造成的 27:
- 如果多个节点同时向同一个目的地发送数据包,网络交换机就必须让这些数据包排队,再逐一送入通往目的地的网络链路(如 图 9-2 所示)。网络链路繁忙时,数据包可能需要等待一段时间才能获得发送机会,这叫作 网络拥塞(network congestion)。如果流入的数据太多,交换机队列被塞满,数据包就会被丢弃,因而必须重发——即使网络本身仍在正常工作。
- 数据包抵达目标机器时,如果所有 CPU 核心都在忙,操作系统会把来自网络的请求放进队列,直到应用程序有能力处理为止。等待时间取决于机器的负载,可能任意之长 28。
- 在虚拟化环境中,当另一台虚拟机占用某个 CPU 核心时,正在运行的操作系统常常会暂停数十毫秒。在此期间,虚拟机无法消费任何网络数据,因此虚拟机监控器会将流入的数据排队(缓冲)29,进一步增大网络延迟的波动。
- 如前所述,为避免网络过载,TCP 会限制数据的发送速率。这意味着数据甚至还没进入网络,就已经在发送方额外排了一次队。

此外,当 TCP 检测到数据包丢失并自动重传时,应用程序虽然不会直接看到丢包,却会感受到由此造成的延迟:先等待超时,再等待重传的数据包得到确认。
一些对延迟敏感的应用程序(例如视频会议和 IP 语音,即 VoIP)使用 UDP 而非 TCP。这是在可靠性与延迟波动之间做出的权衡:UDP 不进行流量控制,也不重传丢失的数据包,因而避开了网络延迟发生波动的部分原因(不过,它仍会受到交换机排队和调度延迟的影响)。
如果数据一旦延迟就不再有价值,UDP 是很好的选择。例如在 VoIP 通话中,等到音频该从扬声器播放时,通常已经来不及重传丢失的数据包。这种情况下,重传毫无意义——应用程序只能用静音填补丢失数据包对应的时间片(听起来就是声音短暂中断),然后继续播放后面的音频。真正的重试发生在人这一层。(“能再说一遍吗?刚才声音断了一下。”)
上述因素都会造成网络延迟的波动。当系统接近最大容量时,排队延迟的变化范围尤其大:拥有充足余量的系统可以轻松排空队列,而利用率很高的系统则可能迅速积起长队。
在公共云和多租户数据中心中,许多客户共享同一批资源:网络链路和交换机是共享的,甚至每台机器的网络接口和 CPU(使用虚拟机时)也是共享的。处理大量数据时,可能耗尽网络链路的全部容量,使其达到 饱和(saturation)。你既无法控制,也无法了解其他客户如何使用共享资源;如果附近某个 吵闹的邻居(noisy neighbor)正在大量消耗资源,网络延迟就可能剧烈波动 30 31。
在这样的环境中,只能通过实验来选择超时时间:长期测量大量机器之间网络往返时间的分布,确定延迟通常会有多大波动。然后结合应用程序自身的特点,在故障检测的延迟与过早超时的风险之间做出适当权衡。
更好的办法是不用固定配置的超时时间,而让系统持续测量响应时间及其波动(抖动,jitter),再根据观测到的响应时间分布自动调整超时。Phi 累积故障检测器 32 就采用了这种方法,Akka 和 Cassandra 等系统都在使用它 33。TCP 的重传超时也以类似方式工作 5。
同步网络与异步网络
如果可以指望网络在某个固定的最长延迟内送达数据包,而且绝不丢包,分布式系统就会简单得多。为什么不能在硬件层面解决这个问题,让网络变得可靠,从而使软件不必再操心呢?
要回答这个问题,不妨把数据中心网络与传统的固定电话网络(非蜂窝网络,也非 VoIP)比较一下。传统电话网络极其可靠:音频帧延迟和通话中断都很罕见。电话通话需要持续保持较低的端到端延迟,并提供足够带宽来传输语音采样。如果计算机网络也能具备类似的可靠性与可预测性,岂不是很好?
通过电话网络拨号时,网络会建立一条 电路(circuit):从一名通话者到另一名通话者的整条路径上,都会为这次通话分配固定且有保证的带宽。这条电路会一直保留到通话结束 34。例如,ISDN 网络以每秒 4,000 帧的固定速率运行。建立通话时,每一帧都会在两个方向上各分配 16 位空间。因此在整个通话期间,双方都能保证每 250 微秒恰好发送 16 位音频数据 35。
这类网络是 同步的(synchronous):即使数据要经过多台路由器,也不会排队,因为每一跳都已经为这次通话预留了 16 位的空间。由于无需排队,网络的最大端到端延迟是固定的。我们称之为 有界延迟(bounded delay)。
我们不能简单地让网络延迟可预测吗?
注意,电话网络中的电路与 TCP 连接截然不同:电路会预留固定数量的带宽,只要电路存在,其他人就无法使用这部分带宽;TCP 连接的数据包则会伺机利用当时可用的任意网络带宽。你可以交给 TCP 一块大小不定的数据(例如电子邮件或网页),它会尽力在最短时间内传输完毕。TCP 连接空闲时不会占用带宽,偶尔发送的保活包除外。
如果数据中心网络和互联网采用电路交换,那么建立电路时就可以同时确定有保证的最大往返时间。但它们并非如此:以太网和 IP 都是分组交换协议,数据包会排队,因而网络延迟没有上界。这些协议中根本没有“电路”的概念。
为什么数据中心网络和互联网要使用分组交换?因为它们针对 突发流量(bursty traffic)做了优化。音频或视频通话在整个通话期间每秒传输的位数相当稳定,很适合使用电路。相比之下,请求网页、发送电子邮件或传输文件并没有特定的带宽要求——我们只希望它们尽快完成。
如果要通过电路传输文件,就必须先猜测应当分配多少带宽。猜得太低,传输速度会慢得毫无必要,同时还有网络容量闲置;猜得太高,电路又无法建立,因为网络不能为无法保证带宽分配的电路放行。因此,用电路传输突发数据既浪费网络容量,又使传输无谓地变慢。TCP 则不同,它会根据可用的网络容量动态调整数据传输速率。
人们也曾尝试构建兼具电路交换与分组交换的混合网络。异步传输模式(Asynchronous Transfer Mode,ATM)在 20 世纪 80 年代曾与以太网竞争,但除了电话网络的核心交换机外,并未得到广泛采用。InfiniBand 与之有些相似 36:它在链路层实现端到端流量控制,减少了网络排队的必要,不过链路拥塞仍可能带来延迟 37。如果谨慎运用 服务质量(quality of service,QoS,即数据包的优先级与调度)和 准入控制(admission control,即对发送方限速),就可以在分组网络上模拟电路交换,或提供统计意义上的有界延迟 27 34。低延迟、低损耗和可伸缩吞吐量(L4S)等新型网络算法,试图从客户端和路由器两端缓解部分排队与拥塞控制问题。Linux 的流量控制器(TC)也允许应用程序为实现 QoS 而重新安排数据包的优先级。
更一般地说,可以把延迟的波动看作动态划分资源的结果。
假设两台电话交换机之间有一条线路,最多可以承载 10,000 路并发通话。经这条线路交换的每条电路都会占用一个通话槽位。因此,可以把这条线路看作一项最多由 10,000 名并发用户共享的资源。资源以 静态 方式划分:即使现在整条线路上只有你这一通电话,另外 9,999 个槽位全部空闲,你的电路得到的仍然只是那份固定带宽,与线路满载时一模一样。
相比之下,互联网以 动态 方式共享网络带宽。发送方彼此争抢,都想尽快把自己的数据包送上线路;网络交换机则随时决定接下来发送哪个数据包,也就是如何分配带宽。这种方法的缺点是会造成排队,优点则是能最大限度地利用线路。线路的成本是固定的,利用率越高,经由它发送的每个字节就越便宜。
CPU 也有类似的情况:如果多个线程动态共享一个 CPU 核心,那么当另一个线程正在运行时,某个线程有时就必须在操作系统的运行队列中等待,因而可能暂停长短不一的时间 38。不过,与给每个线程静态分配固定数量的 CPU 周期相比,这样做能更充分地利用硬件(参见 “响应时间保证”)。为了提高硬件利用率,云平台也会在同一台物理机器上运行来自不同客户的多台虚拟机。
在某些环境中,只要静态划分资源(例如采用专用硬件并分配独占带宽),就能提供延迟保证。但代价是资源利用率降低——换句话说,也就是成本更高。动态划分资源的多租户模式利用率更好,因此更便宜,代价则是延迟会发生波动。
网络延迟会发生波动并非自然法则,只是成本与收益权衡的结果。
不过,多租户数据中心、公共云以及经由互联网的通信,目前都没有启用这样的服务质量机制。现有部署技术无法让我们对网络的延迟或可靠性作出任何保证:必须假定网络拥塞、排队和无界延迟都会发生。因此,超时时间没有所谓“正确”的取值,只能通过实验来确定。
互联网服务提供商之间的对等互联协议,以及通过边界网关协议(BGP)建立路由的方式,比 IP 本身更像电路交换。在这个层面上,确实可以买到专用带宽。不过,互联网路由工作在网络层面,而不是主机之间的单条连接层面,其时间尺度也长得多。
不可靠的时钟
时钟和时间都很重要。应用程序会以各种方式依赖时钟,回答下面这些问题:
- 这个请求已经超时了吗?
- 这项服务响应时间的第 99 百分位数是多少?
- 过去五分钟里,这项服务平均每秒处理多少次查询?
- 用户在我们的网站上停留了多长时间?
- 这篇文章是什么时候发布的?
- 提醒邮件应该在哪一天、什么时间发出?
- 这个缓存条目什么时候过期?
- 日志文件里这条错误消息的时间戳是什么?
问题 1~4 测量的是 持续时间(duration;例如从发送请求到收到响应之间的时间间隔),问题 5~8 描述的则是 时间点(point in time;在某个具体日期和时刻发生的事件)。
在分布式系统中,时间是个棘手的问题,因为通信并非瞬间完成:消息从一台机器经网络传到另一台机器需要时间。收到消息的时刻必然晚于发出消息的时刻,但由于网络延迟会发生变化,我们不知道究竟晚了多少。涉及多台机器时,这一事实有时会让事件的先后顺序难以确定。
此外,网络中的每台机器都有自己的时钟,而且它是实实在在的硬件设备,通常是石英晶体振荡器。这些设备并不十分精确,因此每台机器都各有一套时间观念,可能比其他机器走得稍快或稍慢。时钟可以在一定程度上同步:最常用的机制是网络时间协议(NTP),它根据一组服务器报告的时间来校准计算机时钟 39;这些服务器又从 GPS 接收器等更精确的时间源取得时间。
单调时钟与日历时钟
现代计算机至少配有两种不同的时钟:日历时钟(time-of-day clock)和 单调时钟(monotonic clock)。它们虽然都用来度量时间,却服务于不同目的,因此必须加以区分。
日历时钟
日历时钟所做的,正是人们直觉上认为时钟应该做的事:按照某种历法返回当前日期和时间(也称为 墙上时钟时间,wall-clock time)。例如,Linux 的 clock_gettime(CLOCK_REALTIME) 和 Java 的 System.currentTimeMillis() 返回从 纪元(epoch)至今经过的秒数(或毫秒数):这里的纪元是格里高利历中的 1970 年 1 月 1 日 UTC 零时,且不计闰秒。有些系统使用其他日期作为参考点。(Linux 虽然把这个时钟称为 实时 时钟,但它与实时操作系统毫无关系,参见 “响应时间保证”。)
日历时钟通常会与 NTP 同步,因此理想情况下,一台机器上的某个时间戳与另一台机器上的同一时间戳表示同一时刻。不过,日历时钟也有各种怪异之处,下一节将进一步说明。尤其是当本地时钟比 NTP 服务器快得太多时,它可能会被强制重置,看起来就像突然跳回了过去。这样的跳变,以及闰秒造成的类似跳变,使日历时钟不适合测量已经过去了多长时间 40。
夏令时(DST)开始或结束时,日历时钟也可能跳变;只要始终采用没有夏令时的 UTC 时区,就能避开这类跳变。历史上,日历时钟的分辨率也相当粗糙,例如旧版 Windows 系统的时钟每次会向前跳 10 毫秒 41。在较新的系统上,这已经不是什么大问题。
单调时钟
单调时钟适合测量持续时间(时间间隔),例如超时时间或服务响应时间。Linux 的 clock_gettime(CLOCK_MONOTONIC)、clock_gettime(CLOCK_BOOTTIME) 42 和 Java 的 System.nanoTime() 都属于单调时钟。之所以叫“单调”,是因为它保证只会向前走(日历时钟却可能突然跳回过去)。
你可以先读取一次单调时钟,做些事情,稍后再读一次。两次读数之 差 就是其间经过的时间——它更像秒表,而不是挂钟。不过,单调时钟的 绝对 读数没有任何意义:它可能表示计算机启动至今的纳秒数,也可能采用其他任意起点。尤其不能比较两台不同计算机的单调时钟读数,因为它们表示的并不是同一回事。
在有多个 CPU 插槽的服务器上,每颗 CPU 都可能有自己的计时器,而且未必与其他 CPU 同步 43。操作系统会补偿其中的差异,尽力让应用程序线程看到单调递增的时钟,即使线程在不同 CPU 之间调度也是如此。不过,对这样的单调性保证最好还是有所保留 44。
如果 NTP 发现计算机的本地石英钟比 NTP 服务器走得更快或更慢,可以调节单调时钟向前推进的速率,这叫作对时钟进行 渐进校准(slewing)。默认情况下,NTP 最多可以把时钟速率加快或减慢 0.05%,但不能让单调时钟突然向前或向后跳变。单调时钟通常有很好的分辨率:在大多数系统上,它能测量微秒级甚至更短的时间间隔。
在分布式系统中,用单调时钟测量经过的时间(例如超时)通常没有问题,因为这不要求不同节点的时钟彼此同步,而且测量中的轻微误差也不会造成太大影响。
时钟同步和准确性
单调时钟不需要同步,但日历时钟只有根据 NTP 服务器或其他外部时间源进行校准才有用。遗憾的是,让时钟显示正确时间的手段远没有想象中那么可靠和准确——硬件时钟与 NTP 都可能反复无常。下面只是其中几个例子:
- 计算机里的石英钟并不十分精确,它会发生 漂移(走得比应有速度更快或更慢),而漂移程度又会随机器温度变化。Google 假定其服务器的时钟漂移最高可达 200 ppm(百万分之二百)45。这相当于每隔 30 秒与服务器重新同步一次的时钟会漂移 6 毫秒,而每天才同步一次的时钟会漂移 17 秒。即使其他一切都正常,时钟漂移也会限制所能达到的最佳精度。
- 如果计算机时钟与 NTP 服务器相差太大,它可能拒绝同步,也可能强制重置本地时钟 39。在重置前后观察时间的应用程序,可能会看到时间向后倒退,或突然向前跳跃。
- 如果防火墙意外阻断了节点与 NTP 服务器的通信,这项错误配置可能很长时间都无人察觉;在此期间,漂移不断累积,不同节点的时钟最终可能相差甚远。轶事证据表明,实践中确实发生过这种情况。
- NTP 同步的准确性不可能优于网络延迟。因此,在数据包延迟波动不定的拥塞网络上,它的准确性必然有限。一项实验表明,经由互联网同步所能达到的最小误差为 35 毫秒 46,而网络延迟偶尔出现尖峰时,误差会达到一秒左右。具体取决于配置,过大的网络延迟甚至可能让 NTP 客户端彻底放弃同步。
- 有些 NTP 服务器本身不正确或配置有误,报告的时间能相差数小时 47 48。NTP 客户端会查询多台服务器并忽略离群值,以减轻这类错误的影响。即便如此,把系统的正确性押在某个互联网陌生人报给你的时间上,多少还是令人不安。
- 闰秒会让一分钟变成 59 秒或 61 秒,从而打乱那些设计时未考虑闰秒的系统对时序所作的假设 49。闰秒已经导致许多大型系统崩溃 40 50,足见关于时钟的错误假设有多么容易悄然混入系统。处理闰秒的最佳办法,或许是让 NTP 服务器“撒谎”:把闰秒调整分摊到一整天内逐渐完成,这称为 平滑处理(smearing)51 52;不过实践中各 NTP 服务器的实际行为并不一致 53。好在从 2035 年起将不再使用闰秒,这个问题也会随之消失。
- 在虚拟机中,硬件时钟也是虚拟化的,这给需要精确计时的应用程序带来了额外挑战 54。多个虚拟机共享 CPU 核心时,一台虚拟机运行,其他虚拟机就可能暂停数十毫秒。从应用程序的视角看,这种暂停表现为时钟突然向前跳跃 29。如果虚拟机暂停了几秒,其时钟随后可能比实际时间慢几秒,但 NTP 仍可能报告时钟几乎完全同步 55。
- 如果软件运行在你无法完全控制的设备上(例如移动设备或嵌入式设备),那么设备的硬件时钟可能根本不值得信任。有些用户会故意把硬件时钟设置成错误的日期和时间,例如借此在游戏中作弊 56。因此,时钟可能被设到离谱的过去或未来。
只要足够重视时钟精度,并愿意投入大量资源,的确可以做到非常精确。例如,欧洲针对金融机构的 MiFID II 法规要求所有高频交易基金把时钟与 UTC 的误差控制在 100 微秒以内,以便排查“闪崩”等市场异常,并帮助发现市场操纵 57。
借助专用硬件(GPS 接收器和/或原子钟)、精确时间协议(PTP),再辅以审慎的部署与监控,就能达到这样的精度 58 59。不过,只依赖 GPS 也有风险,因为 GPS 信号很容易受到干扰;某些地方(例如军事设施附近)甚至经常发生这种情况 60。一些云服务商已经开始为虚拟机提供高精度时钟同步 61。即便如此,时钟同步仍需格外小心。如果 NTP 守护进程配置有误,或防火墙阻断了 NTP 流量,漂移造成的时钟误差很快就会变得很大。
对同步时钟的依赖
时钟的问题在于,它看似简单易用,陷阱却多得惊人:一天未必恰好有 86,400 秒,日历时钟可能倒着走,一个节点所认为的时间也可能与另一个节点相差很大。
本章前面讨论过网络丢包和任意延迟数据包的问题。尽管网络在绝大多数时候表现良好,设计软件时仍必须假定网络偶尔会发生故障,并妥善处理这些故障。时钟也是如此:它们大多数时候都走得好好的,但健壮的软件必须做好应对错误时钟的准备。
部分问题在于,错误的时钟很容易无人察觉。如果机器的 CPU 有缺陷,或网络配置有误,它多半会彻底无法工作,因此问题很快就会暴露并得到修复。反之,如果石英钟有缺陷,或 NTP 客户端配置有误,那么即使时钟逐渐偏离现实越来越远,大多数事情看上去仍然运行正常。如果某段软件依赖精确同步的时钟,最终结果更可能是隐蔽而细微的数据丢失,而不是一场惊天动地的崩溃 62 63。
因此,如果使用的软件要求时钟同步,就必须仔细监控所有机器之间的时钟偏差。凡是时钟与其他节点偏离太远的节点,都应当被宣告死亡并移出集群。这样的监控可以确保在故障时钟造成太大破坏以前及时发现它。
用于事件排序的时间戳
下面来看一种很容易让人想依赖时钟、却十分危险的情形:为多个节点上的事件排序 64。例如,两个客户端都向分布式数据库写入时,谁先到达?哪次写入更新?
图 9-3 展示了在采用多主复制的数据库中,日历时钟的一种危险用法(这个例子与 图 6-8 类似)。客户端 A 在节点 1 上写入 x = 1;这次写入复制到节点 3;客户端 B 在节点 3 上将 x 递增(此时 x = 2);最后,两次写入都复制到节点 2。

在 图 9-3 中,写入复制到其他节点时,会按照写入起源节点的日历时钟附上时间戳。这个例子中的时钟同步已经非常好:节点 1 与节点 3 的偏差不到 3 毫秒,实践中恐怕很难达到这么好的水平。
递增操作建立在先前写入的 x = 1 之上,因此我们自然会认为 x = 2 这次写入应该具有更大的时间戳。遗憾的是,图 9-3 中并非如此:写入 x = 1 的时间戳是 42.004 秒,写入 x = 2 的时间戳却是 42.003 秒。
正如 “最后写入者胜(丢弃并发写入)” 所讨论的,解决不同节点并发写入值之间冲突的一种办法是 最后写入者胜(LWW):对于同一个键,只保留时间戳最大的写入,丢弃所有时间戳更早的写入。在 图 9-3 的例子中,节点 2 收到这两个事件后,会错误地断定 x = 1 才是更新的值,并丢弃对 x = 2 的写入,于是递增操作便丢失了。
要避免这个问题,可以确保每当覆盖一个值时,新值的时间戳一定高于被覆盖值,即使这个时间戳已经超前于写入者的本地时钟。不过,这样就要付出额外读取的代价,先找出当前最大的时间戳。Cassandra 和 ScyllaDB 等系统希望在一次往返中写入所有副本,因此它们直接使用客户端时钟生成的时间戳,并采用最后写入者胜策略 62。这种做法存在一些严重问题:
- 数据库写入可能莫名其妙地消失:在走得较慢的节点追上走得较快的节点之前,它无法覆盖后者先前写入的值 63 65。这样可能在不向应用程序报告任何错误的情况下,悄无声息地丢弃任意数量的数据。
- LWW 无法区分短时间内接连发生的顺序写入(在 图 9-3 中,客户端 B 的递增操作显然发生在客户端 A 的写入 之后)与真正的并发写入(两个写入者都不知道对方的写入)。为了避免违反因果关系,还需要版本向量等额外的因果关系追踪机制(参见 “检测并发写入”)。
- 两个节点可能各自独立生成具有相同时间戳的写入,特别是时钟分辨率只有毫秒时。解决这类冲突还需要一个额外的决胜值(简单地取一个很大的随机数即可),但这种做法同样可能违反因果关系 62。
因此,尽管保留最“新”的值并丢弃其他值,看上去是很诱人的冲突解决办法,但必须意识到,“新”的定义取决于本地日历时钟,而它很可能并不准确。即使时钟经过严密的 NTP 同步,也可能出现这样的情况:数据包在时间戳 100 毫秒时发出(按发送方的时钟),却在时间戳 99 毫秒时到达(按接收方的时钟)——看起来仿佛数据包还没发出就已经到达,这当然不可能。
能否把 NTP 同步做得足够精确,从而彻底避免这类错误排序?恐怕不能。除了石英钟漂移等其他误差源以外,NTP 的同步精度本身就受网络往返时间限制。若要保证顺序正确,时钟误差必须显著小于网络延迟,而这是不可能做到的。
所谓的 逻辑时钟(logical clock)66 以递增计数器为基础,而不是以振荡的石英晶体为基础,因此是为事件排序时更安全的选择(参见 “检测并发写入”)。逻辑时钟既不度量一天中的时刻,也不度量经过了多少秒,只记录事件之间的相对顺序(一个事件发生在另一个事件之前还是之后)。与之相对,日历时钟和单调时钟度量真实流逝的时间,因此也称为 物理时钟(physical clock)。我们将在 “ID 生成器和逻辑时钟” 中更详细地讨论逻辑时钟。
带置信区间的时钟读数
机器的日历时钟也许能以微秒甚至纳秒为分辨率提供读数,但测量得如此精细,并不表示读数真的精确到了这个程度。事实上,它大概率没有这么准。前面说过,即使每分钟都与局域网中的 NTP 服务器同步一次,不精确的石英钟也很容易漂移数毫秒。若使用公共互联网上的 NTP 服务器,最理想的精度大概也只有数十毫秒;遇到网络拥塞时,误差很容易飙升到 100 毫秒以上。
因此,不应把一次时钟读数理解为一个精确的时间点,它更像是落在某个置信区间内的一段时间范围。例如,系统也许有 95% 的把握认为当前时刻位于这一分钟的第 10.3 秒至第 10.5 秒之间,但无法知道得更精确 67。如果时间只能确定到 ±100 毫秒,那么时间戳中精确到微秒的那些数字基本毫无意义。
不确定性的边界可以根据时间源来计算。如果计算机直接连接着 GPS 接收器或原子钟,预期误差范围取决于设备本身;对 GPS 而言,还取决于卫星信号的质量。如果从服务器获取时间,不确定性大致等于:自上次与服务器同步以来石英钟的预期漂移,加上 NTP 服务器自身的不确定性,再加上与服务器之间的网络往返时间——这是第一步近似,并且假定服务器值得信任。
遗憾的是,大多数系统都不会暴露这种不确定性。例如,调用 clock_gettime() 时,返回值并不会告诉你时间戳的预期误差,因此你无从得知它的置信区间究竟是五毫秒,还是五年。
不过也有例外:Google Spanner 的 TrueTime API 45 和 Amazon 的 ClockBound 都会明确报告本地时钟的置信区间。查询当前时间时,得到的是两个值:[earliest, latest],分别表示 最早可能 与 最晚可能 的时间戳。根据对不确定性的计算,时钟能够断定真实的当前时间位于这个区间内。区间的宽度取决于多种因素,其中包括本地石英钟距离上次与更精确的时间源同步已经过去多久。
用于全局快照的同步时钟
在 “快照隔离与可重复读” 中,我们讨论了 多版本并发控制(MVCC)。对于既要支持短小快速的读写事务,又要支持大型、长时间运行的只读事务(例如备份或分析)的数据库来说,MVCC 是一项非常有用的功能。它让只读事务能够看到数据库在某个特定时间点的一致状态,也就是一个 快照,同时又不必锁住或干扰读写事务。
一般来说,MVCC 需要单调递增的事务 ID。如果某次写入发生在快照之后(也就是说,写入的事务 ID 大于快照的事务 ID),那么快照事务就看不到这次写入。在单节点数据库上,用一个简单的计数器就足以生成事务 ID。
可是,当数据库分布在许多机器上,甚至横跨多个数据中心时,生成全局单调递增的事务 ID(覆盖所有分片)就很困难,因为这需要协调。事务 ID 还必须反映因果关系:如果事务 B 读取或覆盖了事务 A 先前写入的值,B 的事务 ID 就必须大于 A,否则快照将不一致。面对大量短小快速的事务,在分布式系统中生成事务 ID 会成为难以承受的瓶颈。(我们将在 “ID 生成器和逻辑时钟” 中讨论这类 ID 生成器。)
能不能直接把同步日历时钟产生的时间戳当作事务 ID?如果时钟同步能做得足够好,时间戳的确具备所需属性:越晚的事务,时间戳越大。当然,问题依旧在于时钟精度的不确定性。
Spanner 正是以这种方式实现跨数据中心的快照隔离 68 69。它利用 TrueTime API 报告的时钟置信区间,依据的是下面这个观察:假设有两个置信区间,每个区间都由最早和最晚的可能时间戳组成(A = [A最早, A最晚],B = [B最早, B最晚]),如果两个区间不重叠(即 A最早 < A最晚 < B最早 < B最晚),那么 B 必定发生在 A 之后,不存在任何疑问。只有当两个区间发生重叠时,我们才无法确定 A 与 B 的先后顺序。
为了确保事务时间戳能够反映因果关系,Spanner 会在提交读写事务之前,刻意等待相当于置信区间长度的一段时间。这样一来,任何可能读到该数据的事务都会发生在足够晚的时刻,使它们的置信区间不再重叠。为了尽量缩短等待,Spanner 需要把时钟不确定性控制得尽可能小;为此,Google 在每个数据中心都部署了 GPS 接收器或原子钟,从而把时钟同步误差控制在大约 7 毫秒以内 45。
严格来说,Spanner 并非一定要使用原子钟和 GPS 接收器:真正重要的是获得置信区间,精确的时间源只是帮助缩小这个区间。其他系统也开始采用类似方法。例如,YugabyteDB 在 AWS 上运行时可以利用 ClockBound 70,还有若干系统也开始在不同程度上依赖时钟同步 71 72。
进程暂停
再来看一个在分布式系统中危险使用时钟的例子。假设某个数据库的每个分片都只有一个领导者,而且只有领导者可以接受写入。一个节点怎么知道自己仍是领导者(没有被其他节点宣告死亡),因而可以安全地接受写入呢?
一种办法是由领导者向其他节点取得一份 租约(lease),它类似于带有超时的锁 73。任何时刻只能有一个节点持有租约。因此,节点拿到租约后,就知道自己在租约到期以前的一段时间内仍是领导者。为保持领导者身份,节点必须在租约到期前定期续租。如果节点失效,就会停止续租;租约到期后,另一个节点便可接管。
可以想象,请求处理循环大致如下:
这段代码有什么问题?首先,它依赖同步时钟:租约到期时间由另一台机器设置(例如以当前时间加 30 秒来计算),却要与本地系统时钟比较。如果两台时钟的偏差超过几秒,这段代码的行为就会变得古怪。
其次,即使把协议改成只使用本地单调时钟,仍然存在另一个问题:代码假定读取时间(System.currentTimeMillis())与处理请求(process(request))之间只会经过极短时间。通常这段代码确实运行得很快,预留 10 秒足以确保租约不会在请求处理到一半时过期。
可是,如果程序执行时意外暂停了呢?例如,假设线程执行到 lease.isValid() 附近时停了 15 秒,随后才恢复。在处理请求时,租约很可能早已过期,另一个节点也已接任领导者。然而,没有任何东西会告诉这个线程它刚才停了那么久;直到循环进入下一轮,它才会发现租约已经过期——而在此之前,它可能已经处理了请求,做出了不安全的操作。
认为线程可能暂停这么久,是否合理?很遗憾,完全合理。造成长时间暂停的原因多种多样:
- 多个线程争用锁、队列等共享资源时,线程可能把大量时间花在等待上。换用 CPU 核心更多的机器甚至可能让这类问题进一步恶化,而且争用问题往往很难诊断 74。
- 许多编程语言运行时(例如 Java 虚拟机)都带有 垃圾回收器(GC),偶尔需要停止所有正在运行的线程。过去,这类 STW GC 暂停 有时会让程序停上几分钟 75!现代 GC 算法已经大大缓解了这个问题,但 GC 暂停仍可能相当明显(参见 “限制垃圾回收的影响”)。
- 在虚拟化环境中,虚拟机可以被 挂起(暂停所有进程并把内存内容保存到磁盘),随后再 恢复(还原内存内容并从原处继续执行)。这种暂停可能发生在进程执行的任何时刻,持续时间也没有上限。这个功能有时用于把虚拟机从一台宿主机 实时迁移 到另一台宿主机而无需重启;在这种情况下,暂停多久取决于进程写入内存的速率 76。
- 在笔记本电脑和手机等终端用户设备上,执行也可能随时挂起并恢复,例如用户合上笔记本电脑屏幕时。
- 当操作系统切换到另一个线程,或者虚拟机监控器切换到另一台虚拟机时,当前线程可能停在代码中的任意位置。对虚拟机而言,被其他虚拟机占用的 CPU 时间称为 窃取时间(steal time)。如果机器负载很高——也就是有很长的线程队列在等待运行——暂停的线程可能要过一阵子才能再次获得运行机会。
- 如果应用程序执行同步磁盘访问,线程可能暂停下来,等待缓慢的磁盘 I/O 操作完成 77。在许多语言中,即使代码没有明确读写文件,磁盘访问也可能出人意料地发生。例如,Java 类加载器会在类第一次使用时才加载类文件,而这可能出现在程序执行的任何时刻。I/O 暂停与 GC 暂停甚至可能相互叠加 78。如果所谓的磁盘其实是网络文件系统或网络块设备(例如 Amazon EBS),I/O 延迟还会受到网络延迟波动的影响 31。
- 如果操作系统允许 换页到磁盘(分页),一次简单的内存访问也可能触发缺页错误,必须从磁盘把某个页面载入内存。在这项缓慢的 I/O 操作完成以前,线程会一直暂停。如果内存压力很大,还可能需要先把另一个页面换出到磁盘。极端情况下,操作系统会把大部分时间耗在内存页面的换入换出上,几乎不做实际工作,这叫作 抖动(thrashing)。为了避免这种情况,服务器通常会禁用分页——与其冒着发生抖动的风险,不如终止一个进程来释放内存。
- 向 Unix 进程发送
SIGSTOP信号也会令其暂停,例如在 shell 中按 Ctrl-Z。这个信号会立即停止给进程分配 CPU 周期,直到SIGCONT令它恢复;随后,它会从先前停下的位置继续运行。即使你的环境通常不用SIGSTOP,运维人员也可能不小心发出这个信号。
上述任何事件都可能在任意位置 抢占 正在运行的线程,过一段时间再让它恢复,而线程对此毫无察觉。这个问题类似于保证单机多线程代码的线程安全:不能对时序作任何假定,因为上下文切换与并行执行随时都可能发生。
在单台机器上编写多线程代码时,我们有不少成熟工具来保证线程安全:互斥锁、信号量、原子计数器、无锁数据结构、阻塞队列等等。遗憾的是,这些工具不能直接套用到分布式系统,因为分布式系统没有共享内存,只有经由不可靠网络传递的消息。
分布式系统中的节点必须假定:自己的执行可能在任意时刻暂停很长时间,哪怕正处于函数执行途中。暂停期间,外部世界仍在继续运转,甚至可能因为这个节点迟迟没有响应而宣告它死亡。最终节点恢复运行时,甚至意识不到自己曾经“睡着”,直到稍后再次读取时钟。
响应时间保证
如上所述,在许多编程语言与操作系统中,线程和进程都可能暂停任意长的时间。不过,只要投入足够努力,这些暂停的原因 确实可以 消除。
有些软件运行在这样的环境中:如果不能在规定时间内响应,就可能造成严重损害。控制飞机、火箭、机器人、汽车及其他实体设备的计算机,必须快速而且可预测地响应传感器输入。在这些系统中,软件必须赶在明确规定的 截止时间(deadline)之前响应;错过截止时间,就可能导致整个系统失效。这样的系统称为 硬实时(hard real-time)系统。
在嵌入式系统中,实时 是指系统经过精心设计与测试,能够在任何情况下满足规定的时序保证。这与 Web 领域对 实时 一词较为宽泛的用法形成对比:后者通常只是指服务器向客户端推送数据或进行流处理,并没有严格的响应时间约束(参见 第 12 章)。
例如,汽车的车载传感器检测到碰撞正在发生时,你绝不会希望安全气囊因为控制系统恰好遭遇 GC 暂停而延迟弹出。
要在系统中提供实时保证,需要软件栈的每一层共同支持:需要 实时操作系统(RTOS),保证在规定的时间间隔内为进程分配 CPU 时间;库函数必须说明最坏情况下的执行时间;动态内存分配可能要受到限制,甚至完全禁止(虽然存在实时垃圾回收器,应用程序仍必须确保不会给 GC 安排太多工作);此外还要进行海量测试与测量,验证系统确实满足保证。
这些要求不仅带来大量额外工作,也严重限制了可用的编程语言、库和工具,因为大多数语言和工具都不提供实时保证。正因如此,实时系统开发极其昂贵,最常用于安全攸关的嵌入式设备。而且,“实时”并不等于“高性能”——事实上,实时系统的吞吐量可能更低,因为及时响应必须优先于一切(另见 “延迟与资源利用率”)。
对于大多数服务器端数据处理系统,实时保证既不经济,也不合适。因此,这些系统只能承受非实时环境带来的进程暂停与时钟不稳定。
限制垃圾回收的影响
垃圾回收曾是造成进程暂停的最大原因之一 79。好在 GC 算法已经有了长足进步:如今,经过适当调优的回收器通常只会暂停几毫秒。Java 运行时提供并发标记清除(CMS)、垃圾优先(G1)、Z 垃圾回收器(ZGC)、Epsilon 和 Shenandoah 等回收器,分别针对高频创建对象、大型堆等不同内存使用特征进行优化。相比之下,Go 提供的是一种更简单、尝试自我优化的并发标记清除垃圾回收器。
如果必须彻底避免 GC 暂停,可以选择完全没有垃圾回收器的语言。例如,Swift 使用自动引用计数来判断何时能够释放内存;Rust 和 Mojo 则通过类型系统追踪对象的生命周期,让编译器判断内存需要保留多久。
也可以继续使用带垃圾回收的语言,同时减轻暂停造成的影响。一种做法是把 GC 暂停看作节点短暂的计划内停机:某个节点进行垃圾回收时,由其他节点处理客户端请求。如果运行时能够提前告知应用程序节点即将进行 GC 暂停,应用程序就可以停止向该节点发送新请求,等待它处理完尚未完成的请求,然后趁没有请求进行时执行 GC。这种技巧能向客户端隐藏 GC 暂停,并降低响应时间的高百分位数 80 81。
这种思路还有一个变体:只让垃圾回收器处理容易快速回收的短命对象,并定期重启进程,赶在长期存活对象积累到需要执行一次完整 GC 之前 79 82。每次可以只重启一个节点,并在计划重启前先把流量从该节点迁走,就像滚动升级一样(参见 第 5 章)。
这些措施无法彻底杜绝垃圾回收暂停,却能切实减轻其对应用程序的影响。
知识、真相和谎言
到目前为止,本章已经考察了分布式系统与单机程序的不同之处:系统没有共享内存,只能通过延迟不定的不可靠网络来传递消息,还可能遭遇部分失效、不可靠的时钟和进程暂停。
如果还不熟悉分布式系统,这些问题带来的后果会让人极度迷失方向。网络中的一个节点不可能 确切知道 其他节点的任何事情,只能根据收到(或没有收到)的消息作出猜测。一个节点只有与另一个节点交换消息,才能得知对方处于什么状态,例如存储了哪些数据、是否正常运行等等。如果远程节点没有响应,就无从得知它的状态,因为网络问题与节点自身的问题无法可靠地区分。
关于这类系统的讨论已经近乎哲学:在系统中,我们究竟知道什么为真、什么为假?如果感知和测量的机制都不可靠,我们又能在多大程度上确信自己的认知 83?软件系统是否应当服从我们认为物理世界必然遵循的法则,例如因果律?
好在我们不必一路追问到生命的意义。在分布式系统中,可以明确写出对系统行为所作的假设(即 系统模型,system model),再把实际系统设计成符合这些假设。我们还可以证明某个算法在特定系统模型中能够正确运行。这意味着,即使底层系统模型只提供极少保证,也仍然可以实现可靠的行为。
不过,要让软件在不可靠的系统模型中表现良好,绝非轻而易举。本章余下部分将进一步探讨分布式系统中的知识与真相,帮助我们思考可以作出哪些假设,以及希望提供哪些保证。在 第 10 章 中,我们将继续考察一些分布式算法:它们在特定假设下提供特定保证。
多数派原则
设想一个存在非对称故障的网络:某个节点可以收到发给它的所有消息,但它发出的消息都会丢失或延迟 22。这个节点明明工作得完全正常,也在接收其他节点的请求,可其他节点就是听不见它的回应。等到超时后,其他节点由于一直收不到回复,便宣告它已经死亡。接下来的场面如同噩梦:这个半失联的节点被拖向墓地,一路挣扎高喊“我还没死!”——可惜谁也听不见它的呼喊,送葬队伍仍以坚忍不拔的决心继续前进。
在一个没那么噩梦般的场景中,半失联的节点也许会发现自己发出的消息得不到其他节点的确认,于是意识到网络必定出了故障。然而,其他节点仍会错误地宣告它死亡,而它对此无能为力。
第三种场景是,假设某个节点暂停执行一分钟。在此期间,它既不处理请求,也不发送响应。其他节点一边等待、一边重试,终于失去耐心,宣告它死亡并把它抬上灵车。最终暂停结束,节点的线程若无其事地继续执行。其他节点惊讶地看到,那个本应死去的节点突然精神抖擞地从棺材里探出头来,兴高采烈地与旁人聊天。刚恢复时,这个节点甚至不知道整整一分钟已经过去,也不知道自己已经被宣告死亡——在它看来,距离上次和其他节点交谈仿佛只过了一瞬间。
这些故事告诉我们,节点未必能相信自己对处境的判断。分布式系统不能只依赖某一个节点,因为节点随时可能失效,使系统陷入僵局而无法恢复。因此,许多分布式算法依赖 法定人数(quorum),也就是让多个节点投票(参见 “读写仲裁”):一项决策必须获得若干节点的最低票数,借此减少对任何单个节点的依赖。
宣告节点死亡的决定也是如此。如果达到法定人数的节点宣告另一个节点已经死亡,那么即使它觉得自己还活得好好的,也必须被视为死亡。单个节点必须服从法定人数作出的决定并下台。
最常见的法定人数,是超过节点总数一半的绝对多数,当然也可以采用其他形式。多数法定人数让系统能在少数节点发生故障时继续工作:三个节点可以容忍一个故障节点,五个节点则可以容忍两个。与此同时,它仍然是安全的,因为系统中只能形成一个多数派,不可能同时出现两个作出冲突决定的多数派。我们将在 第 10 章 讨论 共识算法 时,更详细地介绍法定人数的用法。
分布式锁和租约
分布式应用程序中的锁与租约很容易被误用,也是程序缺陷的常见来源 84。下面来看一种具体的出错方式。
在 “进程暂停” 中,我们看到租约是一种会超时的锁:如果原持有者停止响应(可能是因为它崩溃了、暂停太久,或与网络断开),租约就可以交给新的持有者。当系统要求某种东西只能有一个时,就可以使用租约。例如:
- 只允许一个节点担任数据库分片的领导者,以避免脑裂(参见 “处理节点故障”)。
- 只允许一个事务或客户端更新特定资源或对象,以免并发写入将其损坏。
- 一项大型处理作业中的每个输入文件只应由一个节点处理,避免多个节点重复执行同一份工作而白白浪费计算资源。
值得仔细想一想:如果多个节点同时相信自己持有租约——也许是进程暂停造成的——会发生什么?对第三个例子而言,后果不过是浪费一些计算资源,并不严重;但在前两个例子中,数据可能丢失或损坏,严重得多。
例如,图 9-4 展示了锁实现错误导致的数据损坏。(这并非纯粹的理论问题:HBase 曾经就有这个缺陷 85 86。)假设你想确保某个存储服务里的文件一次只能由一个客户端访问,因为多个客户端同时写入会损坏文件。于是,你要求客户端在访问文件之前,先向锁服务取得租约。这类锁服务通常用共识算法来实现,我们将在 第 10 章 中进一步讨论。

问题正是 “进程暂停” 所讨论的情形:持有租约的客户端如果暂停太久,租约就会过期。另一个客户端此时可以取得同一文件的租约,并开始写入。等暂停的客户端恢复后,它误以为自己的租约依然有效,也继续写入文件。于是便出现了脑裂:两个客户端的写入彼此冲突,最终损坏文件。
图 9-5 展示了另一个后果类似的问题。这个例子里没有进程暂停,只有客户端 1 崩溃。就在崩溃前,客户端 1 向存储服务发出了一项写请求,但请求在网络中延迟了很久。(回想 “实践中的网络故障”,数据包有时会延迟一分钟以上。)等写请求抵达存储服务时,租约早已超时,客户端 2 已经取得租约并发出了自己的写入。结果便是类似 图 9-4 的数据损坏。

用栅栏机制隔离僵尸与延迟请求
僵尸(zombie)一词有时用来形容这样的原租约持有者:它还不知道自己已经失去租约,仍把自己当作当前持有者行事。既然无法彻底杜绝僵尸,就必须确保它们不能以脑裂的形式造成任何破坏。这称为用 栅栏机制(fencing)隔离僵尸。
有些系统试图通过关停僵尸来隔离它,例如断开它的网络连接 9、通过云服务商的管理界面关闭虚拟机,甚至直接切断机器电源 87。这种做法称为 STONITH,即“击毙另一个节点”。遗憾的是,它有几个问题:无法防范 图 9-5 所示的超长网络延迟;所有节点可能彼此关停 19;而等到僵尸被发现并关闭时,也许早已为时过晚,数据已经损坏。
图 9-6 展示了一种更健壮的栅栏机制,既能防范僵尸,也能防范延迟请求。

假设锁服务每次授予锁或租约时,还会返回一个 栅栏令牌(fencing token)。这是一个每次授予锁都会增大的数字,例如由锁服务负责递增。随后可以要求客户端每次向存储服务发送写请求时,都必须带上自己当前的栅栏令牌。
在 图 9-6 中,客户端 1 取得租约及令牌 33,随后却长时间暂停,导致租约过期。客户端 2 接着取得租约及令牌 34(令牌值始终递增),并向存储服务发送带令牌 34 的写请求。稍后,客户端 1 恢复执行,也向存储服务发送写请求,其中带着自己的令牌 33。然而,存储服务记得自己已经处理过令牌值更高(34)的写入,因此会拒绝令牌 33 的请求。刚取得租约的客户端必须立刻向存储服务执行一次写入;一旦这次写入完成,所有僵尸都会被隔离在外。
如果使用 ZooKeeper 作为锁服务,可以把事务 ID zxid 或节点版本 cversion 用作栅栏令牌 85。在 etcd 中,修订号与租约 ID 共同发挥类似作用 89。Hazelcast 的 FencedLock API 则会显式生成栅栏令牌 90。
这种机制要求存储服务能够检查写入所携带的令牌是否已经过时。另一种办法是让服务支持类似原子比较并设置(CAS)的写入:只有从当前客户端上次读取对象以后,没有其他客户端写过该对象,写入才会成功。对象存储服务就支持这类检查:Amazon S3 称之为 条件写入(conditional write),Azure Blob Storage 称之为 条件标头(conditional header),Google Cloud Storage 则称之为 请求前置条件(request precondition)。
多副本隔离
如果客户端只需写入一个支持这类条件写入的存储服务,那么锁服务多少有些多余 91 92,因为完全可以直接依托该存储服务来分配租约 93。不过,有了栅栏令牌以后,也可以把它用于多个服务或副本,确保原租约持有者在所有这些服务上都被隔离。
例如,假设存储服务是一个采用最后写入者胜来解决冲突的无主复制键值存储(参见 “无主复制”)。在这样的系统中,客户端直接向每个副本发送写入,各副本根据客户端分配的时间戳,自行决定是否接受写入。
如 图 9-7 所示,可以把写入者的栅栏令牌放在时间戳最高有效的若干位或数字中。这样就能确保,新租约持有者生成的任何时间戳,都大于原租约持有者生成的所有时间戳,即使原持有者的写入实际发生得更晚。

在 图 9-7 中,客户端 2 的栅栏令牌为 34,因此它生成的所有以 34… 开头的时间戳,都大于客户端 1 生成的任何以 33… 开头的时间戳。客户端 2 成功写入达到法定人数的一组副本,但无法连接副本 3。因此,僵尸客户端 1 稍后尝试写入时,这次写入虽然会被副本 1 和 2 忽略,却可能在副本 3 上成功。这不成问题,因为后续的仲裁读会优先选择客户端 2 那个时间戳更大的写入,而读修复或反熵过程最终会覆盖客户端 1 写入的值。
从这些例子可以看出,假定任何时刻只有一个节点持有租约并不安全。所幸,只要稍加留意,就可以借助栅栏令牌防止僵尸与延迟请求造成任何破坏。
拜占庭故障
栅栏令牌能够发现并阻止 无意间 出错的节点,例如尚未发现租约已经过期的节点。但如果节点蓄意破坏系统保证,只需发送带有伪造栅栏令牌的消息,就能轻易绕过这种防护。
本书假定节点虽然不可靠,却是诚实的:它们可能因为故障而响应缓慢或永不响应,也可能因为 GC 暂停或网络延迟而持有过时状态;但我们假定,只要节点 确实 作出响应,它说的就是“真话”——至少据它自己所知,它是在遵守协议规则。
如果节点有可能“撒谎”,也就是发出任意错误或损坏的响应,分布式系统问题就会困难得多。例如,一个节点可能在同一轮选举中投出多张互相矛盾的票。这种行为叫作 拜占庭故障(Byzantine fault),而在这种互不信任的环境中达成共识的问题,则称为 拜占庭将军问题(Byzantine Generals Problem)94。
拜占庭将军问题是所谓 两将军问题(Two Generals Problem)95 的推广。两将军问题设想,两名军队将领必须就作战计划达成一致,但他们驻扎在不同地点,只能派信使互通消息,而信使有时会迟到或失踪(就像网络数据包一样)。我们将在 第 10 章 中讨论这个 共识(consensus)问题。
在拜占庭版本的问题中,需要达成一致的将军有 n 名,但其中混入了若干叛徒。大多数将军都忠诚可靠,会发送真实消息;叛徒却可能发送虚假消息,企图欺骗并迷惑其他人。而谁是叛徒,事先无从得知。
拜占庭原是古希腊的一座城市,后来成为君士坦丁堡,所在地就是今天土耳其的伊斯坦布尔。没有任何历史证据表明,拜占庭的将军比其他地方的将军更爱阴谋诡计。这个名称其实来自 拜占庭式 一词在政治语境中的含义:过度复杂、官僚而诡诈;这种用法远在计算机出现以前便已存在 96。Lamport 想选择一个不会冒犯读者的国籍,而别人劝他最好不要把问题叫作 阿尔巴尼亚将军问题 97。
如果某些节点发生故障、不遵守协议,或有恶意攻击者干扰网络时,系统仍能继续正确运行,就称这个系统具有 拜占庭容错(Byzantine fault tolerance,BFT)能力。这类问题在某些特定情形下确实很重要。例如:
- 在航空航天环境中,辐射可能破坏计算机内存或 CPU 寄存器里的数据,使节点以任意且不可预测的方式回应其他节点。系统失效的代价极其高昂——例如飞机坠毁导致机上人员全部遇难,或火箭撞上国际空间站——因此飞行控制系统必须能够容忍拜占庭故障 98 99。
- 在有多方参与的系统中,部分参与者可能企图欺骗或诈骗其他人。此时,节点不能轻信另一个节点发来的消息,因为消息可能带有恶意。比特币等加密货币及其他区块链,就可以看作一种无需依赖中央权威,让彼此不信任的各方就某笔交易是否发生达成一致的方式 100。
不过,对本书讨论的系统而言,通常可以放心假定不存在拜占庭故障。数据中心里的所有节点都受你的组织控制,因而有望值得信任;辐射水平也足够低,内存损坏并非主要问题——尽管人们已经在考虑把数据中心送入轨道 101。多租户系统中的租户彼此并不信任,但系统依靠防火墙、虚拟化和访问控制策略把租户相互隔离,而不是使用拜占庭容错。让系统实现拜占庭容错的协议代价很高 102,而具有容错能力的嵌入式系统又依赖硬件层面的支持 98。对大多数服务器端数据系统来说,部署拜占庭容错方案的成本高得并不现实。
Web 应用程序的确必须预料到,Web 浏览器等由最终用户控制的客户端可能作出任意乃至恶意的行为。这也正是输入验证、数据清理和输出转义如此重要的原因,例如它们可以防止 SQL 注入与跨站脚本攻击。不过,我们通常不会为此使用拜占庭容错协议,只需让服务器充当权威,决定哪些客户端行为可以接受、哪些不可以。在没有这种中央权威的点对点网络中,拜占庭容错才更为重要 103 104。
软件缺陷也可以看作拜占庭故障,但如果所有节点部署的都是同一套软件,拜占庭容错算法也救不了你。大多数拜占庭容错算法要求超过三分之二的节点正常运行(例如四个节点中最多只能有一个发生故障)。若想用这种办法防范软件缺陷,就得准备同一软件的四种独立实现,并寄希望于缺陷只出现在其中一种实现里。
同理,如果某种协议能保护我们免受漏洞、安全失陷和恶意攻击,当然很有吸引力。遗憾的是,这同样不现实:在大多数系统中,攻击者既然能够攻陷一个节点,多半也能攻陷所有节点,因为它们很可能运行相同的软件。因此,身份认证、访问控制、加密和防火墙等传统机制,仍然是抵御攻击者的主要手段。
弱形式的谎言
虽然我们通常假定节点是诚实的,但仍值得在软件中加入一些机制,防范较弱形式的“撒谎”,例如硬件问题、软件缺陷或配置错误产生的无效消息。这些机制算不上完整的拜占庭容错,因为它们挡不住意志坚定的攻击者;但作为提高可靠性的手段,它们既简单又务实。例如:
- 硬件问题,或操作系统、驱动程序、路由器等组件中的缺陷,确实可能损坏网络数据包。TCP 与 UDP 内置的校验和通常能够发现损坏的数据包,但偶尔也会漏检 105 106 107。一般只需一些简单措施就足以防范这类损坏,例如在应用层协议中加入校验和。TLS 加密连接也能抵御数据损坏。
- 可公开访问的应用程序必须仔细清理所有用户输入,例如检查数值是否落在合理范围内,并限制字符串长度,防止攻击者通过分配巨量内存来造成拒绝服务。防火墙后的内部服务也许可以放宽输入检查,但协议解析器仍然最好保留基本校验 105。
- NTP 客户端可以配置多个服务器地址。同步时,客户端会联系所有服务器,估算各自的误差,并检查大多数服务器是否就某个时间范围达成一致。只要大部分服务器正常,配置有误、报告错误时间的 NTP 服务器就会被识别为离群值,并排除在同步过程之外 39。与只依赖一台服务器相比,使用多台服务器能让 NTP 更加健壮。
系统模型与现实
人们已经设计了许多算法来解决分布式系统问题。例如,我们将在 第 10 章 中考察共识问题的解决方案。要真正有用,这些算法必须能够容忍本章讨论的分布式系统中的各种故障。
算法的编写方式不应过度依赖其运行环境中具体的硬件与软件配置。这又要求我们以某种方式,把系统中预期会发生的故障类型形式化。为此,我们定义 系统模型(system model):它是一种抽象,用来说明算法可以作出哪些假设。
关于时序假设,通常使用以下三种系统模型:
- 同步模型(synchronous model)
- 同步模型假定网络延迟、进程暂停和时钟误差都有上界。这并不表示时钟完全同步,也不表示网络延迟为零;它只是意味着,我们知道网络延迟、暂停和时钟漂移永远不会超过某个固定上限 108。对大多数实际系统而言,同步模型并不现实,因为正如本章所讨论的,无界延迟和暂停确实可能发生。
- 部分同步模型(partially synchronous model)
- 部分同步是指系统在 大多数时候 都像同步系统一样运行,但偶尔会突破网络延迟、进程暂停和时钟漂移的界限 108。对许多系统来说,这是一个现实的模型:绝大多数时候,网络与进程表现得相当规矩,否则我们什么事情也做不成;但我们也必须正视这样一个事实——任何时序假设都有可能偶尔失效。发生这种情况时,网络延迟、进程暂停和时钟误差都可能变得任意之大。
- 异步模型(asynchronous model)
- 在这种模型中,算法不能作出任何时序假设——事实上,它甚至没有时钟可用(因而也不能使用超时)。有些算法可以针对异步模型来设计,但这种模型限制极大。
除了时序问题,还必须考虑节点失效。常见的节点系统模型包括:
- 崩溃停止故障(crash-stop fault)
- 在 崩溃停止(crash-stop,或 故障停止,fail-stop)模型中,算法可以假定节点只有一种失效方式:崩溃 109。节点可能在任意时刻突然停止响应,此后便永远消失,再也不会回来。
- 崩溃恢复故障(crash-recovery fault)
- 我们假定节点可能在任意时刻崩溃,也可能在一段未知时间后重新开始响应。在崩溃恢复模型中,节点拥有能在崩溃后保留数据的稳定存储(即非易失性磁盘存储),但内存中的状态会丢失。
- 性能下降和功能不全
- 除了崩溃与重启,节点还可能变慢:它们也许仍能响应健康检查,却慢得无法完成任何实际工作。例如,千兆网络接口可能因为驱动程序缺陷,吞吐量突然跌到 1 Kb/s 110;面临内存压力的进程可能把大部分时间花在垃圾回收上 111;磨损的 SSD 可能表现得极不稳定;高温、连接器松动、机械振动、电源问题、固件缺陷等也会影响硬件 112。这类情形称为 跛行节点(limping node)、灰色失效(gray failure)或 慢失效(fail-slow)113,甚至可能比彻底失效的节点更难处理。还有一种相关问题:进程不再执行原本应做的某些工作,其他功能却仍在继续,例如后台线程崩溃或死锁时 114。
- 拜占庭(任意)故障
- 节点可能做出任何行为,包括像上一节所述那样欺骗其他节点。
对现实系统建模时,带崩溃恢复故障的部分同步模型通常最有用。它允许出现无界网络延迟、进程暂停和慢节点。那么,分布式算法该如何应对这种模型呢?
定义算法的正确性
要定义一个算法怎样才算 正确,可以描述它必须具备哪些 属性。例如,排序算法的输出具有这样一项属性:对输出列表中的任意两个不同元素,位于左侧的元素都小于位于右侧的元素。这不过是用形式化语言定义一份列表何谓“有序”。
同理,我们也可以列出分布式算法应具备的属性,以此定义它怎样才算正确。例如,如果算法要为锁生成栅栏令牌(参见 “用栅栏机制隔离僵尸与延迟请求”),我们可能要求它满足以下属性:
- 唯一性
- 任意两个栅栏令牌请求都不能返回相同的值。
- 单调序列
- 如果请求 x 返回令牌 t**x,请求 y 返回令牌 t**y,而且 x 在 y 开始以前已经完成,那么 t**x < t**y。
- 可用性
- 请求栅栏令牌且没有崩溃的节点,最终会收到响应。
如果算法在某个系统模型允许发生的所有情形下,始终满足这些属性,就可以说它在该系统模型中是正确的。然而,如果所有节点都崩溃,或所有网络延迟都突然变成无限长,那么任何算法都不可能完成工作。面对一个允许系统彻底失效的模型,我们怎样才能仍然给出有用的保证呢?
安全性与活性
为了说清这个问题,有必要区分两类不同的属性:安全性(safety)与 活性(liveness)。在刚才的例子中,唯一性 和 单调序列 属于安全属性,可用性 则属于活性属性。
两者究竟有何区别?一个明显线索是,活性属性的定义中往往含有“最终”二字。(没错,你已经猜到了:最终一致性 就是一项活性属性 115。)
安全性常被非正式地定义为 坏事不会发生,活性则是 好事最终会发生。不过,最好不要过度解读这种通俗说法,因为“好”与“坏”属于价值判断,并不太适合用来描述算法。安全性与活性的严格定义要精确得多 116:
- 如果安全属性遭到违反,我们可以指出它在哪一个具体时刻被破坏。例如,唯一性遭到破坏时,可以找到返回重复栅栏令牌的那次具体操作。安全属性一旦被违反,就无法撤销——损害已经造成。
- 活性属性正好相反:它在某个时刻可能尚未成立(例如节点已经发出请求,却还没收到响应),但未来始终还有满足它的希望(也就是最终收到响应)。
区分安全属性与活性属性的一个好处,是能帮助我们应对棘手的系统模型。对分布式算法,通常要求安全属性在系统模型允许的所有情形下都 始终 成立 108。也就是说,即使所有节点崩溃,或整个网络失效,算法仍必须保证绝不返回错误结果,因而安全属性依旧得到满足。
而对活性属性,我们可以附加条件。例如,可以规定只有在多数节点尚未崩溃、且网络最终能从中断中恢复时,请求才必须得到响应。部分同步模型的定义要求系统最终回到同步状态——也就是说,任何一次网络中断都只能持续有限时间,随后会得到修复。
将系统模型映射到现实世界
安全属性、活性属性和系统模型,对推理分布式算法的正确性非常有用。但当算法真正落地实现时,现实世界混乱的一面又会回来找麻烦;这时便清楚地看到,系统模型只是对现实的简化抽象。
例如,崩溃恢复模型中的算法通常假定稳定存储里的数据能挺过崩溃。但如果磁盘数据损坏了,或因硬件错误、配置错误而被清空,会发生什么 117?如果服务器存在固件缺陷,重启时明明硬盘连接无误,却无法识别它们,又会怎样 118?
法定人数算法(参见 “读写仲裁”)依赖节点记住自己声称已经存储的数据。如果节点会患上“失忆症”,忘掉先前保存的数据,就会破坏法定人数条件,进而破坏算法的正确性。或许还需要定义一种新系统模型:假定稳定存储在崩溃后通常得以保留,但偶尔也会丢失。可这样的模型又会变得更加难以推理。
算法的理论描述可以直接声明,假定某些事情绝不会发生——在非拜占庭系统中,我们确实必须对哪些故障可能发生、哪些不可能发生作出假设。然而,真实实现也许仍要包含代码,处理那些理论上“不可能”发生的事情;哪怕处理方式归结为 printf("Sucks to be you") 和 exit(666),也就是把残局留给人类运维人员收拾 119。(这正是计算机科学与软件工程之间的一项区别。)
这并不是说理论化、抽象化的系统模型毫无价值——恰恰相反。它们极其有助于把真实系统的复杂性提炼成一组可管理、可推理的故障,使我们能够理解问题,并尝试用系统化的方法加以解决。
形式化方法和随机测试
怎样知道一个算法确实满足所需属性?并发、部分失效与网络延迟会产生海量潜在状态。我们需要保证这些属性在每一种可能状态下都成立,还要确保没有遗漏任何边界情况。
一种办法是对算法进行形式化验证:用数学语言描述算法,再运用证明技术,证明它在系统模型允许的所有情形下都满足所需属性。证明算法正确,并不表示它在真实系统中的 实现 必然始终行为正确。不过,这是非常好的一步,因为理论分析能够发现算法中的隐患;这些问题在真实系统中可能潜伏很久,直到某些异常情况打破了你的假设(例如时序假设)才突然发作。
把理论分析与经验性测试结合起来,验证实现的行为是否符合预期,是一种稳妥做法。基于属性的测试(property-based testing)、模糊测试(fuzz testing)和确定性模拟测试(deterministic simulation testing,DST)等技术都利用随机化,在各种不同情形下测试系统。Amazon Web Services 等公司已经成功地把这些技术组合运用于许多产品 120 121。
模型检查与规范语言
模型检查器(model checker)是帮助验证算法或系统行为是否符合预期的工具。算法规范要用 TLA+、Gallina 或 FizzBee 等专门设计的语言来编写。借助这类语言,我们可以专注于算法行为,不必纠缠于代码实现细节。模型检查器随后会系统地尝试各种可能发生的情况,利用模型来验证不变量是否在算法的所有状态中都成立。
严格来说,模型检查无法证明算法的不变量在每一种可能状态下都成立,因为大多数现实算法的状态空间是无限的。要真正验证所有状态,需要给出形式化证明;这虽然可行,但通常比运行模型检查器困难得多。因此,使用模型检查器时,通常要把算法模型缩减成一个可以完全验证的近似版本,或是给执行设置某种上限(例如限制最多可以发送多少条消息)。这样一来,只会在更长执行过程中出现的缺陷就无法被发现。
即便如此,模型检查器仍然在易用性与发现隐蔽缺陷的能力之间取得了很好的平衡。CockroachDB、TiDB、Kafka 和许多其他分布式系统,都使用模型规范来发现并修复缺陷 122 123 124。例如,研究人员借助 TLA+ 证明,视图戳复制(VR)的文字描述存在歧义,可能导致数据丢失 125。
按照设计,模型检查器运行的并非实际代码,而是一个只描述协议核心思想的简化模型。这样更容易系统地探索状态空间,但也带来规范与实现逐渐偏离的风险 126。可以检查模型与真实实现的行为是否等价,不过这需要在真实实现中插桩 127。
故障注入
许多缺陷只有在机器或网络发生失效时才会触发。故障注入(fault injection)是一种有效(有时也相当吓人)的技术,用来验证系统实现遇到问题时是否仍会按预期工作。思路很简单:向正在运行的系统环境注入故障,再观察系统如何反应。注入的故障可以是网络失效、机器崩溃、磁盘损坏、进程暂停——凡是你能想到的计算机出错方式都可以尝试。
故障注入测试通常在与系统实际生产环境十分相似的环境中运行;有些团队甚至直接在生产环境中注入故障。Netflix 通过 Chaos Monkey 工具推广了这种做法 128。在生产环境中注入故障通常称为 混沌工程(chaos engineering),我们在 “可靠性与容错” 中已经讨论过。
运行故障注入测试时,首先要部署被测系统,以及故障注入协调者和脚本。协调者负责决定注入哪些故障、何时注入;本地或远程脚本则负责让单个节点或进程发生失效。注入脚本会利用许多不同工具来触发故障:Linux 进程可以用 kill 命令暂停或终止,磁盘可以用 umount 卸载,网络连接可以通过防火墙规则中断。检查故障注入期间以及之后的系统行为,就能确认系统是否一如预期。
触发不同失效所需的工具五花八门,因此故障注入测试写起来颇为繁琐。常见做法是采用 Jepsen 之类的故障注入框架来简化流程。这类框架集成了多种操作系统,并预置了大量故障注入器 129。Jepsen 在许多广泛使用的系统中都成功发现过关键缺陷,成效极为显著 130 131。
确定性模拟测试
确定性模拟测试(DST)也已成为模型检查与故障注入的一种流行补充。它探索状态空间的方式与模型检查器相似,但测试的是真实代码,而不是模型。
在 DST 中,模拟器会自动运行大量随机化的系统执行。模拟期间的网络通信、I/O 和时钟计时都由模拟组件取代,让模拟器能够精确控制各种时序与失效场景下,所有事情发生的先后顺序。这样一来,模拟器所能探索的情形远多于手写测试或故障注入。如果测试失败,还可以重新运行,因为模拟器知道触发失败的准确操作序列;故障注入则无法对系统实施如此细粒度的控制。
DST 要求模拟器能够控制网络延迟等一切非确定性来源。要让代码变得确定,通常采用以下三种策略之一:
- 应用程序层
- 有些系统从一开始就以便于确定性执行代码为目标进行构建。例如,FoundationDB 是 DST 领域的先驱之一,它基于名为 Flow 的异步通信库构建。Flow 为开发者提供了一个接入点,可以把确定性网络模拟注入系统 132。类似地,TigerBeetle 是一款原生支持 DST 的在线事务处理(OLTP)数据库。它把系统状态建模为状态机,所有状态变更都在单一事件循环内发生。再配合时钟等确定性模拟原语,这种架构便能以确定性方式运行 133。
- 运行时层
- 带异步运行时与常用库的编程语言,为引入确定性提供了接入点。可以用单线程运行时,强制所有异步代码顺序执行。例如,FrostDB 修改 Go 运行时,让 goroutine 依次执行 134。Rust 的 madsim 库也采用类似方式。Madsim 为 Tokio 异步运行时 API、AWS S3 库、Kafka 的 Rust 库以及许多其他组件提供确定性实现。应用程序只需换用确定性的库与运行时,无需修改自身代码,就能获得确定性的测试执行。
- 机器层
- 除了在运行时修改代码,也可以让整台机器变得确定。这个过程十分精细:机器必须对所有通常带有非确定性的调用给出确定性响应。Antithesis 等工具通过构建定制的虚拟机监控器来实现这一点,由它把通常非确定性的操作替换成确定性操作。从时钟到网络再到存储,一切都必须纳入考虑。完成后,开发者就能在虚拟机监控器内的一组容器中运行整个分布式系统,得到一个完全确定的分布式系统。
DST 的优势不止是能够重放。Antithesis 等工具发现较为罕见的行为时,会把一次测试执行分叉成多个子执行,借此尝试探索应用程序中的更多代码路径。确定性测试通常使用模拟时钟与网络调用,因此运行速度可以快于现实中的时间流逝。例如,TigerBeetle 的时间抽象可以模拟网络延迟和超时,而无需真的等待足够长时间让超时触发。这类技术让模拟器能够以更快速度探索更多代码路径。
非确定性正是本章所讨论的所有分布式系统难题的核心:并发、网络延迟、进程暂停、时钟跳变和崩溃都会以不可预测的方式发生,而且系统每次运行时都可能不同。反过来说,如果能让系统变得确定,许多事情都会大为简化。
事实上,让事物具有确定性是个简单而强大的思想,在分布式系统设计中反复出现。除了确定性模拟测试,前面几章还介绍过多种利用确定性的方式:
- 事件溯源的一项关键优势(参见 “事件溯源与 CQRS”),是可以确定性地重放事件日志,重新构建派生的物化视图。
- 工作流引擎(参见 “持久化执行与工作流”)依赖确定性的工作流定义,借此提供持久化执行语义。
- 我们将在 “使用共享日志” 中讨论的 状态机复制,会在每个副本上独立执行相同的确定性事务序列,从而复制数据。我们已经见过这种思想的两个变体:基于语句的复制(参见 “复制日志的实现”),以及通过存储过程串行执行事务(参见 “存储过程的利弊”)。
不过,要让代码具有彻底的确定性,仍须十分小心。即使已经消除了所有并发,并用确定性模拟替换 I/O、网络通信、时钟和随机数生成器,系统里仍可能残留非确定性。例如,在某些编程语言中,遍历哈希表元素的顺序可能不确定;是否会触及资源上限(内存分配失败、栈溢出)同样具有非确定性。
总结
本章讨论了分布式系统中可能发生的各种问题,其中包括:
- 每当试图通过网络发送数据包时,它都可能丢失或遭到任意延迟。响应同样可能丢失或延迟,因此只要没有收到响应,你就无从知道消息是否送达。
- 即使已经尽力配置 NTP,一个节点的时钟仍可能与其他节点严重不同步,也可能突然向前或向后跳变。依赖这样的时钟非常危险,因为你多半无法可靠估计其置信区间。
- 进程可能在执行到任意位置时暂停很长时间,被其他节点宣告死亡;随后它又恢复执行,却完全没有意识到自己曾经暂停。
可能出现这类 部分失效,正是分布式系统的决定性特征。只要软件试图完成任何涉及其他节点的事情,就可能偶尔失败、莫名其妙地变慢,或彻底不作响应(最终超时)。在分布式系统中,我们会把容忍部分失效的能力构建到软件里,使整个系统即使有部分组件损坏,也能继续工作。
要容忍故障,第一步是 检测 故障,可就连这一步也很困难。大多数系统都没有准确判断节点是否已经失效的机制,因此多数分布式算法只能依靠超时来判断远程节点是否仍然可用。然而,超时无法区分网络失效与节点失效,而网络延迟的波动有时又会让系统错误地怀疑某个节点已经崩溃。跛行节点更加难以处理:它们仍在响应,却慢得根本做不了任何有用的工作。
即便检测出了故障,要让系统容忍它也并不容易:机器之间既没有全局变量,没有共享内存,没有公共知识,也没有其他形式的共享状态 83。节点甚至无法就“现在几点”达成一致,更不用说更深刻的问题了。信息从一个节点流向另一个节点的唯一途径,就是经由不可靠网络发送。重大决策不能安全地交给单个节点,因此我们需要协议来召集其他节点参与,并争取让达到法定人数的节点达成一致。
如果你习惯于在单台计算机那种数学般完美的理想环境中编写软件——同一个操作总会确定性地返回相同结果——那么转向分布式系统混乱的物理现实,难免会感到震惊。反过来,分布式系统工程师常常认为,只要能在单台计算机上解决,问题就微不足道 4。事实上,如今单台计算机确实能做很多事情。如果能够避免打开潘多拉魔盒,把工作简单地留在一台机器上,例如使用嵌入式存储引擎(参见 “嵌入式存储引擎”),通常就值得这样做。
不过,正如 “分布式与单节点系统” 所讨论的,可伸缩性并不是采用分布式系统的唯一理由。容错与低延迟(把数据放在地理位置更靠近用户的地方)同样重要,而这些目标无法靠单个节点实现。分布式系统的力量在于,原则上它可以在服务层面永不停机,因为所有故障和维护都能在节点层面处理。(当然在实践中,如果把一项错误配置发布到所有节点,分布式系统照样会被彻底击垮。)
本章还稍稍岔开话题,探讨网络、时钟和进程的不可靠性是否属于无可避免的自然法则。答案是否定的:网络可以提供硬实时响应保证与有界延迟,只不过代价极其高昂,硬件资源的利用率也会随之降低。大多数非安全攸关系统都会在便宜而不可靠与昂贵而可靠之间选择前者。
本章始终在讨论问题,呈现出一幅颇为黯淡的图景。下一章将转向解决方案,讨论一些专为应对分布式系统问题而设计的算法。
参考文献
Mark Cavage. There’s Just No Getting Around It: You’re Building a Distributed System. ACM Queue, volume 11, issue 4, pages 80-89, April 2013. doi:10.1145/2466486.2482856 ↩︎
Jay Kreps. Getting Real About Distributed System Reliability. blog.empathybox.com, March 2012. Archived at perma.cc/9B5Q-AEBW ↩︎
Coda Hale. You Can’t Sacrifice Partition Tolerance. codahale.com, October 2010. https://perma.cc/6GJU-X4G5 ↩︎
Jeff Hodges. Notes on Distributed Systems for Young Bloods. somethingsimilar.com, January 2013. Archived at perma.cc/B636-62CE ↩︎ ↩︎
Van Jacobson. Congestion Avoidance and Control. At ACM Symposium on Communications Architectures and Protocols (SIGCOMM), August 1988. doi:10.1145/52324.52356 ↩︎ ↩︎
Bert Hubert. The Ultimate SO_LINGER Page, or: Why Is My TCP Not Reliable. blog.netherlabs.nl, January 2009. Archived at perma.cc/6HDX-L2RR ↩︎ ↩︎
Jerome H. Saltzer, David P. Reed, and David D. Clark. End-To-End Arguments in System Design. ACM Transactions on Computer Systems, volume 2, issue 4, pages 277–288, November 1984. doi:10.1145/357401.357402 ↩︎
Peter Bailis and Kyle Kingsbury. The Network Is Reliable. ACM Queue, volume 12, issue 7, pages 48-55, July 2014. doi:10.1145/2639988.2655736 ↩︎ ↩︎
Joshua B. Leners, Trinabh Gupta, Marcos K. Aguilera, and Michael Walfish. Taming Uncertainty in Distributed Systems with Help from the Network. At 10th European Conference on Computer Systems (EuroSys), April 2015. doi:10.1145/2741948.2741976 ↩︎ ↩︎
Phillipa Gill, Navendu Jain, and Nachiappan Nagappan. Understanding Network Failures in Data Centers: Measurement, Analysis, and Implications. At ACM SIGCOMM Conference, August 2011. doi:10.1145/2018436.2018477 ↩︎
Urs Hölzle. But recently a farmer had started grazing a herd of cows nearby. And whenever they stepped on the fiber link, they bent it enough to cause a blip. x.com, May 2020. Archived at perma.cc/WX8X-ZZA5 ↩︎
CBC News. Hundreds lose internet service in northern B.C. after beaver chews through cable. cbc.ca, April 2021. Archived at perma.cc/UW8C-H2MY ↩︎
Will Oremus. The Global Internet Is Being Attacked by Sharks, Google Confirms. slate.com, August 2014. Archived at perma.cc/P6F3-C6YG ↩︎
Jess Auerbach Jahajeeah. Down to the wire: The ship fixing our internet. continent.substack.com, November 2023. Archived at perma.cc/DP7B-EQ7S ↩︎
Santosh Janardhan. More details about the October 4 outage. engineering.fb.com, October 2021. Archived at perma.cc/WW89-VSXH ↩︎
Tom Parfitt. Georgian woman cuts off web access to whole of Armenia. theguardian.com, April 2011. Archived at perma.cc/KMC3-N3NZ ↩︎
Antonio Voce, Tural Ahmedzade and Ashley Kirk. ‘Shadow fleets’ and subaquatic sabotage: are Europe’s undersea internet cables under attack? theguardian.com, March 2025. Archived at perma.cc/HA7S-ZDBV ↩︎
Shengyun Liu, Paolo Viotti, Christian Cachin, Vivien Quéma, and Marko Vukolić. XFT: Practical Fault Tolerance beyond Crashes. At 12th USENIX Symposium on Operating Systems Design and Implementation (OSDI), November 2016. ↩︎
Mark Imbriaco. Downtime last Saturday. github.blog, December 2012. Archived at perma.cc/M7X5-E8SQ ↩︎ ↩︎
Tom Lianza and Chris Snook. A Byzantine failure in the real world. blog.cloudflare.com, November 2020. Archived at perma.cc/83EZ-ALCY ↩︎ ↩︎
Mohammed Alfatafta, Basil Alkhatib, Ahmed Alquraan, and Samer Al-Kiswany. Toward a Generic Fault Tolerance Technique for Partial Network Partitioning. At 14th USENIX Symposium on Operating Systems Design and Implementation (OSDI), November 2020. ↩︎
Marc A. Donges. Re: bnx2 cards Intermittantly Going Offline. Message to Linux netdev mailing list, spinics.net, September 2012. Archived at perma.cc/TXP6-H8R3 ↩︎ ↩︎
Troy Toman. Inside a CODE RED: Network Edition. signalvnoise.com, September 2020. Archived at perma.cc/BET6-FY25 ↩︎
Kyle Kingsbury. Call Me Maybe: Elasticsearch. aphyr.com, June 2014. perma.cc/JK47-S89J ↩︎
Salvatore Sanfilippo. A Few Arguments About Redis Sentinel Properties and Fail Scenarios. antirez.com, October 2014. perma.cc/8XEU-CLM8 ↩︎
Nicolas Liochon. CAP: If All You Have Is a Timeout, Everything Looks Like a Partition. blog.thislongrun.com, May 2015. Archived at perma.cc/FS57-V2PZ ↩︎
Matthew P. Grosvenor, Malte Schwarzkopf, Ionel Gog, Robert N. M. Watson, Andrew W. Moore, Steven Hand, and Jon Crowcroft. Queues Don’t Matter When You Can JUMP Them! At 12th USENIX Symposium on Networked Systems Design and Implementation (NSDI), May 2015. ↩︎ ↩︎
Theo Julienne. Debugging network stalls on Kubernetes. github.blog, November 2019. Archived at perma.cc/K9M8-XVGL ↩︎
Guohui Wang and T. S. Eugene Ng. The Impact of Virtualization on Network Performance of Amazon EC2 Data Center. At 29th IEEE International Conference on Computer Communications (INFOCOM), March 2010. doi:10.1109/INFCOM.2010.5461931 ↩︎ ↩︎
Brandon Philips. etcd: Distributed Locking and Service Discovery. At Strange Loop, September 2014. ↩︎
Steve Newman. A Systematic Look at EC2 I/O. blog.scalyr.com, October 2012. Archived at perma.cc/FL4R-H2VE ↩︎ ↩︎
Naohiro Hayashibara, Xavier Défago, Rami Yared, and Takuya Katayama. The ϕ Accrual Failure Detector. Japan Advanced Institute of Science and Technology, School of Information Science, Technical Report IS-RR-2004-010, May 2004. Archived at perma.cc/NSM2-TRYA ↩︎
Jeffrey Wang. Phi Accrual Failure Detector. ternarysearch.blogspot.co.uk, August 2013. perma.cc/L452-AMLV ↩︎
Srinivasan Keshav. An Engineering Approach to Computer Networking: ATM Networks, the Internet, and the Telephone Network. Addison-Wesley Professional, May 1997. ISBN: 978-0-201-63442-6 ↩︎ ↩︎
Othmar Kyas. ATM Networks. International Thomson Publishing, 1995. ISBN: 978-1-850-32128-6 ↩︎
Mellanox Technologies. InfiniBand FAQ, Rev 1.3. network.nvidia.com, December 2014. Archived at perma.cc/LQJ4-QZVK ↩︎
Jose Renato Santos, Yoshio Turner, and G. (John) Janakiraman. End-to-End Congestion Control for InfiniBand. At 22nd Annual Joint Conference of the IEEE Computer and Communications Societies (INFOCOM), April 2003. Also published by HP Laboratories Palo Alto, Tech Report HPL-2002-359. doi:10.1109/INFCOM.2003.1208949 ↩︎
Jialin Li, Naveen Kr. Sharma, Dan R. K. Ports, and Steven D. Gribble. Tales of the Tail: Hardware, OS, and Application-level Sources of Tail Latency. At ACM Symposium on Cloud Computing (SOCC), November 2014. doi:10.1145/2670979.2670988 ↩︎
Ulrich Windl, David Dalton, Marc Martinec, and Dale R. Worley. The NTP FAQ and HOWTO. ntp.org, November 2006. ↩︎ ↩︎ ↩︎
John Graham-Cumming. How and why the leap second affected Cloudflare DNS. blog.cloudflare.com, January 2017. Archived at archive.org ↩︎ ↩︎
David Holmes. Inside the Hotspot VM: Clocks, Timers and Scheduling Events – Part I – Windows. blogs.oracle.com, October 2006. Archived at archive.org ↩︎
Joran Dirk Greef. Three Clocks are Better than One. tigerbeetle.com, August 2021. Archived at perma.cc/5RXG-EU6B ↩︎
Oliver Yang. Pitfalls of TSC usage. oliveryang.net, September 2015. Archived at perma.cc/Z2QY-5FRA ↩︎
Steve Loughran. Time on Multi-Core, Multi-Socket Servers. steveloughran.blogspot.co.uk, September 2015. Archived at perma.cc/7M4S-D4U6 ↩︎
James C. Corbett, Jeffrey Dean, Michael Epstein, Andrew Fikes, Christopher Frost, JJ Furman, Sanjay Ghemawat, Andrey Gubarev, Christopher Heiser, Peter Hochschild, Wilson Hsieh, Sebastian Kanthak, Eugene Kogan, Hongyi Li, Alexander Lloyd, Sergey Melnik, David Mwaura, David Nagle, Sean Quinlan, Rajesh Rao, Lindsay Rolig, Dale Woodford, Yasushi Saito, Christopher Taylor, Michal Szymaniak, and Ruth Wang. Spanner: Google’s Globally-Distributed Database. At 10th USENIX Symposium on Operating System Design and Implementation (OSDI), October 2012. ↩︎ ↩︎ ↩︎
M. Caporaloni and R. Ambrosini. How Closely Can a Personal Computer Clock Track the UTC Timescale Via the Internet? European Journal of Physics, volume 23, issue 4, pages L17–L21, June 2012. doi:10.1088/0143-0807/23/4/103 ↩︎
Nelson Minar. A Survey of the NTP Network. alumni.media.mit.edu, December 1999. Archived at perma.cc/EV76-7ZV3 ↩︎
Viliam Holub. Synchronizing Clocks in a Cassandra Cluster Pt. 1 – The Problem. blog.rapid7.com, March 2014. Archived at perma.cc/N3RV-5LNL ↩︎
Poul-Henning Kamp. The One-Second War (What Time Will You Die?) ACM Queue, volume 9, issue 4, pages 44–48, April 2011. doi:10.1145/1966989.1967009 ↩︎
Nelson Minar. Leap Second Crashes Half the Internet. somebits.com, July 2012. Archived at perma.cc/2WB8-D6EU ↩︎
Christopher Pascoe. Time, Technology and Leaping Seconds. googleblog.blogspot.co.uk, September 2011. Archived at perma.cc/U2JL-7E74 ↩︎
Mingxue Zhao and Jeff Barr. Look Before You Leap – The Coming Leap Second and AWS. aws.amazon.com, May 2015. Archived at perma.cc/KPE9-XMFM ↩︎
Darryl Veitch and Kanthaiah Vijayalayan. Network Timing and the 2015 Leap Second. At 17th International Conference on Passive and Active Measurement (PAM), April 2016. doi:10.1007/978-3-319-30505-9_29 ↩︎
VMware, Inc. Timekeeping in VMware Virtual Machines. vmware.com, October 2008. Archived at perma.cc/HM5R-T5NF ↩︎
Victor Yodaiken. Clock Synchronization in Finance and Beyond. yodaiken.com, November 2017. Archived at perma.cc/9XZD-8ZZN ↩︎
Mustafa Emre Acer, Emily Stark, Adrienne Porter Felt, Sascha Fahl, Radhika Bhargava, Bhanu Dev, Matt Braithwaite, Ryan Sleevi, and Parisa Tabriz. Where the Wild Warnings Are: Root Causes of Chrome HTTPS Certificate Errors. At ACM SIGSAC Conference on Computer and Communications Security (CCS), pages 1407–1420, October 2017. doi:10.1145/3133956.3134007 ↩︎
European Securities and Markets Authority. MiFID II / MiFIR: Regulatory Technical and Implementing Standards – Annex I. esma.europa.eu, Report ESMA/2015/1464, September 2015. Archived at perma.cc/ZLX9-FGQ3 ↩︎
Luke Bigum. Solving MiFID II Clock Synchronisation With Minimum Spend (Part 1). catach.blogspot.com, November 2015. Archived at perma.cc/4J5W-FNM4 ↩︎
Oleg Obleukhov and Ahmad Byagowi. How Precision Time Protocol is being deployed at Meta. engineering.fb.com, November 2022. Archived at perma.cc/29G6-UJNW ↩︎
John Wiseman. gpsjam.org, July 2022. ↩︎
Josh Levinson, Julien Ridoux, and Chris Munns. It’s About Time: Microsecond-Accurate Clocks on Amazon EC2 Instances. aws.amazon.com, November 2023. Archived at perma.cc/56M6-5VMZ ↩︎
Kyle Kingsbury. Call Me Maybe: Cassandra. aphyr.com, September 2013. Archived at perma.cc/4MBR-J96V ↩︎ ↩︎ ↩︎
John Daily. Clocks Are Bad, or, Welcome to the Wonderful World of Distributed Systems. riak.com, November 2013. Archived at perma.cc/4XB5-UCXY ↩︎ ↩︎
Marc Brooker. It’s About Time! brooker.co.za, November 2023. Archived at perma.cc/N6YK-DRPA ↩︎
Kyle Kingsbury. The Trouble with Timestamps. aphyr.com, October 2013. Archived at perma.cc/W3AM-5VAV ↩︎
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 ↩︎
Justin Sheehy. There Is No Now: Problems With Simultaneity in Distributed Systems. ACM Queue, volume 13, issue 3, pages 36–41, March 2015. doi:10.1145/2733108 ↩︎
Murat Demirbas. Spanner: Google’s Globally-Distributed Database. muratbuffalo.blogspot.co.uk, July 2013. Archived at perma.cc/6VWR-C9WB ↩︎
Dahlia Malkhi and Jean-Philippe Martin. Spanner’s Concurrency Control. ACM SIGACT News, volume 44, issue 3, pages 73–77, September 2013. doi:10.1145/2527748.2527767 ↩︎
Franck Pachot. Achieving Precise Clock Synchronization on AWS. yugabyte.com, December 2024. Archived at perma.cc/UYM6-RNBS ↩︎
Spencer Kimball. Living Without Atomic Clocks: Where CockroachDB and Spanner diverge. cockroachlabs.com, January 2022. Archived at perma.cc/AWZ7-RXFT ↩︎
Murat Demirbas. Use of Time in Distributed Databases (part 4): Synchronized clocks in production databases. muratbuffalo.blogspot.com, January 2025. Archived at perma.cc/9WNX-Q9U3 ↩︎
Cary G. Gray and David R. Cheriton. Leases: An Efficient Fault-Tolerant Mechanism for Distributed File Cache Consistency. At 12th ACM Symposium on Operating Systems Principles (SOSP), December 1989. doi:10.1145/74850.74870 ↩︎
Daniel Sturman, Scott Delap, Max Ross, et al. Roblox Return to Service. corp.roblox.com, January 2022. Archived at perma.cc/8ALT-WAS4 ↩︎
Todd Lipcon. Avoiding Full GCs with MemStore-Local Allocation Buffers. slideshare.net, February 2011. Archived at https://perma.cc/CH62-2EWJ ↩︎
Christopher Clark, Keir Fraser, Steven Hand, Jacob Gorm Hansen, Eric Jul, Christian Limpach, Ian Pratt, and Andrew Warfield. Live Migration of Virtual Machines. At 2nd USENIX Symposium on Symposium on Networked Systems Design & Implementation (NSDI), May 2005. ↩︎
Mike Shaver. fsyncers and Curveballs. shaver.off.net, May 2008. Archived at archive.org ↩︎
Zhenyun Zhuang and Cuong Tran. Eliminating Large JVM GC Pauses Caused by Background IO Traffic. engineering.linkedin.com, February 2016. Archived at perma.cc/ML2M-X9XT ↩︎
Martin Thompson. Java Garbage Collection Distilled. mechanical-sympathy.blogspot.co.uk, July 2013. Archived at perma.cc/DJT3-NQLQ ↩︎ ↩︎
David Terei and Amit Levy. Blade: A Data Center Garbage Collector. arXiv:1504.02578, April 2015. ↩︎
Martin Maas, Tim Harris, Krste Asanović, and John Kubiatowicz. Trash Day: Coordinating Garbage Collection in Distributed Systems. At 15th USENIX Workshop on Hot Topics in Operating Systems (HotOS), May 2015. ↩︎
Martin Fowler. The LMAX Architecture. martinfowler.com, July 2011. Archived at perma.cc/5AV4-N6RJ ↩︎
Joseph Y. Halpern and Yoram Moses. Knowledge and common knowledge in a distributed environment. Journal of the ACM (JACM), volume 37, issue 3, pages 549–587, July 1990. doi:10.1145/79147.79161 ↩︎ ↩︎
Chuzhe Tang, Zhaoguo Wang, Xiaodong Zhang, Qianmian Yu, Binyu Zang, Haibing Guan, and Haibo Chen. Ad Hoc Transactions in Web Applications: The Good, the Bad, and the Ugly. At ACM International Conference on Management of Data (SIGMOD), June 2022. doi:10.1145/3514221.3526120 ↩︎
Flavio P. Junqueira and Benjamin Reed. ZooKeeper: Distributed Process Coordination. O’Reilly Media, 2013. ISBN: 978-1-449-36130-3 ↩︎ ↩︎
Enis Söztutar. HBase and HDFS: Understanding Filesystem Usage in HBase. At HBaseCon, June 2013. Archived at perma.cc/4DXR-9P88 ↩︎
SUSE LLC. SUSE Linux Enterprise High Availability 15 SP6 Administration Guide, Section 12: Fencing and STONITH. documentation.suse.com, March 2025. Archived at perma.cc/8LAR-EL9D ↩︎
Mike Burrows. The Chubby Lock Service for Loosely-Coupled Distributed Systems. At 7th USENIX Symposium on Operating System Design and Implementation (OSDI), November 2006. ↩︎
Kyle Kingsbury. etcd 3.4.3. jepsen.io, January 2020. Archived at perma.cc/2P3Y-MPWU ↩︎
Ensar Basri Kahveci. Distributed Locks are Dead; Long Live Distributed Locks! hazelcast.com, April 2019. Archived at perma.cc/7FS5-LDXE ↩︎
Martin Kleppmann. How to do distributed locking. martin.kleppmann.com, February 2016. Archived at perma.cc/Y24W-YQ5L ↩︎
Salvatore Sanfilippo. Is Redlock safe? antirez.com, February 2016. Archived at perma.cc/B6GA-9Q6A ↩︎
Gunnar Morling. Leader Election With S3 Conditional Writes. www.morling.dev, August 2024. Archived at perma.cc/7V2N-J78Y ↩︎
Leslie Lamport, Robert Shostak, and Marshall Pease. The Byzantine Generals Problem. ACM Transactions on Programming Languages and Systems (TOPLAS), volume 4, issue 3, pages 382–401, July 1982. doi:10.1145/357172.357176 ↩︎
Jim N. Gray. Notes on Data Base Operating Systems. in Operating Systems: An Advanced Course, Lecture Notes in Computer Science, volume 60, edited by R. Bayer, R. M. Graham, and G. Seegmüller, pages 393–481, Springer-Verlag, 1978. ISBN: 978-3-540-08755-7. Archived at perma.cc/7S9M-2LZU ↩︎
Brian Palmer. How Complicated Was the Byzantine Empire? slate.com, October 2011. Archived at perma.cc/AN7X-FL3N ↩︎
Leslie Lamport. My Writings. lamport.azurewebsites.net, December 2014. Archived at perma.cc/5NNM-SQGR ↩︎
John Rushby. Bus Architectures for Safety-Critical Embedded Systems. At 1st International Workshop on Embedded Software (EMSOFT), October 2001. doi:10.1007/3-540-45449-7_22 ↩︎ ↩︎
Jake Edge. ELC: SpaceX Lessons Learned. lwn.net, March 2013. Archived at perma.cc/AYX8-QP5X ↩︎
Shehar Bano, Alberto Sonnino, Mustafa Al-Bassam, Sarah Azouvi, Patrick McCorry, Sarah Meiklejohn, and George Danezis. SoK: Consensus in the Age of Blockchains. At 1st ACM Conference on Advances in Financial Technologies (AFT), October 2019. doi:10.1145/3318041.3355458 ↩︎
Ezra Feilden, Adi Oltean, and Philip Johnston. Why we should train AI in space. White Paper, starcloud.com, September 2024. Archived at perma.cc/7Y3S-8UB6 ↩︎
James Mickens. The Saddest Moment. USENIX ;login, May 2013. Archived at perma.cc/T7BZ-XCFR ↩︎
Martin Kleppmann and Heidi Howard. Byzantine Eventual Consistency and the Fundamental Limits of Peer-to-Peer Databases. arxiv.org, December 2020. doi:10.48550/arXiv.2012.00472 ↩︎
Martin Kleppmann. Making CRDTs Byzantine Fault Tolerant. At 9th Workshop on Principles and Practice of Consistency for Distributed Data (PaPoC), April 2022. doi:10.1145/3517209.3524042 ↩︎
Evan Gilman. The Discovery of Apache ZooKeeper’s Poison Packet. pagerduty.com, May 2015. Archived at perma.cc/RV6L-Y5CQ ↩︎ ↩︎
Jonathan Stone and Craig Partridge. When the CRC and TCP Checksum Disagree. At ACM Conference on Applications, Technologies, Architectures, and Protocols for Computer Communication (SIGCOMM), August 2000. doi:10.1145/347059.347561 ↩︎
Evan Jones. How Both TCP and Ethernet Checksums Fail. evanjones.ca, October 2015. Archived at perma.cc/9T5V-B8X5 ↩︎
Cynthia Dwork, Nancy Lynch, and Larry Stockmeyer. Consensus in the Presence of Partial Synchrony. Journal of the ACM, volume 35, issue 2, pages 288–323, April 1988. doi:10.1145/42282.42283 ↩︎ ↩︎ ↩︎
Richard D. Schlichting and Fred B. Schneider. Fail-stop processors: an approach to designing fault-tolerant computing systems. ACM Transactions on Computer Systems (TOCS), volume 1, issue 3, pages 222–238, August 1983. doi:10.1145/357369.357371 ↩︎
Thanh Do, Mingzhe Hao, Tanakorn Leesatapornwongsa, Tiratat Patana-anake, and Haryadi S. Gunawi. Limplock: Understanding the Impact of Limpware on Scale-out Cloud Systems. At 4th ACM Symposium on Cloud Computing (SoCC), October 2013. doi:10.1145/2523616.2523627 ↩︎
Josh Snyder and Joseph Lynch. Garbage collecting unhealthy JVMs, a proactive approach. Netflix Technology Blog, netflixtechblog.medium.com, November 2019. Archived at perma.cc/8BTA-N3YB ↩︎
Haryadi S. Gunawi, Riza O. Suminto, Russell Sears, Casey Golliher, Swaminathan Sundararaman, Xing Lin, Tim Emami, Weiguang Sheng, Nematollah Bidokhti, Caitie McCaffrey, Gary Grider, Parks M. Fields, Kevin Harms, Robert B. Ross, Andree Jacobson, Robert Ricci, Kirk Webb, Peter Alvaro, H. Birali Runesha, Mingzhe Hao, and Huaicheng Li. Fail-Slow at Scale: Evidence of Hardware Performance Faults in Large Production Systems. At 16th USENIX Conference on File and Storage Technologies, February 2018. ↩︎
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 ↩︎
Chang Lou, Peng Huang, and Scott Smith. Understanding, Detecting and Localizing Partial Failures in Large System Software. At 17th USENIX Symposium on Networked Systems Design and Implementation (NSDI), February 2020. ↩︎
Peter Bailis and Ali Ghodsi. Eventual Consistency Today: Limitations, Extensions, and Beyond. ACM Queue, volume 11, issue 3, pages 55-63, March 2013. doi:10.1145/2460276.2462076 ↩︎
Bowen Alpern and Fred B. Schneider. Defining Liveness. Information Processing Letters, volume 21, issue 4, pages 181–185, October 1985. doi:10.1016/0020-0190(85)90056-0 ↩︎
Flavio P. Junqueira. Dude, Where’s My Metadata? fpj.me, May 2015. Archived at perma.cc/D2EU-Y9S5 ↩︎
Scott Sanders. January 28th Incident Report. github.com, February 2016. Archived at perma.cc/5GZR-88TV ↩︎
Jay Kreps. A Few Notes on Kafka and Jepsen. blog.empathybox.com, September 2013. perma.cc/XJ5C-F583 ↩︎
Marc Brooker and Ankush Desai. Systems Correctness Practices at AWS. Queue, Volume 22, Issue 6, November/December 2024. doi:10.1145/3712057 ↩︎
Andrey Satarin. Testing Distributed Systems: Curated list of resources on testing distributed systems. asatarin.github.io. Archived at perma.cc/U5V8-XP24 ↩︎
Jack Vanlightly. Verifying Kafka transactions - Diary entry 2 - Writing an initial TLA+ spec. jack-vanlightly.com, December 2024. Archived at perma.cc/NSQ8-MQ5N ↩︎
Siddon Tang. From Chaos to Order — Tools and Techniques for Testing TiDB, A Distributed NewSQL Database. pingcap.com, April 2018. Archived at perma.cc/5EJB-R29F ↩︎
Nathan VanBenschoten. Parallel Commits: An atomic commit protocol for globally distributed transactions. cockroachlabs.com, November 2019. Archived at perma.cc/5FZ7-QK6J ↩︎
Jack Vanlightly. Paper: VR Revisited - State Transfer (part 3). jack-vanlightly.com, December 2022. Archived at perma.cc/KNK3-K6WS ↩︎
Hillel Wayne. What if the spec doesn’t match the code? buttondown.com, March 2024. Archived at perma.cc/8HEZ-KHER ↩︎
Lingzhi Ouyang, Xudong Sun, Ruize Tang, Yu Huang, Madhav Jivrajani, Xiaoxing Ma, Tianyin Xu. Multi-Grained Specifications for Distributed System Model Checking and Verification. At 20th European Conference on Computer Systems (EuroSys), March 2025. doi:10.1145/3689031.3696069 ↩︎
Yury Izrailevsky and Ariel Tseitlin. The Netflix Simian Army. netflixtechblog.com, July, 2011. Archived at perma.cc/M3NY-FJW6 ↩︎
Kyle Kingsbury. Jepsen: On the perils of network partitions. aphyr.com, May, 2013. Archived at perma.cc/W98G-6HQP ↩︎
Kyle Kingsbury. Jepsen Analyses. jepsen.io, 2024. Archived at perma.cc/8LDN-D2T8 ↩︎
Rupak Majumdar and Filip Niksic. Why is random testing effective for partition tolerance bugs? Proceedings of the ACM on Programming Languages (PACMPL), volume 2, issue POPL, article no. 46, December 2017. doi:10.1145/3158134 ↩︎
FoundationDB project authors. Simulation and Testing. apple.github.io. Archived at perma.cc/NQ3L-PM4C ↩︎
Alex Kladov. Simulation Testing For Liveness. tigerbeetle.com, July 2023. Archived at perma.cc/RKD4-HGCR ↩︎
Alfonso Subiotto Marqués. (Mostly) Deterministic Simulation Testing in Go. polarsignals.com, May 2024. Archived at perma.cc/ULD6-TSA4 ↩︎