分布式理论面试
分布式理论面试
分布式常识
扩展
- Designing Data-Intensive Applications - 《数据密集型应用系统设计》,Martin Kleppmann 著,目前公认最优秀的分布式系统入门到进阶书籍,第一、二、八、九章全面覆盖分布式基础理论。
- Distributed Systems: Principles and Paradigms - Andrew S. Tanenbaum 著,分布式系统经典教材。
- A Note on Distributed Systems - 1994 年经典论文,阐述了"远程交互不能像本地对象那样"这一关键认知。
- The Eight Fallacies of Distributed Computing - Peter Deutsch 总结的分布式系统新手常犯的 8 个误区。
【简单】什么是分布式系统?它和集中式系统有什么区别?⭐
- 分布式系统:由多个独立节点通过网络连接,协同完成同一任务,对外表现为单一系统的集群架构,核心是多节点协作、去中心化 / 弱中心化。
- 集中式系统:由单个核心节点处理所有任务,数据存储、计算、请求响应全由该节点完成,核心是单节点独占、强中心化。
二者核心差异对比如下:
| 对比维度 | 集中式系统 | 分布式系统 |
|---|---|---|
| 架构特征 | 单节点处理所有任务,强中心化 | 多节点协作,去中心化 / 弱中心化 |
| 可扩展性 | 受限于单机硬件瓶颈,垂直扩展成本高 | 可水平扩展,理论上无上限 |
| 可用性 | 单点故障即全系统故障 | 部分节点故障不影响整体(需冗余) |
| 一致性 | 天然强一致(数据只有一份) | 需通过协议保证一致性(CAP 权衡) |
| 性能瓶颈 | CPU、内存、磁盘、网络 IO | 网络延迟、节点间协调开销 |
| 典型代表 | 传统关系型数据库单机部署 | HDFS、Kafka、ZooKeeper、TiDB |
【简单】分布式系统有哪些核心特征?⭐
要点
分布式系统的核心特征可概括为:透明性、可扩展性、高可用性、容错性、一致性、并发性。
| 特性 | 说明 |
|---|---|
| 透明性 | 用户感知不到系统的分布式特性,使用时如同使用单机系统。包括访问透明、位置透明、迁移透明、复制透明、并发透明、故障透明等 |
| 可扩展性 | 通过增加节点即可线性提升系统能力,包括垂直扩展(提升单机配置)和水平扩展(增加节点数量) |
| 高可用性 | 系统中部分节点故障不影响整体服务,通过冗余部署、故障转移等机制保障 |
| 容错性 | 系统能在部分节点故障、网络分区等异常情况下继续运行 |
| 一致性 | 多个数据副本之间保持一致的状态,根据强度分为强一致性、弱一致性、最终一致性 |
| 并发性 | 系统中多个节点可同时处理请求,通过并发控制机制协调 |
【中等】分布式系统面临哪些核心挑战?⭐
要点
分布式系统的核心挑战源于不可靠的网络和不可靠的时钟,由此衍生出部分失效、网络分区、时钟漂移、一致性、共识等难题。
分布式系统相比单机系统,主要面临以下挑战:
部分失效(Partial Failure):分布式系统中某些组件可能故障而其他组件正常工作,这种"部分失效"具有不确定性,是最难处理的问题。例如,请求未收到响应时,无法区分是请求丢失、节点崩溃还是响应丢失。
不可靠的网络:互联网及数据中心内部网络(以太网)都是异步网络,消息可能丢失、延迟、乱序。常见的网络问题包括:
- 请求丢失(网线被拔、网络拥塞)
- 请求排队等待(接收方过载)
- 远程节点崩溃或暂时无法响应(如 GC 暂停)
- 响应丢失或延迟
不可靠的时钟:不同节点的物理时钟无法完全同步,即使使用 NTP 校准仍存在误差。时钟问题影响超时检测、事件排序、缓存过期等场景。
网络分区:网络故障导致集群被分割为多个互不可达的子集,各子集独立运作可能产生数据冲突。
数据一致性:多副本环境下,如何保证各副本数据一致是核心难题,涉及 CAP 权衡。
分布式共识:如何在可能故障的节点间就某个值达成一致,是分布式系统的基石问题。
【中等】什么是分布式计算的八大谬误?⭐
要点
分布式计算八大谬误是 Peter Deutsch 等人在 Sun Microsystems 总结的新手常犯的认知错误,这些"想当然"的假设在分布式系统中都不成立。
1994 年,Peter Deutsch 等人在 Sun Microsystems 总结了分布式系统新手常犯的 8 个谬误:
| 序号 | 谬误 | 说明 |
|---|---|---|
| 1 | 网络是可靠的 | 网络会丢包、断连,必须处理消息丢失 |
| 2 | 延迟为零 | 网络延迟不为零且不可预测,远程调用不能等同本地调用 |
| 3 | 带宽是无限的 | 带宽有限,大对象传输需要考虑压缩、分片 |
| 4 | 网络是安全的 | 网络不安全,需要认证、加密、防重放 |
| 5 | 拓扑不会变化 | 网络拓扑会变化(节点增减、路由变化),系统需适应动态拓扑 |
| 6 | 只有一个管理员 | 多团队、多组织协作,配置和策略可能冲突 |
| 7 | 传输代价为零 | 数据序列化、反序列化、网络传输都有成本 |
| 8 | 网络是同构的 | 异构网络(不同协议、版本、硬件)普遍存在 |
这八大谬误的本质是:在分布式系统中,故障是不可避免的。构建可靠的分布式系统必须建立容错机制,明确软件在故障场景下的行为。
逻辑时钟
扩展
- Time, Clocks, and the Ordering of Events in a Distributed System,译文,解读 - Lamport 介绍 happened before、偏序关系(partial ordering)、逻辑时钟(Logical Clocks)概念,提出解决分布式系统中区分事件发生的时序问题的方法。
- Virtual Time and Global States of Distributed Systems,解读 - 逻辑时钟无法描述事件的因果关系。本文提出了向量时钟,这种算法利用了向量这种数据结构将全局各个进程的逻辑时间戳广播给各个进程,通过向量时间戳就能够比较任意两个事件的因果关系。
- 逻辑时钟 - 如何刻画分布式中的事件顺序
【简单】为什么需要逻辑时钟?⭐
要点
不同节点的物理时钟即使校准(NTP)也无法完全保持一致。
为什么需要逻辑时钟?分布式系统中以系统时间来确定事件顺序有什么问题吗?
不同节点的物理时钟无法完全保持一致。即使引入一个全局时钟(例如:NTP)来进行校准,由于网络通信延迟的不确定性,以及时钟计时的偏差,无法保证每个节点的时间完全一致。
在分布式系统中,由于网络通信延迟的不确定性, 仅仅以接收顺序作为整个分布式系统中事件的发生顺序是不可取的。
【中等】什么是偏序?什么是全序?⭐
全序和偏序是数学上的术语,按照数学内容阐述比较晦涩,简单来说:
- 偏序是部分可比较的有序关系。
- 全序是在偏序基础上,要求全部元素必须可比较的有序关系。
【困难】什么是逻辑时钟?⭐⭐
要点
- Lamport 逻辑时钟构建了一个全序时钟来描述事件顺序。
- Lamport 逻辑时钟的缺陷是无法描述同时发生的事件。
1978 年,Lamport 在 Time, Clocks, and the Ordering of Events in a Distributed System 中提出了逻辑时钟的概念,来解决分布式系统中区分事件发生的时序问题。
逻辑时钟并不度量时间本身,仅区分事件发生的前后顺序。
分布式系统中按是否存在节点交互可分为三类事件,一类发生于节点内部,二是发送事件,三是接收事件。Lamport 时间戳原理如下:

