TOPIC #61Intermediate 9 min read

Sharding Strategies (Range, Hash, Directory, Geolocation)

CSD
CompleteSystemDesign Editorial
Report an issue
Key takeawayCore 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
Interactive Lab · 🗂️ 4 Sharding Strategies RacerFull lab guide

Range vs Hash vs Directory vs Geo Sharding

Auto-incrementing orders flow through each strategy; measure write skew and range-scan broadcast cost.

Shard 1

0 rows

0.0% of writes

Shard 2

0 rows

0.0% of writes

Shard 3

0 rows

0.0% of writes

Shard 4

0 rows

0.0% of writes

Orders placed

0

Range scan WHERE id BETWEEN recent

4 shards touched

Write skew

max 0% balanced

Murmur/FNV(key) mod N: perfectly uniform writes, but range scans scatter to every shard. Under Range sharding on auto-increment IDs, 100% of fresh writes hit the newest shard while history sits idle; Hash fixes that but converts scans into broadcasts; Directory pays one lookup-hop for surgical VIP isolation; Geo buys data-sovereignty compliance at the price of uneven regional growth.

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%
Touchpad: Pinch to zoom • Drag to pan
Rendering visual architecture flowchart...

01.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 Check

Exercise 1 of 2 • Test your architectural comprehension.

Exercise 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 ChapterFeedback & Rating

How clear and actionable was this distributed systems breakdown?

Interactive Engineering Workbenches: