Sharding Strategies

⭐ Interview Importance: MEDIUM
⏱️ Revision Time: 3 min

Concept

If you make the decision to Shard your database, you must choose a mathematical strategy to distribute the data.
The column you choose to divide the data by (e.g., user_id or country) is called the Shard Key.
Choosing the wrong Shard Key will destroy your database performance.

1. Range-Based Sharding

You divide the data into sequential chunks based on the Shard Key.

  • Shard A: user_id 1 to 10,000
  • Shard B: user_id 10,001 to 20,000

Pros: Extremely easy to implement. Range queries (WHERE user_id BETWEEN 5000 AND 6000) are blazingly fast because they only hit a single shard.
Cons: The Hotspot Problem. If users are created sequentially, all brand-new users (IDs 20,001+) will be dumped exclusively into Shard C. Shard A and Shard B will sit idle, while Shard C gets crushed by 100% of the write traffic, completely defeating the purpose of horizontal scaling.

2. Hash-Based Sharding

You take the user_id and run it through a Hash Function (e.g., user_id % 4).

  • ID 1 % 4 = 1 (Goes to Shard 1)
  • ID 2 % 4 = 2 (Goes to Shard 2)
  • ID 5 % 4 = 1 (Goes to Shard 1)

Pros: Perfect, even distribution. It completely eliminates Hotspots. Write traffic is mathematically distributed perfectly across all servers.
Cons: Range queries are completely destroyed. If you want to SELECT * WHERE user_id BETWEEN 1 AND 10, the router must execute the query across all 4 shards simultaneously (Scatter-Gather), wait for all 4 servers to respond, and sort the data in memory. Furthermore, if you need to add a 5th server, the modulo math (% 5) changes, meaning you must physically move almost every single row to a new server (unless you use Consistent Hashing).

3. Directory-Based Sharding (Lookup Table)

Instead of doing math, you create a dedicated “Lookup Database” that acts as a central map.

  • “Where is User A?” -> Check Lookup Table -> “Shard 3”.
  • “Where is User B?” -> Check Lookup Table -> “Shard 1”.

Pros: Absolute flexibility. You can manually move extremely active users to empty shards without changing any math.
Cons: The Lookup Database becomes a massive Single Point of Failure and a latency bottleneck. Every single query requires two database hits (one to check the map, one to fetch the data).

4. Geo-Based Sharding (Multi-Region)

You partition the data based on geography to reduce physical light-speed latency.

  • Shard A (Deploy in Frankfurt): European Users
  • Shard B (Deploy in Tokyo): Asian Users

Pros: Data Localization (GDPR compliance) and massive latency reduction (users in Tokyo hit a server 10 miles away instead of 5,000 miles away).
Cons: The “Celebrity Problem”. If a European user creates a post that goes viral in Asia, millions of Asian users are forced to fetch data from the Frankfurt shard, crushing the transcontinental network cable and slowing the application to a crawl.

Interview Questions

Q: You are building a Slack clone. You must choose a Shard Key for the messages table. You can shard by message_id (Hash) or by workspace_id (Directory/Hash). Which do you choose?
A: You must shard by workspace_id.
If you shard by message_id, the messages for a single chat room will be scattered randomly across 10 different database servers. To load a single chat channel, your Node.js server must query all 10 shards, wait for them to respond, and stitch the messages together.
If you shard by workspace_id, all data relating to a specific company (users, channels, messages) is guaranteed to live on the exact same physical server. This allows you to use standard JOIN queries and keeps latency incredibly low. (This pattern is called Tenant-Based Sharding).