- 每个事件对应一个 Lamport 计数器,初始值为 0
- 如果事件在节点内发生,计数器加 1
- 如果事件属于发送事件,计数器加 1 并在消息中带上该计数器
- 如果事件属于接收事件,计数器 = Max(本地计数器,消息中的计数器) + 1
综上,Lamport 逻辑时钟构建了一个全序时钟来描述事件顺序。Lamport 逻辑时钟的缺陷是无法描述同时发生的事件。
【困难】什么是向量时钟?⭐⭐
要点
- 向量时钟在逻辑时钟基础上改进:不仅记录了本节点的时间戳,还记录了其他节点的时间戳。
- 其本质在于将逻辑时钟的全序计数器改造为向量时钟的偏序大小关系:向量有序,则事件有序;向量平行,则事件并发。
- 向量时钟可以发现数据冲突,但不能解决数据冲突。
向量时钟其实是在逻辑时钟的基础上进行了演进,算法逻辑类似,只是不仅记录了本节点的时间戳,还记录了其他节点的时间戳。其本质在于将逻辑时钟的全序计数器改造为向量时钟的偏序大小关系:向量有序,则事件有序;向量平行,则事件并发。

向量时钟可以发现数据冲突,但不能解决数据冲突。
【困难】什么是版本向量时钟?⭐
要点
版本向量时钟只有在更新数据的时候做向量自增。
在向量时钟算法中, 消息传播后,发送方的向量一定会小于接收者的向量, 是因为接收者对齐了发送者的原因。
版本向量在此基础上,做了一点加强:消息传播后,发送方也对齐接收者的向量,也就是双向对齐,在版本向量中,叫做同步。
发送消息和接收消息的时候不再自增向量中的自己的计数器,而是只做双方的向量对齐操作。 也就是,只有在更新数据的时候做向量自增。
【困难】逻辑时钟、向量时钟、版本向量时钟有什么差异?⭐
| 特性维度 | 逻辑时钟 | 向量时钟 | 版本向量时钟 |
|---|---|---|---|
| 本质 | 单一递增计数器 | 节点ID → 计数器的映射 | 向量时钟的特化 |
| 作用 | 判断事件因果关系 | 精确判断事件是因果还是并发关系 | 数据版本冲突检测 |
| 能否识别并发 | 不能 | 能 | 能(主要用途) |
| 能否判断因果 | 部分能(可能漏判) | 精确能 | 能(专用于数据因果) |
| 数据开销 | 最小(传递1个整数) | 较大(传递整个向量,大小=节点数) | 同向量时钟(大小=副本数) |
| 典型应用场景 | 分布式锁、事件日志排序 | 因果一致性消息系统、分布式故障检测 | 多主/无主复制数据库的冲突检测(如Cassandra、Riak) |
| 类比说明 | 电影院取票机(只给递增票号) | 多窗口叫号系统(各窗口独立号,可见全局进度) | 协同编辑历史(记录每人编辑版本,用于合并冲突) |
【困难】什么是混合逻辑时钟(HLC)?⭐
要点
- HLC(Hybrid Logical Clock,混合逻辑时钟) 结合了物理时钟和逻辑时钟的优点:既能逼近物理时间,又能保证因果一致性。
- HLC 由 Kulkarni 等人在 2014 年论文 Logical Physical Clocks 中提出。
- 典型应用:CockroachDB、MongoDB、YugabyteDB 等全球分布式数据库使用 HLC 实现分布式事务排序。
物理时钟的问题:NTP 同步存在误差,不同节点的时间戳可能偏离几十毫秒甚至秒级,无法用于严格的事件排序。
逻辑时钟的问题:与物理时间脱节,无法回答"这个事件发生在几点几分"这类问题,对依赖超时、TTL 的场景不友好。
HLC 的核心思想:每个节点维护一个混合时间戳 (physical_time, logical_counter),规则如下:
- 本地事件:取
max(本地物理时间, HLC 的物理部分),若物理部分未增长则逻辑计数器 +1,否则逻辑计数器清零。 - 发送事件:先按本地事件规则更新 HLC,将 HLC 随消息发送。
- 接收事件:物理部分取
max(本地物理时间, 消息 HLC 物理部分, 本地 HLC 物理部分);若物理部分相同则逻辑计数器取max(本地计数器, 消息计数器) + 1。
HLC 的优势:
- 保证因果一致性:若事件 A 因果先于事件 B,则
HLC(A) < HLC(B)。 - 逼近物理时间:HLC 的物理部分始终接近真实时间,便于与超时、TTL 配合。
- 有界偏差:当物理时钟偏差有界时,HLC 与物理时间的偏差也有界。
一致性
扩展
- Brewer’s Conjecture and the Feasibility of Consistent, Available, Partition-Tolerant Web Services,解读 - 经典的 CAP 定理,即:在一个分布式系统中,当发生网络分区时,那么强一致性和可用性只能二选一。
- CAP Twelve Years Later: How the “Rules” Have Changed, 解读 - CAP 定理的新解读,并阐述 CAP 定理的一些常见误区。
- BASE: An Acid Alternative,译文 - BASE 定理是对 CAP 中一致性和可用性的权衡,提出采用适当的方式来使系统达到最终一致性。
【简单】什么是强一致性?什么是弱一致性?什么是最终一致性?⭐⭐
一致性(Consistency)指的是多个数据副本是否能保持一致的特性。
数据一致性又可以分为以下几点:
- 强一致性:数据更新操作结果和操作响应总是一致的,即操作响应通知更新失败,那么数据一定没有被更新,而不是处于不确定状态。通俗的说,分布式系统在执行写操作成功后,如果所有用户都能够读取到最新的值,该系统就被认为具有强一致性。
- 弱一致性:系统在写入数据成功后,不承诺立即能读到最新的值,也不承诺什么时候能读到,但是过一段时间之后用户可以看到更新后的值。那么用户读不到最新数据的这段时间被称为"不一致窗口时间"。
- 最终一致性:最终一致性作为弱一致性中的特例,强调的是所有数据副本,在经过一段时间的同步后,最终能够到达一致的状态,不需要实时保证系统数据的强一致性。
【中等】什么是线性一致性、顺序一致性、因果一致性?⭐⭐
要点
一致性模型从强到弱大致为:线性一致性 > 顺序一致性 > 因果一致性 > 读己之写 / 单调读等会话一致性 > 最终一致性。一致性越强,对客户端越直观,但性能和可用性代价越大。
这是分布式系统中最常被混淆的几个一致性模型,它们的强度依次递减:
| 一致性模型 | 核心约束 | 是否需要物理时钟 | 典型应用 |
|---|---|---|---|
| 线性一致性(Linearizability) | 所有操作看起来像在单一时间点原子完成,且符合真实物理时间的先后顺序。即:若操作 A 在操作 B 开始前完成,则所有节点都必须看到 A 先于 B | 需要(依赖全局物理时间) | etcd、ZooKeeper 的线性读、Spanner 的 TrueTime |
| 顺序一致性(Sequential Consistency) | 所有节点看到操作的全局顺序一致,但该顺序不需要符合物理时间。只要同一节点内的操作顺序保留即可 | 不需要 | 多核 CPU 内存模型、ZooKeeper 默认读 |
| 因果一致性(Causal Consistency) | 保证有因果关系的操作顺序一致,但并发操作(无因果关系)的顺序在不同节点可以不同 | 不需要(可用向量时钟实现) | 评论系统(回复必须在被回复之后可见) |
| 读己之写(Read-Your-Writes) | 客户端总能读到自己刚写入的值 | 不需要 | 会话级缓存、CDN |
| 单调读(Monotonic Reads) | 客户端不会读到比之前更旧的数据 | 不需要 | 社交网络时间线 |
| 最终一致性(Eventual Consistency) | 停止写入后,所有副本最终收敛到一致状态,但不保证过程中的可见性顺序 | 不需要 | DNS、Cassandra、DynamoDB |
关键区别说明:
- 线性一致性 vs 顺序一致性:线性一致性额外要求操作顺序符合物理时间(实时性),顺序一致性只要求全局顺序一致即可。例如,节点 A 写入 x=1 完成后,节点 B 立刻读 x——线性一致性要求 B 必须读到 1,顺序一致性不保证(因为不要求符合物理时间)。
- 顺序一致性 vs 因果一致性:顺序一致性要求所有操作有一个全局一致顺序,因果一致性只要求有因果关系的操作顺序一致,并发操作可乱序。因果一致性弱于顺序一致性。
- CAP 中的 C 指的是线性一致性,这是最强的一致性模型,因此 CAP 的约束才如此严格。
【简单】什么是 ACID?⭐⭐⭐
那么,什么是 ACID 特性呢?ACID 是数据库事务正确执行的四个基本要素的单词缩写:
- 原子性(Atomicity)
- 原子是指不可分解为更小粒度的东西。事务的原子性意味着:事务中的所有操作要么全部成功,要么全部失败。
- 回滚可以用日志来实现,日志记录着事务所执行的修改操作,在回滚时反向执行这些修改操作即可。
- ACID 中的原子性并不关乎多个操作的并发性,它并没有描述多个线程试图访问相同的数据会发生什么情况,后者其实是由 ACID 的隔离性所定义。
- 一致性(Consistency)
- 数据库在事务执行前后都保持一致性状态。
- 在一致性状态下,所有事务对一个数据的读取结果都是相同的。
- 一致性本质上要求应用层来维护状态一致(或者恒等),应用程序有责任正确地定义事务来保持一致性。这不是数据库可以保证的事情。
- 隔离性(Isolation)
- 同时运行的事务互不干扰。换句话说,一个事务所做的修改在最终提交以前,对其它事务是不可见的。
- 持久性(Durability)
- 一旦事务提交,则其所做的修改将会永远保存到数据库中。即使系统发生崩溃,事务执行的结果也不能丢失。
- 可以通过数据库备份和恢复来实现,在系统发生奔溃时,使用备份的数据库进行数据恢复。
一个支持事务(Transaction)中的数据库系统,必需要具有这四种特性,否则在事务过程(Transaction processing)当中无法保证数据的正确性。
- 只有满足一致性,事务的执行结果才是正确的。
- 在无并发的情况下,事务串行执行,隔离性一定能够满足。此时只要能满足原子性,就一定能满足一致性。
- 在并发的情况下,多个事务并行执行,事务不仅要满足原子性,还需要满足隔离性,才能满足一致性。
- 事务满足持久化是为了能应对系统崩溃的情况。
扩展
- Brewer’s Conjecture and the Feasibility of Consistent, Available, Partition-Tolerant Web Services,解读 - 经典的 CAP 定理,即:在一个分布式系统中,当发生网络分区时,那么强一致性和可用性只能二选一。
- CAP Twelve Years Later: How the “Rules” Have Changed, 解读 - CAP 定理的新解读,并阐述 CAP 定理的一些常见误区。
- BASE: An Acid Alternative,译文 - BASE 定理是对 CAP 中一致性和可用性的权衡,提出采用适当的方式来使系统达到最终一致性。
【中等】什么是 CAP 定理?⭐⭐
CAP 要点
- CAP 就是取 Consistency、Availability、Partition Tolerance 的首字母而命名。
- CAP 定理提出:Consistency、Availability、Partition Tolerance 三者不可兼得。
- 在分布式系统中,分区容错不可避免,因此,CAP 定理实际上是要在可用性(A)和一致性(C)之间做权衡。
CAP 定理解读
CAP 定理提出:分布式系统有三个指标,这三个指标不能同时做到:
- 一致性(Consistency):这里的 C 特指线性一致性(Linearizability),即任何读操作都能读到最近一次写操作的结果,所有节点看到的数据视图如同单一节点。注意:CAP 的 C 与 ACID 的 C 含义不同。
- 可用性(Availability):对每一个非故障节点的请求都必须返回(非错误)响应,但不保证返回的是最新数据。注意:这里要求的是"每个请求都能得到响应",而不是"系统能处理多少请求"。
- 分区容错性(Partition Tolerance):当网络发生分区(节点间消息丢失或延迟)时,系统仍能继续运行。
CAP 就是取 Consistency、Availability、Partition Tolerance 的首字母而命名。

