Contents

Architecture & System Design › System Design Fundamentals · also in Database Operations

Sharding

Splitting data across databases by key.

Also known as: database sharding, horizontal partitioning, shards, shard, data sharding

Sharding splits a dataset across multiple independent databases (shards), each holding a subset of the data, so that no single machine has to store everything or handle all the traffic. It’s the main way to scale writes and data size beyond one server. (Replicas scale reads, sharding scales writes and storage.)

user_id → shard function
 shard 0: users A (ids hash to 0)      shard 1: users B      shard 2: users C
 each shard: its own server(s), its own storage, its own slice of the data

Choosing a shard key and strategy

The shard key (such as user_id or tenant_id) decides where each row lives. Common strategies:

StrategyHowProsCons
RangeKey ranges per shard (ids 1-1M, 1M-2M)Range queries are efficientHot spots (newest range gets all writes)
Hashhash(key) % N or consistent hashingEven distributionRange queries scatter, and resharding is harder (consistent hashing)
Directory / lookupA mapping table says which shard holds each keyFlexible, easy to move dataThe directory is a critical component
Geographic / tenant-basedWhole tenants or regions per shardData locality, isolationUneven sizes

A good shard key spreads data and traffic evenly, and keeps most queries on a single shard.

What gets harder

  • Cross-shard queries and joins are slow, complex or unsupported. You design so common queries include the shard key.
  • Transactions across shards need distributed coordination, or you avoid them (distributed transactions, sagas).
  • Unique constraints and auto-increment IDs can’t rely on a single database (distributed ID generation).
  • Hot shards: one celebrity user or tenant overloads a shard (hot partitions).
  • Rebalancing and resharding: adding shards means moving data while serving traffic, which is operationally difficult (rebalancing).
  • Operational complexity: backups, schema migrations, monitoring and failure handling for many databases.
  • Application complexity: routing logic, shard-aware code (or a proxy or middleware that hides it).

Alternatives to try first

Sharding is costly and mostly irreversible. Before it, exhaust:

When you do shard, pick the key very carefully, since changing it later means migrating all the data.