Skip to main content

Distributed queries

By the end of this page you will know what actually happens when a query touches data spread across more than one node, which query shapes stay cheap and which force data across the network, and how to see what a query actually did.

Connecting​

Connect to any node, exactly as described in Cluster overview: there is no separate coordinator, router, or proxy process to find or configure. Whichever node you connect to coordinates your statement: it plans it, runs the parts it can serve locally, and forwards or fetches the rest from whichever nodes actually hold or lead the data involved.

A client needs to know nothing cluster-specific to get this. No routing hint, no "which node owns this row" lookup, no special driver or connection string beyond the ordinary PostgreSQL ones you'd use against a single node. Learner nodes are the one case worth naming explicitly: the node you connect to always coordinates your statement, but the pieces of a distributed read run on learners first, then followers, and on a shard group's leader last (see Where fragments run), whichever node you connected to.

What changes, and what doesn't​

The SQL is the same SQL, and the wire protocol is the same wire protocol; nothing about how you write a query or which client library you use changes because a table's rows happen to live on more than one node. What's actually different:

  • Isolation is unchanged. A cluster transaction is still fully serializable, exactly as on one node; a conflicting transaction is aborted with SQLSTATE 40001, and your existing retry-on-serialization-failure logic handles it unmodified. See Cluster overview and Failover for the full contract.
  • A statement can now involve more than one node's worth of work. Single-node ScramDB never has to move a row over a network to answer a query; a cluster sometimes does. That is the entire added cost surface this page is about, and it is why two queries that look equally simple can have very different latency on a cluster even though neither would on a single node.
  • A statement confined to one bucket's worth of data still takes a fast, local path. ScramDB detects when everything a statement touches lives in a single shard group and, in that case, commits it without paying for the general cross-node transaction protocol. This is exactly why the guidance below centers on staying inside one bucket wherever you can.

How a distributed query actually executes​

At the level that matters to you: a query plan's individual operations, a scan, a filter, a join, an aggregation, run wherever the data they need already is, and intermediate results move between nodes only when a later step needs data that no single node already holds all of. Three shapes cover almost everything:

  • A scan can be pruned to the nodes that actually own the relevant buckets. If your filter identifies a specific value of a table's distribution key, the engine can compute exactly which bucket that value hashes to and only visit the node or nodes that own it, never touching the rest of the table's data at all.
  • A join or aggregation that needs rows from more than one node's buckets brought together picks between shipping the small side once, or repartitioning both sides by the join key. Both are genuine network transfers; the difference is how much data moves and to how many nodes. See the next section for exactly when each applies.
  • The engine can adapt some of this at runtime, not just at planning time, when you opt in. Under the materialize exchange mode (SET exchange_mode = 'materialize'; the default is streaming, which skips all of this), small transfers into the same destination can be coalesced into fewer, larger ones; a partition that turns out far larger than the others once execution is under way, a skewed key, can be split rather than left to bottleneck one node; a fragment running unusually slowly can get a speculative backup dispatched elsewhere. Each of these three behaviors has its own opt-in setting, covered in Practical guidance below. None of it changes a query's result, only how quickly it arrives, and none of it runs unless you turn it on.

The practical consequence: latency on a cluster is not just "how much work," it's "how much work, plus how many network hops the data involved needs to make, plus the size of what crosses on each hop." A query whose plan needs one hop over the low-latency network between nodes in the same cluster costs little beyond that hop's round trip; a query whose plan needs several rounds of data movement, or moves a large table across the network, costs accordingly.

Data distribution, as you experience it​

A new table with a primary key is distributed automatically the moment it's created in cluster mode, no special syntax required:

CREATE TABLE orders (
order_id bigint PRIMARY KEY,
customer_id bigint NOT NULL,
total numeric NOT NULL,
created_at timestamptz NOT NULL DEFAULT now()
);

ScramDB hashes each row's primary key to place it in one of [cluster] default_buckets buckets (8 by default), and each bucket is owned by as many nodes as your replication_factor calls for (see Multi-zone). Both the hash key and the bucket count are cluster-wide settings today, not something you choose per table: CREATE TABLE accepts a WITH (distribution_key = ..., buckets = ...) clause without error, but it does not currently change where a table's rows actually land in a running cluster, so don't rely on it. Treat the primary-key-hashed, default_buckets-wide default as the real, current behavior. A table with no primary key is not distributed across buckets at all.

