Sharding Strategies (Range, Hash, Directory, Geolocation)
Compare the 4 primary sharding partitioning schemes: Range-Based, Hash-Based (Modulo/Consistent), Directory-Based (Lookup Table), and Geolocation-Based Sharding.
The 4 Core Database Sharding Strategies & Tradeoff Matrix 🗂️
Structural breakdown of Range-Based (ordered lookups), Hash-Based (uniform distribution), Directory-Based (flexible lookup table), and Geolocation-Based (GDPR compliance) sharding.
01.1. The 4 Fundamental Sharding Strategies
Database architects choose between four primary sharding algorithms based on query access patterns, range scan requirements, and regulatory data compliance:
1. Range-Based Sharding
Data is partitioned based on contiguous ranges of an ordered key (e.g., user_id 1 - 1,000,000 on Shard 1; 1,000,001 - 2,000,000 on Shard 2; or timestamps by year/month).
- Pros: Range queries (
WHERE order_date BETWEEN '2026-01-01' AND '2026-01-31') are routed to a single physical shard. - Cons: Severe Hotspots. If sharded by auto-incrementing ID or timestamp, 100% of all current write traffic hits the newest shard, leaving older shards completely idle.
2. Hash-Based Sharding (Algorithmic Sharding)
Applies a hash function (MD5, MurmurHash3) to the shard key:
Target Shard = Murmur3(user\_id) \pmod N
- Pros: Uniform, random distribution of rows across all shards, completely eliminating write hotspots.
- Cons: Range scans cannot be localized (must execute expensive Scatter-Gather across all shards). Adding or removing a shard node invalidates previous modulos (
N → N+1), requiring resharding the entire dataset unless Consistent Hashing is used.
3. Directory-Based Sharding (Lookup Table)
Maintains a centralized, dynamic mapping service or lookup table (stored in Redis, etcd, or ZooKeeper) that explicitly maps entity IDs to specific physical shard IDs:
textUser_42 -> Shard_1 (10.0.1.10) User_1009 -> Shard_2 (10.0.1.11) VIP_Customer -> Shard_9 (Dedicated High-Memory Server)
- Pros: Extreme flexibility. You can migrate a single massive "celebrity" or enterprise customer to a dedicated isolated shard without moving any other customer data.
- Cons: The lookup directory adds a network hop to every query and can become a single point of failure if not cached aggressively.
4. Geolocation-Based Sharding
Partitions data based on the physical geographic region or nationality of the user (e.g., European users on EU shards in Frankfurt, American users on US shards in Virginia).
- Pros: Sub-millisecond local read/write latency and strict compliance with GDPR, HIPAA, and Data Sovereignty regulations (ensuring European user data never leaves EU physical borders).
- Cons: Uneven data growth across regions; complex cross-region queries for international interactions.
⚖️Architectural Trade-offs & Production Realities
Architectural Advantages
- Hash sharding guarantees balanced CPU and disk utilization across all cluster nodes.
- Directory sharding allows surgical rebalancing of VIP tenants to dedicated hardware.
- Geo-sharding provides both low latency and statutory data residency compliance.
Trade-offs & Constraints
- Range sharding creates massive write hotspots on the latest timestamp/ID partition.
- Hash sharding makes range queries scatter-gather broadcasts across all shards.
Pinterest uses Hash-Based Sharding across thousands of MySQL instances, hashing `user_id` via MurmurHash to pin all of a user's boards and pins to a single shard. Discord uses Hash Sharding on `guild_id` (Server ID) in ScyllaDB to ensure high-velocity chat streams are evenly distributed across its distributed cluster.
🎯 Staff+ Engineering Takeaways
- Range: Fast range scans; vulnerable to write hotspots on new IDs/timestamps.
- Hash: Uniform write distribution; scatter-gather for range scans.
- Directory: Central lookup service; extreme flexibility for tenant migration.
- Geo: Partitions by country/region; delivers low latency and GDPR compliance.
Topic Knowledge Assessment 🧠
Step through 2 scenario questions to test your staff-level grasp.
Why does Range-Based Sharding on an auto-incrementing Primary Key (e.g. 1-1000 on Shard 1, 1001-2000 on Shard 2) create a severe write hotspot in production?
How clear and staff-actionable was this system breakdown?