《Designing Data-Intensive Applications》(Martin Kleppmann, O’Reilly 2017) 第二部分「分布式数据」读书笔记。用我自己的话重新梳理每章的核心概念、关键权衡和例子。第一部分讲的是单机上的数据系统;第二部分开始,数据被分散到多台机器上,于是所有问题都变了味——网络会断、时钟不准、节点会毫无征兆地卡住。这一部分就是教你在这种”没有一样东西靠得住”的世界里,如何依然保证数据正确。

本篇目录

  1. 第 5 章 · 复制 (Replication)
  2. 第 6 章 · 分区 (Partitioning)
  3. 第 7 章 · 事务 (Transactions)
  4. 第 8 章 · 分布式系统的麻烦
  5. 第 9 章 · 一致性与共识

为什么要有第二部分:把数据放到多台机器上,通常出于三个动机——可扩展性(数据量/负载超过单机)、容错与高可用(一台挂了另一台顶上)、降低延迟(把数据放到离用户近的地方)。而实现”多机”有两种正交的手段:复制 (replication)(同一份数据存多份副本)和分区 (partitioning)(把数据切成块分散存)。这两者通常一起用。第 5、6 章分别讲它们,第 7~9 章处理由此产生的一致性难题。


第 5 章 · 复制 (Replication)

复制是什么、为什么难

复制 = 在多个节点上保存同一份数据的多个副本 (replica)。 目的:离用户更近(降延迟)、部分节点挂了系统还能用(高可用)、多副本分摊读请求(提升读吞吐)。

如果数据永远不变,复制就太简单了——拷贝一份就完事。真正的难点全在于”数据会变”:一个写操作怎么可靠地传播到所有副本?主流有三类方案。

一、单主复制 (Single-Leader)

也叫主从复制 (master-slave)。这是最常见的方案(MySQL、PostgreSQL 默认):

  1. 指定一个副本为主库 (leader)。所有写必须发给它。
  2. 主库把数据变更写进本地,然后把复制日志 (replication log) 发给所有从库 (follower)。
  3. 从库按相同顺序应用这些变更。
  4. 读可以从任意副本(主或从)进行。

同步 vs 异步:一个核心权衡

主库把变更发给从库后,要不要等从库确认?

  • 同步 (synchronous):等从库确认才算写成功。好处是从库一定有最新数据;坏处是只要一个从库慢/挂,写就被卡住。
  • 异步 (asynchronous):发出去就算成功,不等确认。写很快、可用性好;但如果主库在数据传出去之前崩溃,这些写就永久丢了。

💡 权衡 tip:全同步不现实(一个从库故障会拖垮整个写入),全异步有丢数据风险。实践中常用半同步 (semi-synchronous)——只要求一个从库同步确认,其余异步。这样既保证”至少有两份最新数据”,又不会被单个慢节点拖死。这是”持久性 vs 性能/可用性”这条永恒权衡的典型折中。

节点故障怎么办

  • 从库挂了:好办。从库重启后,看自己日志里最后处理到哪,向主库要之后的变更补上即可(追赶恢复 catch-up)。
  • 主库挂了:麻烦,这叫故障转移 (failover)。要选一个从库升为新主库,让其他从库和客户端都改认它。这个过程充满陷阱:
    • 丢写:异步复制下,旧主库还没传出去的写会丢。
    • 脑裂 (split brain):两个节点都以为自己是主库,同时接受写 → 数据冲突甚至损坏。
    • 超时设定难:判定主库”死了”的超时设太短会误判(把慢当死),太长则恢复慢。

💡 运维 tip:故障转移看着自动化很美,但它是分布式系统里最容易出大事故的环节之一。很多成熟团队宁愿在某些场景下用手动故障转移——因为自动化的误判和脑裂造成的数据损坏,往往比多停机几分钟更可怕。

复制日志怎么实现(了解即可)

  • 基于语句 (statement-based):把 SQL 语句本身发过去。坑多(NOW()、随机数、自增在不同副本结果不同)。
  • 传输 WAL (write-ahead log shipping):直接传底层的物理日志。紧凑,但和存储引擎版本绑死,难做滚动升级。
  • 逻辑日志 (logical / row-based):描述”哪一行发生了什么变化”,与存储引擎解耦,也方便被外部系统(如数仓、缓存)消费——CDC (变更数据捕获) 就基于此。
  • 基于触发器:灵活但开销大、易出错。

