Sharding
Splits data across multiple database nodes by a shard key so write/read capacity scales beyond one machine.
Interview tip Lead with a 30-second definition, then one real system example and name 2–3 designs where Sharding is non-negotiable.
① What it is (30 seconds)
Splits data across multiple database nodes by a shard key so write/read capacity scales beyond one machine.
② How it works in system design
Choose shard key (user_id, tenant_id). hash(key) → shard index. Cross-shard queries expensive — design access patterns per shard. Rebalancing moves ranges when adding nodes.
Typical placement
Client→Edge / Gateway→Sharding→Services→Data stores
③ Concrete system design example
Scenario: Twitter timelines: tweets sharded by user_id. Write tweet to author shard; fan-out workers read follower lists and write to follower shard inboxes.
④ Important interview Q&A
| Question | Answer |
|---|---|
| Shard key selection? | High cardinality, even distribution, aligns with query patterns — avoid cross-shard joins. |
| Hot shard? | Split range, add read replicas, or reshard by sub-key. |
| Shard vs read replica? | Replica scales reads on same data; sharding splits data for write scale. |
⑤ Seen in these system designs
- Twitter — timeline shards
- WhatsApp & Chat — message shards
- Distributed Search — index shards
In interviews, after explaining the concept, say: "This shows up directly in …" and link two designs.
⑥ Revision checklist
- Shard key justified
- Cross-shard query plan
- Rebalance strategy
- Hot shard mitigation