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 ranges | Data 本来就落在 ranges 里;range scans 留在一个 shard |
| Hash + consistent hashing | hash(key) 落在 ring 上 | Default。均匀铺开,加一个 shard 不会把一切重洗 |
| Directory | Lookup table 说 key 住在哪 | 必须挪走一个 hot key,又不想 rehash 整个 cluster |
1. 为什么需要 sharding
一台大的 RDS Postgres 能撑很久。~70 TB。~10k writes per second。先升级。更大的箱子仍然在发火,再 shard。
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 不是。
| Partitioning | Sharding | |
|---|---|---|
| Data 住在哪 | 一个 Postgres | 许多 Postgres instances |
| 什么在 scale | Vacuum、pruning、drop-a-month | Writes、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。
| Key | Why |
|---|---|
社交 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。 |
- 让相关 rows 住在一起。Comments 与对应的 post 在同一 shard。Detail page 是一个 database,不是跨机器的 join。
- Tenant-centric 产品常常按
organization_idshard。这套 stack 上的 isolation 故事仍是 row 上的 RLS。用 Hono、Better Auth、Drizzle 与 Postgres RLS 打造 Multi-Tenant 后端。
Failure: 按 created_at shard,于是所有 writes 打在今天的 shard;或 posts 与 comments 用不同 keys,于是每个 detail page 都是 cross-shard join。
4. 组怎么落到机器上
Key 把 rows 分组。Strategy 把组放到机器上。三个选项。面试默认是 hash 加 consistent hashing。
Range
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
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_idshard」已经暗示它,除非你是 mid-level 而被问 how。
Directory
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_idshard。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 不该是这样。
GET /users/123 → one shard
GET /posts/trending → all shards, then merge- Alignment 是第一道防线。如果大多数 queries 是 user-scoped,就按
user_idshard。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 跨不过去。
-- 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 拆开。证明它。
如果数字要求拆开:
- Key 来自 access patterns. 社交 app,user-centric reads —— posts、followers、likes 都 scoped 到一个 user —— 按
user_idshard。 - Distribution. Hash-based 加 consistent hashing。均匀铺开。加一个 shard 只动一小部分,不是整个 cluster。Mid-level:说出来。Senior:这是默认。
- Trade-off. Global queries 变贵。Trending 是 cache 或预计算 job,不是每次 request 都 scatter-gather。
- Growth. 一开始就留够 shards 去长。Consistent hashing 是你加更多、又不用搬走一切的办法。
Failure: bottleneck 还没出现就画出 shards,或点名一个没有 access pattern 的 key,或跳过 global-query 的代价,逼面试官自己把它挖出来。