1. What It Is
You reach for sharding when a single primary can't keep up with write QPS or disk growth β but it's expensive to undo, so we only shard after caches and read replicas are exhausted.
What:
Horizontal partitioning (sharding) splits single tables horizontally into separate database servers; vertical partitioning splits table columns into separate tables.
Primary purpose:
Scale write throughput and storage capacity beyond the physical hardware limits of a single master server.
Usually used for:
High-volume transactional databases, heavy social timeline layers, and multi-tenant SaaS backends.
2. Core Mental Model
The shard key is the most important decision you'll make β it determines which queries hit one node and which fan out to every shard:
π The Sharding Key Rule
Select a key (e.g. `user_id` or `tenant_id`) that routes over 90% of your common queries to a single shard, entirely avoiding cross-shard lookups.
βοΈ Even Distribution
Choose keys with high cardinality to distribute storage and write QPS uniformly across your database node pool.
𧬠Application Routing
Move routing logic either into an application middleware library (e.g., Vitess) or deploy a dedicated proxy tier (e.g., Citus).
In the room
Interviewers love asking "how would you query across shards?" If your shard key is user_id but the product needs a global leaderboard, you have a problem. Name the shard key early and explain which queries stay single-shard vs which need a secondary index or denormalized table.
3. Why It Matters in HLD
Sharding is the last resort after vertical scaling, caching, and read replicas β but when we need it, the shard key is everything. Three lenses:
Needed When:
Database write volume outgrows physical disk write I/O speeds, or total storage requirements exceed single SSD arrays.
Avoids:
Primary database CPU locks during writes, slow sequential table scans, index cache churns, and massive hardware migration costs.
Optimizes For:
Write concurrency scaling, distributed storage limits, data locality latencies, and transaction boundary isolations.
4. Architecture & Data Flow
Walk partition strategies as interview steps. Step 1 β Range: contiguous key ranges per shard β great for time-series, risky for hot ranges. Step 2 β Hash: hash(key) mod N or consistent hash ring β even spread, cross-range queries expensive. Step 3 β Directory: lookup service maps keys to shards β flexible moves, extra hop. Step 4 β Geo: pin data to regions for GDPR or latency. Step 5 β Name the hot query: pick the strategy that keeps 90%+ of queries single-shard.
In the room
Name your shard key in the first five minutes. If it is user_id but the product needs a global leaderboard, say how you would denormalize or scatter-gather β do not wait for the interviewer to catch the mismatch.
5. Key Characteristics
Each strategy determines whether queries hit one shard or fan out to all of them β we compare before committing:
- Partition strategy determines whether queries hit one shard or fan out to all of them:
| Strategy | Routing Model | Primary Benefit | Severe Cost |
|---|---|---|---|
| Range-Based | Map key intervals (e.g. A-E, F-L) to specific shard nodes. |
| Creates severe hotspots (e.g. all updates on active suffix keys). |
| Hash-Based | Apply hash function to partition key modulo node count: hash(key) % N. | Ensures uniform data distribution across nodes. | Adding/removing nodes remaps almost 100% of keys (requires Consistent Hashing). |
| Directory (Lookup) | Maintain a routing table mapping keys or ranges to shard IDs. | Flexible rebalancing and hot-key relocation without changing hash function. |
|
| Geo-Located | Route user tables to nodes inside their geographic region. |
| Complex cross-region backups and failovers. |
6. Strategic Tradeoffs
Horizontal partition buys write scale at the cost of cross-shard complexity β we state both:
| Benefit | Cost |
|---|---|
| Horizontal Scalability (bypasses the physical scale limit of a single server, scaling database QPS linearly) |
|
| Small Blast Radius (if shard node A crashes, only 10% of users are degraded, keeping the remaining pool alive) | Loss of ACID Atomicity (transactions span multiple shards, requiring complex 2-Phase Commit protocols) |
7. Failure / Bottleneck Awareness
Hot shards and scatter-gather are the classic sharding interview traps β we volunteer fixes:
Problem: Sharding by `user_id` works uniformly until a celebrity account (e.g. 100M followers) executes massive read/write activities. The destination shard node hosting that celebrity collapses under the QPS load, while other shards remain idle.
Mitigation: Implement a hybrid key model: append random salt suffixes to hot user keys (e.g. `celebrity_123:salt_0` to `celebrity_123:salt_9`) to spread their data across multiple shards.
Problem: Executing queries without specifying the partition sharding key forces the database coordinator proxy to query all shards concurrently and merge the results. This destroys throughput and maximizes tail latency.
Mitigation: Ensure all high-volume transactional queries include the sharding key, or maintain secondary indexes mapped in a dedicated lookup directory table.
8. Common HLD Usage
These production shard keys show how real teams co-locate data with access patterns:
| Production System | Sharding Key Selection | Architectural Rationale |
|---|---|---|
| Twitter Timeline Shards | user_id (Hash-sharded) |
|
| E-Commerce Ledger | tenant_id (Lookup-sharded) | Multi-tenant B2B platforms isolate customers completely, routing all tenant requests to dedicated databases. |
| URL Shortener | short_code (Directory lookup) | A directory table maps vanity codes to shard IDs so hot codes can be relocated without re-hashing the key space. |
9. Decision Signals
Shard only when measured write QPS or disk size exceeds single-node limits β not because DAU sounds big:
- Your database storage volume exceeds 2 Terabytes, or write requirements exceed 10,000 write QPS.
- You can easily isolate users/tenants data into disjoint sets with zero overlap requirements.
- You need to geographical locate data in specific datacenters to comply with data privacy policies (GDPR).
11. Deep Dive (Optional)
Vertical Partitioning vs Horizontal Sharding
It is critical to distinguish between these two sizing dimensions:
1. Horizontal Sharding
Splits table **rows** across multiple database servers. The schema remains identical across all shards, but each shard stores a unique subset of rows. Ideal for scaling QPS and disk capacities.
2. Vertical Partitioning
Splits table **columns** into separate tables stored on different disks or servers. For example, moving a large text column (`user_bio`) or blob column (`profile_photo_blob`) out of the core `users` table into a secondary `user_profiles` table.
**Benefit**: Reduces the size of the primary table's memory-mapped pages, significantly boosting scan speeds for common narrow queries (e.g. `SELECT username FROM users`).
When Sharding Is Premature
Sharding is an operational tax β cross-shard JOINs disappear, distributed transactions get harder, and rebalancing shards is a multi-week project. Before partitioning horizontally, exhaust vertical scaling and simpler wins: read replicas, connection pooling, query/index tuning, caching hot keys, and archiving cold data. A single well-tuned PostgreSQL instance often carries millions of DAU; interviewers reward candidates who shard only when measured QPS or disk size forces it, not because the problem mentions "scale."
Cross-Shard Operations
Once sharded, queries that span keys on different partitions become expensive:
- Scatter-gather: Fan out to all shards, merge results in the app layer β acceptable for rare admin reports, deadly for user-facing search at high QPS.
- Global secondary indexes: Maintain a separate index shard keyed differently (e.g., email β user_id lookup) with async replication β adds write amplification.
- Co-locate related data: Choose shard key so hot queries hit one partition β user_id shards user posts, orders, and settings together.
Design the shard key around your hottest access path, not around entity nouns in isolation.
Review
How helpful was this walkthrough?
Click a star to rate. We actively use this feedback to refine and update our system design content.
Discussion
Share your thoughts, ask questions, or help others.