跳至主要內容
返回

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 核心概念