在分布式系统中,分区容错性是一个既定的事实:因为分布式系统总会出现各种各样的问题,如由于网络原因而导致节点失联;发生机器故障;机器重启或升级等等。因此,CAP 定理实际上是在可用性(A)和一致性(C)之间做权衡。
选择 CP 还是 AP?
选择 CP 还是 AP,应该视具体业务场景而定
- 选择 AP 模式,偏向于保证服务的高可用性。用户访问系统的时候,都能得到响应数据,不会出现响应错误;但是,当出现分区故障时,相同的读操作,访问不同的节点,得到响应数据可能不一样。
- 选择 CP 模式,一旦因为消息丢失、延迟过高发生了网络分区,就会影响用户的体验和业务的可用性。因为为了防止数据不一致,系统将拒绝新数据的写入。
一个最具代表性的问题是:服务注册中心应该选择 AP 还是 CP?
在微服务架构下,服务注册和服务发现机制中主要有三种角色:
- 服务提供者(RPC Server / Provider)
- 服务消费者(RPC Client / Consumer)
- 服务注册中心(Registry)
注册中心负责协调服务注册和服务发现,显然它是核心中的核心。主流的注册中心有很多,如:ZooKeeper、Nacos、Eureka、Consul、etcd 等。在针对注册中心进行技术选型时,其 CAP 设计也是一个比较的维度。
- CP 模型代表:ZooKeeper、etcd。系统强调数据的一致性,当数据一致性无法保证时(如:正在选举主节点),系统拒绝请求。
- AP 模型代表:Nacos、Eureka。系统强调可用性,牺牲一定的一致性(即服务节点上的数据不保证是最新的),来保证整体服务可用。
对于服务注册中心而言,即使不同节点保存的服务注册信息存在差异,也不会造成灾难性的后果,仅仅是信息滞后而已。但是,如果为了追求数据一致性,使得服务发现短时间内不可用,负面影响更严重。所以,对于服务注册中心而言,可用性比一致性更重要,一般应该选择 AP 模型。
【中等】CAP 定理真的正确吗?⭐⭐
要点
CAP 理论模型局限性很大,未考虑网络分区以外的各种故障、异常情况,因此对于复杂的分布式系统场景指导意义不足。
CAP 定理在分布式系统领域大名鼎鼎,以至于被很多人视为了真理。然而,CAP 定理真的正确吗?
网络分区是一种故障,不管喜欢还是不喜欢,它都可能发生,所以无法选择或逃避分区的问题。在网络正常的时候,系统可以同时保证一致性(线性化)和可用性。而一旦发生了网络故障,必须要么选择一致性,要么选择可用性。因此,对 CAP 更准确的理解应该是:当发生网络分区(P)的情况下,可用性(A)和一致性(C)二者只能选其一。
CAP 定理所描述的模型实际上局限性很大,它只考虑了一种一致性模型和一种故障(网络分区故障),而没有考虑网络延迟、节点失效等情况。因此,它对于指导一个具体的分布式系统设计来说,没有太大的实际价值。
值得一提的是,在 CAP 定理提出十二年之后,其提出者也发表了一篇文章 CAP Twelve Years Later: How the “Rules” Have Changed,来阐述 CAP 定理的局限性。
【中等】什么是 BASE 定理?⭐⭐
要点
BASE 是 基本可用(Basically Available)、软状态(Soft State) 和 最终一致性(Eventually Consistent) 三个短语的缩写。
BASE 核心思想是:要求最终一致性,通过牺牲强一致性来达到可用性。
BASE 是 基本可用(Basically Available)、软状态(Soft State) 和 最终一致性(Eventually Consistent) 三个短语的缩写。BASE 定理是对 CAP 定理中可用性(A)和一致性(C)权衡的结果。
BASE 定理的核心思想是:即使无法做到强一致性,但每个应用都可以根据自身业务特点,采用适当的方式来使系统达到最终一致性。
ACID 要求强一致性,通常运用在传统的数据库系统上。而 BASE 要求最终一致性,通过牺牲强一致性来达到可用性,通常运用在大型分布式系统中。

