Day 8 我們學到 Database Replication:把相同資料複製到其他 Database
Server,主要改善 Read Scaling 與 Availability。
但如果 Primary 有 100 TB 資料,每個 Replica 仍然需要保存大量相同資料;而
Single Primary 的 Write 也可能成為 Bottleneck。
今天要解決的是:
如果一台 Database 已經放不下所有資料,能不能把資料拆到不同
Database?
這就是 Database Sharding。
**Partition(分區)**就是把一大份資料拆成數個較小部分。
例如 1 億 Users:
Partition A → User 1 ~ 25M
Partition B → User 25M ~ 50M
Partition C → User 50M ~ 75M
Partition D → User 75M ~ 100M
可以想成一個倉庫太大,所以把商品分散到多個倉庫。
按照 Column 拆:
users_basic
├── id
├── name
└── email
users_profile
├── user_id
├── profile_image
├── bio
└── address
按照 Row 拆:
Users 1 ~ 1M → Partition A
Users 1M ~ 2M → Partition B
Users 2M ~ 3M → Partition C
Database Sharding 通常可以理解成一種把 Horizontal Partitions 分散到不同
Database Server 的方式。
每一份被分散出去的資料叫做 Shard:
Shard #1 → Users 1 ~ 1M
Shard #2 → Users 1M ~ 2M
Shard #3 → Users 2M ~ 3M
所以 Sharding 的核心是:
把大型 Dataset 拆成多份,分散到不同 Database Server。
Replication:
DB A → A B C D
DB B → A B C D
DB C → A B C D
重點是 Copy Data,常用來改善 Read
Scaling、Availability、Redundancy。
Sharding:
Shard A → A B
Shard B → C D
Shard C → E F
重點是 Split Data,常用來 Scale Storage 與 Write。
一句話:
Replication = Copy
Sharding = Split
Shard Key 是用來決定一筆資料應該去哪個 Shard 的欄位或 Key。
例如使用:
user_id
系統看到:
user_id = 123
就利用 Sharding Rule 決定:
→ Shard #2
可以把 Shard Key 想成郵件上的 Postal
Code:它幫助系統判斷資料應該送到哪個區域。
Range 就是「範圍」。
例如:
Shard #1 → user_id 1 ~ 1,000,000
Shard #2 → user_id 1,000,001 ~ 2,000,000
Shard #3 → user_id 2,000,001 ~ 3,000,000
user_id = 1,500,000 很容易知道要去 Shard #2。
優點:
Simple
Easy to Understand
Range Query Friendly
但可能出現 Hot Shard。
Hot Shard 指某一個 Shard 承受遠高於其他 Shard 的 Data 或 Traffic。
如果 User ID 一直增加,最新 User 都落到最後一個 Range:
Shard #1 → Low Traffic
Shard #2 → Low Traffic
Shard #3 → 🔥🔥🔥
Traffic Distribution 不平均,就失去 Sharding 原本想分散 Load 的目的。
Cardinality 在這裡可以簡單理解成:
一個欄位有多少不同的 Value。
例如 country 只有有限數量的值;user_id 可能有一億個不同 Value。
Shard Key 通常希望有足夠高的 Cardinality,讓資料比較容易分散。
但 High Cardinality 不代表一定是好 Shard Key,還要看:
Data Distribution
Traffic Distribution
Query Pattern
例如用 country Sharding,如果 80% User 都在同一國,某個 Shard
還是可能特別大。
先理解 Hash Function。
它可以先簡單想成:
把 Input 經過固定計算,轉換成另一個數值。
例如:
hash(user_id)
接著:
hash(user_id) % number_of_shards
決定資料去哪個 Shard。
例如:
hash(123) % 3 → Shard #1
hash(456) % 3 → Shard #2
如果 Hash Distribution 良好,資料通常比簡單 Range 更平均。
Trade-off:
Range-based
→ Range Query Friendly
→ Hotspot Risk
Hash-based
→ Better Distribution
→ Range Query Harder
Backend 必須知道 User 123 在哪個 Shard。
負責根據 Shard Key 決定目的 Shard 的 Routing Logic,可以概念化成 Shard
Router:
Request
↓
Backend
↓
Shard Key
↓
Shard Router
↓
Shard #2
因此 Sharding 不只是 Database 的改變,也會增加 Application Architecture
的 Complexity。
如果要找所有 age > 30 的 User,而資料散落三個 Shards:
Query Shard #1
Query Shard #2
Query Shard #3
↓
Merge Results
這就是 Cross-Shard Query:一個 Query 需要存取多個 Shards。
通常會增加:
Network Calls
Latency
Computation
Complexity
如果 JOIN 的資料也位於不同 Shards,就可能形成 Cross-Shard
JOIN,比單一 Database JOIN 更麻煩。
因此選 Shard Key 時要思考:
哪些資料經常一起被 Query?
假設原本 3 個 Shards 快滿了,因此新增 Shard #4。
這時需要把部分資料從舊 Shards 移到新 Shard。
這個過程叫:
Rebalancing
而系統搬資料時通常還需要繼續 Serving Reads / Writes,所以要處理 Data
Migration、Routing 與 Consistency。
假設:
shard = hash(user_id) % 3
增加 Shard 後變成:
shard = hash(user_id) % 4
很多 Key 的結果會改變,因此可能需要搬動大量資料。
這就引出 Consistent Hashing。
今天先理解它的目的,不需要先背演算法:
當 Server / Shard 數量改變時,盡量只重新分配一部分資料。
普通 Modulo:
3 Shards → 4 Shards
↓
大量 Keys 位置改變
Consistent Hashing 希望:
3 Shards → 加入 Shard #4
↓
只有部分 Keys 搬家
可以想像新增一個倉庫時,只搬部分商品,而不是把所有倉庫清空重新分配。
假設 Transfer:
User A → Shard #1
User B → Shard #2
需要:
Shard #1 扣錢
Shard #2 加錢
如果第一步成功、第二步失敗,就不能像單一 Database Transaction 那麼簡單。
這種跨多個 Database / Shard 的 Transaction 叫 Distributed
Transaction。
它可能進一步涉及:
Two-Phase Commit
Saga
Idempotency
Compensation
今天先建立核心觀念:
Sharding 可能把原本簡單的 Database Transaction 變成 Distributed
System Problem。
可以:
Shard #1
├── Primary
├── Replica
└── Replica
Shard #2
├── Primary
├── Replica
└── Replica
Shard #3
├── Primary
├── Replica
└── Replica
兩者分工:
Sharding → Split Data → Scale Storage / Write
Replication → Copy Data → Scale Read / Availability
如果面試官問:
系統有 1 億 User,一台 Database 快放不下,怎麼設計?
不要直接說 Sharding。
先確認 Bottleneck:
Storage?
Read Throughput?
Write Throughput?
Query Latency?
如果真正問題是 Dataset 太大與 Write Traffic 太高,再考慮 Sharding。
例如把 user_id 當 Shard Key 候選,再比較 Range-based 與
Hash-based,並主動說明:
Hot Shard
Cross-Shard Query
Rebalancing
Distributed Transaction
這才是 System Design 思考。
可以問:
1. Cardinality 高嗎?
2. Data Distribution 平均嗎?
3. Traffic Distribution 平均嗎?
4. 常見 Query 是否包含這個 Key?
5. 能否讓相關資料留在同一 Shard?
6. Rebalancing 容易嗎?
7. 容易產生 Hot Shard 嗎?
重點:
Data Distribution 和 Traffic Distribution 都要看。
1. Partition 是什麼?
2. Vertical vs Horizontal Partitioning?
3. Database Sharding 是什麼?
4. Shard 是什麼?
5. Shard Key 是什麼?
6. Replication vs Sharding?
7. Range-based Sharding?
8. Hash-based Sharding?
9. Hot Shard 是什麼?
10. Cardinality 是什麼?
11. 為什麼 country 可能不是好的 Shard Key?
12. Cross-Shard Query 是什麼?
13. Cross-Shard JOIN 為什麼麻煩?
14. Rebalancing 是什麼?
15. 為什麼 hash(key) % N 在新增 Shard 時很麻煩?
16. Consistent Hashing 想解決什麼?
17. Distributed Transaction 是什麼?
18. Replication 和 Sharding 可以一起使用嗎?
19. 1 億 User、一台 DB 放不下時怎麼分析?
20. Shopping System 的 Orders 可以怎麼選 Shard Key?
Sharding:
Large Dataset
↓
Split Data
↓
Multiple Shards
可以改善:
Storage Scaling
Write Scaling
Large Dataset Distribution
但帶來:
Shard Key Selection
Hot Shard
Cross-Shard Query
Cross-Shard JOIN
Rebalancing
Distributed Transaction
Routing Complexity
最重要的一句:
Sharding 不是隨便把 Database 切成很多份,而是根據 Query Pattern 與
Traffic Pattern,讓資料與負載合理分散。
下一篇:
Day 10|CDN:為什麼網站可以讓全球 User 都快速載入圖片與 Static
Files?
在進入 CDN 前,我們會先從零解釋新名詞:
Latency
Static Content
Origin Server
Edge Server
Edge Location
再進入:
CDN
Cache Hit / Miss
Origin Pull
TTL
Cache Invalidation
並比較:
CDN vs Redis Cache
繼續把 System Design 往全球規模推進。