System design2 min

Sharding vs Partitioning

When a database becomes too large to fit on a single machine, or when the read/write load overwhelms the CPU, you must scale horizontally. Breaking a massive database into smaller, manageable pieces is known as Partitioning.

There are two main ways to partition a database.

1. Vertical Partitioning

Vertical partitioning means splitting a table's columns into separate tables.

If you have a massive Users table, you might put the frequently accessed columns (ID, Username, PasswordHash) on one server, and the rarely accessed columns (Bio, AvatarBlob, LastLoginData) on a completely different server.

Pros: Reduces the amount of I/O required to fetch hot data. Cons: If a query needs data from both partitions, the application must make two network calls and join the data in memory.

2. Horizontal Partitioning (Sharding)

Sharding means splitting a table's rows across multiple different database servers (called shards). All shards have the exact same schema, but hold completely different subsets of data.

For example, you could put all users whose names start with A-M on Shard 1, and N-Z on Shard 2.

Sharding Strategies

  1. Algorithmic Sharding (Hash Sharding): You apply a hash function to the partition key (e.g., hash(UserID) % NumberOfShards).

    • Pros: Data is distributed very evenly.
    • Cons: Adding or removing shards is a nightmare (requires massive data migrations). Consistent Hashing is often used to mitigate this.
  2. Range Sharding: Data is divided based on specific ranges (e.g., UserIDs 1 to 1,000,000 go to Shard A; 1,000,001 to 2,000,000 go to Shard B).

    • Pros: Excellent for range queries.
    • Cons: High risk of "hotspots". If the newest users are the most active, Shard B will be crushed with traffic while Shard A sits idle.
  3. Directory Sharding: A separate "Lookup Database" keeps track of exactly which shard holds which piece of data.

    • Pros: Incredibly flexible. You can move data around easily.
    • Cons: The lookup table becomes a single point of failure and a bottleneck.

The Pain of Sharding

Sharding is often considered a last resort because it introduces immense complexity:

  • No Cross-Shard Joins: You cannot easily JOIN data if Table A is on Shard 1 and Table B is on Shard 2.
  • Complex Transactions: Maintaining ACID properties across distributed nodes is notoriously difficult and slow (e.g., Two-Phase Commit).
  • Uneven Growth: Some shards will inevitably grow faster than others, requiring complex re-balancing operations.