Partitioning and sharding
Split one dataset across many machines, so that each machine holds and serves only part of it. The shard key decides which machine owns a row.
- 1Shard last. Replicas, a bigger machine and partitioned tables come first.
- 2The shard key decides everything: the common query must hit one shard, and writes must spread evenly.
- 3Hash mod N moves almost every key when N changes. A ring or fixed slots move only the new node's share.
- 4Hot keys, cross-shard queries and cross-shard transactions are the cost. Design them out with the key.
- Shard when one primary cannot take the writes, or cannot hold the data.
- Each shard owns a subset of keys. A router maps a key to its shard.
- The placement rule decides how much data moves when the cluster grows.
I never place keys with hash mod N on a cluster that grows. Adding one node would move most of the data.
| strategy | range scan | add a node | fits |
|---|---|---|---|
| Range of the key | Approved | Split a range | Ordered keys, time within a source |
| Hash mod N | Not approved | Not approved | A shard count that never changes |
| Hash into fixed slots, slot map | Not approved | Move slots | Most key-value and OLTP shards |
| Consistent hash ring, virtual nodes | Not approved | Moves 1/(N+1) | Peer clusters with no central map |
| Directory: key to shard table | Per entry | Move one key | Big tenants, uneven key sizes |
| Geographic: region in the key | In a region | Per region | Data residency, local latency |
Fixed slots and the ring both move only the new node's share. Slots keep an explicit map; a ring derives the map from token positions.
I use ranges when queries scan ordered keys, and hashing into slots or a ring when writes must spread evenly.
- cardinality
- Many more values than shards. A country column has about 200 values; most traffic sits in a few.
- even spread
- No single value takes a large share of rows or writes.
- alignment
- The common query names the key, so it touches one shard.
- growth
- New data does not all go to one value. A timestamp fails this test.
| workload | shard key | the common query | status |
|---|---|---|---|
| SaaS app | tenant_id | Everything for one tenant | Big tenants move out |
| Hotel booking | hotel_id | Hold and confirm nights at one hotel | Approved |
| Feed, profile | user_id | One user's posts and settings | Approved |
| Chat | conversation_id | Messages of one conversation, in order | Approved |
| Logs, metrics | created_at | Recent lines | One hot shard |
| Logs, metrics | (source_id, day) | One source over a time range | Approved |
| Orders | order_id | A customer's orders | Scatters by customer |
| Orders | customer_id | A customer's orders | Approved |
I pick the key that the most common query names. Then I check that no value takes a large share of writes, and that new data spreads too.
Step 1: Find the key
- Take the shard key from the request: tenant_id = 42.
- Hash it, or keep it as is for range sharding.
If it fails
The request has no shard key. Then it must go to every shard (step 4).
The router finds the shard from a cached map. A query with the key touches one shard; a query without it touches all of them.
| tool | capability | what it gives this design | also used for |
|---|---|---|---|
| Postgres | Declarative partitioning: PARTITION BY RANGE, LIST or HASH | One logical table, many physical tables on one server. The step before sharding. | Time series, multi-tenant tables |
| Postgres | Partition pruning | A query with the partition key reads one partition. Without it, the plan reads all of them. | Date-bounded reports |
| Postgres | DETACH PARTITION, DROP TABLE | Remove a whole month at once, with no row-by-row DELETE. | Retention for logs and events |
| Postgres | Unique indexes on a partitioned table | Limit They must contain the partition key. Global uniqueness needs another plan. | |
| Postgres | Logical replication with row filters (version 15 and later) | Copy only one tenant's rows to a new server while it stays live. | Major-version upgrades, change data capture |
| Postgres | postgres_fdw: a partition can be a foreign table | A partition can live on another server and still answer through the parent table. | Federated reporting |
| Citus | Distributed tables by a distribution column, co-located tables, reference tables | Sharded Postgres: joins on the same key run on one worker. | Multi-tenant SaaS, real-time analytics |
| Redis | Cluster: 16,384 hash slots, CRC16 of the key | Fixed slots and a slot map. Resharding moves slots while the cluster serves. | Any Redis data set larger than one node |
| Redis | Hash tags {user:42} and MOVED, ASK replies | Keys with one tag share a slot. A client with a stale map is told where the slot went. | Multi-key scripts in a cluster |
| Cassandra | Token ring with virtual nodes (num_tokens, 16 by default) | A consistent hash ring. Partition key hashed with Murmur3; replicas follow the ring. | Write-heavy wide tables |
| DynamoDB | Partition key hashed to partitions; partitions split as they grow | Managed hash sharding. A hot partition key still hits one partition. | Serverless key-value |
| Service | A router library with a cached shard map | Finds the shard in memory. Reloads the map when a shard refuses a key. | Any shard-aware service |
On one Postgres I partition tables and let pruning skip partitions. Across machines I use slots or a ring, and I move data with logical replication or double writes.
build(nodes, vnodes): // once per membership change
FOR EACH node IN nodes:
FOR i IN 0 .. vnodes-1:
put token hash(node + "#" + i) on the ring1
sort tokens by position
owner(key):
h = hash(key)
t = first token at or after h2 // binary search
IF there is none: t = the first token // wrap around3
RETURN t.node
add(node):
put the vnodes tokens of node on the ring
keys just before each new token move to node4- 1A virtual node is one token. Its position depends only on the node name and the index, so other nodes never move.
- 2Binary search over the sorted tokens: O(log of the token count).
- 3The ring is a circle. A hash past the last token belongs to the first one.
- 4Only these keys move, and all of them move to the new node.
Tested source Go: the ring
// Ring is a consistent hash ring. Each node owns vnodes points (tokens) on the ring; a key
// belongs to the first token at or after its hash, wrapping at the end.
type Ring struct {
tokens []Token // sorted by Pos
}
// Token is one virtual node: a point on the ring and the node that owns it.
type Token struct {
Pos uint64
Node string
}
// NewRing places vnodes tokens for each node. A token's position depends only on its node and
// index, so adding a node never moves another node's tokens.
func NewRing(nodes []string, vnodes int) *Ring {
r := &Ring{}
for _, n := range nodes {
r.tokens = append(r.tokens, NodeTokens(n, vnodes)...)
}
sort.Slice(r.tokens, func(i, j int) bool { return r.tokens[i].Pos < r.tokens[j].Pos })
return r
}
// NodeTokens returns the tokens of one node.
func NodeTokens(node string, vnodes int) []Token {
out := make([]Token, vnodes)
for v := range out {
out[v] = Token{Pos: Hash(fmt.Sprintf("%s#%d", node, v)), Node: node}
}
return out
}
// Owner implements Placer: binary search for the first token at or after the key's hash.
func (r *Ring) Owner(key string) string {
h := Hash(key)
i := sort.Search(len(r.tokens), func(i int) bool { return r.tokens[i].Pos >= h })
if i == len(r.tokens) {
i = 0 // past the last token: wrap to the first
}
return r.tokens[i].Node
}
A key belongs to the first token clockwise from its hash. A new node only takes the arcs just before its own tokens.
Going from 4 to 5 nodes, hash mod N moves about 80% of the keys. Slots and a ring with enough tokens move about 20%, the new node's fair share.
- A's ranges
- key on A
- key moved by the last change
Fair share 25.0% (dashed). Fullest node: 1.37× its fair share.
Six nodes, any subset, five token counts: 315 rings, each simulated over 20,000 keys. Dots are 48 of those keys.
With one token per node the arcs are uneven and a new node takes keys from one neighbour. With a few hundred tokens each node owns many small arcs, so load and moves both even out.
| fix | status | cost |
|---|---|---|
| Hash a high-cardinality column first, time second | Approved | A time-range query over all sources scatters. |
| Split the hot range | Range sharding | The newest range is hot again after the split. |
| Cache the hot reads | Read-hot keys | Does nothing for writes. |
| Salt the hot key | Write-hot keys | Every read of that key reads k shards. |
| Own shard for the hot tenant | With a directory | One more entry to keep in the map. |
A timestamp as the key sends every new write to the newest range. I hash the source first and keep time inside the shard.
write(key, value):
IF key is hot:
key = key + "#" + (count MOD k1) // k sub-keys
shard(key).write(value)
read(key):
IF key is hot:
RETURN merge(2shard(key + "#" + i).read FOR i IN 0 .. k-1)
RETURN shard(key).read- 1Round robin over k sub-keys. Each sub-key hashes to its own shard, so the writes spread.
- 2The cost: one read becomes k reads. Salt only keys that are hot for writes.
Tested source Go: salt
// Salt spreads the writes of one hot key over k sub-keys ("user:0#0" to "user:0#k-1"), round
// robin. Each sub-key hashes to its own shard. A read of the hot key must read all k sub-keys.
func Salt(writes []string, hot string, k int) []string {
out := make([]string, len(writes))
next := 0
for i, w := range writes {
if w != hot {
out[i] = w
continue
}
out[i] = fmt.Sprintf("%s#%d", hot, next%k)
next++
}
return out
}
Sub-keys hash like any key, so two can land on one shard. The whisker shows that spread.
One key with a quarter of the writes made its shard take 34.4% of all writes. Split into 8 sub-keys, the hottest shard took about 17.3%.
-- One logical table, split into 8 partitions by a hash of tenant_id.
CREATE TABLE events (
tenant_id bigint NOT NULL,
event_id bigint NOT NULL,
created_at timestamptz NOT NULL,
body text NOT NULL,
PRIMARY KEY (tenant_id, event_id)2
) PARTITION BY HASH (tenant_id)1;
CREATE TABLE events_p0 PARTITION OF events
FOR VALUES WITH (MODULUS 8, REMAINDER 03);
CREATE TABLE events_p1 PARTITION OF events
FOR VALUES WITH (MODULUS 8, REMAINDER 1);
-- ... events_p2 to events_p7- 1Postgres hashes tenant_id to pick the partition. A sharded cluster does the same across servers.
- 2The key must contain tenant_id. Each partition checks uniqueness only for its own rows.
- 3Hash mod 8. To grow, detach one partition, replace it with two of modulus 16, and copy its rows back.
| method | status | cost |
|---|---|---|
| Scatter-gather to every shard | Rare queries | Latency of the slowest shard; load on all. |
| Global secondary index: value to shard key | Approved | A second write, usually asynchronous. |
| Denormalised copy sharded by the other key | Hot lookups | Two copies to keep in step. |
| Shard number inside the ID | Approved | The ID format is fixed for good. |
A query without the shard key goes to every shard. For a common lookup I keep a global index or put the shard number in the ID.
| method | status | why |
|---|---|---|
| Co-locate: both rows share the key | Approved | A normal local transaction. |
| Saga with compensations, from an outbox | Approved | Each step commits locally. Others can see the middle state. |
| Two-phase commit | Few shards, low rate | A prepared transaction holds its locks until the coordinator decides. |
| Two independent commits, no plan | Not approved | A crash between them leaves money debited and not credited. |
I choose the key so that one transaction touches one shard. When it cannot, I run a saga with compensations, not two-phase commit.
| method | status | note |
|---|---|---|
| Unique index without the shard key | Not approved | Postgres refuses it on a partitioned table. |
| UUID, random or time-ordered | Approved | No coordination. 16 bytes; time order keeps inserts local in the index. |
| 64-bit ID: time, worker, sequence | Approved | Each process needs a unique worker number. |
| Shard number inside the ID | Approved | A lookup by ID finds its shard with no index. |
| Claim table sharded by the value: email to user_id | Unique emails | Insert the claim first. A duplicate fails there. |
| One central sequence | Low rate | Every insert waits on one server. |
The lab error, verbatim: "unique constraint on partitioned table must include all partitioning columns".
A unique index works only inside one shard. I generate IDs that are unique by design, and I claim unique values like emails in a table sharded by that value.
move(range, old, new):
1 turn on double writes:
write old shard // version = version + 1
write new shard IF its version is older1
2 backfill: FOR EACH batch of rows IN old:
copy to new, SKIP rows new already has2
3 verify: rows that differ == 03
4 switch reads of the range to new4
5 stop writes to old, then delete the range from old- 1Two writers can reach the new shard in either order. The version guard keeps the newer row.
- 2ON CONFLICT DO NOTHING. A double write made after the backfill read is newer than the copy.
- 3A full outer join of the two shards. Switch nothing until it is zero.
- 4One entry in the shard map changes. Routers that still use the old map are refused and reload.
Tested source Go: double write · Go: backfill · SQL: writes and copy
// Write changes one account during a move: old shard first, then the new shard. The version
// guard on the new shard makes the two writes safe in any order across writers.
func Write(ctx context.Context, db *pgxpool.Pool, id, delta int64) error {
var balance, version int64
if err := db.QueryRow(ctx, reshard["write_old"], id, delta).Scan(&balance, &version); err != nil {
return fmt.Errorf("write old shard, account %d: %w", id, err)
}
if _, err := db.Exec(ctx, reshard["write_new"], id, balance, version); err != nil {
return fmt.Errorf("write new shard, account %d: %w", id, err)
}
return nil
}
// Backfill copies the old shard to the new one in batches of size, in id order, while the
// service keeps double writing.
func Backfill(ctx context.Context, db *pgxpool.Pool, size int) error {
var after int64
for {
batch, err := ReadBatch(ctx, db, after, size)
if err != nil {
return err
}
if len(batch) == 0 {
return nil
}
if err := CopyBatch(ctx, db, batch, false); err != nil {
return err
}
after = batch[len(batch)-1].ID
}
}
-- The service writes the old shard first. The version counts every change to the row.
UPDATE accounts_old
SET balance = balance + $2, version = version + 1
WHERE id = $1
RETURNING balance, version;
-- Then it writes the same row to the new shard. An older version never replaces a newer one.
INSERT INTO accounts_new (id, balance, version) VALUES ($1, $2, $3)
ON CONFLICT (id) DO UPDATE
SET balance = excluded.balance, version = excluded.version
WHERE accounts_new.version < excluded.version;
-- The backfill skips a row the new shard already has. A double write put it there, and it is newer.
INSERT INTO accounts_new (id, balance, version)
SELECT * FROM unnest($1::bigint[], $2::bigint[], $3::bigint[])
ON CONFLICT (id) DO NOTHING;
-- Rows that differ between the shards, or exist on one side only.
SELECT count(*)
FROM accounts_old o
FULL JOIN accounts_new n USING (id)
WHERE (o.balance, o.version) IS DISTINCT FROM (n.balance, n.version);- Tested with 8 writers and 3,200 double writes during a backfill of 2,000 rows: 0 rows differ.
- A delete during the move needs a tombstone. Otherwise the backfill copies the deleted row back.
- With a ring or slots, only 1/(N+1) of the data takes this path.
I double write with a version guard, backfill rows the new shard does not have, verify, and only then switch the map. The backfill never overwrites a row.
| event | result | why it is safe | saved by |
|---|---|---|---|
| A shard primary fails | Its keys fail; other shards serve. | Its replica is promoted. Only 1/N of users notice. | Replica |
| A router has a stale map after a move | It sends the key to the old shard. | The old shard refuses keys it no longer owns. The router reloads and retries. | Map version |
| The backfill copies a row older than a double write | The new shard already has the row. | The copy skips rows that exist. | ON CONFLICT |
| Two writers reach the new shard out of order | The older version arrives second. | The upsert changes the row only when its version is older. | Version guard |
| A row is deleted during the move | The backfill can copy it back. | Delete with a tombstone until the move ends. Verify finds any leftover. | Tombstone |
| One tenant grows to 30% of the load | Its shard runs hot. | A directory entry moves that tenant to its own shard. | Directory |
| One shard is slow during a scatter query | The whole query waits. | A timeout returns partial results, marked as partial. | Timeout |
| A transfer fails between two shards | Debited, not credited. | The saga runs the refund from its outbox record. | Saga |
| step | add | it handles | move up when you see |
|---|---|---|---|
| 1 | One Postgres with indexes. Scale it up first: more cores, memory, faster disks. | About 50,000 single-row inserts a second in the lab, and every query on one machine. | Reads slow the primary. |
| 2 | Read replicas, as on the replication sheet. | Reads grow with each replica. Writes do not. | Tables grow so large that indexes and retention slow down. |
| 3 | Partitioned tables on the same server. | Smaller indexes, fast retention by DROP, pruning by key. | One primary cannot take the write rate, even on the largest machine. |
| 4 | Shards by key, with a router and a shard map. | Writes and storage grow with the shard count. The common query stays on one shard. | One shard runs hot, or the cluster needs more nodes. |
| 5 | Reshard: move ranges or slots to new shards online. | Growth with no downtime. Only 1/(N+1) of the data moves per new node. | Top of the ladder. |
Demand example: 1 billion events a day is 1,000,000,000 / 86,400 ≈ 11,574 a second. Shard capacity is the lab measurement times the shard count; real shards scale a little below that. Postgres 16 on an 8-core laptop, 32 clients.
I add replicas for reads, a bigger machine and partitions for size. I shard only when one primary cannot take the writes, and I pick the key from the main query.
0 of 10 known
You go from 4 to 5 shards with hash mod N. About how many keys move?
You add node E to a consistent hash ring. Which keys move, and to where?
Why give each node many virtual nodes?
You shard logs by timestamp. What goes wrong?
A SaaS app shards by tenant_id. One tenant sends 30% of the traffic. What do you do?
Orders are sharded by customer_id. How do you find an order by order_id?
Why does Postgres refuse UNIQUE (event_id) on a table partitioned by tenant_id?
During a backfill, why skip rows that already exist on the new shard?
A transfer moves money between users on two shards. How do you make it safe?
Why does Redis Cluster use 16,384 fixed slots instead of a ring?
- mod N
- N to N+1 nodes moves N/(N+1) of keys: 4 to 5 moved 79.6%.
- ring, slots
- Move about 1/(N+1): 19.9% with slots, 20.0% with a 256-token ring.
- tokens
- Fullest of 8 nodes: 2.42× fair share with 1 token, 1.09× with 256.
- hot key
- A key with 25.0% of writes: its shard takes 34.4% of 8 shards' writes.
- Redis
- 16,384 hash slots (cluster specification).
- Cassandra
- 16 tokens per node by default (configuration reference).
- one node
- About 50,000 single-row inserts a second, measured.
Simulations are seeded and repeat exactly. Postgres 16 on an 8-core laptop, 32 clients; a server that flushes each commit to durable storage commits slower.