跳到主要内容
返回

System Design 里的 Sharding

系统设计

为什么 writes 与 storage 要离开一个 node —— shard keys、hash vs range、hotspots、cross-shard work,以及随之而来的 failure modes

Sharding 把一份 dataset 拆到彼此独立的 databases 上,好让没有一台机器独自扛 writes、disk、或 connection pool。每个 shard 是自己的 Postgres:CPU、memory、storage、connections。合在一起才是那份 dataset。Vertical scale 是第一步。Sharding 是机器已经到顶之后才伸手去拿的东西。

这篇 note 依 Evan 的 walkthrough。Schema 见 System Design 里的 Data Modeling。同一个 Postgres 里的 partitioning 见 SQL 核心概念。Denormalized 的全局副本住在 System Design 里的 Caching。Compensating steps 见 System Design 里的 Message Queues。这篇 note 讲的是拆开。



Pattern Map

Pattern一行如何找到一台机器Reach for it when
Range连续的 key rangesData 本来就落在 ranges 里;range scans 留在一个 shard
Hash + consistent hashinghash(key) 落在 ring 上Default。均匀铺开,加一个 shard 不会把一切重洗
DirectoryLookup table 说 key 住在哪必须挪走一个 hot key,又不想 rehash 整个 cluster


1. 为什么需要 sharding

一台大的 RDS Postgres 能撑很久。~70 TB。~10k writes per second。先升级。更大的箱子仍然在发火,再 shard。


text
One Postgres  →  bigger Postgres  →  shards
   (default)        (vertical)       (horizontal)

  • Traffic 涨。Reads 涨。Writes 涨。第一反应是更大的 instance —— 更多 CPU、更多 disk、140 TB、每秒几万次 writes。大多数公司永远不离开这台箱子。
  • Ceiling 是真的。Storage 满了。Write throughput 饱和。Backups 要做一辈子。Working set 再也装不进,queries 变慢。更大的机器不存在。
  • 这时 sharding 是唯一还能加容量的动作:把 rows 拆到多台机器。Storage 与 throughput 靠加第四个、第五个、第六个 shard 往上走。
  • 代价是 operational。你现在要选 key、route queries、忍受 hot shards、再 rebalance。这份复杂度,就是为什么画 boxes 之前先算数。

Failure: 一份 2.5 TB 的 dataset 因为「我们总是 shard」就被拆开,而单台 Postgres 本来装得下。面试要的是数字,不是拆分。



2. Partitioning vs sharding

Partitioning 是在 同一个 database 里 整理一张表。Sharding 是 多台 databases。口头上经常是同一个词。Failure domain 不是。


PartitioningSharding
Data 住在哪一个 Postgres许多 Postgres instances
什么在 scaleVacuum、pruning、drop-a-monthWrites、disk、connections
Failure domain一个 primary、一份 WAL、一套 backup每个 shard 是自己的故障
App 要改什么通常不用 —— parent table 仍是 query按 key routing。Cross-shard 工作现在是你的

  • 一个 node 上 500M 行的 orders 表是 partition 问题:巨大的 indexes、整张表的 autovacuum、一条按日期的 query 扫过好几年。按月拆。Primary 仍是一台机器。
  • Horizontal partition:同样的 columns,每片更少 rows。Vertical partition:同样的 rows,每片更少 columns。两者都不加第二台 primary。
  • 一个 shard 有自己的 CPU、memory、storage、pool。没有一台箱子装着整份 dataset。这就是 writes 与 storage 如何越过一台机器。
  • Statement 模型、以及 partition key 何时真的 prune,见 SQL 核心概念。Key 必须服务的 schema 见 System Design 里的 Data Modeling

Failure: 把 table partitions 叫成「sharding」,然后发现 unique constraints、joins、transactions 仍活在一份 WAL 上 —— 或把第二台 RDS instance 当成 planner 看得见的 partition。



3. Shard key

两个决定。按什么 分组 —— shard key。这些组如何 落到机器上 —— 下一节。面试里,点名 column 以及为什么。Key 通常是永久的。