在实际的分布式场景中,不同业务单元和组件对一致性的要求是不同的,因此 ACID 和 BASE 往往会结合在一起使用。
共识
扩展
- Paxos
- Raft
- ZAB
【困难】Paxos 的工作原理是什么?⭐⭐
Paxos 算法要点
Paxos 是一种基于消息传递且具有容错性的共识性(consensus)算法。Paxos 算法运行在允许宕机故障的异步系统中,不要求可靠的消息传递,可容忍消息丢失、延迟、乱序以及重复。Paxos 利用多数派 (Majority) 机制保证了一定的容错能力,即 N 个节点的系统最多允许 N / 2 - 1 个节点同时出现故障。
Paxos 的核心思想是【两阶段提交】和【多数派决议】。
Paxos 算法包含 2 个部分:
- Basic Paxos 算法
- 要点:多节点之间如何就某个值达成共识。
- 实现:通过二阶段提交的方式来达成共识。
- Multi Paxos 思想
- 要点:执行多个 Basic Paxos,就一系列值达成共识。
Basic Paxos 算法
Basic Paxos 是通过二阶段提交的方式来达成共识的。
Paxos 将分布式系统中的节点分 Proposer、Acceptor、Learner 三种角色。
- 提议者(Proposer):发出提案(Proposal),用于投票表决。Proposal 信息包括提案编号 (Proposal ID) 和提议的值 (Value)。在绝大多数场景中,集群中收到客户端请求的节点,才是提议者。这样做的好处是,对业务代码没有入侵性,也就是说,我们不需要在业务代码中实现算法逻辑。
- 接受者(Acceptor):对每个 Proposal 进行投票,若 Proposal 获得多数 Acceptor 的接受,则称该 Proposal 被批准。一般来说,集群中的所有节点都在扮演接受者的角色,参与共识协商,并接受和存储数据。
- 学习者(Learner):不参与接受,从 Proposers/Acceptors 学习、记录最新达成共识的提案(Value)。一般来说,学习者是数据备份节点,比如主从架构中的从节点,被动地接受数据,容灾备份。
Paxos 算法有 3 个阶段,其中,前 2 个阶段负责协商并达成共识:
- 准备(Prepare)阶段:Proposer 向 Acceptors 发出 Prepare 请求,Acceptors 针对收到的 Prepare 请求进行 Promise 承诺。
- 接受(Accept)阶段:Proposer 收到多数 Acceptors 承诺的 Promise 后,向 Acceptors 发出 Propose 请求,Acceptors 针对收到的 Propose 请求进行 Accept 处理。
- 学习(Learn)阶段:Proposer 在收到多数 Acceptors 的 Accept 之后,标志着本次 Accept 成功,决议形成,将形成的决议发送给所有 Learners。
Multi Paxos 思想
Basic Paxos 有以下问题,导致它不能应用于实际:
- Basic Paxos 算法只能对一个值形成决议。
- Basic Paxos 算法会消耗大量网络带宽。Basic Paxos 中,决议的形成至少需要两次网络通信,在高并发情况下可能需要更多的网络通信,极端情况下甚至可能形成活锁。如果想连续确定多个值,Basic Paxos 搞不定了。
Multi Paxos 基于 Basic Paxos 做了两点改进:
- 针对每一个要确定的值,运行一次 Paxos 算法实例(Instance),形成决议。每一个 Paxos 实例使用唯一的 Instance ID 标识。
- 在所有 Proposer 中选举一个 Leader,由 Leader 唯一地提交 Proposal 给 Acceptor 进行表决。这样没有 Proposer 竞争,解决了活锁问题。在系统中仅有一个 Leader 进行 Value 提交的情况下,Prepare 阶段就可以跳过,从而将两阶段变为一阶段,提高效率。
【困难】Raft 的工作原理是什么?⭐⭐⭐
Raft 算法要点
Raft 是一种管理日志复制的分布式共识性算法。
从本质上说,Raft 算法是通过集中制(一切以领导者为准),实现一系列值的共识和各节点日志的一致。
Raft 出现之前,Paxos 一直是分布式共识性算法的标准。Paxos 难以理解,更难以实现。Raft 的设计目标是简化 Paxos,使得算法既容易理解,也容易实现。
Raft 将一致性问题分解成了三个子问题:
- 选举 Leader:
- Leader 心跳:Leader 定时向 Follower 发心跳以续活;Follower 超时未收到心跳,视其为下线。
- 多数派决议
- Follower 判断 Leader 下线后,发起 Leader 选举,成为 Candidate。
- 得到大多数选票的 Candidate 当选为 Leader。
- 随机超时时间:
- 每个 Follower 都设置一个随机的竞选超时时间,该时间范围内,未收到 Leader 的心跳,就视为当前 Term 无 Leader,再次 Leader 选举。
- 之所以是随机时间,是为了避免重复出现相同投票结果,导致始终选不出 Leader 的情况(一种典型的活锁)。
- 日志复制:Leader 负责处理所有客户端读写请求;Follower 只负责同步 Leader 的日志,并更新本地的日志状态机(Offset)
- 安全性
- 选举限制:拥有最新的已提交的日志条目的 Follower 才有资格成为 Leader。
- 提交旧任期的日志条目:Raft 永远不会通过计算副本数目的方式去提交一个之前 Term 内的日志条目。
- 日志压缩:Raft 采用对整个系统进行快照来解决,快照之前的日志都可以丢弃。以此,避免日志无限膨胀,导致故障恢复过久。
选主
(1)服务器角色
在 Raft 中,任何时刻,每个服务器都处于这三个角色之一 :
Leader:领导者,通常一个系统中是一主(Leader)多从(Follower)。Leader 负责处理所有的客户端请求。Follower:跟随者,不会发送任何请求,只是简单的 响应来自 Leader 或者 Candidate 的请求。Candidate:参选者,选举新 Leader 时的临时角色。

