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 的代價,逼面試官自己把它挖出來。