Cheap: a query answerable from one bucket.

SELECT * FROM orders WHERE order_id = 482913;

An equality filter on the primary key lets the engine compute the one bucket that row lives in and go straight to the node or nodes that own it. The same is true of a single-row UPDATE or DELETE keyed the same way, and it's why point lookups and point writes on the primary key stay fast on a cluster.

Costly: a query with nothing to narrow the search.

SELECT * FROM orders WHERE total > 1000;

Filtering on a column other than the primary key gives the engine nothing to compute a bucket from, so it has to visit every bucket the table has. This is a fully-supported query; it costs proportionally more than the point lookup above, exactly as a full scan costs more than an index lookup on a single node, plus a cluster's extra step of gathering results from every node holding a piece of the table.

Grouping and aggregating is often cheaper than the scan it sits on top of. For an aggregate whose math can be computed in pieces and merged, sum, count, avg, min, max, and similar, each node computes its own partial result locally and ships back only that small partial result, not the raw rows; the coordinator does one final merge.

SELECT customer_id, sum(total)
FROM orders
GROUP BY customer_id;

This still pays for the full scan above, since nothing narrows it to one bucket, but the network cost past that scan is proportional to the number of distinct customer_id values that come back, not the number of rows in orders. An aggregate that can't be split this way is still answered, exactly, at a higher network cost. count(distinct col) repartitions the rows by col across the nodes, so each value is counted on one node, and when col is the distribution key nothing moves at all. string_agg, array_agg and any other aggregate called with DISTINCT stream the rows the query reads, after its filter and only the columns it names, to the node that received the query, which computes the aggregate there and sets rows aside on disk past its memory. Filter such a query down to a single bucket when you can, and it costs no more than a local one.

Joins move data when the two sides aren't already together, which is the normal case. Joining a large distributed table against a small one is cheap: the engine ships the small side to every node holding a piece of the large one, once, rather than moving the large table at all. Joining two large, distributed tables is the expensive case: both sides have to be repartitioned by the join key over the network so matching rows end up together, real traffic proportional to both tables' filtered size. There is no way today to make two independently-distributed tables share a placement so a join between them costs nothing, even when they're distributed on identically-named or identically-typed columns, so plan for a large-to-large join to move data rather than trying to design around it.

How a read picks its route​

Underneath the scan, join, and aggregate shapes above sits one more decision, and it runs on every read the cluster can route: a cost-based choice between a point lane, a single replica, and a full scatter, cheapest first. Read freshness and routing covers the freshness side of that choice in full, learner_read and max_staleness; this section covers the mechanics every route shares.

A point lookup takes one bounded hop, never the general path. A SELECT that resolves to exactly one row of one distribution key is a point lookup. When the node you are connected to leads that key's bucket, it answers directly off its own index. When another node leads it, this node forwards the request over a dedicated internal channel and waits, bounded, for the leader's answer, instead of running the statement through the full parse, plan, and scatter machinery a less targeted query needs. A send failure, a timeout, or a malformed reply steps aside to the general path automatically rather than hang or guess, so the shortcut is never a way to get a wrong answer, only occasionally a slower one. It only ever carries reads: a point write keeps the general commit path unconditionally, since a write's atomicity comes from the same machinery every other write already uses.

Above a size ceiling, nothing rides a single node, ever. However small a statement's own answer looks, if the data it has to read to produce that answer is large, it scatters across the replicas that hold the table's buckets the same way the full-scan case above does: a large query never becomes one replica's problem to carry alone.

Every one of these decisions renders as one line in EXPLAIN ANALYZE's Distributed Route: output and increments exactly one counter in /metrics, covered next.

Where fragments run​

A scatter runs one fragment per replica it uses, and each bucket of the table is read by exactly one replica, so the pieces add up to the whole table once. Any replica that holds a bucket can read it, and every one of them reads at the statement's one read position, so the answer is the same wherever a piece runs; only where the work lands changes. For each bucket the engine picks, in this order:

  • a replica in the region of the node you connected to, before one in another region;
  • within that, a learner, then a follower, then the shard group's leader last, so analytical scans leave the leaders' cores to the transactions they commit;
  • among equals, the replica carrying the fewest of this table's buckets so far, the node you connected to first on a tie.

