Database Sharding:行按什么切、查询怎么找到它、搬数据和热 key 怎么办
比较 Range Sharding with Auto Split、Hash Slots、Directory by Tenant 与 Write Sharding of Hot Keys 四种分片拓扑:哪个值决定行落在哪个 shard、请求怎么找到那个 shard、加减 shard 时数据怎么搬、一个 key 或一段区间变热时会发生什么、路由的真相归谁,以及一次打错 shard 或一次广播能波及多远、怎么收住。
分片就是把一张逻辑表拆到几个数据库上,每个库只放一部分行、只扛一部分负载。所有分片方案都要回答四个问题:哪个值决定一行落在哪个 shard、一个请求怎么找到那个 shard、加减 shard 时已有的行怎么办、一个 key 或一段区间拿走大部分流量时会发生什么。这四个答案,而不是用的哪家产品,才定义了拓扑。
想亲手走一遍四种拓扑、注入故障再恢复,打开互动 Lab:/system-design-lab/database-sharding-architectures。
有约束的设计问题
一家公司先后要给四类数据分片:一张按客户和时间做范围扫描、每月都在涨的业务表;一个点查为主、节点常增减、扩缩容不能停服的 KV 缓存;一个每个请求都属于某个租户、业务表之间频繁 join、个别大租户要单独安置的多租户 SaaS;以及一份所有写都落在「今天」这一个分区键上、超过单分区吞吐上限的时序数据。每一类该怎么切?
四种 topology signature
| 架构 | 什么决定 shard | 请求怎么找到它 | 怎么搬数据 | 热点表现 | State owner | 适合 | 主要代价 |
|---|---|---|---|---|---|---|---|
| Range Sharding with Auto Split | 有序 shard key 落在哪段区间 | 路由器查集群元数据里的区间 → shard 映射 | 区间按大小分裂,均衡器搬区间并更新元数据 | 单调 key 把所有插入送到最后一段;低基数造 jumbo chunk | 区间 → shard 的元数据 | 按 key 范围扫描、前缀查询、数据持续增长 | 路由依赖元数据新鲜度;搬区间有开销;单调 key 要哈希或复合 key |
| Hash Slots | hash(key) mod 槽数 | 客户端缓存槽表,节点用 MOVED / ASK 重定向 | 槽逐个在节点间搬,两边同时服务 | 天然均匀,除非一个 key 本身热;多 key 操作只能在同一槽内 | 集群达成一致的槽 → 节点表 | 点查为主的 KV 和缓存,节点弹性增减 | 没有范围扫描;跨槽操作不支持;热 key 仍在一个节点 |
| Directory by Tenant | 查表:租户 → shard | 网关先解析 id 再路由;同址的表本地 join | 改目录并复制该租户的行 | 大租户就是热 shard;按租户搬而不是按行搬 | 目录(lookup vindex、协调者元数据) | 多租户 SaaS,查询按租户,join 要留在本地 | 插入多写一条、查询多查一次;不带租户 id 就广播;目录是每个请求的依赖 |
| Write Sharding of Hot Keys | 应用给逻辑 key 算出的后缀 | 写时算后缀;读查遍所有后缀再合并 | 不需要:后缀数量就是并行度 | 专为一个热逻辑 key(日期、计数器)设计 | N 个物理分区 | 只有一个热 key 的时序和计数器 | 读要 N 次查询加合并;随机后缀让单条读不可能,除非后缀可算 |
规范里的 hashed shard key(MongoDB 的哈希分片)是先哈希再走范围拓扑;meta ranges(寻址区间的两级索引)是范围拓扑的路由元数据;全局二级索引 是目录拓扑对非 key 查找的答案。它们在这里解释,不单独画。
1. Range Sharding with Auto Split:区间切块,元数据指路,均衡器搬块
MongoDB 把定义写得很直接:「shard key 是一个带索引的字段或复合索引覆盖的多个字段,决定集合文档在集群 shard 之间的分布」;「MongoDB 把 shard key 值的范围(或哈希后的值)切成不重叠的区间。每个区间对应一个 chunk,MongoDB 尽量把 chunk 均匀分布到集群的 shard 上」。路由这一侧,mongos「通过缓存 config server 的元数据来跟踪哪些数据在哪个 shard」,「可以把包含 shard key 或复合 shard key 前缀的查询路由到特定的 shard」,否则「把查询广播给该集合的所有 shard,除非 mongos 能确定哪个 shard 或哪些 shard 存有这些数据」;「多文档更新操作永远是广播操作」。
搬数据的是均衡器:「一个后台进程,监控每个分片集合在每个 shard 上的数据量」,到达迁移阈值时「尝试自动在 shard 之间迁移数据,在遵守 zone 的前提下让每个 shard 的数据量均衡」;它「运行在 config server 副本集的主节点上」;过程「对用户和应用层完全透明,但进行期间可能有性能影响」;「长期禁用均衡器会降低集群性能」,可以设置均衡窗口避开生产流量。热点的来源在故障排查页写得清楚:单调递增的 key 让「新文档通常写入同一个 shard 和 chunk。接收写入的 shard 和 chunk 称为热 shard 和热 chunk」;出现 jumbo chunk 说明「shard key 的基数不够,或者 shard key 值的频率分布不均」;办法是「给现有 key 加后缀字段来 refine 以提高基数」、改用哈希分片,或「用基数更高的 shard key 重新分片」。CockroachDB 是同一个形状的无运维版:key 空间「被划分成我们称为 range 的连续 key 空间块,所以每个 key 总能在单个 range 里找到」;「range 一到默认大小就分裂成两个」;小 range 会合并;寻址靠「key 空间开头的两级索引,称为 meta range」;系统「在节点上线或下线时自动重新平衡 range 的分布」。
典型事故:团队用自增的 orderId 做 shard key,大促期间所有插入都落进最大的那段区间,那段一直在分裂但新的一段仍在同一个 shard 上,写吞吐等于一台机器——对单调 key 用哈希分片或高基数非单调的复合前缀。新加的「按手机号查订单」没带 shard key,每次都广播全部 shard 再合并,功能一热整个集群一起忙——带上 key 或先查一张「手机号 → shard key」的查找表。均衡器把区间搬走后某台路由器还拿着旧映射打到老 shard——shard 校验映射版本,旧了返回「配置过期」让路由器刷新重试,而不是静默返回空结果。
2. Hash Slots:key 哈希到固定的槽,搬数据就是搬槽
Redis Cluster 规范:HASH_SLOT = CRC16(key) mod 16384;hash tag 让「{user1000}.following 和 {user1000}.followers 两个 key 哈希到同一个槽,因为只有子串 user1000 参与哈希」;节点收到不属于自己的槽的请求会回 -MOVED 3999 127.0.0.1:6381;搬槽时对源节点执行 CLUSTER SETSLOT 8 MIGRATING B、对目标节点执行 CLUSTER SETSLOT 8 IMPORTING A,期间「如果收到 ASK 重定向,只把被重定向的这一条查询发给指定节点,后续查询继续发给老节点」;「向集群加节点就是加入一个空节点,再把一批哈希槽从现有节点搬过去」,「再平衡就是在节点之间搬一批槽」,「移除节点就是把它的槽搬到其他节点」。可用性方面,集群「能在多数主节点可达、且每个不可达的主节点至少有一个可达副本的分区中存活」;写安全是尽力而为,「存在已确认的写被丢失的小窗口」,「客户端在少数派分区时这个窗口更大」。
典型事故:应用把 user1000 的关注列表和粉丝列表存成两个没有 hash tag 的 key,想一条命令一起读,集群回 CROSSSLOT——给它们同一个 hash tag,或改成一个 key 下的哈希结构。一场直播的弹幕计数器是一个 key,每秒几十万次自增,它所在的节点 CPU 打满,同节点几千个槽一起变慢,搬槽也没用,因为热的是一个 key——在应用层给这个 key 做写分片,读时把 N 份加总。
3. Directory by Tenant:查表定位租户,同租户的表放一起,维表复制到处
Vitess:「Vindex 提供把列值映射到 keyspace ID 的方式」;「表的 Primary Vindex 类似数据库主键」,「给定输入值必须产生唯一的 keyspace ID」;「Functional Vindex 的列值到 keyspace ID 映射是预先确定的,通常通过算法函数」,而「Lookup Vindex 提供在一个值和 keyspace ID 之间建立关联并在需要时找回的能力」;代价是「插入和删除多一次对查找表的写,读多一次查找」;一致性 lookup vindex「用谨慎的加锁和事务顺序保证一致性而不用 2PC」;VTGate 让「更新某个用户信息的查询可以只发往一个 shard」,而「查几个商品的查询可能发往一个或多个 shard」;重新分片时「Vitess 在新 shard 上复制、校验并保持数据最新,现有 shard 继续服务实时读写」。
Citus 给出目录拓扑的另一种实现:「每个集群有一个特殊节点叫协调者」;「应用把查询发给协调者,协调者转给相关的 worker 并汇总结果」;它的元数据表「跟踪 worker 节点的 DNS 名和健康状态,以及数据在节点上的分布」;「shard 到 worker 的映射称为 shard placement」;「Citus 用分布式表的分布列把行分配到 shard」,选它「是最重要的建模决策之一」;「分布列值相同的行总在同一台机器上,即使跨不同的表」,所以「限定同一个 account_id 时,Accounts 和 Campaigns 的 join 在一个节点上就有全部需要的数据」;「reference table 是一种分布式表,全部内容集中在一个 shard 里并复制到每个 worker」;建议「用共同的 tenant_id 列对分布式表分区」;没有分布列的查询要付「查询每个 shard、运行多条查询的开销」和「分多步写查询再合并结果的成本」。
典型事故:运营后台加了「按 email 域名查所有用户」,条件里没有 tenant_id,网关无法定位,发给全部 shard 各跑一遍,所有租户的 shard 同时被这一个查询占用——所有业务表按 tenant_id 分布,共享维表做 reference table,确实要按非租户列查的少数场景加 lookup vindex。存 lookup vindex 的库磁盘满了,每个请求都要先查它,所有租户的请求在网关失败尽管每个 shard 都健康——目录多副本、网关缓存最近的映射,缓存里没有的租户明确报错而不猜。
4. Write Sharding of Hot Keys:给热 key 加后缀摊开写,读查遍后缀再合并
DynamoDB 的分区键最佳实践:「你应该把应用设计成在表及其二级索引的所有分区键上活动均匀」;每个分区「设计为最多提供每秒 3,000 个读单元和 1,000 个写单元」。摊开热 key 的办法是「在分区键值末尾加一个随机数」,于是「每天对表的写入被均匀分散到多个分区」,但「要读某一天的全部条目,你必须查询所有后缀的条目再合并结果」;用计算后缀时「你可以轻松地对特定条目和日期执行 GetItem,因为你能为特定的 OrderId 算出分区键值」,而整天读仍要查「每个 2014-07-09.N 的 key」再合并;「好处是避免了单个『热』分区键值承担全部负载」。
典型事故:第一版用了 1 到 200 的随机后缀,写得很均匀,上线后客服要按 orderId 查一笔订单,没人知道它当时拼的是哪个后缀,只能 200 个 key 全查一遍——后缀改成从 orderId 算出来,旧数据按新公式回填一次。把 N 从 8 调成 10 之后,读整天的循环仍然只查到 date.8,存储对 date.9 和 date.10 一无所知也不会报错,日报连续几天少了两片订单——把 N 版本化成写读共同引用的配置,读时校验返回的片数等于 N。
四种拓扑共同的底线
- 说清路由的真相归谁。 区间元数据、槽表、目录,还是逻辑 key 的 N 个分区。
- 把 shard key 放进每一条能放的查询里。 不带 key 的查询就是广播,广播的比例要单独计数和告警。
- 先排练搬迁再等它变紧急。 均衡器、搬槽、搬租户各有窗口和开销,都要在第一个 shard 满之前跑过一次。
- 热点先看每 shard 的数据量和读写速率。 远高于别人的那个 shard 就是热点或 jumbo range。
- 路由元数据要高可用、有版本。 拿着旧映射的路由器必须收到明确的「配置过期」,而不是错的数据。
故障与恢复
| 架构 | 故障 | 用户看到什么 | 恢复 |
|---|---|---|---|
| Range Sharding | 单调 shard key 把所有插入送到一个热 shard | 一台机器的写吞吐 | 哈希或复合 key;refine 或 reshard |
| Range Sharding | 不带 shard key 的查询广播 | 全集群一起忙 | 带上 key 或前缀;查找表先定位 |
| Range Sharding | 搬区间后路由器元数据过期 | 打到老 shard | 「配置过期」错误触发刷新重试 |
| Hash Slots | 跨槽多 key 命令被拒 | CROSSSLOT | hash tag 让相关 key 同槽;或改单 key 结构 |
| Hash Slots | 客户端不更新槽表 | 一直被 MOVED | 智能客户端改表;ASK 只跟一次 |
| Hash Slots | 一个热 key 打满节点 | 同节点所有槽变慢 | 应用层写分片;只读热 key 走副本 |
| Directory | 不带 tenant_id 的查询广播 | 所有租户延迟变差 | 全表按 tenant_id 分布;reference table;lookup vindex |
| Directory | 目录不可用 | 全部请求失败 | 目录多副本、网关缓存;缓存没有就明确报错 |
| Write Sharding | 随机后缀让单条读不可能 | 单条查询等于整天扫描 | 从可查询属性算后缀;整天读才查遍后缀 |
怎么选
- 按 key 范围扫描、数据持续增长:Range Sharding with Auto Split,key 用高基数非单调的复合前缀,查询都带 key,均衡窗口避开峰值,元数据当作必须保鲜的东西。
- 点查为主、节点弹性增减:Hash Slots,相关 key 用 hash tag,智能客户端在
MOVED时刷新槽表,槽逐个搬,并且明白哈希摊开的是 key 不是一个 key 的流量。 - 多租户、查询按租户、join 要本地:Directory by Tenant,所有表按租户列分布,小共享表做 reference table,少数非租户查找加 lookup vindex,大租户改目录搬到独占 shard。
- 一个逻辑 key 天生热:Write Sharding,用可算的后缀让单条读仍然可能,N 刚好够用,整 key 读必须扇出并合并。
- 无论哪种:写下路由真相归谁、key 进每条查询、搬迁路径在第一个 shard 满之前排练过。
回到开头的四类数据:业务表走 Range Sharding;KV 缓存走 Hash Slots;多租户 SaaS 走 Directory by Tenant;时序数据走 Write Sharding。
面试时这样回答
- 先复述约束:范围扫描还是点查、节点会不会增减、查询是否都带租户 id、有没有天生的热 key。
- 说路由路径:点名图上的边,例如「路由器查元数据后只打一个 shard,均衡器后台搬区间」「客户端算槽,
MOVED后重试,搬槽两边同时服务」「网关查目录再路由,同址 join 加本地 reference table」「拼后缀写入,读整天扇出 N 次再合并」。 - 说路由真相归谁:区间元数据、槽表、目录,还是 N 个分区。
- 说代价与一个故障:例如自增 orderId 造成热 shard,修法是哈希分片或非单调复合 key,已上线的集合先 refine 再 reshard。
一手证据
- MongoDB:Shard Keys
- MongoDB:Sharded Cluster Balancer
- MongoDB:Troubleshoot Shard Keys
- MongoDB:mongos(query router)
- Redis:Redis cluster specification
- Vitess:Vindexes
- Vitess:Sharding
- Citus:Concepts
- Citus:Data Modeling
- Amazon DynamoDB:Partition key design best practices
- Amazon DynamoDB:Write sharding
- CockroachDB:Distribution Layer