(2)任期

Raft 把时间分割成任意长度的 任期(Term),任期用连续的整数标记。每一段任期从一次选举开始。Raft 保证了在一个给定的任期内,最多只有一个领导者。
任期在 Raft 算法中充当逻辑时钟的作用,使得服务器节点可以查明一些过期的信息(比如过期的 Leader)。每个服务器节点都会存储一个当前任期号,这一编号在整个时期内单调的增长。当服务器之间通信的时候会交换当前任期号。
(3)选举 Leader 流程
领导者心跳消息:Raft 使用一种心跳机制来触发 Leader 选举。Leader 需要周期性的向所有 Follower 发送心跳消息,以此维持 Leader 身份。
随机的竞选超时时间:每个 Follower 都设置了一个随机的竞选超时时间,一般为 150ms ~ 300ms,如果在竞选超时时间内没有收到 Leader 的心跳消息,就会认为当前 Term 没有可用的 Leader,并发起选举来选出新的 Leader。开始一次选举过程,Follower 先要增加自己的当前 Term 号,并转换为 Candidate。
Candidate 会并行的向集群中的所有服务器节点发送投票请求(RequestVote RPC),它会保持当前状态直到以下三件事情之一发生:
- 自己成为 Leader
- 其他的服务器成为 Leader
- 没有任何服务器成为 Leader
Raft 算法通过:领导者心跳消息、随机选举超时时间、得到大多数选票才通过原则、任期最新者优先、先来先服务等投票原则,保证了一个任期只有一位领导,也极大地减少了选举失败的情况。
日志复制