二、复制延迟带来的一致性问题

异步复制下,从库总是稍微落后于主库,这段差距叫复制延迟 (replication lag)。绝大多数时候只有几毫秒,但负载高时可能到几秒甚至几分钟。这就是所谓的最终一致性 (eventual consistency)——只要停止写入,副本”最终”会追上。问题是”最终”到底多久没保证,中间用户会看到各种诡异现象。作者列了三个必须防住的:

  • 读己之写 (read-your-own-writes):你刚发了条评论(写到主库),刷新页面(从一个还没同步的从库读)却看不到自己的评论——体验极差。解决:让用户读自己刚改过的东西时走主库,或读足够新的副本。
  • 单调读 (monotonic reads):你连刷两次,第一次从较新的从库读到了评论,第二次从较旧的从库读,评论又”消失”了——时光倒流。解决:让同一个用户始终读同一个副本。
  • 一致前缀读 (consistent prefix reads):如果因果相关的写以错误顺序被看到——比如先看到回答”大约十秒”,再看到问题”你能看多远的未来”——因果就乱了。解决:保证有因果关系的写被同一分区按序处理。

💡 认知 tip:这三个”异常”的根子是——“最终一致”这四个字掩盖了中间那段窗口里用户会遇到的真实困惑。开发时如果假装”数据总是最新的”,就会写出这类 bug。第二部分后面讲的”更强一致性”(如线性一致性),本质就是用性能换取消除这些异常。

三、多主复制 (Multi-Leader)

允许多个节点都能接受写。典型用途:跨多数据中心(每个数据中心一个主库,本地写本地)、离线客户端(手机日历,离线也能改,联网再同步)、协同编辑(多人同时编辑文档)。

好处是写延迟低、离线可用。但引入了单主复制没有的大麻烦——写冲突 (write conflict):两个数据中心可能同时改了同一条记录的同一个字段。解决冲突的办法:

  • LWW (Last Write Wins):给每个写打时间戳,保留”最后”的。简单但会丢数据(且依赖不可靠的时钟)。
  • 让写不冲突:比如把某条记录的写都路由到同一个主库。
  • 合并 / CRDT:设计能自动合并的数据结构(如集合的并集),或让应用层介入解决。

四、无主复制 (Leaderless)

没有主库,客户端直接把写发给多个副本(或经一个协调节点转发)。源自 Amazon Dynamo,代表:Cassandra、Riak、Voldemort。

核心机制是法定人数 (quorum):设共 n 个副本,每次写要成功写入 w 个,每次读要从 r 个副本读。只要满足 w + r > n,读和写的副本集合必然有交集,从而保证读到至少一个最新值。

  • 常见配置 n=3, w=2, r=2。
  • 怎么修复落后的副本?读修复 (read repair)(读到旧值时顺手更新它)和反熵 (anti-entropy)(后台进程持续对比、同步差异)。
  • 宽松法定人数 (sloppy quorum) + 提示移交 (hinted handoff):网络分区时,写先临时存到能连上的其他节点,等原节点恢复再转交,提升可用性。
  • 检测并发写:无主/多主都需要判断两个写是否”并发”(没有因果先后)。用版本向量 (version vectors) 记录”谁基于什么版本改的”,识别出并发冲突交给应用处理。这背后是happens-before(先发生)关系——分布式系统里判断因果的基本工具,后面第 9 章还会深入。

💡 架构 tip:单主 / 多主 / 无主,本质是在“一致性的简单程度”和”可用性/写灵活性”之间取不同的点。单主一致性最好推理但主库是瓶颈和单点;无主可用性最高但要应用层直面并发冲突。没有最好的复制方案,只有匹配你对一致性/延迟/可用性要求的方案。


第 6 章 · 分区 (Partitioning)

为什么分区

当数据量或吞吐大到单机存不下、扛不住时,就把数据切成分区 (partition,也叫 shard 分片),分散到多个节点。每条数据属于且仅属于一个分区,每个节点可以存多个分区。目标是可扩展性——把数据和查询负载均摊到多台机器。分区几乎总是和复制一起用:每个分区再复制到多个节点上,兼顾扩展性和容错。

一、键值数据怎么分区