好的 key 有三个性质:


  • High cardinality. 许多 distinct values,data 才能真正铺开。Boolean 是两组。你已经把自己封在两个 shards。
  • Even distribution. Values 不该堆到一台机器。user_id 通常能铺开。Country,如果 90% 的 users 在一个国家,铺不开。
  • Query alignment. Hot path 应该打中 一个 shard。GET /users/{id}/posts 要的是 user_id。对不上 APIs 的 key,会把每次常见 read 变成 scatter-gather。

KeyWhy
社交 app 上的 user_id几百万个 values。Profiles、posts、likes 都是 user-scoped。一个 user,一个 shard。
Checkout 上的 order_id几百万笔 orders。Create、status、receipt 是一笔 order。
is_premium两个 values。两个 shards。然后你卡住了。
Write path 上的 created_at每次 insert 都落在 今天。那个 shard 着火。更旧的 shards 闲着。Time-range 属于 archival,不属于 OLTP。


Failure:created_at shard,于是所有 writes 打在今天的 shard;或 posts 与 comments 用不同 keys,于是每个 detail page 都是 cross-shard join。



4. 组怎么落到机器上

Key 把 rows 分组。Strategy 把组放到机器上。三个选项。面试默认是 hash 加 consistent hashing。


Range

text
Shard 1 → user_id 0 – 10M
Shard 2 → user_id 10M – 20M
Shard 3 → user_id 20M – 30M

  • 简单。落在一个 interval 里的 range scan 只打一个 shard。
  • 早期只有 shard 1 有 users。后来新 ids 是单调递增的 —— 每个新 user、每次 write,都落在最高的 range。那个 shard 吃掉热度。
  • 当不同 tenants 本来就 query 不同 ranges 时,range 能用。对一份在涨的 user_id,它不是 production 默认。

Hash

text
shard = hash(user_id) % N

user 42  → hash(42)  % 3 = shard 1
user 99  → hash(99)  % 3 = shard 2
user 123 → hash(123) % 3 = shard 0

  • Hash 把输入打散。新 users 均匀铺开。这是赢的地方。
  • N 一变就是输。% 3 改成 % 4,几乎每个 key 都 remap。大部分 dataset 要搬家。Operationally 是噩梦。
  • Consistent hashing 把 keys 与 shards 放上一个 ring。顺时针走到下一个 shard。加一个 node,只有一小部分 keys 移动 —— 邻居,不是整个 cluster。Virtual nodes 避免 ring 结块。这是 industry default。面试里,「按 user_id shard」已经暗示它,除非你是 mid-level 而被问 how。

Directory

text
user_to_shard
---------------
15  → shard 1
87  → shard 4
204 → shard 2

  • Lookup 说 key 住在哪。改一行就能把 Messi 挪到自己的 shard。Rebalance 不用 rehash。
  • 每个 request 现在是两跳:directory,然后 shard。Directory 是 single point of failure。它挂了,健康的 shards 也够不着 —— 你不知道 data 在哪。
  • Production 里,对少数 hot keys,这份灵活性你付得起。几乎从来不是面试答案。它会把后半小时拖进 lookup 的 HA。

Failure: hash % N 却没有 N+1 的计划,或一上来就 directory,于是剩下的一小时都在讲 lookup service。



5. Hotspots

好的 key 仍有 outliers。Hashing 铺开的是 keys,不是 traffic。一个 key 可以是大部分负载。


  • Celebrity problem.user_id shard。Messi 落在 shard 1。每次 profile view、like、comment、DM 都打那台。其它的闲着。Cluster 看起来没事。一台 primary 着火。
  • Time-range 是另一种形状:所有新 writes 都去最新的 shard。同一场事故,不同原因。
  • 从 shard metrics 发现它 —— latency、CPU、RPS —— 不是从「keyspace 很均匀」。

两种应对:


  • Compound key. Hash user_id + n,或 user_id + date,让一个 celebrity 的 posts 铺到多个 shards。曾经是一个 shard 的 reads 现在要 fan-out。你用这个 user 上更小的 scatter-gather,换来了 write 铺开。
  • Dedicated celebrity shard. 找出 hot keys。把它们挪到自己的硬件。一小层 directory overlay:如果是 celebrity,走那台 shard;否则 hash。大多数系统永远不需要这个。一张有 Messi 的 social graph 需要。

Cache 一条热 read 是另一半。Profile 前面的 Redis 不会让 write path 变成无限,但它能拦住 read storm 同时变成 database storm。System Design 里的 Caching


