Database Sharding - Complete Deep Dive
Prerequisites: Database Concepts, Consistent Hashing Used in: Key-Value Store, Chat System, URL Shortener, any system beyond single-DB scale
What is Sharding?
Sharding splits a large database into smaller pieces (shards), each stored on a different server. Each shard holds a subset of the data.
Real-world analogy: A library with 1 million books. Instead of one giant room (single DB), you build 10 rooms (shards), each holding 100K books organized by authorβs last name initial (A-C in Room 1, D-F in Room 2, etc).
Without sharding:
1 billion rows β 1 database server β CPU maxed, disk full, queries slow
With sharding (4 shards):
Shard 1: rows where userId 0-249M (Server 1)
Shard 2: rows where userId 250M-499M (Server 2)
Shard 3: rows where userId 500M-749M (Server 3)
Shard 4: rows where userId 750M-1B (Server 4)
When Do You Need Sharding?
NOT yet: If your database is < 1TB and < 10K queries/sec, a single server with read replicas is enough. Donβt shard prematurely.
Time to shard when:
- Single DB exceeds disk capacity (> 1-5 TB)
- Write throughput exceeds single server (> 10-50K writes/sec)
- Read replicas arenβt enough (write-heavy workload)
- Query latency grows despite indexes
Rule of thumb: Scale reads with replicas first. Shard only when writes are the bottleneck.
Sharding Strategies
1. Hash-Based Sharding
Hash the shard key and mod by number of shards.
shard_id = hash(userId) % num_shards
hash("user_001") % 4 = 2 β Shard 2
hash("user_002") % 4 = 0 β Shard 0
hash("user_003") % 4 = 3 β Shard 3
Pros: Even distribution (if hash function is good). No hotspots from sequential keys. Cons: Range queries impossible (canβt ask βall users with ID 100-200β β theyβre scattered). Adding shards requires rehashing (use consistent hashing to fix this).
Best for: Key-value lookups, user data, sessions.
2. Range-Based Sharding
Assign contiguous ranges to each shard.
Shard 0: userId 0 - 999,999
Shard 1: userId 1,000,000 - 1,999,999
Shard 2: userId 2,000,000 - 2,999,999
Pros: Range queries work naturally (βget all users from 1M to 1.5Mβ). Simple to understand. Cons: Hotspots if data isnβt evenly distributed (new users all hit the last shard). Uneven shard sizes over time.
Best for: Time-series data (shard by month), geographic data (shard by region).
3. Directory-Based Sharding
A lookup table maps each key to its shard.
Lookup table:
user_001 β Shard 2
user_002 β Shard 0
user_003 β Shard 1
Query: look up shard in directory, then query that shard.
Pros: Flexible β can move any key to any shard. Rebalancing is easy (update directory). Cons: Directory is a single point of failure. Extra hop for every query.
Best for: When you need fine-grained control over data placement.
4. Geographic Sharding
Shard by userβs geographic region.
Shard "US": all US users (servers in us-east-1)
Shard "EU": all EU users (servers in eu-west-1)
Shard "APAC": all Asia-Pacific users (servers in ap-south-1)
Pros: Data locality (low latency for users near their shard). Compliance (EU data stays in EU for GDPR). Cons: Uneven distribution (US shard might be 5x larger). Cross-region queries are expensive.
Choosing a Shard Key
The shard key determines which shard a row goes to. Itβs the most important decision in sharding.
Good shard key properties:
- High cardinality (many unique values β userId is good, country is bad)
- Even distribution (no single value dominates)
- Matches query patterns (queries include the shard key β single-shard lookup)
| Shard Key | Good For | Bad For |
|---|---|---|
| userId | User-centric apps (get all data for one user) | Cross-user queries |
| orderId | Order lookups | βAll orders for user Xβ (need to scan all shards) |
| timestamp | Time-series | Hot partition (all writes hit current time shard) |
| country | Geo-partitioning | Uneven (US has 60% of traffic) |
| hash(userId) | Even distribution | Range queries, debugging |
Hot Partition Problem
The problem: One shard gets disproportionately more traffic than others.
Examples:
- Celebrity user (10M followers β their shard gets hammered on every post)
- Time-based shard key (all current writes go to βtodayβ shard)
- Popular product (flash sale β one productβs shard overwhelmed)
Solutions:
| Solution | How |
|---|---|
| Add salt/suffix to key | shard_key = userId + "_" + random(0,9) β spreads across 10 sub-shards. Reads must fan-out. |
| Dedicated shard for hot keys | Detect hot keys β move to a dedicated high-capacity shard |
| Caching in front | Cache hot data aggressively β most reads donβt hit the shard |
| Further split the hot shard | Break one shard into 4 smaller ones |
Cross-Shard Queries (The Pain)
The biggest downside of sharding: Queries that span multiple shards are expensive.
"Get the top 10 orders across all users sorted by amount"
β Query ALL shards β each returns its top 10 β merge β pick global top 10
β Scatter-gather pattern (slow, expensive)
How to handle:
| Pattern | When | Trade-off |
|---|---|---|
| Denormalize | Store duplicate data to avoid cross-shard JOINs | Write amplification |
| Application-side join | Query both shards, merge in app code | Complex, higher latency |
| Scatter-gather | Fan out query to all shards, aggregate results | Latency = slowest shard |
| Secondary index (Elasticsearch) | Sync data to a search index that isnβt sharded the same way | Extra infra, lag |
Best practice in interviews: βIβd shard by userId so all of a userβs data is on one shard. For cross-user queries (leaderboards, analytics), Iβd use a separate denormalized read store or Elasticsearch.β
Rebalancing
When shards become uneven (one grows too large), you need to move data between shards.
Approaches:
| Strategy | How | Downtime? |
|---|---|---|
| Fixed partitions | Pre-create many partitions (e.g., 1000), assign groups to nodes. Rebalance = move partition groups. | Minimal |
| Dynamic splitting | When a shard exceeds size threshold, split into two. | Zero (if background) |
| Consistent hashing | Add virtual nodes for new server, only ~K/N keys move. | Zero |
DynamoDB: Auto-splits partitions when they exceed 10GB or 3000 RCU/1000 WCU. You donβt manage this manually.
Cassandra: Uses consistent hashing with virtual nodes. Adding a node = automatic rebalancing.
Sharding vs Replication
| Β | Sharding | Replication |
|---|---|---|
| Purpose | Scale writes + storage | Scale reads + availability |
| Data | Each shard has DIFFERENT data | Each replica has SAME data |
| Failure | Shard down = that data unavailable | Replica down = other replicas serve |
| Complexity | High (routing, cross-shard queries) | Medium (replication lag, failover) |
| When | Write-heavy, large data | Read-heavy, high availability |
You usually need BOTH: Shard for writes/storage, replicate each shard for reads/HA.
Shard 1: Primary β Replica 1A, Replica 1B
Shard 2: Primary β Replica 2A, Replica 2B
Shard 3: Primary β Replica 3A, Replica 3B
Common Interview Questions
Q: βHow would you shard this database?β A: βIβd shard by [entity]Id using hash-based sharding for even distribution. All data for one [entity] lives on one shard, so most queries are single-shard. For cross-entity queries, Iβd use a secondary index (Elasticsearch) synced via CDC.β
Q: βWhatβs the risk of sharding by timestamp?β A: βHot partition. All current writes go to the βnowβ shard. Fix: compound key = timestamp + random suffix, or hash-based sharding with separate time-series index for time queries.β
Q: βWhen would you NOT shard?β A: βWhen your data fits on one server (< 1TB), when writes are < 10K/sec, or when your queries frequently need cross-entity JOINs (sharding makes JOINs very expensive).β
Q: βHow do you handle transactions across shards?β A: βYou canβt do ACID transactions across shards easily. Options: 1) Design so transactions are single-shard (shard by the transactional entity). 2) Use saga pattern for cross-shard operations. 3) Use two-phase commit (slow, avoid).β
| β Back to Fundamentals | β Consistent Hashing | Next: CDN β |