核心目标是均匀分布,避免热点 (hot spot)(负载集中在个别分区,其他机器闲着)。两种主流策略:

  • 按键的范围分区 (by key range):像百科全书按字母分卷。分区内有序,支持高效的范围查询(”查 1 月所有数据”)。缺点:容易产生热点——比如用时间戳当 key,当天所有写全砸在同一个分区。
  • 按键的哈希分区 (by hash of key):对 key 求哈希再分区。分布非常均匀,避开热点。代价是丧失了范围查询能力(相邻的 key 被打散到各处)。

💡 权衡 tip:这里又是一个”鱼和熊掌”——范围分区给你有序和范围扫描但可能热点;哈希分区给你均匀但牺牲范围查询。有些系统(如 Cassandra)用组合键折中:一部分字段做哈希决定分区,另一部分在分区内做范围排序。另外注意:即便哈希分区,也挡不住超级热点 key(如某明星账号被疯狂访问)——这种”名人问题”通常得靠应用层给 key 加随机后缀等手段手动打散。

二、分区与二级索引

主键分区好办,但如果要按”非主键字段”查询(二级索引,如”查所有红色的车”),索引也得分区,有两种做法:

  • 本地索引 / 按文档分区 (local / document-partitioned):每个分区维护自己那部分数据的索引。写只碰一个分区(简单);但读要问遍所有分区再汇总(scatter/gather),读慢。
  • 全局索引 / 按词条分区 (global / term-partitioned):索引本身按被索引的值来分区。读只需查一个分区(快);但一次写可能要更新多个分区的索引(因为一条记录的不同字段落在不同索引分区),写变复杂、常需异步更新。

三、再平衡 (Rebalancing)

集群加/减机器时,分区要在节点间重新分配,这叫再平衡。原则:搬动的数据尽量少、期间还能继续服务。

  • 千万别用 hash mod N:节点数 N 一变,几乎所有 key 的归属都变,得搬几乎全部数据。
  • 固定数量的分区:一开始就建远多于节点数的分区(如 1000 个),加节点时只是把一些整分区搬过去。简单常用。
  • 动态分区:分区太大就分裂、太小就合并(类似 B-tree),适合数据量变化大的场景。
  • 手动 vs 自动:全自动再平衡很方便,但和自动故障检测一起可能”火上浇油”(把慢节点误判为死 → 触发大规模搬迁 → 更慢)。很多系统让再平衡需要人确认。

四、请求路由 (Request Routing)

客户端怎么知道某个 key 该找哪个节点?这是服务发现问题。几种方案:让客户端随便连一个节点由它转发;加一个路由层;或让客户端直接知道分区分布。分区到节点的映射会变,需要一个可靠的地方维护这份”权威信息”——很多系统用 ZooKeeper 这类协调服务来存储集群元数据并通知变更(这又引出第 9 章的共识话题)。


第 7 章 · 事务 (Transactions)

事务解决什么问题

现实中一次操作常涉及多次读写(转账要”A 减钱 + B 加钱”)。中途任何一步可能失败:程序崩溃、网络断、并发操作互相干扰。挨个处理这些错误极其繁琐。事务 (transaction) 就是把一组读写打包成一个逻辑单元,给你一个简单的承诺:

要么全部成功 (commit),要么全部不做 (abort)、可以安全重试——中间状态永远不对外可见。

有了它,应用层就不用操心一大堆部分失败的中间态。事务是一种简化编程模型的抽象,代价是数据库内部要做很多工作,且并非没有性能成本。

ACID:常被误用的四个字母

  • A 原子性 (Atomicity):不是”并发”的意思,而是可中止性——出错时能把已做的部分全部撤销,就像什么都没发生。
  • C 一致性 (Consistency):指数据满足应用定义的不变量(如”账户总额守恒”)。作者点破:这其实主要是应用的责任,数据库只能帮你保证它力所能及的部分。(C 其实是硬凑进 ACID 的。)
  • I 隔离性 (Isolation):并发执行的事务互不干扰,效果如同一个接一个串行执行。这是本章重点。
  • D 持久性 (Durability):一旦提交,数据不会丢(哪怕断电)。

💡 认知 tip:ACID 这个词营销味很重,不同数据库对它的实现差别巨大,尤其是隔离性 I——号称”支持事务”的系统,隔离级别可能弱得超出你想象。别看标签,要看它具体提供哪个隔离级别、防住了哪些并发异常。

