Database Sharding
Concept
Replication solves the problem of too many Reads, but all Writes still hit a single Primary server. When your application scales to millions of writes per second, or the dataset becomes larger than a single hard drive (e.g., 50 Terabytes), you must use Sharding (Horizontal Partitioning).
Sharding splits your massive database into smaller, independent databases (shards), each holding a specific slice of the data.
Mental Model
How It Works
- The Shard Key: You must choose a column (like
user_idorregion) to act as the Shard Key. - The Routing Logic: When the app wants to insert or find a user, it passes the Shard Key to a router. The router uses a mathematical algorithm to determine which physical server holds that data.
- The Shards: Shard 1 only knows about Users 1-1000. It has no idea that Shard 2 or User 2000 exists.
Trade-Offs
Sharding is the most complex database operation you can perform. It should be avoided until absolutely necessary.
Pros:
- Infinite Scalability: You can scale writes and storage linearly simply by adding more commodity servers.
Cons:
- No Cross-Shard JOINs: If you shard the
Userstable byuser_id, and theOrderstable byorder_id, you cannot perform a SQLJOINto find all orders for a user. The data lives on two entirely different machines. - Complex Transactions: ACID transactions across multiple shards require complex protocols like Two-Phase Commit (2PC), which are incredibly slow and error-prone.
- Resharding: If Shard 1 fills up its hard drive, you have to split it into Shard 1A and Shard 1B. Moving terabytes of data between servers without downtime is extremely dangerous.
Real-World Usage
- NoSQL: Databases like Cassandra and MongoDB are designed to be sharded out of the box. They handle the routing and data rebalancing automatically.
- SQL (PostgreSQL / MySQL): Native sharding is historically very difficult. Companies often implement sharding at the application layer (writing code to route connections to
db_shard_1ordb_shard_2). Today, tools like Vitess (for MySQL) or Citus (for Postgres) provide transparent sharding for relational DBs.
Interview Questions
Q: You decide to shard your multi-tenant SaaS application (like Slack) by tenant_id (Company ID). What is the potential danger of this Shard Key?
A: The danger is the “Celebrity Problem” (Hotspotting). If you shard by tenant_id, all data for a specific company lands on one single server. If 99 of your clients are small startups with 10 employees, but your 100th client is Microsoft with 100,000 employees, the shard holding Microsoft’s data will be overwhelmed with traffic and run out of disk space, while the other shards sit idle. You have failed to distribute the load evenly.
Q: How would you fix the Celebrity Problem in the previous example?
A: Instead of using a logical partition like tenant_id, I would use a completely randomized algorithmic shard key, such as applying Consistent Hashing to the message_id. This guarantees that Microsoft’s 100,000 messages are perfectly distributed across all the shards, preventing any single server from becoming a hotspot. (The trade-off is that fetching all messages for Microsoft now requires a scatter-gather query to all shards).