Failure: hash 了 user_id 就宣布 cluster 均衡,因为每个 shard 的 user 数一样,而其中一个 user 就是产品本身。



6. Cross-shard 工作

Rows 一旦住在许多机器上,任何需要超过一个 shard 的 query 都是 fan-out:查 N 台,等,merge。Hot path 不该是这样。


text
GET /users/123            → one shard
GET /posts/trending       → all shards, then merge

  • Alignment 是第一道防线。如果大多数 queries 是 user-scoped,就按 user_id shard。Global top-10 于是是例外,不是 feed。
  • 你消不掉 global queries。Trending、leaderboards、「有多少 users」。第一次昂贵的 scatter-gather cache 起来。Trending page 五分钟 stale 是产品选择。用 job 预计算,让 request 永远不用 fan-out。System Design 里的 Caching
  • Denormalize,让相关 facts 住在一起。把 hot read 需要的 fields 抄到已经有这个 user 的 shard。Writes 去两个地方。Reads 留在一个。Schema 上的取舍见 System Design 里的 Data Modeling
  • 偶发的 admin totals 可以打每一个 shard。一条常见的、面向 user 的路径如果总是 scatter-gather,说明 key 错了

Failure: 把「我们会 query 所有 shards 再 aggregate」当成 homepage 的计划。那是该换 key、cache merge、或预计算的信号 —— 不是把 N 次 round-trips 当成产品。



7. Consistency

一个 Postgres 让一笔 transfer 是一次 transaction。两个 shards 让它变成两台彼此不认识的 databases。ACID 跨不过去。


text
-- one database
BEGIN;
  UPDATE accounts SET balance = balance - 5 WHERE id = bob;
  UPDATE accounts SET balance = balance + 5 WHERE id = alice;
COMMIT;

  • Bob 在 shard 3,Alice 在 shard 1:扣款可以成功,入账可以失败。钱没了。反过来是双记一笔。没有覆盖两台箱子的 ROLLBACK
  • Two-phase commit 先让每个 shard prepare,再 commit。正确、慢、脆。Coordinator 或某个 shard 在半路挂掉,会留下不容易解开的 locks。Production 避开它。
  • 金科玉律:不要有这场 distributed transaction。把一个 user 的 balance、history、profile 留在那个 user 的 shard。Collocation 就是 consistency 策略。Key 才是设计。
  • 当两个 shards 上的两个 users 必须转账,用 saga:扣 Bob,给 Alice;入账失败就 refund Bob。Compensation 不是 undo。它是一笔新的、durable、可重试的 write,带着自己的 audit。System Design 里的 Message Queues
  • Follower counts 与 denormalized tallies 可以 eventually consistent。几秒的不一致,比 2PC 便宜。钱不是 tally。

Failure: 因为教科书点了名就把 2PC 当默认,或把 saga compensation 当成 rollback —— 钱已经动了;refund 是第二件 business event。



8. 面试里怎么讲

Sharding 出现在 deep dive,当你在满足一条 scaling 的 non-functional。先算数。一台调好的 Postgres 能走很远。


  • Storage. 500M users × 5 KB = 2.5 TB。一台 instance 装得下。把这话讲出来。10× 或 100× 再 shard。
  • Write throughput. Peak 5 万 writes/s 已经越过一台舒服的 primary。现在你有理由。
  • Read throughput. Replicas 能带你走很远。1 亿 DAU、每人好几次 queries,仍可能需要把 read load 拆开。证明它。

如果数字要求拆开:


  1. Key 来自 access patterns. 社交 app,user-centric reads —— posts、followers、likes 都 scoped 到一个 user —— 按 user_id shard。
  2. Distribution. Hash-based 加 consistent hashing。均匀铺开。加一个 shard 只动一小部分,不是整个 cluster。Mid-level:说出来。Senior:这是默认。
  3. Trade-off. Global queries 变贵。Trending 是 cache 或预计算 job,不是每次 request 都 scatter-gather。
  4. Growth. 一开始就留够 shards 去长。Consistent hashing 是你加更多、又不用搬走一切的办法。

Failure: bottleneck 还没出现就画出 shards,或点名一个没有 access pattern 的 key,或跳过 global-query 的代价,逼面试官自己把它挖出来。



Recap Q&A

阅读下一篇笔记
SQL 核心概念