弱隔离级别(因为最强的太贵)

理想的隔离是可串行化 (serializable)——效果完全等同于事务一个个串行跑。但它性能开销大,所以现实中数据库提供了一系列较弱的隔离级别,各自防住一部分并发异常:

  • 读已提交 (Read Committed):防两件事——脏读(读到别人还没提交的数据)和脏写(覆盖别人还没提交的写)。最常见的默认级别。
  • 快照隔离 (Snapshot Isolation) / 可重复读:每个事务从一个一致性快照读数据——它看到的是事务开始那一刻的整个数据库状态,不受期间其他事务影响。防住读偏差 (read skew),比如你在转账过程中查两个账户余额,不会看到”钱转出了但还没转入”的中间态。实现靠 MVCC (多版本并发控制):保留数据的多个版本,让读不阻塞写、写不阻塞读。
  • 防丢失更新 (lost update):两个事务同时”读-改-写”同一个值(如计数器 +1),一个覆盖了另一个。解决:原子操作 (UPDATE ... SET x = x + 1)、显式加锁 (FOR UPDATE)、或CAS (compare-and-set)。
  • 写偏差 (write skew) 与幻读 (phantom):更微妙的一类。两个事务各自读了数据、基于读到的结果做决定、再写入,单看都没问题,合起来却破坏了不变量(经典例子:两名医生同时确认”至少还有一名同事在值班”于是都请假了,结果没人值班)。这类问题连快照隔离也防不住,往往需要显式锁或物化冲突。

可串行化:三种实现

如果你需要最强保证,有三条路:

  • 真的串行执行 (Actual Serial Execution):干脆单线程一个个跑事务(如 VoltDB、Redis)。前提是数据能放进内存、事务短小、用存储过程提交。在合适场景下反而简单高效——现代内存足够大让这变得可行。
  • 两阶段锁 (2PL, Two-Phase Locking):悲观锁。读加共享锁、写加排他锁,还需谓词锁/索引范围锁来防幻读。正确但性能差、易死锁——传统关系库的可串行化多用它。
  • 可串行化快照隔离 (SSI, Serializable Snapshot Isolation):乐观并发控制。让事务先在快照上跑,提交时才检查有没有和别的事务发生真正的冲突,有冲突就中止重试。冲突不多时性能远好于 2PL,是较新的、有前途的方案(如 PostgreSQL 的可串行化级别)。

💡 实践 tip:99% 的应用不需要手写并发控制——选对隔离级别就够了。关键是搞清楚你用的数据库默认是哪个级别(很多默认只是”读已提交”或”快照隔离”),以及你的业务逻辑是否存在写偏差这类连快照隔离都防不住的隐患。当你写出”先查一下,再根据结果决定写什么”的代码时,就要警惕了。


第 8 章 · 分布式系统的麻烦

这一章是全书基调最”悲观”的一章,也是理解分布式系统的思想钢印。它的核心信息只有一句:

在分布式系统里,任何东西都可能以你想不到的方式出错,而且你常常无法确定到底出没出错。

单机上,一个操作要么成功要么失败,干脆利落。分布式系统的本质困难是部分失败 (partial failure)——系统的一部分坏了,另一部分好好的,而且这种失败是非确定性的:同样的操作,有时成有时败。

一、不可靠的网络

分布式系统通过异步网络通信。发一个请求出去,可能:请求丢了、请求在排队、远程节点挂了、响应丢了、响应在排队……当你迟迟收不到回复时,你根本无法区分到底是哪种情况。

唯一的工具是超时 (timeout):等一段时间没回复就当它失败。但超时设多久是个两难——太短会把”只是慢”误判为”死了”(可能导致一个操作被执行两次),太长则故障恢复迟缓。而且网络延迟高度多变,没有”正确”的超时值。

二、不可靠的时钟

每台机器有自己的石英钟,会漂移,靠 NTP 同步但同步本身也不可靠。作者区分两种时钟:

  • 墙上时钟 (time-of-day clock):返回真实日期时间,但会跳变(NTP 校正、闰秒),甚至倒退。绝不能用它来测量时间间隔。
  • 单调时钟 (monotonic clock):只保证一直往前走,适合测量”过了多久”(如超时),但它的绝对值没有意义。