A replica that cannot serve the read position yet (a learner or follower behind its group, or one still receiving a full copy of the group's rows) says so before it reads anything, and its buckets are read by another replica in the same statement; the engine then passes that replica over for that group until it catches up. A piece that ran and failed fails the statement, which you retry: a partial answer is never returned as a complete one. [cluster] fragment_any_replica = false keeps every piece on the buckets' owners instead. EXPLAIN ANALYZE shows the mix:

Distributed Route: scatter(4 fragments, strong); buckets on learners 3, followers 1, leaders 0

Practical guidance​

  • Favor point lookups and point writes on the primary key for latency-sensitive paths; they take the fast, single-bucket path described above.
  • Put the small side of a join on the small side. A join between a large fact table and a small dimension table is cheap however the optimizer orients it, so you don't have to hint it, but a query written so the optimizer can tell which side is actually small, filtered, narrow, helps it pick well.
  • Expect a JOIN between two large distributed tables, or a full scan with no primary-key filter, to move real data, and size your expectations accordingly; this isn't a bug or a missing optimization, it's the real cost of that query shape on data that isn't already together.
  • Reach for sum, count, avg, min, and max over count(distinct ...), string_agg, or array_agg on a distributed table where you have the choice: the former group stays cheap past the scan it sits on, while the latter group is answered by moving rows (or, for count(distinct ...), by repartitioning them), which costs network traffic in proportion to the rows read.
  • A transaction confined to rows in one bucket is materially cheaper than one that touches several. If your workload can be structured so a single transaction's writes usually land in one primary-key value's bucket, a single order's rows, a single tenant's rows, it consistently takes the fast local path instead of the general cross-node one.
  • Turn on the adaptive shuffle behaviors for a join- or aggregation-heavy workload that's genuinely moving data. SET exchange_mode = 'materialize' is the prerequisite (the default, streaming, skips all of it); aqe_coalesce, aqe_skew_split, and straggler_backup (each 'on' or 'off', all default 'off') opt into the specific behaviors from How a distributed query actually executes on top of it. None of them change a query's result. Their thresholds (the coalesce target, the skew factor and floor, the straggler delay) are the [cluster.distributed_join] keys of the config file.

Observing what happened​

EXPLAIN and EXPLAIN ANALYZE work exactly as they do on single-node ScramDB, and show the same plan tree: scan, filter, join, and aggregate operators, plus, for ANALYZE, planning time, execution time, and rows returned. Verify this before relying on it for cluster diagnosis: today's EXPLAIN output does not annotate which node a step runs on, or estimate how many bytes a join or aggregation will move across the network. It tells you the shape of the plan, genuinely useful for predicting whether a query is a point lookup or a full scan, but not which physical nodes did the work.

EXPLAIN ANALYZE (never plain EXPLAIN, which does not execute the statement and so does not know yet) adds one more line after the plan and the timings: Distributed Route: ..., naming the route the statement actually took as one of oltp, point(group N), local(strong) / local(stale), replica(NODE,strong) / replica(NODE,stale), or scatter(N fragments, strong); buckets on learners L, followers F, leaders D (every shape that scatters: a scan, an aggregate, a join, a k-NN search, a shuffle). See How a read picks its route above and Read freshness and routing for what each of those means and when a read qualifies for which.

The log is where the actual cross-node join decision shows up. Every join the engine places across the cluster logs one structured line naming which strategy it chose and the estimated bytes for every strategy it considered, for example:

distributed join placement: chosen=broadcast left=orders right=customers N=3 est_left=48000 est_right=1200 colocate=None broadcast=2400 shuffle=49200 preserved=[] alignment_excluded_broadcast=false

Grep your node logs for distributed join placement: to see exactly what a specific join did and why, including the estimated cost of the strategies it did not pick. The estimator always considers a same-node option too (colocate in that log line); for two independently created tables it is never actually viable today, so expect colocate=None and a choice between broadcast and shuffle in practice.

Prometheus metrics (/metrics, default port 9090) cover the cluster transport and the runtime adaptive behavior described above. The shuffle_aqe_* and shuffle_straggler_backups rows only move off zero once you've opted into the matching session setting from Practical guidance; seeing zeros there on a cluster that has never enabled them is expected, not a sign anything is broken:

MetricWhat it tells you
scramdb_cluster_bytes_sent_total / scramdb_cluster_bytes_received_totalTotal cross-node traffic, cumulative.
scramdb_cluster_frames_sent_total / scramdb_cluster_frames_received_totalMessage counts, cumulative.
scramdb_cluster_send_queue_bytesCurrent outbound backlog; sustained growth here means a peer isn't draining fast enough.
scramdb_cluster_routes_total{route="..."}Routing decisions from How a read picks its route: oltp, point, local_replica, forward_replica, or scatter; one per statement, and one per scatter for a statement that scatters more than once (each outer row of a k-NN LATERAL join). Fixed cardinality by construction, always exactly one of those five label values, never a node id or a group id, so this counter can never grow with the size of the cluster.
scramdb_cluster_exchange_placed_learner_buckets_total / _placed_follower_buckets_total / _placed_leader_buckets_totalBuckets of distributed reads this node placed on learners, followers and shard group leaders (see Where fragments run).
scramdb_cluster_exchange_placement_skips_total, scramdb_cluster_exchange_fragment_redispatches_totalReplicas passed over because they could not serve the read position, and fragments read by another replica after one could not serve them.
scramdb_cluster_exchange_fragments_served_learner_total, scramdb_cluster_exchange_fragments_refused_totalFragments this node served from its learner store, and fragments it refused because its replica could not serve the read position.
scramdb_cluster_shuffle_consumer_fragments_totalHow many shuffle transfers a join or aggregation actually dispatched.
scramdb_cluster_shuffle_aqe_coalesce_groups_totalHow often small transfers were merged into fewer, larger ones.
scramdb_cluster_shuffle_aqe_skew_split_flagged_total / _acted_totalHow often a disproportionately large partition was detected, and how often the engine actually split it in response.
scramdb_cluster_shuffle_aqe_broadcast_demote_flagged_totalHow often a broadcast side over its size budget was flagged as a candidate for a different strategy. Detection only today; no automatic correction follows yet.
scramdb_cluster_shuffle_straggler_backups_totalHow often a slow fragment got a speculative backup dispatched elsewhere.
scramdb_cluster_handshake_latency_secondsPeer connection setup latency; not per-query, but a useful health signal for the transport every distributed query rides on.
scramdb_cluster_dml_statements_distributed_totalUPDATE and DELETE statements on a distributed table that read rows from every owner because this node holds only part of the table (see UPDATE and DELETE across nodes).
scramdb_cluster_dml_candidates_remote_total / _local_total, scramdb_cluster_dml_candidate_batches_totalRows such statements read from other nodes and from this one, and the batches they arrived in.
scramdb_cluster_dml_candidate_fetch_secondsTime a statement waited for its next batch of rows.
scramdb_cluster_dml_candidate_spilled_bytes_totalRow bytes set aside on disk while a statement's memory budget was spoken for.
scramdb_cluster_dml_owners_pruned_totalOwners a statement did not ask because its WHERE named one bucket.
scramdb_cluster_dml_locking_selects_totalLocking SELECT statements (FOR UPDATE, FOR NO KEY UPDATE, FOR SHARE, FOR KEY SHARE) that read a distributed table's rows without locking them (see Row locks on distributed tables).
scramdb_cluster_dml_locking_read_rows_totalRows such statements took and registered one by one for their commit to check (a SERIALIZABLE statement without SKIP LOCKED has its scan checked instead, and adds nothing here unless a locked table has row-level security).
scramdb_cluster_dml_own_writes_reads_totalQueries that read a distributed table together with their transaction's own changes to it.
scramdb_cluster_own_writes_spilled_bytes_totalBytes of such a transaction's view of a table a query set aside on disk because the execution memory budget had no room for them.
scramdb_cluster_rows_moved_totalRows an UPDATE moved to a new primary key or distribution-key bucket.
scramdb_cluster_unique_probes_total, scramdb_cluster_unique_violations_totalNew keys an INSERT or key-changing UPDATE checked with the nodes that hold their buckets, and the ones refused because another row already held them.
scramdb_cluster_upsert_remote_conflicts_totalINSERT ... ON CONFLICT rows whose conflicting row is stored on another node.
scramdb_cluster_read_stability_retries_total, scramdb_cluster_read_stability_failures_totalReads re-run because a replica applied writes while they scanned, and reads that failed with 40001 because it kept doing so.
scramdb_cluster_serialization_retries_totalAutocommit statements sent again, and COPY votes voted again, after a serialization failure.

The log line and the metrics above show what a distributed query did. Where a table's buckets live is SHOW SHARDS, or the same rows as the scram.shards table (see Inspecting the cluster).

UPDATE and DELETE across nodes​

When every node holds every bucket of a table (a replication factor equal to the node count), an UPDATE or DELETE runs exactly as it does on a single node: it finds its rows in the local store and commits through the cluster. When the node you are connected to holds only some of a table's buckets, the same statement finds its rows everywhere instead:

  • Every bucket owner reads its own rows. The statement's WHERE travels to each owner as far as every node evaluates it identically (comparisons, IN, BETWEEN, LIKE and IS tests over numeric, boolean, text and uuid columns); a condition that calls a function, casts, or compares a date or time value is applied on the node you are connected to instead, which only costs rows on the wire, never a row the statement should have reached. A WHERE that fixes the whole distribution key to one value, UPDATE orders SET ... WHERE order_id = 42, asks only the owner of that value's bucket.
  • Rows stream back in batches. The node you are connected to holds only the batches it has not processed yet, charged to the execution memory budget and set aside on disk when that budget is spoken for, so a statement over a table of any size runs in bounded memory.
  • No row is locked, and nothing waits. A second writer on any node that reaches a row another transaction has changed is not held: its statement goes ahead, and the two are told apart when they commit. The first to commit wins; the other fails at its COMMIT with SQLSTATE 40001 for your retry logic, at every isolation level, READ COMMITTED included, and leaves nothing behind. A first writer that rolls back leaves the row to the second. See Row locks on distributed tables.
  • A key change moves the row. An UPDATE that changes a row's primary key or its distribution column moves the row to the bucket its new key belongs to, inside the same transaction, so a read by the new key finds it on every node and a read by the old key finds nothing. A new key that another row already holds, on any node, fails the statement with SQLSTATE 23505 and changes nothing.
  • Everything else is the statement you wrote. RETURNING, UPDATE ... FROM, DELETE ... USING, row triggers (once per changed row) and statement triggers, all fired on the node you are connected to, CHECK and foreign-key rules, and partitioned tables whose partitions are distributed all behave as on one node, and the affected-row count is the cluster-wide count.
  • A transaction reads its own changes. Inside a transaction that changed a table this node holds only part of, a later query of that table sees the changes: the table's committed rows are read from their owners and combined with the transaction's own changes on this node. A query over that table alone takes its WHERE along; a query joining it reads it whole, holding in memory what the execution memory budget has room for and setting the rest aside on disk, so such a join answers whatever the table's size.

A node that holds every bucket keeps the local path above at no added cost, and its statements move none of the scramdb_cluster_dml_* counters in the metrics table above.

COPY into a distributed table​

A COPY ... FROM into a distributed table is all or nothing, on every node at once, however large the input. The node you are connected to reads and checks the input a batch at a time exactly as an INSERT of those rows would (constraints, keys, foreign keys), so a bad row fails the COPY at once and names its line. It sends each batch's rows to the shard groups that own them, where every replica, learners included, writes them outside the table, invisible to every reader. At the end of the input the COPY commits once, and every replica turns its rows on at the same commit position:

  • All or nothing. A bad row at line 50,001, a key that two lines of the input (or two COPY statements of one transaction block) both take, a cancel, a lost client, or a crash of the coordinating node before the commit leaves none of the rows on any node; a replica that crashes mid-COPY leaves every node agreeing with the answer the COPY got, none of the rows if it failed and all of them if it committed. A COPY inside BEGIN commits or rolls back with the rest of the transaction.
  • None or all, for every reader. A query on any node sees none of the COPY's rows or all of them, never part.
  • Unique values hold. A COPY and a concurrent INSERT or COPY on another node that take the same primary key or unique value never both commit: the second to commit fails with 40001, as two INSERTs would.
  • Bounded by disk, not memory. The staged rows wait on each replica's disk; a COPY larger than the execution memory budget commits whole.
  • Nothing left behind. Rows a COPY staged and never committed are discarded when it fails, and a group leader discards what a coordinating node that went silent left behind after [cluster.dilith] copy_stage_idle_timeout (by default ten election timeouts).
  • A leader change does not stall it. When a shard group changes leader while a COPY, or a transaction whose rows are too large for one commit message, is sending its rows, whatever the old leader had not answered is sent again to the new leader: at once when the sending node holds a replica of the group and sees the change itself, otherwise on the commit protocol's own re-send interval ([cluster.dilith] retry_round_trips round trips), not after a whole commit wait. A COPY whose rows make no progress at all for [cluster.transactions] commit_answer_timeout fails with 40001 and leaves nothing.
  • Slow links are waited for. The re-send interval of each round of rows counts the time its bytes take at the rate the node's earlier rounds moved theirs, so rows that are only slow to cross the network are waited for rather than sent again, up to half of [cluster.dilith] copy_stage_idle_timeout (a group's leader drops staged rows whose sender stays silent that long, so a slower round is sent again to keep them); a link that suddenly becomes much slower than before can also see a round sent twice. scramdb_cluster_dilith_stage_rounds_resent_total counts the rounds sent again.