- Leader 负责处理所有客户端的请求。
- Leader 把请求作为日志条目加入到它的日志中,然后并行的向其他服务器发送
AppendEntries RPC请求,要求 Follower 复制日志条目。 - Follower 复制成功后,返回确认消息。
- 当这个日志条目被半数以上的服务器复制后,Leader 提交这个日志条目到它的复制状态机,并向客户端返回执行结果。
安全性
- 选举限制:拥有最新的已提交的日志条目的 Follower 才有资格成为 Leader。
- 提交旧任期的日志条目:Raft 永远不会通过计算副本数目的方式去提交一个之前 Term 内的日志条目。
- 日志压缩:Raft 采用对整个系统进行快照来解决,快照之前的日志都可以丢弃。以此,避免日志无限膨胀,导致故障恢复过久。
【困难】ZAB 的工作原理是什么?⭐
ZAB 协议要点
ZAB 协议是 Zookeeper 专门设计的一种支持故障恢复的原子广播协议。ZAB 协议是 ZooKeeper 的数据一致性和高可用解决方案。
ZAB 协议定义了两个可以无限循环的流程:
选主:用于故障恢复,从而保证高可用。采用多数派决议选举 Leader。原子广播:用于主从同步,从而保证数据一致性。- Leader 负责所有读写;每次更新事务会有一个唯一标识(ZXID)
- Leader 会广播同步数据给 Follower,当半数以上 Follower 更新成功,Leader 才会提交数据。
选主
ZooKeeper 集群采用一主(称为 Leader)多从(称为 Follower)模式,主从节点通过副本机制保证数据一致。
- 如果 Follower 节点挂了:ZooKeeper 集群中的每个节点都会单独在内存中维护自身的状态,并且各节点之间都保持着通讯,只要集群中有半数机器能够正常工作,那么整个集群就可以正常提供服务。
- 如果 Leader 节点挂了:如果 Leader 节点挂了,系统就不能正常工作了。此时,需要通过 ZAB 协议的选举 Leader 机制来进行故障恢复。
ZAB 协议的选举 Leader 机制简单来说,就是:基于过半选举机制产生新的 Leader,之后其他机器将从新的 Leader 上同步状态,当有过半机器完成状态同步后,就退出选举 Leader 模式,进入原子广播模式。
原子广播
ZooKeeper 通过副本机制来实现高可用。
那么,ZooKeeper 是如何实现副本机制的呢?答案是:ZAB 协议的原子广播。