💡 警示 tip:最危险的错误之一,是用墙上时钟的时间戳给分布式事件排序(比如前面提到的 LWW 冲突解决)。因为不同机器的时钟有偏差,”时间戳大”不代表”真的更晚发生”——你可能悄无声息地丢掉本该保留的数据。Google Spanner 的解法很有意思:TrueTime 不返回一个确切时间,而是返回一个置信区间 [最早, 最晚],并在提交时主动等待这个区间过去,从而确保顺序正确——用”承认不确定性”来换取正确性。

三、进程暂停

一个更可怕的问题:你的进程可能在任意一行代码之间被暂停任意长的时间而毫不知情——GC 停顿(Stop-The-World)、虚拟机被挂起、笔记本合盖、线程被抢占……醒来后它以为只过了一瞬,实际可能已过去几十秒。

这对”自己以为还是主库”的节点是致命的:它 GC 停顿了 30 秒,期间集群已经选了新主库,它醒来却继续以主库身份写数据 → 数据损坏。

四、真相由多数决定 + 防护令牌

既然单个节点的判断(包括它对”我是不是主库”的判断)都不可信,分布式系统的解法是:真相由多数节点(法定人数)共同决定,不信任任何单个节点。

针对上面”假主库”的问题,有个精巧的工具——防护令牌 (fencing token):每次授予锁/主库身份时发一个递增的编号,节点每次写存储时都带上自己的令牌号;存储系统只接受令牌号更大的写,拒绝更小的。这样即使旧主库从暂停中醒来想写,它那个过期的小令牌也会被拒绝。

此外还有拜占庭故障 (Byzantine faults)——节点不只是崩溃,还会”撒谎”、发送错误或恶意信息(如航天、区块链场景)。防御它非常昂贵,所以大多数数据中心系统假设节点不作恶(非拜占庭),只防崩溃和网络问题。

💡 心智模型 tip:这一章要建立的核心直觉是——在分布式系统里,”我知道现在的真实状态”本身就是一种奢望。 你不知道消息到没到、节点死没死、自己暂停了多久。好的分布式设计不是假装这些问题不存在,而是用多数决、令牌、超时等机制,让系统在充满不确定性的前提下依然能推导出可靠结论。 这正是下一章”共识”要正面解决的。


第 9 章 · 一致性与共识

前一章讲了各种”坏消息”,这一章讲”怎么办”——如何在不可靠的基础上构建可靠的、有强保证的抽象。这是整个第二部分的理论高峰。

一、线性一致性 (Linearizability)

也叫强一致性、原子一致性。直觉定义:

让整个系统看起来只有一份数据,每个操作都在某个瞬间原子地生效。一旦某个客户端读到了新值,之后所有客户端也必须读到新值(或更新的值)——不允许时光倒流。

它是一种新近性保证 (recency guarantee):你读到的一定是最新的。这消除了第 5 章那些复制延迟异常。典型用途:主库选举、分布式锁、唯一性约束——这些场景必须所有人对”当前状态”达成一致。

⚠️ 别把它和可串行化搞混:可串行化是关于事务隔离(多个对象、多个操作的并发效果等价于串行);线性一致性是关于单个对象上读写的新近性和顺序。两个不同维度的概念,虽然常被一起讨论。

代价——CAP 定理:当发生网络分区 (Partition) 时,你只能在一致性 (Consistency, 指线性一致) 和可用性 (Availability) 之间二选一。要线性一致,分区期间少数派那边就得停止服务(不可用);要一直可用,就得放弃线性一致。此外,即使没有分区,线性一致本身也慢(要跨副本协调),所以很多系统主动选择不要它。

二、顺序保证

一致性和”顺序”深度相关。关键概念是因果顺序 (causal order):如果事件 A 导致了事件 B(B 依赖 A),那所有人都必须先看到 A 再看到 B。因果顺序是一种偏序 (partial order)——有因果关系的事件有先后,并发的事件则无所谓顺序。

  • 因果一致性:保证因果顺序,是比线性一致性更弱、但更便宜的保证(分区时仍可用),对很多应用已经足够。
  • Lamport 时间戳:一种给所有事件生成全序编号的巧妙方法,能保证与因果一致(但无法单凭它判断两个事件是否并发)。
  • 全序广播 (Total Order Broadcast):让所有节点以完全相同的顺序收到同一批消息。这是个极强的原语——它被证明和”共识”等价。有了它,就能实现可串行化事务、主库复制等。

