Skip to content
Prompt Words
History

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.

    Rows 3 3 3

    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

    • Notion: Split its Postgres data into 480 logical shards spread over 32 physical databases, with the workspace ID as the shard key.
    • MongoDB: A sharded cluster spreads a collection's documents across shards by a shard key, and the mongos router sends each query to the right shard.

    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.