ZAB 协议的原子广播要求:
所有的写请求都会被转发给 Leader,Leader 会以原子广播的方式通知 Follow。当半数以上的 Follow 已经更新状态持久化后,Leader 才会提交这个更新,然后客户端才会收到一个更新成功的响应。这有些类似数据库中的两阶段提交协议。
在整个消息的广播过程中,Leader 服务器会每个事务请求生成对应的 Proposal,并为其分配一个全局唯一的递增的事务 ID(ZXID),之后再对其进行广播。
【困难】Paxos、Raft、ZAB 有什么区别与联系?⭐⭐
要点
- 三者都是非拜占庭容错的共识算法(只容忍崩溃故障,不容忍恶意节点),都基于**多数派(Quorum)**机制,都要求
2f + 1个节点容忍f个故障。 - Paxos 是理论基石,Raft 是工程友好的 Paxos 简化,ZAB 是 ZooKeeper 专用的主备原子广播协议。
- Raft 和 ZAB 都采用"强 Leader"模型,而 Paxos(Basic Paxos)允许多 Proposer 并发。
| 对比维度 | Paxos | Raft | ZAB |
|---|---|---|---|
| 提出者 / 年份 | Lamport / 1998 | Ongaro & Ousterhout / 2014 | Yahoo / 2008 |
| 设计目标 | 理论完备的共识算法 | 易理解、易实现的共识算法 | ZooKeeper 专用的原子广播协议 |
| 角色模型 | Proposer / Acceptor / Learner | Leader / Follower / Candidate | Leader / Follower / Observer |
| Leader 机制 | Basic Paxos 无固定 Leader;Multi-Paxos 选 Leader | 强 Leader,所有请求经 Leader | 强 Leader,所有写经 Leader |
| 日志复制 | Multi-Paxos 需自行实现日志机制 | 明确定义日志复制流程 | 基于 ZXID 的原子广播 |
| 选举 | 多数派投票,可能活锁 | 随机超时选举,避免活锁 | 基于 ZXID(epoch + counter)选举 |
| 顺序保证 | 不保证日志顺序(需上层封装) | 保证日志顺序 | 保证全局顺序(ZXID 单调递增) |
| 日志编号 | Instance ID | (Term, Index) | ZXID = (epoch, counter) |
| 典型实现 | Chubby(Google)、Spanner | etcd、Consul、TiKV | ZooKeeper |
| 复杂度 | 难理解、难实现 | 易理解、易实现 | 中等,专用于主备场景 |
| 适用场景 | 通用共识、理论参考 | 通用共识、配置管理 | 主备数据同步、配置管理 |
联系:
- 三者都解决分布式共识问题,核心思想都是多数派决议。
- Raft 和 ZAB 都可视为 Multi-Paxos 的特化和工程优化,都引入了强 Leader来简化协议、避免活锁。
- 三者都只容忍崩溃故障(Crash Fault),不容忍拜占庭故障;要容忍拜占庭故障需使用 PBFT、PoW 等 BFT 算法。
拜占庭容错
扩展
- The Byzantine Generals Problem - Lamport 经典论文,提出拜占庭将军问题。
- Practical Byzantine Fault Tolerance - Castro & Liskov 提出 PBFT 算法,将 BFT 带入实用。
- Bitcoin: A Peer-to-Peer Electronic Cash System - 中本聪论文,提出 PoW 解决开放网络的拜占庭共识。
- Bitcoin and Cryptocurrency Technologies - 普林斯顿大学比特币教材,系统讲解区块链共识。
【困难】什么是拜占庭将军问题?它与普通共识问题有什么区别?⭐
要点
- 拜占庭将军问题由 Lamport 提出,研究在**存在恶意节点(可能发送错误、伪造信息)**的情况下如何达成共识。
- 与普通共识(Crash Fault)的区别:普通共识假设节点只会崩溃,不会作恶;拜占庭共识假设节点可能任意作恶(伪造、篡改、选择性转发)。
- 容错能力:非拜占庭(崩溃)容错需要
2f + 1个节点容忍f个故障;拜占庭容错需要3f + 1个节点容忍f个故障。
拜占庭将军问题借用一个故事来阐述:一群拜占庭将军各领一支军队围困一座城市,他们只能通过信使互相通信来协商"进攻"或"撤退"。将军中可能有叛徒,叛徒可能向不同将军发送不同的命令,或伪造其他将军的命令。问题是:忠诚的将军如何在这种条件下达成一致的行动策略?
映射到分布式系统:将军 = 节点,信使 = 通信网络,叛徒 = 故障或恶意节点。
与普通共识的区别:
| 对比维度 | 普通共识(Crash Fault) | 拜占庭共识(Byzantine Fault) |
|---|---|---|
| 故障类型 | 节点只会崩溃或停止响应 | 节点可能发送任意错误信息、伪造数据、选择性转发 |
| 容错算法 | Paxos、Raft、ZAB | PBFT、PoW、PoS、HotStuff |
| 节点要求 | 2f + 1 个节点容忍 f 个故障 | 3f + 1 个节点容忍 f 个故障 |
| 通信复杂度 | 较低 | 较高(PBFT 为 O(n²)) |
| 适用场景 | 分布式数据库、协调服务(可信内网) | 区块链、跨机构金融系统(不可信网络) |
【困难】PBFT 算法的工作原理是什么?⭐
要点
- PBFT(Practical Byzantine Fault Tolerance,实用拜占庭容错) 由 Castro 和 Liskov 在 1999 年提出,是第一个实用的拜占庭容错算法。
- PBFT 需要
3f + 1个节点容忍f个拜占庭节点,通过**三阶段协议(Pre-Prepare、Prepare、Commit)**达成共识。 - PBFT 适用于许可链(节点身份已知)场景,如 Hyperledger Fabric。
PBFT 的核心是通过三阶段消息交换确保所有诚实节点以相同顺序执行相同请求,即使存在至多 f 个拜占庭节点。
角色:
- Primary(主节点):接收客户端请求,排序并广播给副本节点。Primary 可被替换(View Change)。
- Replica(副本节点):参与共识,所有节点(含 Primary)都是 Replica。
三阶段协议:
- Pre-Prepare(预准备):Primary 收到请求后,分配序号
n,向所有副本发送<PRE-PREPARE, view, n, request>。副本验证后进入 Prepare 阶段。 - Prepare(准备):每个副本向所有节点广播
<PREPARE, view, n, digest>。当收到2f个 Prepare 消息(含自己共2f + 1)时,说明该请求已"准备好",进入 Commit 阶段。 - Commit(提交):每个副本向所有节点广播
<COMMIT, view, n, digest>。当收到2f个 Commit 消息(含自己共2f + 1)时,执行请求并回复客户端。
为什么需要 3f + 1?:在 3f + 1 个节点中,至多 f 个拜占庭节点,至少 2f + 1 个诚实节点。需要从 3f + 1 个节点收集 2f + 1 个消息,其中诚实节点的 2f + 1 条占多数,可覆盖拜占庭节点的 f 条错误信息。
View Change(视图更换):当副本节点检测到 Primary 异常(超时未响应)时,发起 View Change,选举新的 Primary,保证系统在 Primary 故障或作恶时仍能继续运行。
PBFT 的局限:通信复杂度为 O(n²),节点数超过 100 时性能急剧下降,因此不适用于大规模公链。
【中等】PoW 和 PoS 如何解决拜占庭共识问题?⭐
要点
- PoW(Proof of Work,工作量证明):通过算力竞争记账权,攻击者需控制 51% 算力才能篡改,用于 Bitcoin。
- PoS(Proof of Stake,权益证明):通过质押代币竞争记账权,攻击者需控制 1/3 质押代币才能破坏共识,用于 Ethereum 2.0。
- 两者都通过经济成本提高 Sybil 攻击门槛,使拜占庭节点难以超过阈值。
开放网络(公链)中,节点身份未知,攻击者可低成本创建大量虚假节点(Sybil 攻击),使拜占庭节点比例超过 1/3,破坏 PBFT 等算法。PoW 和 PoS 通过经济成本解决这一问题:
| 对比维度 | PoW(工作量证明) | PoS(权益证明) |
|---|---|---|
| 核心思想 | 通过计算哈希难题竞争记账权,算力越强越可能获胜 | 通过质押代币竞争记账权,质押越多越可能获胜 |
| 记账权获取 | 解出 SHA-256(prefix=0...) 的 nonce | 随机按质押比例选取(含随机数算法) |
| 攻击成本 | 控制 51% 算力(需大量硬件和电力) | 控制 1/3 质押代币(需大量资金) |
| 作恶惩罚 | 产出的无效区块被网络拒绝,算力白费 | 质押代币被罚没(Slashing) |
| 能源消耗 | 极高 | 极低(无需大量计算) |
| 出块速度 | 慢(Bitcoin 约 10 分钟/块) | 快(Ethereum 约 12 秒/块) |
| 典型应用 | Bitcoin | Ethereum 2.0、Cardano |
PoW 解决拜占庭问题的原理:
- 所有节点竞争解哈希难题,最先解出的节点获得记账权(相当于成为"司令")。
- 其他节点验证答案正确性(极快),然后在其上继续延展。
- 最长链规则:节点总是选择最长的有效链。篡改历史需重新计算所有后续区块,成本极高且追不上主链增长。
PoS 解决拜占庭问题的原理:
- 验证者质押代币获得记账资格。
- 按质押比例随机选取出块者。
- 若出块者作恶(双签、无效块),其质押代币被罚没(Slashing),经济上不可行。
同步
【困难】Gossip 的工作原理是什么?⭐⭐
扩展
要点
Gossip 是一种用于分布式系统节点间信息交换的协议。
Gossip 的设计思想基于去中心化和最终一致性。
Gossip 协议的工作原理:
- 周期性、成对通信:每个节点每隔一段时间就随机选择集群中的另一个节点(这个节点称为“邻居”)。
- 交换信息:两个节点连接后,会互相交换自己拥有的信息(例如,其他节点的状态、存储的数据等)。
- 感染式传播:接收到新信息的节点,会将这些新信息融入到自己的信息库中。在下一次周期中,它又会成为传染源,将(包含新信息的)所有信息再次传播给其他随机节点。
- 最终一致性:不需要中央协调,经过一段时间后,通过这种“八卦”式的传播,集群中的所有节点最终都会拥有完全相同的信息。
Gossip 也叫 Epidemic Protocol (流行病协议),这个协议基于最终一致性以及去中心化设计思想。主要用于分布式节点之间进行信息交换和数据同步,这种场景的一个最大特点就是组成的网络的节点都是对等节点,是非结构化网络(去中心化)。
Gossip 过程是由种子节点发起,当一个种子节点有状态需要更新到网络中的其他节点时,它会随机的选择周围几个节点散播消息,收到消息的节点也会重复该过程,直至最终网络中所有的节点都收到了消息。这个过程可能需要一定的时间,由于不能保证某个时刻所有节点都收到消息,但是理论上最终所有节点都会收到消息,因此它是一个最终一致性协议。
Gossip 过程是异步的,也就是说发消息的节点不会关注对方是否收到,即不等待响应;不管对方有没有收到,它都会每隔 1 秒向周围节点发消息。异步是它的优点,而消息冗余则是它的缺点。
Goosip 协议的信息传播和扩散通常需要由种子节点发起。整个传播过程可能需要一定的时间,由于不能保证某个时刻所有节点都收到消息,但是理论上最终所有节点都会收到消息,因此它是一个最终一致性协议。