An autocommit COPY whose commit meets a concurrent writer of the same keys is retried on the server a few times before it fails with 40001. Two current limits are listed on Limitations: a COPY that stages for a long time while its shard groups commit many other transactions, and reading your own COPY's rows inside the same transaction block.

Changing a table on a cluster​

Every form of ALTER TABLE works on a cluster as it does on one node, issued through any node. The change reaches every node's catalog, and every copy of the table (each replica, a branch, a restore) takes it in order with the table's rows, so a row is always read as what it was written under.

  • Existing rows are checked once, for the whole cluster. SET NOT NULL, a NOT NULL column added with no default, a CHECK, a UNIQUE or primary key constraint and a foreign key judge every row of the table before the change is made, with PostgreSQL's error codes. NOT VALID and VALIDATE CONSTRAINT work as in PostgreSQL.
  • Writers retry around a change. While a distributed table's definition changes, and for a transaction that began before the change reached its node, a COMMIT that writes the table fails with SQLSTATE 40001; retry the transaction, as for any serialization failure. The change itself waits for transactions that hold the table as long as lock_timeout allows, on every node, and fails with SQLSTATE 55P03 when it runs out.
  • A column's default is read once. A column added with a default such as now() or random() gives every row already stored the same value on every node, the value the statement read; rows written later evaluate the default as usual.
  • The key that places the rows cannot change. Dropping a distributed table's primary key or a column of it, or changing a key column's type, is refused with SQLSTATE 0A000: create a new table with the key you want and copy the rows into it.
  • Dropping a table waits for its commits. DROP TABLE of a distributed table takes the table's lock on every node, waits for the commits already writing it, and has every replica apply them (each node waits at most [cluster] truncate_node_timeout) before the table goes. No row of the dropped table ever reaches a table created later under the same name. A node whose replicas do not catch up in time, or that restarted meanwhile, fails the statement with a message asking you to retry, and the table stays.
  • Rolling upgrades. While the cluster's nodes run different versions, ALTER TABLE of a distributed table is refused, and DROP TABLE of one cannot run either, since the nodes cannot agree on its drop across versions; upgrade every node first.

