Sharding
Also called: horizontal partitioning, partitioning, shards.
Splitting one big dataset across several databases by a shard key, so each database (a shard) holds only its part. All reads and writes for one key go to one shard, which spreads both the data and the write load. Queries that need many shards get slower and harder to write.
hash(user_id) % 6 picks a logical shard, and a lookup table says which database (shard) holds it.
Shard 0: 3 users · Shard 1: 3 users · Shard 2: 3 users
9 users on 3 shards, placed by hash(user_id) % 6 and a lookup table.
Say it in a prompt
Shard the orders table by tenant_id: logical_shard = hash(tenant_id) % 64, and a lookup table maps the 64 logical shards onto 8 Postgres databases. Every query must include tenant_id and go through one shard-router module. Moving a shard means copying it and updating the lookup table, with no rehashing. Vague vs precise prompt
Vague prompt
our orders table is huge, split it up so it scales Typical resultCreates orders_2024 and orders_2025 tables on the same server. The disk and the write load stay on one machine.
Precise prompt
Shard orders by tenant_id: logical_shard = hash(tenant_id) % 64, with a lookup table mapping the 64 logical shards onto 8 Postgres databases. All queries include tenant_id and go through one shard-router module. Typical resultEach database holds about an eighth of the tenants and their writes, one tenant's queries touch one database, and a shard can be moved by editing the lookup table.
Seen on
You might describe it as
- split the customers across several databases
- each database only holds some of the users
- one table is too big for one server
Not to be confused with
- Read replica
Sharding gives each database a different slice of the data; a read replica is a full copy that only serves reads.