专业书籍精读 · DDIA · 第 6 章
Designing Data-Intensive Applications · Ch 6 · Martin Kleppmann · 2017
上一章说的是把同一份数据抄好几份(复制)。这一章解决的是另一个问题:数据本身太大了,大到一台机器根本装不下、也扛不住。淘宝上亿条商品、微信几十亿条消息——没有哪台机器的硬盘和 CPU 顶得住。办法很直白:把一大坨数据切成很多小块,每块交给一台机器管。这个「切开分放」的动作,就叫「分区」(也叫分片)。
想象一座巨大的图书馆,书多到一个书架放不下。你有两种分法:一种是按书名首字母——A 到 F 放 1 号柜、G 到 M 放 2 号柜……另一种是给每本书算个「暗号」随机打散——同一个书名永远算出同一个暗号,按暗号丢进对应的柜子。前者找「所有 D 开头的书」很方便(都在一起);后者每个柜子塞得一样满、谁也不会特别挤。数据库切数据,用的就是这两招。
不切开,问题是死结:一台机器存不下那么多数据、也算不过来那么多请求。而一旦切开,新麻烦立刻冒出来——切得不匀会「堵车」。要是所有热门请求恰好都落在同一个柜子(比如某个明星发了条动态,几千万人同时来看),那台机器照样被压垮,其他机器却闲着。这种「一个格子被挤爆、其他格子空荡荡」的情况,是这一章反复要对付的敌人。
① 怎么切才匀? 按首字母切,简单、还方便按顺序找一段;但容易堆热点(大家都查最新的)。按暗号打散,最匀;但没法「按顺序取一段」了(相邻的书被丢到了天南海北)。各有各的甜头和苦头。
② 加了新柜子,书怎么搬? 图书馆扩建、加了几个新书架,总不能把全馆的书推倒重排——那得关门好几天。聪明的做法是只从每个旧柜挪一小部分书到新柜,动得越少越好。
③ 要找一本书,怎么知道它在哪个柜? 得有个「前台」或一张「目录」告诉你:你要的这条数据,归哪台机器管。柜子搬动后,这张目录还得及时更新。
不是,但它俩是黄金搭档。复制是「同一块内容,抄给几台机器各存一份」(防坏、就近、分摊读);分区是「把整个大数据集切成不同的块,各机器只管其中一块」(为了装得下、算得动)。真实系统几乎两个一起上:先切成很多块,每块再各抄三份放到不同机器上——既装得下,又坏得起。
分区 = 把一个大到单机装不下的数据集切成很多小块,每块交一台机器管,为的是能扩展。切法有两种:按顺序切(方便按段找、但易堵)、按暗号打散切(最匀、但没法按段找)。真正的难点是别让某一块被挤爆(热点)、加机器时少搬数据、以及找得到每条数据在哪。它和上一章的复制通常配套使用。
想进到具体机制、二级索引与再平衡策略、示意图? → 切到精读版
分区(partitioning,也叫分片 / sharding)就是把一个大数据集切成若干互不重叠的小块,每块(一个分区)放在一台节点上,让存储和查询负载能横向摊到多台机器——这是数据「扩展」的核心手段。本章的全部张力集中在四件事上:怎么切才均匀(键范围 vs 哈希)、怎么防某一分区被压垮(热点 / 偏斜)、二级索引怎么办(本地索引 vs 全局索引)、以及加减节点时怎么把分区重新摊匀又少搬数据(再平衡)。
本章是 Part II「分布式数据」的第二章,紧接第 5 章复制。两者是一对互补的横向扩展手段,也常常叠加:复制是「同一份数据整份多放几处」,分区是「把大数据集切片、每片各放一处」。真实系统几乎都是「先分区、每个分区再复制几份」(见图 4)。分区又直接牵出下一层难题:一次操作若跨多个分区,就需要跨分区的一致性与协调——这正是第 7 章事务、第 9 章一致性与共识要接手的。本章可看作把第 5 章的「多副本」问题,换成了「多分片」问题。
复制解决了「一台挂了怎么办」和「读太多怎么办」,但它没解决数据本身太大——每台副本仍存着全量数据。设想一个社交应用:10 亿用户、几十 TB 数据、写入峰值每秒几十万条。这个量级下,单机的两道墙同时撞上来:硬盘装不下几十 TB;单机 CPU / 内存 / 磁盘 IO 扛不住几十万 QPS 的写。加副本没用——每份副本都得存全量、都得吃下全部写。
出路是把数据切开:每台只存 1/N 的数据、只处理落在自己这块上的请求,理论上加机器就能线性扩容。但切开会带来单机时代没有的三类新问题,本章逐一解决:① 按什么规则切,才能让数据和负载都摊得均匀、不出热点?② 二级索引(非主键查询)怎么在切开的数据上工作?③ 集群加减节点时,怎么把分区重新摊匀,还尽量少搬数据、且搬的时候不停服? 切不好,扩容带来的不是性能,而是一个被挤爆的分区拖垮全局。
分区的唯一敌人是偏斜(skew):数据或请求不成比例地压在某些分区上,那些分区成了热点(hot spot),而扩容的初衷正是别让任何单点过载。最傻的「完全随机分」能摊得极匀,但代价是读的时候不知道数据在哪、只能问遍所有分区——不可用。所以真正的方案要在「摊得匀」和「找得到」之间取平衡,于是有了下面两种主流切法。
是什么:把键按顺序划成一段段连续区间,每个分区负责一段——像纸质百科全书按字母分卷(A–C 一卷、D–F 一卷)。区间边界可以人工设,也可由数据库自动选(让每段数据量大致相当)。分区内部按键有序存储(配合 LSM-tree / SSTable),于是范围扫描非常高效:查「某传感器在某天的全部读数」只需扫一个分区里连续的一段。用它的系统:Bigtable、HBase、RethinkDB、早期 MongoDB。
致命弱点是热点:如果键本身带有递增趋势——最典型的是用时间戳当键——那么所有新写入都砸向「今天」这一个分区,昨天的分区闲着、今天的被压垮。DDIA 给的解法是给键加前缀打散:比如把键从「时间戳」改成「传感器名 + 时间戳」,写入就按传感器名先散开;代价是查「所有传感器某时段的数据」时要对每个传感器各扫一段、再合并。
是什么:先用哈希函数把键搅成一个均匀分布的数字,再按这个哈希值的范围分区。好的哈希能把「偏斜的键」变成「均匀的哈希值」——哪怕大量键长得很像,打散后也均匀落在各分区。用它的系统:Cassandra、MongoDB(hashed sharding)、Voldemort。
代价是丢掉了范围查询:相邻的键(如连续的时间戳)被哈希打散到天南海北,「取一段连续键」就得问遍所有分区。Cassandra 的折中很经典:用复合主键(compound primary key)——第一列(分区键)拿去哈希决定落哪个分区,其余列在分区内部按序存储、可做范围扫描。于是「同一个用户的所有帖子按时间排」可以在一个分区里高效范围查,但跨用户就不行。这也是「先按 A 打散、再在每个 A 里按 B 有序」的通用套路。
顺带一提:常被提到的一致性哈希(consistent hashing)是 Karger 等人 1997 年为 CDN / Web 缓存提出的、把哈希值排成一个环来减少节点增减时的数据迁移;DDIA 特别提醒,数据库里说的「哈希分区」多半不是那篇论文严格意义上的一致性哈希,这个词用得挺乱,别被绕进去。
哈希摊匀的是不同的键;可现实里常有同一个键被疯狂访问——名人账号发了条动态,几千万人同时读写同一条记录。这条记录的键哈希后只会落到一个分区,那台机器照样被打爆,哈希毫无帮助。DDIA 坦言:今天的系统大多还没法自动补偿这种「热键」,得靠应用层动手——最常见的招是给热键加盐(salting):在键后面拼一个小随机数(比如 0–99),把一条热记录人为拆成 100 个键、散到 100 个分区去分担写。代价是读的时候要把这 100 份都读回来合并,还得自己记住「哪些键被拆过」——是个需要额外簿记的权宜之计。
前面都在讲按主键切数据。但真实查询常问的是「所有红色的车」「作者是张三的所有帖子」——这靠二级索引。麻烦在于:二级索引没法像主键那样干净地跟着分区走,因为「红色的车」可能散落在每一个分区里。DDIA 给了两种方案,是本章最容易考、也最容易搞混的一对:
方案 A · 本地索引 / 按文档分区(local / document-partitioned index):每个分区只给自己那份数据建索引,各管各的。写很简单——改一条记录只动它所在的那一个分区(数据和它的索引都在一起)。但按二级索引读很贵:因为「红色的车」散在所有分区,你必须把查询发给每一个分区、再把结果汇总——这叫分散 / 聚集(scatter / gather),读延迟受最慢那个分区拖累,尾延迟容易被放大。用它的系统:MongoDB、Cassandra、Elasticsearch、SolrCloud、Riak、VoltDB。
方案 B · 全局索引 / 按词条分区(global / term-partitioned index):把二级索引本身也切开,但不是按文档所在分区切,而是按被索引的值(词条 term)切——比如「所有 color=red 的索引项」集中放在索引分区 1,「color=blue」放索引分区 2。读快:查红色车只需问「持有 red 词条」的那一个索引分区,不用问遍全部。但写变复杂:改一条记录可能同时牵动好几个索引分区(它的颜色、价格、品牌词条各在不同分区),一次写就变成跨多个分区的分布式写——所以现实中全局索引的更新常常是异步的,意味着写完后短时间内索引可能还没跟上。用它的系统:DynamoDB 的全局二级索引(GSI)。
本章的每个设计点都是一组取舍,且彼此独立、可自由组合(切法 × 索引方式 × 再平衡策略)。下面三张表把「什么场景选什么」摊开。
表 1 · 键范围分区 vs 哈希分区
| 按键范围 range | 按哈希 hash | |
|---|---|---|
| 分布均匀度 | 取决于键分布,易偏斜 | 天然均匀(好哈希摊平偏斜键) |
| 范围查询 | 高效(分区内有序,扫一段即可) | 差(相邻键被打散,须查所有分区) |
| 写热点 | 递增键(时间戳)会全砸末尾分区 | 不同键无写热点;但单个热键仍打爆一分区 |
| 典型场景 | 时序 / 需要范围扫描(传感器、日志区间) | 点查为主、要均匀(用户资料、KV) |
| 代表系统 | HBase、Bigtable、RethinkDB | Cassandra、MongoDB(hashed)、DynamoDB |
表 2 · 本地(文档分区)索引 vs 全局(词条分区)索引
| 本地索引 local | 全局索引 global | |
|---|---|---|
| 索引按什么切 | 跟着文档所在分区 | 按被索引的值(词条)另行切分 |
| 写 | 简单:只动一个分区 | 复杂:可能跨多分区,常异步 |
| 按二级索引读 | 贵:scatter/gather 问遍所有分区 | 快:只问持有该词条的分区 |
| 一致性 | 与主数据同步 | 异步更新时可能短暂读到旧索引 |
| 代表系统 | MongoDB、Elasticsearch、Cassandra、Riak | DynamoDB 全局二级索引(GSI) |
表 3 · 四种再平衡策略(加减节点时怎么重新摊分区)
| 策略 | 怎么做 | 代价 / 坑 |
|---|---|---|
| hash mod N | 分区 = hash(键) % 节点数 | 绝不要用:节点数一变,N 变,几乎所有键都要搬家 |
| 固定分区数 | 一开始就建远多于节点的分区(如 10 节点建 1000 分区),每节点管一批;加节点就从每个旧节点匀走几个整分区 | 只搬整分区、不重算键;但分区数建库时定死、难改,选大了有开销、选小了限制扩容上限 |
| 动态分区 | 分区数随数据量自动增减:某分区超过阈值就分裂,缩小就合并(类似 B-tree) | 适配数据量;但空库只有一个分区、初期全压一台(可预分裂缓解) |
| 按节点比例 | 每节点固定数量的分区;新节点加入时随机分裂已有分区、抢走一半 | 分区数随节点线性增长;随机分裂可能切得不够均 |
三条实操心法:① 别用 hash mod N——它是再平衡的头号反面教材,一加节点就全体搬迁。② 分区数要「多于节点」但别太多:固定分区数方案里,分区是搬迁的最小单位,太少则扩容受限、太多则每个分区元数据 / 开销累积。③ 再平衡最好留个人工闸门——全自动再平衡叠加自动故障检测很危险:一个节点只是变慢被误判为「挂了」,触发再平衡 → 搬数据加重负载 → 更多节点显得像挂了 → 级联雪崩。让人「点一下确认」能挡掉大多数这类事故。
切完还剩最后一问:客户端要读某个键,怎么知道该连哪台节点?(分区还会因再平衡而搬动,这个映射是变化的。)DDIA 归纳三种做法:① 客户端随便连一个节点,那个节点要么自己处理、要么转发给对的节点;② 加一层路由层(routing tier),所有请求先到它、由它按分区映射转发(它自己不处理数据);③ 客户端自己知道分区映射、直连目标节点。核心难点都一样——分区到节点的映射一变,谁来、怎么通知到做决策的那一方。许多系统靠一个独立的协调服务(如 ZooKeeper)专门维护这份集群元数据、变更时通知订阅者(HBase、SolrCloud、早期 Kafka 走这条);Cassandra、Riak 则用节点间的gossip 协议互相同步、不依赖外部协调者;MongoDB 用专门的 config 服务器 + mongos 路由进程。
分区是所有「大到单机装不下」的数据系统的地基。你在 MongoDB 里设 shard key、在 Cassandra 里定 partition key、在 Kafka 里给 topic 分 partition、在 DynamoDB 里挑 partition key——本质都是在本章的框架里做选择:range 还是 hash、本地还是全局索引、分区键会不会造成热点。面试里的高频题——「Kafka 分区数怎么定」「Cassandra 的 partition key 和 clustering key 有什么区别」「怎么避免热分区 / 热 key」「一致性哈希是什么、和取模有什么区别」「分库分表后二级查询怎么做」——答案全在这一章。选错分区键几乎是分布式数据库最常见、也最难补救的架构失误:它决定了数据怎么散、热点在哪、能不能范围查,而且一旦上线、数据填进去,重选分区键往往意味着一场浩大的数据迁移。
(channel_id, bucket) 作复合分区键——bucket 是一段静态时间窗(约 10 天)——既把同一频道的消息按时间聚在一起,又防止热门频道的消息把单个分区撑爆;这是「复合主键 + 分桶」对抗热点与无界分区的教科书案例。后来他们把存储从 Cassandra 迁到 ScyllaDB 缓解尾延迟。来源:Discord Engineering, "How Discord Stores Trillions of Messages"(2023)hash mod N 定分区:看着简单,扩容时几乎全量数据搬家。要用「固定分区数」或「动态分区」等再平衡友好的方案。① 分区 = 把大到单机装不下的数据集切成互不重叠的小块,每块一台节点,为的是横向扩展存储与吞吐;与复制正交、通常叠加。
② 两种切法:按键范围(有序、利于范围扫描,但递增键易出写热点)、按哈希(分布均匀,但丢失范围查询);Cassandra 用复合主键各取一半。
③ 哈希摊匀的是不同的键,挡不住单个热键——名人 / 爆款要靠应用层加盐拆分,代价是读时合并 + 额外簿记。
④ 二级索引两条路:本地索引(写省心、读要 scatter/gather 问遍全分区)、全局索引(读高效、写跨分区且常异步)——本质是把代价放在读还是写。
⑤ 再平衡策略:hash mod N 绝不能用(全量搬迁);用固定分区数 / 动态分区 / 按节点比例,只搬整分区、少动数据。
⑥ 全自动再平衡叠加自动故障检测易引发级联失效,生产中常留人工确认闸门。
⑦ 请求路由要解决「键在哪台节点」:靠 ZooKeeper 等协调服务维护映射、或 gossip 节点互告;映射随再平衡而变。
⑧ 分区键是最难更改的决定:它定死了热点、范围能力和扩展性,选错常意味着大规模数据迁移——上线前务必压测。