Deep Engineering
Advanced·Published·45 MIN

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 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_capable and parallel_commit are 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. With parallel_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

model with assumptionsbench/sharding/placement.py — an exact computation of the layout; only the keys are random, and they are fixed by a seed. The same keys and hash as in bench/hashring/ring.py.

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

— Valkey — Cluster specification

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

model with assumptionsbench/sharding/placement.py, blocks 2–4. The Zipf exponent is a parameter of the model, not a property of anyone's load; the result is shown for three of its values.

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

measured observationbench/sharding/routing.py, PostgreSQL 16.13. The shards are databases on one server behind postgres_fdw, each behind its own delaying relay. The conclusions are stated in network round trips and do not depend on the number of machines.

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

measured observationbench/sharding/crossshard.py, PostgreSQL 16.13. The server's answers are printed verbatim; the shard's refusal is a deferred trigger.

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

— PostgreSQL 16 — 5.11. Table Partitioning

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

measured observationbench/sharding/resharding.py, PostgreSQL 16.13, 200,000 rows on eight shards behind postgres_fdw. Rows are counted, not time: the time of a transfer depends on the machine, the number of rows moved does not.

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.

— PostgreSQL 16 — CREATE TABLE
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 from bench/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.

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_capable and parallel_commit are 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

Claim

Hashing the key spreads the load over the shards evenly

Actually

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).

Claim

Consistent hashing both moves little and spreads evenly

Actually

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).

Claim

With async_capable a query to every shard costs as much as a query to one

Actually

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).

Claim

If the transaction returned an error, nothing was written

Actually

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).

Claim

parallel_commit only speeds up the commit

Actually

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).

Claim

To add a shard in PostgreSQL, adding a partition is enough

Actually

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).

Claim

Sharding is the same as replicas, only harder

Actually

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

Question 1 of 5

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

  1. 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
  2. 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
  3. 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
  4. 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
  5. 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/
  6. 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
  7. 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