Distributed transactions and locking​

A transaction that touches rows on more than one node commits with the same serializable guarantee as a single-node transaction, and the same SQLSTATE 40001 retry contract on conflict; see Cluster overview for that guarantee stated in full and Failover for what happens if the node coordinating a commit fails partway through. Underneath, a distributed commit is one round: each shard group the transaction touched logs its vote, and the transaction is committed once every one of them holds a durable yes vote. Each vote names every participant, so if the coordinating node crashes mid-commit, another node finishes the transaction from the votes alone rather than leaving affected rows stuck; no lock lease or timeout decides the outcome.

Two real, SQL-level locking primitives are cluster-wide, not just local to the node you happen to be connected to:

  • LOCK TABLE, used inside a transaction block exactly as in PostgreSQL, takes a lock that every node in the cluster respects: a conflicting statement on any node, not just the one that issued the LOCK TABLE, waits for it (or fails immediately with NOWAIT), following the same lock-mode conflict rules PostgreSQL uses (a plain read takes ACCESS SHARE, a write ROW EXCLUSIVE, TRUNCATE / ALTER TABLE / DROP / CREATE INDEX take ACCESS EXCLUSIVE, and so on).
  • pg_advisory_lock and the rest of that family (pg_try_advisory_lock, pg_advisory_unlock, and their transaction-scoped variants) are cluster-wide the same way: a lock taken on one node is respected by every other node.

