Limited Offer

30% OFF Lifetime Access ($139) with code SYSTEM30

TOPIC #61Intermediate 9 min read

Sharding Strategies (Range, Hash, Directory, Geolocation)

💡
Core Architecture Summary

Compare the 4 primary sharding partitioning schemes: Range-Based, Hash-Based (Modulo/Consistent), Directory-Based (Lookup Table), and Geolocation-Based Sharding.

Key Glossary Concepts in this TopicAll Glossary Terms

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.

The 4 Core Database Sharding Strategies & Tradeoff Matrix 🗂️
100%
Rendering visual architecture flowchart...

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:

text
User_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.
Production Implementation in Big Tech
Pinterest & Discord• Hash-Based Sharding and Multi-Tenant Isolation

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.

Question 1 of 20 answered
#1

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?

Rate This Architecture Chapter4.9 / 5.0 (38 ratings)

How clear and staff-actionable was this system breakdown?