Database Sharding and Partitioning
How a partition key splits a table, why sharding the pieces across machines makes that key decide every query, and how to reshard without downtime.
Partitioning splits one logical table into pieces by the value of a partition key; sharding puts those pieces on different machines, so each database owns a subset of the rows and none of them sees the whole table. From then on the key decides everything: a query that includes it goes to one shard, a query without it has to ask every shard, a key whose new values all land in one place creates a hot shard, and changing the key later means moving all of the data.
Context
For most of the 2000s, scaling a database meant buying a bigger server. The large web companies hit the ceiling first and split their MySQL databases by user ID by hand: Flickr, Facebook, and YouTube, whose sharding layer became Vitess (built 2010, open-sourced 2012). Instagram did the same on PostgreSQL in 2012 with thousands of logical shards packed onto a few physical servers. In parallel, the systems built to be distributed from day one made partitioning internal: Google Bigtable (2006) splits tables into row ranges, Amazon Dynamo (2007) hashes keys onto a ring, and MongoDB added sharding in version 1.6 (2010). Today Spanner, CockroachDB and DynamoDB shard automatically, and PostgreSQL has had built-in declarative partitioning on a single server since version 10 (2017), with hash partitioning since 11.
You have met it as the "partition key" you must choose when creating a DynamoDB table, as sh.shardCollection() in MongoDB, as Kafka topic partitions, and in the simplest form, as a PostgreSQL table split by month:
CREATE TABLE events (
id bigint NOT NULL,
created_at timestamptz NOT NULL,
payload jsonb
) PARTITION BY RANGE (created_at);
CREATE TABLE events_2026_09 PARTITION OF events
FOR VALUES FROM ('2026-09-01') TO ('2026-10-01');
-- the planner reads only events_2026_09 (partition pruning)
SELECT count(*) FROM events WHERE created_at >= '2026-09-15';
-- retention: drop a whole month instantly instead of DELETE + VACUUM
DROP TABLE events_2025_09;- Partition key / shard key
- The column (or columns) whose value decides which piece a row belongs to. In DynamoDB and Cassandra it is literally called the partition key.
- Partition vs shard
- A partition is a piece of a table; a shard is a piece that lives on its own server. Partitioning can happen inside one database; sharding always crosses machines.
- Horizontal vs vertical
- Horizontal partitioning splits rows (users 1-1M here, 1M-2M there). Vertical partitioning splits columns or whole tables onto different databases.
- Router
- Whatever maps a key to a shard: application code, a proxy such as Vitess or mongos, or the database itself.
- Scatter-gather
- Sending a query to every shard and merging the answers, because it does not name a shard key.
- Resharding
- Changing how keys map to shards, for example from 4 to 8 servers, which means physically moving rows.
Why it matters
A single database server eventually runs out of something: write throughput, disk, memory for the working set, or maintenance time, when an index rebuild or a vacuum on a multi-terabyte table takes all weekend. Partitioning on one server buys a lot of headroom cheaply. Sharding buys almost unlimited headroom and costs you joins, transactions and unique constraints across shards, plus a key choice that is very expensive to undo. Teams that shard too early pay the cost without needing the capacity; teams that pick the key badly get a cluster where one shard is on fire and the rest are idle.
How a query finds its shard
Once rows are spread over several databases, something has to map each query to the right one. The router takes the shard key from the query, computes which shard owns it, and sends the query there. If the query does not contain the shard key, the router has no way to narrow it down and must send it to every shard, wait for all of them, and merge the results. That scatter-gather query is only as fast as the slowest shard, and its cost grows with the number of shards.
Most production routers do not map keys straight to servers. They hash the key into a fixed, generous number of logical shards and keep a small table from logical shard to physical database. Instagram used several thousand logical shards as PostgreSQL schemas; Notion started in 2021 with 480 logical shards on 32 physical databases. Rebalancing then means moving whole logical shards and editing the map, never rehashing individual rows. It is the same idea as Redis Cluster's 16,384 slots.
const LOGICAL_SHARDS = 1024 // chosen once, never changed
// logical shard -> physical database. Rebalancing = editing this map.
const shardMap: string[] = await loadShardMap() // 1024 entries
const pools: Record<string, Pool> = connectAll(shardMap)
export function logicalShard(tenantId: string): number {
return murmur3(tenantId) % LOGICAL_SHARDS
}
export function dbFor(tenantId: string): Pool {
return pools[shardMap[logicalShard(tenantId)]]
}
// every query for one tenant hits one database
const db = dbFor(req.tenantId)
await db.query('SELECT * FROM invoices WHERE tenant_id = $1',
[req.tenantId])Range, hash or directory
There are three ways to turn a key into a shard, and they fail in different places. Range partitioning keeps neighbouring keys together, which makes range scans cheap, and turns any key that only grows (a timestamp, an auto-increment ID) into a single hot shard that receives every new write. Hash partitioning spreads writes evenly and scatters any query over a range of keys. A directory keeps an explicit lookup table, which can place anything anywhere at the price of an extra lookup.
| Strategy | Row → shard | Good at | Watch out for |
|---|---|---|---|
| Range | Key ranges: [a, m) here, [m, z) there | Range scans, time-based retention, splitting a busy range | Ever-increasing keys all write to one shard |
| Hash | hash(key) mod logical shards | Even write load with no tuning | Range queries and cross-key sorting fan out |
| Directory | Lookup table key → shard | Moving one big tenant, mixed placement | An extra lookup per request, which must be cached |
| Geo / tenant | Region or tenant column | Data residency, keeping a customer together | Tenants differ wildly in size |
Choosing the shard key
The key is chosen from the queries, not from the schema. Write down the queries the product runs most, and pick the key that makes the hot ones single-shard while spreading writes evenly. Worked example: a chat app with messages, conversations and users.
- 1List the hot paths. Opening a conversation loads its latest 50 messages; sending writes one message; the inbox lists a user's conversations; search finds messages by text.
- 2Shard by message_id (hash): perfectly even writes, but loading one conversation touches every shard. Rejected: the most frequent read becomes scatter-gather.
- 3Shard by created_at (range): all current writes hit the newest shard. Rejected: one hot shard.
- 4Shard by sender user_id: a conversation's messages come from several senders, so reading it still fans out. Rejected.
- 5Shard by conversation_id (hash), sorted by time inside: loading and sending are single-shard, writes spread across conversations. Accepted for messages.
- 6The inbox needs "conversations for user X", which conversation_id cannot answer. Keep a separate table keyed by user_id (a second sharded table, written alongside) instead of fanning out. Search goes to a dedicated search index, not the shards.
-- shard key = conversation_id
-- single shard: the key is in the WHERE clause
SELECT * FROM messages
WHERE conversation_id = $1
ORDER BY created_at DESC LIMIT 50;
-- scatter-gather: no shard key, every shard must answer
SELECT * FROM messages WHERE sender_id = $1;
-- fix: a second table sharded by the key this query has
SELECT conversation_id FROM user_conversations
WHERE user_id = $1 ORDER BY last_message_at DESC;Resharding and what you lose across shards
With logical shards, adding capacity is a data move, not a rehash: pick some logical shards, copy them to the new server while they keep taking writes, and switch the map. The sequence below is roughly what Vitess, Citus and hand-rolled setups like Notion's do.
- 1Choose which logical shards move from the busy database to the new one, for example 120 of the 240 on
pg-03. - 2Copy them to the new database: a snapshot first, then stream ongoing changes with logical replication or change data capture, so the copy keeps up while the old copy still serves traffic.
- 3Wait for the replication lag to approach zero, then pause writes to just those logical shards for a few seconds, confirm the copy has caught up, and switch their entries in the shard map.
- 4Routers pick up the new map version; requests for those tenants now go to the new database. Keep the old copy read-only until you have verified row counts and checksums, then drop it.
What stops working across shards
A single database guarantees a lot that sharding takes away. A join between rows on different shards becomes application code or a fan-out. A transaction across shards needs 2PC (two-phase commit), which systems like Spanner and CockroachDB run internally at a latency cost, or a saga with compensating actions. A UNIQUE constraint only holds within one shard, so global uniqueness (an email address) needs its own table keyed by that value. And auto-increment IDs collide, so IDs come from a generator that embeds time and origin: Twitter's Snowflake (2010) uses 41 bits of milliseconds, 10 bits of machine ID and a 12-bit sequence; Instagram embeds the logical shard ID instead, so the ID itself says where the row lives.
Pitfalls
- An ever-increasing key under range partitioning
Timestamps, auto-increment IDs and time-ordered UUIDs (version 7) all sort to the end, so under range partitioning every insert goes to the last shard and the others idle. Hash the key, prefix it with something well distributed, or use range partitioning only for retention inside each shard.
- Choosing the key before listing the queries
The "obvious" key from the schema (user_id, message_id) often turns the most frequent read into scatter-gather. Since changing the key later means rewriting every row, start from the top queries by volume and check each one against the candidate key before committing.
- Fan-out queries on the hot path
A query that asks all N shards waits for the slowest one, so its tail latency gets worse as you add shards, and every shard spends work on it. Keep fan-out for rare, offline or admin queries, and serve frequent lookups by other keys from a second table or index keyed the right way.
- Assuming constraints still hold globally
Unique constraints, foreign keys and transactions only hold within a shard. In PostgreSQL even a single-server partitioned table can only enforce a unique constraint that includes the partition key. Duplicate emails or orphaned rows appear quietly unless global invariants get their own enforcement.
- Sharding straight onto physical servers
Mapping
hash(key) mod serversmeans adding a server rehashes nearly every row. Hash into a large fixed number of logical shards (or use consistent hashing) so that growth moves whole logical shards and touches only the data that has to move.
Interview questions
Q1What is the difference between partitioning and sharding?
Partitioning splits a table into pieces by a key; sharding is partitioning where the pieces live on different servers. Partitioning on one server keeps joins, transactions and constraints and mainly helps pruning and maintenance. Sharding adds write and storage capacity but gives up cross-shard joins, transactions and global constraints.
Q2How would you choose a shard key for a chat application?
Conversation ID, hashed, with messages sorted by time inside the partition. Loading and sending messages are the hottest paths and both name the conversation, so they stay single-shard, and writes spread across many conversations. The inbox, which asks by user, gets its own table sharded by user ID, and search goes to a search index, so neither fans out.
Q3Walk me through growing from 4 to 8 database shards without downtime.
If keys map to many logical shards, I move half of each server's logical shards to the new servers. For each batch: snapshot and stream changes to the new server, wait until lag is near zero, pause writes for those logical shards for seconds, switch their entries in the shard map, and let routers load the new map. The old copies stay read-only until counts and checksums match. If the system was sharded by mod 4 directly, I would first introduce logical shards, because otherwise doubling the server count moves about half of all rows at once.
Q4What happens when a query does not include the shard key?
It becomes scatter-gather: the router sends it to every shard and merges the results. It works, but its latency is the slowest shard's, its cost grows with the shard count, and sorting or limiting has to be redone in the router. If it is a frequent query, I add a second table or index keyed by what the query does have.
Q5How do you generate unique IDs across shards?
With a generator that does not need coordination, such as Snowflake-style 64-bit IDs combining a millisecond timestamp, a machine or shard ID and a per-millisecond sequence, or UUIDs. Time-ordered IDs keep B-tree inserts sequential and sort roughly by creation, and embedding the logical shard ID, as Instagram does, lets you route by ID alone. Per-shard auto-increment collides as soon as you have two shards.
Q6When would you choose range partitioning over hash?
When the dominant queries are ranges over the key or when data ages out by that key: time-series retention, "orders from last week", or alphabetical scans. The price is hot spots when new keys all sort to one end, which systems like Bigtable and CockroachDB handle by splitting busy ranges. Hash is the default when the access pattern is point lookups and writes must spread evenly.
Q7How do you handle an operation that must update rows on two shards?
First I try to make it a single-shard operation by co-locating the data, for example sharding everything by tenant. If it truly spans shards, the options are a distributed transaction with two-phase commit, which databases like Spanner and CockroachDB provide at a latency cost, or a saga: each step commits locally and publishes an event, and failures trigger compensating steps. Sagas give up isolation, so the intermediate state must be acceptable to readers.
Q8When should you not shard?
When a bigger server, read replicas, caching and single-server partitioning still cover the load. Sharding costs joins, transactions, constraints and operational simplicity, and it is very hard to undo. I shard when sustained write throughput or data size is approaching what the largest practical instance can handle, or when data residency requires splitting by region.
- Partitioning splits a table by a key; sharding puts the pieces on different servers and gives up cross-shard joins, transactions and constraints.
- The shard key decides everything: queries that name it hit one shard, queries that do not scatter to all of them.
- Choose the key from the hottest queries, not the schema, and check that writes spread evenly. Big tenants need an escape hatch.
- Range keeps neighbours together but makes ever-increasing keys hot; hash spreads writes but scatters range queries.
- Hash into many logical shards and map those to servers, so resharding moves whole logical shards instead of rehashing every row.
- Partition inside one database first. Shard when one machine truly cannot hold the writes or the data.