Both block by default with no built-in deadline, the same as PostgreSQL itself; bound the wait with your client's statement_timeout if you need one. Rows of a distributed table are not locked at all, as the next section explains.

Row locks on distributed tables​

PostgreSQL runs on one node, and its row locks are local: a transaction that wants a row another transaction has locked (with SELECT ... FOR UPDATE or FOR SHARE, or by changing it) waits until that one ends. A distributed table has no such locks, on any node. FOR UPDATE, FOR NO KEY UPDATE, FOR SHARE and FOR KEY SHARE read the rows they select, and UPDATE and DELETE change theirs, without locking them and without waiting for anyone. What each transaction read is checked when it commits:

  • The first to commit wins. Of two transactions that change the same row, on the same node or on two, whether or not they selected it FOR UPDATE first, the one that commits first commits; the other fails at its COMMIT with SQLSTATE 40001 (could not serialize access) and changes nothing. This holds at every isolation level, READ COMMITTED included, so no update is ever lost. A row a transaction only selected FOR UPDATE counts as read: if another transaction changed it and committed first, a transaction that writes anything is refused the same way, while one that writes nothing commits, having read one consistent state.
  • NOWAIT refuses a row whose writer is committing right now, at once, with SQLSTATE 55P03 (could not obtain lock on row in relation ...). A row changed by a transaction that has not reached its COMMIT yet is still returned, as the first rule describes. Under SERIALIZABLE, a row the statement read as it stood before a commit in flight, without taking it, counts as read: if that commit lands, the transaction fails at its COMMIT with 40001 when it writes anything.
  • SKIP LOCKED leaves out rows whose writer is committing right now, without waiting for that commit, and a row it left out never fails its own commit. A row whose writer has not reached its COMMIT yet is not left out: two workers can take the same row, and the second to commit gets 40001.
  • SKIP LOCKED and NOWAIT never wait, at every isolation level and outside a transaction block alike. A commit counts as in flight until the replica serving the read has applied it, so a row whose commit was acknowledged a moment ago can still be left out or refused where PostgreSQL would return it.
  • A foreign key holds without a lock. A child row's check of its parent, and a parent's check for children still naming it, count as reads of those keys at every isolation level. Of a transaction that deletes a parent row or changes its key and a concurrent one that inserts or updates a child naming it, one fails at its COMMIT with 40001, so no child is ever left without its parent; on a single node PostgreSQL would make one of them wait instead.
  • lock_timeout has no effect, and no row lock takes part in a deadlock, so 40P01 never comes from a row of a distributed table.

What differs from a single node is when the loser finds out and what it gets. On a single node the second writer waits, and then under READ COMMITTED applies its change on top of the first writer's committed row, while REPEATABLE READ and SERIALIZABLE fail with 40001 at the statement. On a distributed table the second writer never waits, and its COMMIT fails with 40001 whatever the isolation level. Code written for PostgreSQL's 40001 retry loop needs no change; code that relies on READ COMMITTED quietly applying on top of a concurrent change, or on SKIP LOCKED handing two workers disjoint rows before either commits, sees 40001 at commit instead. A table a cluster node keeps to itself (one outside the cluster's shard map) keeps PostgreSQL's waiting row locks, bounded by lock_timeout.

Next​