Sharding: the key decides where every query goes
Spreading the data over nodes is half the job; the other half is what stops working afterwards. Computed here and measured on a live PostgreSQL: going from eight shards to nine, hash-mod placement moves 88.9% of the keys where 11.1% is the minimum; on a postgres_fdw stand a query without the shard key pays five exchanges to every shard, and two settings overlap only two of them; a transaction through postgres_fdw that touched two shards left a row behind after the client's error in three cases out of four.
Full technical treatment
TL;DR
- Sharding is one big database spread over several machines, with every row belonging to exactly one of them. Which row lives where is decided by one rule and one column — the key.
- Everything good and everything bad follows from the key. A query that has the key goes to one machine. A query without it goes to all of them, unless there is a separate "what lies where" directory, and pays each one.
- What adding a machine costs depends on the rule. The simplest rule — the remainder of a division — moves 88.9% of the data when eight machines become nine, although 11.1% would have been enough. Smarter rules stay near the minimum.
- An operation that touches two machines loses the guarantees of an ordinary database unless a separate mechanism is added for it. Uniqueness on a column outside the key is not checked by itself, and a transaction across two machines through postgres_fdw can commit halfway — and report an error while doing so.
Why split a database into parts
Picture a card index that no longer fits in one cabinet. You could put copies of the cabinet next to it — those are replicas: handy when cards are read often, and a rescue if one cabinet burns down. But every copy still holds all the cards and receives every new one, so no room is gained. Or you could put several different cabinets in a row and spread the cards across them. That is sharding: every card belongs to exactly one cabinet, and a cabinet holds only its own part. One does not exclude the other: each of the different cabinets can have copies of its own.
The cabinets here are separate machines with a database; they are called shards. To find a card, you need to know which cabinet it is in. For that the card index has a rule: say, "surnames A–F in the first cabinet, G–L in the second". The column the rule works on is called the shard key.
The key decides which cabinet
The placement rule looks only at the key. If the cards are filed by surname, the question "where is Smith's card" is answered at once: open one cabinet. The question "where is the card with this phone number" the rule cannot answer — it knows nothing about phone numbers. Without a separate "phone number → cabinet" directory, every cabinet has to be opened in turn. Such a directory can be set up, but it has to be kept in order itself.
This is the main thing to remember about sharding. The key is chosen not by which column is "the main one" but by which questions are asked of the card index most often. If nearly every question is about a particular person, "person" is a good key. If half the questions are a search by phone number, this key will be expensive.
When there are more cabinets
The card index grows, and eight cabinets are no longer enough. The ninth cabinet has to receive its ninth of the cards — moving fewer is impossible. How many actually move depends on the rule:
The simplest rule is "cabinet number = the remainder of the card number divided by the number of cabinets". It spreads cards evenly, but when the number of cabinets changes, the remainder changes for almost every card at once. This is easy to check yourself:
import hashlib
def owner(key: str, shards: int) -> int:
h = int.from_bytes(hashlib.sha1(key.encode()).digest()[:8], "big")
return h % shards
keys = [f"user-{i}" for i in range(100_000)]
moved = sum(owner(k, 8) != owner(k, 9) for k in keys)
print(f"{moved / len(keys):.1%}")The program will print that almost nine cards in ten changed cabinets. Smarter rules — a ring, "slots" cut in advance and a couple of clever formulas — move close to the minimum: the new cabinet takes a little from each old one, and the other cards stay where they are. How the ring works is covered in detail in the lesson on consistent hashing.
There is also a problem no rule can fix. Cards can be spread evenly, yet they are accessed unevenly: one famous card may be asked for more often than all the others in its cabinet put together. Such a card lies in one cabinet, and the whole stream of requests for it goes there. In a model where the most popular card gets 19.6% of all requests, its cabinet works 2.14 times harder than the average. That one card alone gives its cabinet 1.57 times the average load — and no placement rule can fix that.
A query with the key and without it
Back to the question "where is the card with this phone number". On a live PostgreSQL this can be measured. The main database — the coordinator, on our stand PostgreSQL with the postgres_fdw extension — spends five short network exchanges on every cabinet it has to open: begin, find, fetch, close, confirm. If the network answers in 2 ms, a query with the key takes about 13 ms, however many cabinets there are. A query without the key pays five exchanges to every cabinet:
PostgreSQL has two settings that allow talking to the cabinets at the same time. They help: with eight cabinets a query without the key gets cheaper, from 104.0 to 67.9 ms. But only two of the five exchanges run simultaneously — "fetch" and "confirm" — and three still go cabinet by cabinet. So the more cabinets there are, the dearer a query without the key, however it is tuned.
When an operation touches two cabinets
While an operation touches one cabinet, everything works as in an ordinary database. As soon as it touches two, two familiar guarantees no longer hold by themselves. They can be brought back, but by a separate mechanism, and it is paid for in speed and complexity.
Uniqueness. Picture the rule "one phone number — one card" while the cards are filed by surname. Each cabinet can check that it holds no two cards with the same phone number, but it cannot look into the next one. In this situation PostgreSQL honestly refuses to create a unique index at all.
All or nothing. An operation puts one card into each of two cabinets, and one cabinet refuses at the last moment. In an ordinary database a refusal would mean nothing was stored. On our stand the cabinets confirm one after the other, and if the one that confirmed first is not the one that refused, its card stays:
The worst part is that the program gets an error all the same. The usual reaction to an error is to retry, and the retry puts a second identical card into the cabinet that kept the first.
What to choose
- Questions first, then the key. Choose the key so that the most frequent queries and all important operations touch one cabinet.
- Cabinets with room to spare. Adding them one at a time is expensive, so the data is split in advance into many small parts, and when it grows, whole parts are moved.
- A move — with a plan. While cards are being carried into a new cabinet and not yet removed from the old one, they lie in two places at once: on the stand the table counted 212,670 rows instead of 200,000. Decide in advance who is responsible for a card in transit.
- Popular records need separate care. They can be split into several copies, but reads pay for it.
- Operations across two cabinets — without counting on "all or nothing". Uniqueness outside the key is kept in a separate table, and a write to two cabinets is done in two steps, so that a retry creates no duplicates.
How it was measured
bench/sharding/placement.py— how much data moves under each rule and what uneven access does.bench/sharding/routing.py— a query with the key and without it on a live PostgreSQL.bench/sharding/crossshard.py— uniqueness and a transaction across two shards.bench/sharding/resharding.py— adding a ninth shard to eight.
TL;DR
- Sharding is the choice of a key, not the slicing of a table. The key decides which shard a row goes to, and with it splits every operation into two kinds: those that have the key go to one shard, the rest go to all of them, unless the router has separate information about where to look. Everything else in the article follows from this split.
- What adding a shard costs depends on the strategy. Eight shards become nine: hash mod N moves 88.9% of the keys against a minimum of 11.1%, evenly redrawn ranges move 50.0%. A ring, fixed slots, rendezvous and jump hash stay near the minimum, but the ring leaves a 1.41-fold skew.
- Even keys are not even load. With access following a Zipf law of exponent 1.2, the most popular key is 19.6% of all requests, and its shard carries 2.14 times the average load. This key alone is 1.57 times the average shard's share, and no strategy takes its shard below that: the key lives whole in one shard.
- A query without the key pays every shard. On a PostgreSQL stand with postgres_fdw the coordinator makes five exchanges with each shard it touches. At a 2 ms network round trip a query with the key costs about 13 ms at any shard count — 13.2 with four and 13.4 with eight — without the key 52.1 and 104.0.
- The settings that promise parallelism do not overlap everything.
async_capableandparallel_commitare off by default; together they overlap the fetch and the commit, while three of the five exchanges still go shard by shard. A query without the key still gets dearer with every shard: 67.9 ms on eight. - Two shards in one transaction through postgres_fdw are not atomic.
postgres_fdw does not use two-phase commit. When one shard refused at
COMMIT, the client got an error in all four experiments, and a row stayed on the other shard in three of them. Withparallel_commit, a row stayed whichever of the two refused. Another distributed database may offer a different contract, but atomicity across shards always needs a separate mechanism and is paid for in latency and complexity. - PostgreSQL itself does not enforce uniqueness outside the shard key. On a table whose partitions are shards, PostgreSQL creates no unique index at all.
- Resharding is moving data while it is live. While rows were copied to the
new shard and not yet deleted from the old one,
count(*)showed 212,670 rows instead of 200,000. A plain “copy, then delete” leaves a window in which the data is visible twice.
The basics: a shard, a key, and whoever knows where to go
One machine with a database hits a ceiling: disk size, the number of writes per second its log can take, the memory for the working set. Replicas do not raise that ceiling: each holds all the data and receives all the writes. They give something else — they scale reads, survive a node failure, keep the data when a machine is lost — but not capacity and not write throughput. Sharding splits the data itself: every row belongs to exactly one shard, and a shard holds only its share.
Sharding and replication answer different questions and are usually used together: each shard has its own replicas.
key → shard 3
├ primary: takes the writes
├ replica
└ replica
Four terms the rest of the article relies on:
- A shard is a part of the rows on its own node, with its own connections and its own transactions. To the application a shard is a separate database.
- The shard key is the column (or columns) whose value decides which shard
a row lives in. Here it is
user_id. - A placement strategy is the rule "key value → shard": hash modulo, ranges, a ring, a slot table.
- The router is whoever applies the rule to each query. It can be a library in the application, a separate proxy or, as on this article's stand, a PostgreSQL coordinator with foreign tables.
Hence the main property everything else follows from. The router can send
a query to one shard only if it can derive the shard from the query. The
simplest way is the key in the condition: a query by user_id goes where that
user's rows are. A query by email the placement rule cannot address — it
knows nothing about email. If the router has no other information, the query
goes to every shard. Such information can be added — a directory
email → shard, a global secondary index — but that is a separate structure
with its own price, and the article comes back to it.
The same holds for operations: uniqueness, a transaction, a join work as usual while they touch one shard. Touching two, they do not vanish but stop being local: the same guarantees across several shards need a separate mechanism — a global index, a distributed commit, a join coordinator — and are paid for in latency, availability and complexity.
So a key is chosen not by which column is "natural" but by which queries and
which transactions must stay inside one shard. The key user_id is good for a
system where almost everything happens within one user, and bad for a system
where half the queries are lookups by email.
How a key becomes an address: five strategies and one price
While membership does not change, any sensible strategy lays the keys out evenly. They differ in what happens when membership changes. The most common change is adding a node: eight shards become nine. For the new node to get its fair share — a ninth — at least 11.1% of the keys have to move. That is the minimum under three conditions: the nodes are equal, a key has one owner, and the layout before the change was even. Anything beyond it is moved for nothing, and moving less means the new node is left short of its share.
1. THE SAME CHANGE UNDER FIVE STRATEGIES: 8 NODES BECOME 9
----------------------------------------------------------
strategy keys moved min share max share max/min
modulo: hash mod N 88.9% 10.9% 11.3% 1.04x
ranges, all boundaries redrawn evenly 50.0% 10.9% 11.3% 1.03x
ranges, one range split in half 6.4% 6.3% 12.6% 1.99x
ring, 128 virtual nodes per node 10.7% 9.4% 13.3% 1.41x
slots: 1024, handed over whole 11.0% 11.0% 11.3% 1.03x
slots: 16384, handed over whole 11.2% 11.0% 11.2% 1.02x
Hash modulo — hash(key) mod N — lays out evenly and costs one
operation. Its trouble is that a key's owner turns out to be a property not
of the key but of the current number of nodes. Change the divisor and the
answer changes for almost every key at once: 88.9% against a minimum of 11.1%,
eight times what is needed.
Ranges keep keys ordered, and a query for "all orders of the week" or "users 1000 to 2000" goes to one or two shards, not to all. On a change of membership they have two paths, and both are unpleasant. Redraw all the boundaries evenly — and half the keys move: shifting each boundary touches each range. Cut one range in half — and less than the minimum moves, 6.4%, but the new node gets not its share but half of someone else's: two nodes hold half as much as the other seven.
A ring ties the owner to the key itself: a key goes to the nearest point on a circle, and a new node takes only the arcs before its own points. 10.7% moves — slightly below the minimum, and that breaks nothing: everything moved went to the new node, and it got slightly less than a ninth. The price is the randomness of the arcs: even with 128 points per node the busiest node holds 1.41 times the least busy one. Why many points are needed and what they cost is covered in the lesson on consistent hashing; the ring's numbers here and there come from one and the same computation.
Fixed slots separate two decisions the other strategies mix. A key maps to one of a number of slots fixed in advance — that decision never changes. A slot belongs to a node — that decision changes with membership, but a whole slot at a time: the new node gets a few slots from each old one, and nothing else moves. Hence 11.0% with 1024 slots and an even layout. This is how a Valkey cluster is built:
HASH_SLOT = CRC16(key) mod 16384
Slots are taken with a margin, many times more than nodes: a slot is the unit of transfer, and the layout cannot be more even than its size allows. This is the scheme of "many logical shards on a few physical nodes": the number of slots is chosen once and for a long time, and nodes hand them to each other whole.
Of the five strategies in the table, only slots both move little and stay even. They pay with a "slot → node" table that every router has to hold and keep current. But slots are not the only way: two more rules do without a table. They are computed on the same membership change and the same keys:
5. TWO MORE RULES WITHOUT A SLOT TABLE: RENDEZVOUS AND JUMP HASH
----------------------------------------------------------------
the same change as in block 1: 8 nodes become 9, the same keys
strategy keys moved min share max share max/min
rendezvous (HRW), score per node 11.3% 10.9% 11.3% 1.03x
jump consistent hash 10.9% 10.9% 11.3% 1.04x
Rendezvous, or HRW, gives every node its own score for a key — a hash of the “key, node” pair — and hands the key to the node with the highest one. A new node takes exactly the keys where its score beats all the old ones. Jump consistent hash (Lamping, Veach, 2014) computes the node number with a short piece of arithmetic on the key's hash and stores nothing. Both stayed near the minimum and are as even as slots. They pay with something else: rendezvous scores every node on every lookup, and jump hash needs the nodes numbered in a row — only the last one can be removed. There is also a third question besides moving and evenness: which queries stay inside one shard. On that one only ranges win — a scan over a span of keys goes to one or two shards, while under a hash it goes to all of them.
Even keys are not even load
The table above counts keys. Load is made by requests, and they reach
keys unevenly: a popular product, a celebrity on a social network, a shared
counter. The classic model of such unevenness is Zipf's law: the share of
requests to the key of popularity rank r is proportional to 1 / r^s.
The keys below are placed by the hash mod eight shards, and placed evenly: the shards hold from 12.3 to 12.7% of the keys.
2. EVEN KEYS ARE NOT EVEN LOAD: ACCESS FOLLOWS A ZIPF LAW
---------------------------------------------------------
keys placed by hash mod 8: key shares per shard 12.3% to 12.7%,
max/min 1.03x - an even placement of keys
key of rank r gets a share of requests proportional to 1 / r^s;
'top key alone' is the top key's share over the average shard's
s top key, share of all requests hottest shard vs average top key alone
0.8 2.2% 13.9% 1.11x 0.18x
1.0 8.3% 18.3% 1.46x 0.66x
1.2 19.6% 26.8% 2.14x 1.57x
At s = 1.2 the most popular key takes 19.6% of all requests, and the shard it
lives in carries 2.14 times the average load. Part of this skew depends on
where the other keys landed, but not all of it: this key alone is 1.57 times
the average shard's share, and no placement strategy fixes that. The rule
maps a key to exactly one owner, and all traffic to the key goes there.
The only way to spread a hot key is to stop storing it as one key: split it
into K copies key#0 .. key#K-1 and place the copies by their own hash. It
helps exactly as far as the copies land on different shards:
3. SPLITTING ONE HOT KEY INTO SUB-KEYS
--------------------------------------
s = 1.2; the top key is stored as K copies key#0 .. key#K-1,
each copy placed by its own hash, requests spread evenly over copies
K shards the copies hit hottest shard vs average
1 1 26.8% 2.14x
2 2 20.9% 1.67x
4 3 22.5% 1.80x
8 5 20.0% 1.60x
16 7 18.8% 1.50x
Four copies turned out worse than two: they landed on three shards, one of them next to other heavy keys. The copies are placed by the same hash as everything else, and their shards cannot be chosen. And reading pays for the split: every access to the key now picks a copy — or reads all of them, if the value is a counter that has to be summed.
A separate case of unevenness is an increasing key with ranges. An autoincrement, a timestamp, a time-ordered identifier: each new value is larger than all previous ones and lands in the last range.
4. AN INCREASING KEY AND RANGES: WHERE NEW WRITES GO
----------------------------------------------------
ids 1..100000 already stored, ranges of 12500 ids each, 8 shards;
the next 10000 inserts get ids 100001..110000
placement share of new writes on the busiest shard
ranges of the id 100.0%
hash of the id, mod 8 13.1%
All new writes go to one shard of eight, and the other seven get none. Hashing the identifier spreads inserts evenly — and gives up the ordered range scan that was the reason to use ranges.
Where the query goes
The stand is built with PostgreSQL's own means. The orders table on the
coordinator is partitioned by hash of user_id, and each of its partitions is a
postgres_fdw foreign table pointing into its shard. Partitioning decides where
a row goes, postgres_fdw how to get there. Specialised extensions do the same
job with their own code; here the built-in mechanism is what is wanted, because
its behaviour is described in the documentation and can be checked.
What the router does is visible right in the query plan:
1. WHERE THE QUERY GOES: THE PLAN WITH AND WITHOUT THE SHARD KEY
----------------------------------------------------------------
where user_id = 42 (the table is partitioned by hash of user_id)
Foreign Scan on orders_2 orders
where email = 'user42@example.com'
Append
-> Foreign Scan on orders_0 orders_1
-> Foreign Scan on orders_1 orders_2
-> Foreign Scan on orders_2 orders_3
-> Foreign Scan on orders_3 orders_4
With the key in the condition — one shard: the planner pruned the other partitions before execution. Without the key — all four. What each touched shard costs becomes visible when what the coordinator sends over the wire is decoded. The stand's relay understands the PostgreSQL protocol and prints the messages:
2. WHAT THE COORDINATOR SENDS TO A SHARD FOR ONE QUERY
------------------------------------------------------
decoded from the wire by the relay; connections already open
query without the shard key, shard 0:
1. Q: START TRANSACTION ISOLATION LEVEL REPEATABLE READ
2. P+B+D+E+S: DECLARE c1 CURSOR FOR
3. Q: FETCH 100 FROM c1
4. Q: CLOSE c1
5. Q: COMMIT TRANSACTION
messages per shard, all shards: [5, 5, 5, 5]
query with the shard key, messages per shard: [0, 0, 5, 0]
Five exchanges per touched shard: open a remote transaction, declare a cursor,
fetch the rows, close the cursor, commit. Each exchange is a network round
trip. A query with the key pays five round trips, a query without it five per
shard. Five as long as a shard returns no more than a hundred rows:
FETCH 100 is postgres_fdw's default batch size, and rows beyond it cost more
exchanges.
By default postgres_fdw visits the shards in turn. Two foreign-server settings
promise to change that, and both are off: async_capable
allows foreign tables to be scanned concurrently for asynchronous execution
, and parallel_commit
commits, in parallel, remote transactions opened on a foreign server
.
The time of one query with a 1 ms delay each way:
3. TIME PER QUERY, 4 SHARDS, 1 MS EACH WAY
------------------------------------------
configuration median best round worst round spread
with the shard key 13.2 ms 13.0 ms 13.6 ms 4%
no key, shards in turn 52.1 ms 51.3 ms 53.3 ms 4%
no key, async_capable 44.4 ms 43.9 ms 47.3 ms 8%
no key, async + parallel_commit 36.6 ms 36.4 ms 37.5 ms 3%
Both settings help, and help consistently: together they beat going in turn in all nine series of runs. But how much they help shows only in the growth with the number of shards, and here the measurement departed from expectation. Visiting the shards at once, it would seem, should make the price almost flat: all shards are asked together, so the bill should be that of one. That did not happen:
6. WHAT GROWS WITH THE NUMBER OF SHARDS, 1 MS EACH WAY
------------------------------------------------------
configuration 4 shards 8 shards 8 / 4
with the shard key 13.2 ms 13.4 ms 1.01x
no key, shards in turn 52.1 ms 104.0 ms 2.00x
no key, async_capable 44.4 ms 85.6 ms 1.93x
no key, async + parallel_commit 36.6 ms 67.9 ms 1.86x
Twice the shards — and even with both settings nearly twice the price. The explanation comes from counting the round trips paid strictly one after another. Each configuration was measured at two delays, 1 and 3 ms; the growth in time divided by the growth in round trip is that count. Constant costs cancel in the subtraction.
5. HOW MANY ROUND TRIPS ARE PAID ONE AFTER ANOTHER
--------------------------------------------------
each configuration timed at 1 ms and at 3 ms each way;
(time at 3 - time at 1) / (6 - 2 ms) = round trips in sequence
configuration 4 shards expected 8 shards expected
with the shard key 5.3 5 5.2 5
no key, shards in turn 20.8 20 41.1 40
no key, async_capable 17.6 17 34.2 33
no key, async + parallel_commit 14.6 14 27.0 26
The expected column is not a fit but a consequence of decoding the protocol:
of the five exchanges, async_capable overlaps only the fetch and
parallel_commit only the commit, and an overlapped exchange costs one round
trip for all shards instead of one per shard. The differences between the rows
confirm it regardless of the remainder: "in turn" minus async_capable is 3.2
round trips with four shards and 6.9 with eight, async_capable minus
async + parallel_commit is 3.0 and 7.2. That is N − 1, three and seven, to
within 0.2 of a round trip. The measured count is everywhere 0.2–1.2 above
the decoded one; decoding the protocol does not explain where the remainder
comes from.
Opening the remote transaction, declaring the cursor and closing it are still done shard after shard. Three round trips per shard remain under any setting, and in postgres_fdw a query without the key gets dearer linearly with the number of shards, however it is tuned. The settings are worth turning on — on eight shards they remove a third of the price, 104.0 ms against 67.9 — but only the key in the condition fixes the situation.
Five exchanges per shard are a property of this router, not of sharding in general: another coordinator may talk to a shard differently. The general rule is wider than the stand: the price of a query without the key grows with the number of shards touched, as far as the work with them cannot be done at the same time.
What stops working by itself when an operation touches two shards
Uniqueness
A unique email on a table sharded by user_id is an ordinary requirement:
there must not be two accounts with one address. Here is the server's answer:
1. A UNIQUE EMAIL ON A TABLE SHARDED BY THE USER ID
---------------------------------------------------
the sharded table: partitions are foreign tables on the shards
alter table orders add unique (email)
ERROR: unique constraint on partitioned table must include all partitioning columns
DETAIL: UNIQUE constraint on table "orders" lacks column "user_id" which is part of the partition key.
alter table orders add unique (user_id, email)
ERROR: cannot create unique index on partitioned table "orders"
DETAIL: Table "orders" contains partitions that are foreign tables.
the same table with ordinary local partitions
alter table accounts add unique (email)
ERROR: unique constraint on partitioned table must include all partitioning columns
DETAIL: UNIQUE constraint on table "accounts" lacks column "user_id" which is part of the partition key.
alter table accounts add unique (user_id, email)
OK
On the sharded table there is no unique index at all — not even with the shard key included. On an ordinary partitioned table uniqueness is possible, but only one that includes the key, and the documentation explains why:
the constraint's columns must include all of the partition key columns. This limitation exists because the individual indexes making up the constraint can only directly enforce uniqueness within their own partitions
Each shard checks only its own rows, and a duplicate email under another
user_id lives in another shard. The usual way out is a second table, sharded
by email: email → user_id. Uniqueness is then held by a primary key on
each shard of that table — on the coordinator it cannot be created, as shown
above. It holds because the router sends every row of one address to one
shard, and nothing but the router checks that rule. Registration, meanwhile,
turns into two writes to two different shards — that is, into the next
problem.
Such a table is no longer an auxiliary one but a secondary index that correctness depends on. Keeping it consistent takes care on every operation: registration writes to both tables; deleting an account deletes from both; changing an email deletes the old row and inserts a new one, that is, two shards again; a retry after an error must not create a second entry; and after a partial write — the row is in one table and not in the other — a procedure is needed to find such pairs and finish or roll them back. The second table solves uniqueness, but it is paid for in consistency.
Atomicity
A transaction inserts a row into shard 0 and one into shard 1 and commits. One
of the shards refuses exactly at COMMIT — after both inserts have gone
through. This is not exotic: that is how deferred constraints fire, and a lost
connection to a shard at the wrong moment can look the same.
2. ONE TRANSACTION, TWO SHARDS, ONE SHARD REFUSES AT COMMIT
-----------------------------------------------------------
insert one row for user 1 (shard 0) and one for user 3 (shard 1),
then COMMIT; the refusing shard has a deferred trigger that raises at commit
parallel_commit refusing shard client sees row on shard 0 row on shard 1
false 0 an error none kept
false 1 an error none none
true 0 an error none kept
true 1 an error kept none
the client got an error every time; a row survived anyway in 3 of 4 cases
The client got an error all four times, and a row survived it in three. The
documentation describes the mechanism in one sentence:
The remote transaction is committed or aborted when the local transaction commits or aborts
.
At the local commit postgres_fdw commits the remote transactions in turn. If
the refusing shard goes first, the others are rolled back and nothing is left.
If the other one goes first, it is already committed, and there is nothing that
can undo a commit. The documentation does not specify which shard goes first;
in this run shard 1 went first.
parallel_commit, which speeds up the queries of the previous section, makes
things worse in this failure experiment: every COMMIT goes out at once, and
the shard that did not refuse committed in both experiments, whichever of
the two refused. The luck of the order goes away together with the order. That
is not an argument against the setting in general: it speeds up an ordinary
commit and changes only the outcome of a failure.
The most dangerous part is not the partial commit itself but that the client sees an error. The usual reaction to an error is to retry, and the retry creates a duplicate in the surviving shard.
Two-phase commit
Atomicity across several participants is what two-phase commit gives: first
each shard promises it will commit, and only when all have promised are they
told to commit. PostgreSQL can do it, but postgres_fdw does not use it:
Note that it is currently not supported by postgres_fdw to prepare the remote transaction for two-phase commit
.
And on the shards themselves it is off by default:
3. TWO-PHASE COMMIT ON A SHARD, DEFAULT SETTINGS
------------------------------------------------
max_prepared_transactions = 0
prepare transaction 'order-42'
ERROR: prepared transactions are disabled
HINT: Set max_prepared_transactions to a nonzero value.
Even turned on, it does not solve the problem by itself:
Two-phase transactions are intended for use by external transaction management systems
.
Someone outside has to remember which shards have promised and see it through
after its own crash. In practice, therefore, operations across shards are
designed so that atomicity is not needed: the key is chosen so that a
transaction stays in one shard, and where that is impossible, the write to the
second shard is made a separate step — through
an outbox written in the same transaction
and a retry that is safe against duplicates.
A change of membership on a live PostgreSQL
The model at the start of the article says how many keys move. A live server adds a question the model does not have: whether it allows the change at all. PostgreSQL's hash partitioning is modulo placement, and a first attempt to add a ninth shard to eight looks like this:
1. A NINTH PARTITION WITH MODULUS 9 NEXT TO EIGHT WITH MODULUS 8
----------------------------------------------------------------
create foreign table orders_8 partition of orders
for values with (modulus 9, remainder 8) server s8 ...
ERROR: every hash partition modulus must be a factor of the next larger modulus
DETAIL: The new modulus 9 is not divisible by 8, the modulus of existing partition "orders_7".
The rule from the documentation:
every modulus which occurs among the partitions of a hash-partitioned table is a factor of the next larger modulus
.
Eight and sixteen get along, eight and nine do not. With one partition per
shard, an equal ninth shard cannot be added at all. One of two things can: split one shard in two, or lay
out all the rows again.
The documentation describes the first path itself — and this is exactly the procedure the stand ran:
You can detach one of the modulus-8 partitions, create two new modulus-16 partitions covering the same portion of the key space... and repopulate them with data.
2. SPLITTING SHARD 0 IN TWO: MODULUS 8 BECOMES 16 FOR THAT SHARD ONLY
---------------------------------------------------------------------
rows in shard 0 before the split 24770
rows copied to the new shard 8 12670
share of all rows that travelled 6.3%
count(*) over the table before 200000
count(*) after copying, before deleting 212670
rows deleted from shard 0 12670
count(*) after deleting 200000
rows per shard after, shards 0..8: [12100, 25430, 25700, 24150, 24950, 25050, 24810, 25140, 12670]
largest / smallest shard 2.12x
Three observations: two match the model, and the third is not in the model at all.
6.3% of the rows travelled — as much as the model gave for a split range, and with the same result: two shards hold half as much as the rest, a skew of 2.12. When shards are added one at a time, hash partitioning behaves like ranges, not like slots.
While the rows travel, the table counts them twice. Copy and delete are two
steps, and between them count(*) shows 212,670 rows instead of 200,000: the
12,670 moved rows are visible both in the old shard and in the new one. The new
partition of modulus 16 and remainder 0 points at the whole table of shard 0,
and PostgreSQL does not check what a foreign partition contains: until the
moved rows are deleted from shard 0, the coordinator sees them both there and
in the new shard. The documentation warns of this in advance:
it is then the user's responsibility that the contents of the foreign table satisfy the partitioning rule
.
This observation is wider than PostgreSQL. Resharding is moving data while it is live, and the transition state needs a protocol of its own: who owns a row while it is in transit, and what a reader sees. The usual answers are a layout version number known both to the router and to the shards; reads and writes to both places for the duration of the move; switching the owner in one operation once the copy has caught up with the original. A plain “copy, then delete” without such a protocol leaves a window in which the data is visible twice — the stand showed it with one number.
The second path — an even layout over nine — costs what the model predicts for modulo. Computed with PostgreSQL's own hash function, row by row:
3. WHAT AN EVEN NINE-SHARD LAYOUT WOULD MOVE, BY POSTGRESQL'S OWN HASH
----------------------------------------------------------------------
rows 200000
rows whose shard changes, modulus 8 -> 9 177100
share 88.5%
the new shard fair share, 1 / 9 11.1%
88.5% of the rows against 88.9% of the keys in the model: a different hash, different keys, and the same price, because it is the price of modulo itself.
What to choose
Queries first, then the key. Write down the queries and transactions that must stay fast and atomic, and choose the key so that they touch one shard. Whatever is left without the key will pay every shard — on the stand that is five round trips per shard, and no setting brings them down to one.
The number of shards — once and with a margin. With one partition per shard, PostgreSQL's hash partitioning does not let a shard be added evenly: either a twofold skew, or moving almost everything. The scheme of "many logical shards on a few nodes" — like Valkey's slots — separates the decision "key → shard" from the decision "shard → node" and moves whole shards on a change of membership, near the minimum. With hash partitioning this means more partitions than shards: 72 partitions, nine on each of eight shards, go over to nine shards of eight each, and 8 partitions of 72 move — the same 11.1%. The stand did not run such a change: this is arithmetic, not a measurement.
Ranges — for range scans, and only while keeping an eye on increasing keys. They are the only ones that keep an ordered scan within one or two shards, but an increasing key piles all new writes onto the last shard, and adding a node offers a choice between half the data in transit and a twofold skew.
Resharding — with a protocol, not a script. Moving between layouts is moving data while it is live. Decide in advance who owns a row in transit and what a reader sees: a layout version number, writes to both places during the move, a switch in one operation. Without it, “copy, then delete” shows the data twice — 212,670 rows instead of 200,000 on the stand.
A hot key is a separate problem. No layout will spread it: a key has one owner. Splitting it into copies helps while the copies land on different shards, and shifts the cost to reads.
Operations across shards — without atomicity by default. Uniqueness outside the key — with a separate table sharded by that column. A write to two shards — in two steps with an outbox and a retry safe against duplicates, not in one transaction: postgres_fdw may commit it partially and report an error at the same time.
How it was measured
The first script computes without a server, the other three run against a live
PostgreSQL; its address is set by the DE_BENCH_DSN variable. Each opens from
here, together with the record of its run.
bench/sharding/placement.py— the layout when going from eight nodes to nine under five strategies and two more without a slot table (rendezvous, jump hash), uneven access, a hot key, an increasing key with ranges. Keys and hash come frombench/hashring/ring.py.bench/sharding/routing.py— the query plan, the protocol between the coordinator and a shard decoded, a query's time and the number of round trips in sequence for 4 and 8 shards.bench/sharding/crossshard.py— uniqueness outside the shard key, a transaction across two shards with a refusal at commit, two-phase commit by default.bench/sharding/resharding.py— adding a ninth shard to eight: modulus 9 refused, one shard split with the rows in transit counted, the price of an even layout by PostgreSQL's own hash.
The common part of the last three is bench/sharding/stand.py: the
coordinator, the shards and the delaying relay between them.
The version is PostgreSQL 16.13. The shards on the stand are databases of one server, so the absolute milliseconds here speak of the link's delay, not of the hardware; the conclusions are stated in network round trips and carry over to any delay by multiplication.
This is neither a retelling nor a separate text: everything below is taken from the article itself — its own summary, the section headings, the “actually” column and the version table. Which is why these theses cannot drift from the article.
The gist
- Sharding is the choice of a key, not the slicing of a table. The key decides which shard a row goes to, and with it splits every operation into two kinds: those that have the key go to one shard, the rest go to all of them, unless the router has separate information about where to look. Everything else in the article follows from this split.
- What adding a shard costs depends on the strategy. Eight shards become nine: hash mod N moves 88.9% of the keys against a minimum of 11.1%, evenly redrawn ranges move 50.0%. A ring, fixed slots, rendezvous and jump hash stay near the minimum, but the ring leaves a 1.41-fold skew.
- Even keys are not even load. With access following a Zipf law of exponent 1.2, the most popular key is 19.6% of all requests, and its shard carries 2.14 times the average load. This key alone is 1.57 times the average shard's share, and no strategy takes its shard below that: the key lives whole in one shard.
- A query without the key pays every shard. On a PostgreSQL stand with postgres_fdw the coordinator makes five exchanges with each shard it touches. At a 2 ms network round trip a query with the key costs about 13 ms at any shard count — 13.2 with four and 13.4 with eight — without the key 52.1 and 104.0.
- The settings that promise parallelism do not overlap everything.
async_capableandparallel_commitare off by default; together they overlap the fetch and the commit, while three of the five exchanges still go shard by shard. A query without the key still gets dearer with every shard: 67.9 ms on eight. - Two shards in one transaction through postgres_fdw are not atomic. postgres_fdw does not use two-phase commit. When one shard refused at
COMMIT, the client got an error in all four experiments, and a row stayed on the other shard in three of them. Withparallel_commit, a row stayed whichever of the two refused. Another distributed database may offer a different contract, but atomicity across shards always needs a separate mechanism and is paid for in latency and complexity. - PostgreSQL itself does not enforce uniqueness outside the shard key. On a table whose partitions are shards, PostgreSQL creates no unique index at all.
- Resharding is moving data while it is live. While rows were copied to the new shard and not yet deleted from the old one,
count(*)showed 212,670 rows instead of 200,000. A plain “copy, then delete” leaves a window in which the data is visible twice.
In fact
- It spreads the keys evenly, while load comes from requests, and they access keys unevenly. In a model with a Zipf law of exponent 1.2 the most popular key gets 19.6% of all requests, and its shard carries 2.14 times the average load — under an even layout of keys. This key alone is 1.57 times the average shard's share, and no strategy removes that: a placement rule maps a key to one owner, and the whole stream for the key goes there (
bench/sharding/placement.py, block 2). - It moves little: going from eight nodes to nine, 10.7% of the keys against a minimum of 11.1%. But the ring's arcs are random, and even with 128 points per node the most loaded node holds 1.41 times as much as the least loaded one. Both even and near the minimum are fixed slots — 11.0% moved, a skew of 1.03 — and two rules without a table: rendezvous (11.3% and 1.03) and jump consistent hash (10.9% and 1.04) (
bench/sharding/placement.py, blocks 1 and 5). - Of the five exchanges with a shard, this setting overlaps only the fetch. Opening the remote transaction, declaring the cursor, closing it and — without
parallel_commit— the commit still go shard by shard. On eight shards with both settings a query without the key costs 67.9 ms against 13.4 ms with the key, and 1.86 times as much as on four (bench/sharding/routing.py, blocks 5 and 6). - For a transaction across two shards this is false. When one shard refused at
COMMIT, the client got an error in all four experiments, and a row on the other shard stayed in three. postgres_fdw commits shards without two-phase commit, and there is nothing to undo what is already committed (bench/sharding/crossshard.py, block 2). - It also changes the outcome of a refusal. Without it, a partial commit depends on which shard commits first; with it, every
COMMITgoes out at once, and the shard that did not refuse always commits — whichever of the two refused (bench/sharding/crossshard.py, block 2). - With one partition per shard, the server will not let you add a modulus-9 partition to eight modulus-8 ones: every modulus must divide the next larger one. What remains is to cut one shard in two — 6.3% of the rows in transit and a skew of 2.12 — or to lay everything out anew: under PostgreSQL's own hash 88.5% of the rows would change shards (
bench/sharding/resharding.py). - They solve different problems. A replica stores all the data and accepts all the writes: it scales reads and survives a node failure, but adds neither capacity nor write throughput. A shard stores its own part of the rows — and that is exactly why operations that touch two shards stop getting uniqueness and atomicity for free: they have to be brought back by a separate mechanism. The two are usually combined: each shard has its own replicas.
What is covered
- The basics: a shard, a key, and whoever knows where to go
- How a key becomes an address: five strategies and one price
- Even keys are not even load
- Where the query goes
- What stops working by itself when an operation touches two shards
- A change of membership on a live PostgreSQL
- What to choose
- How it was measured
Take to work
- When choosing a shard key, first write down the queries and transactions that must stay fast and atomic, and take the key under which they touch one shard. Everything left without the key pays every shard: on postgres_fdw that is five exchanges with each.
- Treat a query without the shard key on a hot path as a schema defect, not something to tune.
async_capableandparallel_commitare worth turning on — on eight shards they remove a third of the price — but three exchanges per shard remain, and the price keeps growing with the number of shards. - Do not trust the error of a transaction that touched two shards: postgres_fdw commits shards without two-phase commit, and one of them may stay committed. A retry of such an operation must be safe against duplicates — through an idempotency key or a check before the write.
- Settle the number of shards once and with room to spare — many logical shards or slots on few nodes — instead of adding them one at a time. With one partition per shard, PostgreSQL hash partitioning does not let you add a shard evenly: either a twofold skew or moving almost every row. And design the move itself as moving data while it is live, with a protocol for the transition state: without one, “copy, then delete” shows the data twice — on the stand
count(*)saw 212,670 rows instead of 200,000. - Keep uniqueness on a column outside the shard key in a separate table sharded by that column: on a table whose partitions are shards, PostgreSQL will not create a unique index at all. Registration then becomes a write to two shards — design it as two steps, and design email changes, deletion and retries after an error so that the two tables do not drift apart.
Common misconceptions
Hashing the key spreads the load over the shards evenly
It spreads the keys evenly, while load comes from requests, and they access keys unevenly. In a model with a Zipf law of exponent 1.2 the most popular key gets 19.6% of all requests, and its shard carries 2.14 times the average load — under an even layout of keys. This key alone is 1.57 times the average shard's share, and no strategy removes that: a placement rule maps a key to one owner, and the whole stream for the key goes there (bench/sharding/placement.py, block 2).
Consistent hashing both moves little and spreads evenly
It moves little: going from eight nodes to nine, 10.7% of the keys against a minimum of 11.1%. But the ring's arcs are random, and even with 128 points per node the most loaded node holds 1.41 times as much as the least loaded one. Both even and near the minimum are fixed slots — 11.0% moved, a skew of 1.03 — and two rules without a table: rendezvous (11.3% and 1.03) and jump consistent hash (10.9% and 1.04) (bench/sharding/placement.py, blocks 1 and 5).
With async_capable a query to every shard costs as much as a query to one
Of the five exchanges with a shard, this setting overlaps only the fetch. Opening the remote transaction, declaring the cursor, closing it and — without parallel_commit — the commit still go shard by shard. On eight shards with both settings a query without the key costs 67.9 ms against 13.4 ms with the key, and 1.86 times as much as on four (bench/sharding/routing.py, blocks 5 and 6).
If the transaction returned an error, nothing was written
For a transaction across two shards this is false. When one shard refused at COMMIT, the client got an error in all four experiments, and a row on the other shard stayed in three. postgres_fdw commits shards without two-phase commit, and there is nothing to undo what is already committed (bench/sharding/crossshard.py, block 2).
parallel_commit only speeds up the commit
It also changes the outcome of a refusal. Without it, a partial commit depends on which shard commits first; with it, every COMMIT goes out at once, and the shard that did not refuse always commits — whichever of the two refused (bench/sharding/crossshard.py, block 2).
To add a shard in PostgreSQL, adding a partition is enough
With one partition per shard, the server will not let you add a modulus-9 partition to eight modulus-8 ones: every modulus must divide the next larger one. What remains is to cut one shard in two — 6.3% of the rows in transit and a skew of 2.12 — or to lay everything out anew: under PostgreSQL's own hash 88.5% of the rows would change shards (bench/sharding/resharding.py).
Sharding is the same as replicas, only harder
They solve different problems. A replica stores all the data and accepts all the writes: it scales reads and survives a node failure, but adds neither capacity nor write throughput. A shard stores its own part of the rows — and that is exactly why operations that touch two shards stop getting uniqueness and atomicity for free: they have to be brought back by a separate mechanism. The two are usually combined: each shard has its own replicas.
Knowledge check
The orders table is sharded by a hash of user_id over eight shards. Which query will the coordinator send to one shard?
Sources & further reading
7 SOURCES
- PostgreSQL 16 — F.38. postgres_fdwOfficial documentation. How a transaction that touches a remote server works: “The remote transaction is committed or aborted when the local transaction commits or aborts.” And the limit behind the section on atomicity: “Note that it is currently not supported by postgres_fdw to prepare the remote transaction for two-phase commit.” The same page describes both settings measured in this article and their defaults: async_capable “allows foreign tables to be scanned concurrently for asynchronous execution”, parallel_commit “commits, in parallel, remote transactions opened on a foreign server”; for both, “The default is false.”https://www.postgresql.org/docs/16/postgres-fdw.html
- PostgreSQL 16 — 5.11. Table PartitioningOfficial documentation. The rule that forbids a unique email on a table sharded by user_id, and the reason for it: “the constraint's columns must include all of the partition key columns. This limitation exists because the individual indexes making up the constraint can only directly enforce uniqueness within their own partitions.” And the caveat about foreign partitions, that is, about shards: “it is then the user's responsibility that the contents of the foreign table satisfy the partitioning rule.”https://www.postgresql.org/docs/16/ddl-partitioning.html
- PostgreSQL 16 — CREATE TABLE, hash partitionsOfficial documentation. Why a ninth shard cannot be added to eight with the same hash: “every modulus which occurs among the partitions of a hash-partitioned table is a factor of the next larger modulus.” And exactly the procedure the article runs on the stand: “You can detach one of the modulus-8 partitions, create two new modulus-16 partitions covering the same portion of the key space... and repopulate them with data.”https://www.postgresql.org/docs/16/sql-createtable.html
- PostgreSQL 16 — 74.4. Two-Phase TransactionsOfficial documentation. Who two-phase commit in PostgreSQL is for: “Two-phase transactions are intended for use by external transaction management systems.” The server can promise to commit, but someone outside has to decide when to carry the promise out — and postgres_fdw does not.https://www.postgresql.org/docs/16/two-phase.html
- Valkey — Cluster specificationOfficial documentation. Fixed slots in a production system: “HASH_SLOT = CRC16(key) mod 16384.” A key maps to a slot, not to a node, and on a change of membership nodes hand each other whole slots — which is why the article's model has a separate row for 16,384 slots.https://valkey.io/topics/cluster-spec/
- A Fast, Minimal Memory, Consistent Hash AlgorithmSource. John Lamping, Eric Veach (Google), 2014. Jump consistent hash: the node number is computed by a short piece of arithmetic on the key's hash, with no table and no memory, at the price of the nodes being numbered in a row. The implementation in
bench/sharding/placement.py, block 5, follows the paper; rendezvous (HRW) is computed there on the same keys.https://arxiv.org/abs/1406.2294 - The measurements of this article: a placement model and a PostgreSQL standSource. Four scripts: an exact computation of the layout under a change of membership and under uneven access, and three runs against a live PostgreSQL 16.13 — a query's route with the protocol decoded, operations across two shards, and resharding. The shards are databases on one server behind postgres_fdw, each behind its own delaying relay; why that is enough for the article's conclusions is explained in the description of the measurements./en/bench/sharding/routing.py