三、分布式事务与共识

共识 (consensus) 是分布式系统的圣杯,一句话:让多个节点就某件事达成一致(比如”谁是主库”“这个事务提没提交”)。看似简单,在有节点故障和网络延迟时却极难做对。

先看一个具体问题——原子提交:两阶段提交 (2PC)

一个事务跨多个节点,怎么保证要么所有节点都提交、要么都回滚?2PC 引入一个协调者 (coordinator):

  1. 准备阶段 (prepare):协调者问所有参与者”你们能提交吗?”参与者若回答”能”,就必须保证之后一定能提交(做出承诺,不能反悔)。
  2. 提交阶段 (commit):只要所有人都说”能”,协调者就广播”提交”;否则广播”回滚”。

它的致命弱点:如果协调者在关键时刻崩溃,那些已经承诺”能提交”、正在等指令的参与者就卡死 (blocking) 了——它们既不能自己提交也不能自己回滚,只能一直锁着资源干等。这是 2PC 最被诟病的地方。

共识算法

更健壮的做法是用真正的共识算法——Paxos、Raft、Zab、VSR 等。一个正确的共识算法要满足:所有节点对同一个值达成一致 (uniform agreement)、只对真正被提议过的值达成一致 (validity/integrity)、并且最终一定能得出结论 (termination)——只要多数节点存活。它们能容忍协调者/节点故障,不像 2PC 那样一崩就卡死。

前面说了,共识 ⟺ 全序广播 ⟺ 线性一致的 CAS——它们理论上互相等价,是同一个难题的不同面孔。

💡 实践 tip:你几乎永远不该自己手写共识算法——它们细节繁多、极易出错。正确的做法是用现成的协调服务,如 ZooKeeper、etcd。它们把共识封装成好用的功能:主库选举、分布式锁(配防护令牌)、服务发现、配置管理。第 6 章说的”用 ZooKeeper 存分区元数据”、第 5 章的”故障转移选新主库”,底层都是靠它们的共识实现的。认出”这是个需要共识的问题”,然后把它交给 ZooKeeper/etcd,这才是工程师该有的判断。

不是所有问题都需要共识这么重的武器。很多系统宁愿用更弱的保证(如因果一致)换取更好的性能和可用性。理解每种一致性保证的强弱和代价,在具体场景里做出恰当选择——这就是第二部分教给我们的核心能力。


第二部分小结

章 主题 一句话核心
Ch5 复制 同一份数据存多副本 单主/多主/无主三方案,难点全在”数据会变”;异步复制带来延迟异常(读己之写、单调读、一致前缀)
Ch6 分区 把数据切块分散 范围分区 vs 哈希分区(有序 vs 均匀);二级索引本地 vs 全局;再平衡别用 hash mod N
Ch7 事务 把多个读写打包成全有或全无 ACID 名不副实,重点看隔离级别;从读已提交→快照隔离→可串行化,各防不同的并发异常
Ch8 麻烦 分布式的本质困难 部分失败 + 不确定性:网络、时钟、进程暂停都不可靠,”你无法知道真实状态”
Ch9 一致性与共识 在不可靠之上建可靠 线性一致(强但慢,受 CAP 约束);共识=全序广播,交给 ZooKeeper/etcd

贯穿第二部分的思想:分布式带来了扩展性、容错和低延迟,但代价是一致性变成了一个需要精心权衡的选择。从弱到强——最终一致 → 因果一致 → 线性一致 → 共识,越强的保证越贵。好的系统设计,是清楚每种保证的强弱与成本,然后为每个具体问题选最恰当(而非最强)的那一档。

留给自己的思考题

  1. 我的系统用的是哪种复制方案?有没有踩过”读己之写”“时光倒流”这类复制延迟坑?
  2. 我的数据库默认隔离级别是什么?业务里有没有”先查再改”、可能触发写偏差的逻辑?
  3. 我有没有在代码里(错误地)用时间戳给分布式事件排序?
  4. 我系统里哪些地方其实是”需要共识”的问题(选主、锁、唯一约束)?是不是该交给 ZooKeeper/etcd 而非自己造轮子?