Gossip 有两种类型:
- Anti-Entropy(反熵):以固定的概率传播所有的数据。反熵时通讯成本会很高,可以通过引入校验和等机制,降低需要对比的数据量和通讯消息等。反熵不适合动态变化或节点数比较多的分布式环境。
- Rumor-Mongering(谣言传播):仅传播新到达的数据。谣言传播模型指的是当一个节点有了新数据后,这个节点变成活跃状态,并周期性地联系其他节点向其发送新数据,直到所有的节点都存储了该新数据。在谣言传播模型下,消息可以发送得更频繁,因为消息只包含最新 update,体积更小。而且,一个谣言消息在某个时间点之后会被标记为 removed,并且不再被传播,因此,谣言传播模型下,系统有一定的概率会不一致。而由于,谣言传播模型下某个时间点之后消息不再传播,因此消息是有限的,系统开销小。
【中等】Gossip 协议的收敛性如何?有哪些应用场景?⭐
要点
- Gossip 的消息传播速度为 O(log N),即经过
O(log N)轮传播后,集群中所有节点都能收到消息(N 为节点数)。 - Gossip 适用于对一致性要求不严格、节点规模大、容错性要求高的场景。
收敛性分析:
假设每个节点每轮联系 fanout 个邻居(典型值为 3),经过 k 轮后,理论上最多有 fanout^k 个节点收到消息。要覆盖 N 个节点,需 fanout^k ≥ N,即 k ≥ log_fanout(N)。因此,Gossip 的传播复杂度为 O(log N),收敛速度很快。例如,N=10000、fanout=3 时,约需 8-9 轮即可覆盖全网。
应用场景:
| 应用场景 | 说明 | 代表系统 |
|---|---|---|
| 集群成员管理 | 节点加入、退出、故障检测的信息传播 | Cassandra、Consul、Redis Cluster |
| 状态同步 | 节点间数据副本的最终一致同步 | Cassandra、Riak、DynamoDB |
| 服务发现 | 服务实例上下线信息的传播 | Consul、Serf |
| 广播协议 | 集群范围内的元数据变更通知 | Redis Cluster 的 PONG 消息 |
Gossip 的优缺点:
- 优点:去中心化(无单点)、可扩展(O(log N) 收敛)、容错(节点故障不影响传播)、实现简单。
- 缺点:消息冗余(同一消息被多次传播)、收敛延迟(非即时一致)、不适合强一致场景。
分布式系统设计模式
扩展
- Dynamo: Amazon’s Highly Available Key-value Store - Amazon Dynamo 论文,首次系统阐述了 Quorum、Read Repair、Hinted Handoff、向量时钟等无主复制技术。
- Designing Data-Intensive Applications 第 5 章:复制 - Martin Kleppmann 系统讲解了领导者复制、多主复制、无主复制及冲突处理。
【中等】什么是 Quorum 机制?⭐⭐
要点
- Quorum(法定人数)机制是无主复制(及部分主从复制)中平衡一致性与可用性的核心机制。
- 核心公式:
W + R > N时保证读取能读到最新写入的数据,其中 N 为副本数,W 为写成功副本数,R 为读成功副本数。 - 常见配置:N=3, W=2, R=2(强一致);N=3, W=2, R=1(高可用读);N=3, W=1, R=1(高可用,最终一致)。
在多副本系统中,写入时不需要所有副本都成功,读取时也不需要查询所有副本。Quorum 机制通过约束读写副本数来保证一致性:
- N:每个数据的副本总数。
- W:写操作要求成功响应的副本数(写 Quorum)。
- R:读操作要求成功响应的副本数(读 Quorum)。
关键约束:当 W + R > N 时,读 Quorum 和写 Quorum 必有交集,因此读操作一定能读到至少一个包含最新写入的副本,保证强一致性。
| 配置 | W | R | 一致性 | 可用性 | 适用场景 |
|---|---|---|---|---|---|
| W=N, R=1 | N | 1 | 强一致 | 低(任一副本故障即写失败) | 读多写少,强一致 |
| W=1, R=N | 1 | N | 强一致 | 低(任一副本故障即读失败) | 写多读少,强一致 |
| W+R>N | — | — | 强一致 | 中 | 通用强一致 |
| W+R≤N | — | — | 最终一致 | 高 | 高可用,容忍部分副本不一致 |
Sloppy Quorum(松散 Quorum):当部分副本不可用时,写入会临时落到其他健康节点(见 Hinted Handoff),保证可用性但牺牲临时一致性。
【中等】什么是 Read Repair(读修复)?⭐
要点
- Read Repair(读修复) 是无主复制数据库在读操作时检测并修复副本数据不一致的机制。
- 当读取多个副本发现数据版本不一致时,客户端(或协调者)将最新数据回写到过期的副本。
工作流程:
- 客户端向 R 个副本发起读请求。
- 收到 R 个响应后,比较数据的版本(如时间戳、向量时钟)。
- 若发现某些副本的数据过期,将最新数据回写到这些过期副本。
适用场景:
- 适合读频繁的数据——每次读取都可能触发修复,使副本逐渐趋于一致。
- 不适合冷数据(很少被读取的数据)——因为读修复依赖读取触发,冷数据可能长期不一致,需配合 Anti-Entropy(反熵修复,定期全量对比)。
代表系统:Cassandra、Riak、DynamoDB 的无主复制模式均使用 Read Repair。
【中等】什么是 Hinted Handoff(提示移交)?⭐
要点
- Hinted Handoff(提示移交) 是无主复制数据库在目标副本暂时不可用时,将写入临时转交给其他健康节点暂存的机制。
- 当原副本恢复后,暂存节点将数据"移交"回原副本,保证数据不丢失。
工作流程:
- 协调者需要将写入复制到节点 A、B、C(假设 N=3, W=2)。
- 节点 C 因故障不可用,协调者将本应写入 C 的数据连同"hint"(提示:这份数据属于 C)写入另一个健康节点 D。
- 节点 D 标记这份数据为"hinted"(不属于自己,暂存)。
- 节点 C 恢复后,节点 D 检测到 C 上线,将暂存数据移交回 C,然后删除本地暂存。
作用:
- 保证可用性:即使部分副本故障,写入仍可成功(Sloppy Quorum),不会因副本不足而拒绝写入。
- 保证数据不丢失:故障副本恢复后能补齐期间错过的写入。
与 Quorum 的关系:Hinted Handoff 常与 Sloppy Quorum 配合——Sloppy Quorum 允许写入临时落到非目标节点,Hinted Handoff 负责后续移交。
代表系统:Cassandra、Riak、DynamoDB。