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
materializeexchange mode (SET exchange_mode = 'materialize'; the default isstreaming, 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
JOINbetween 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, andmaxovercount(distinct ...),string_agg, orarray_aggon 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, forcount(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, andstraggler_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:
| Metric | What it tells you |
|---|---|
scramdb_cluster_bytes_sent_total / scramdb_cluster_bytes_received_total | Total cross-node traffic, cumulative. |
scramdb_cluster_frames_sent_total / scramdb_cluster_frames_received_total | Message counts, cumulative. |
scramdb_cluster_send_queue_bytes | Current 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_total | Buckets 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_total | Replicas 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_total | Fragments 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_total | How many shuffle transfers a join or aggregation actually dispatched. |
scramdb_cluster_shuffle_aqe_coalesce_groups_total | How often small transfers were merged into fewer, larger ones. |
scramdb_cluster_shuffle_aqe_skew_split_flagged_total / _acted_total | How often a disproportionately large partition was detected, and how often the engine actually split it in response. |
scramdb_cluster_shuffle_aqe_broadcast_demote_flagged_total | How 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_total | How often a slow fragment got a speculative backup dispatched elsewhere. |
scramdb_cluster_handshake_latency_seconds | Peer connection setup latency; not per-query, but a useful health signal for the transport every distributed query rides on. |
scramdb_cluster_dml_statements_distributed_total | UPDATE 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_total | Rows such statements read from other nodes and from this one, and the batches they arrived in. |
scramdb_cluster_dml_candidate_fetch_seconds | Time a statement waited for its next batch of rows. |
scramdb_cluster_dml_candidate_spilled_bytes_total | Row bytes set aside on disk while a statement's memory budget was spoken for. |
scramdb_cluster_dml_owners_pruned_total | Owners a statement did not ask because its WHERE named one bucket. |
scramdb_cluster_dml_locking_selects_total | Locking 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_total | Rows 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_total | Queries that read a distributed table together with their transaction's own changes to it. |
scramdb_cluster_own_writes_spilled_bytes_total | Bytes 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_total | Rows an UPDATE moved to a new primary key or distribution-key bucket. |
scramdb_cluster_unique_probes_total, scramdb_cluster_unique_violations_total | New 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_total | INSERT ... ON CONFLICT rows whose conflicting row is stored on another node. |
scramdb_cluster_read_stability_retries_total, scramdb_cluster_read_stability_failures_total | Reads 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_total | Autocommit 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
WHEREtravels to each owner as far as every node evaluates it identically (comparisons,IN,BETWEEN,LIKEandIStests 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. AWHEREthat 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
COMMITwithSQLSTATE 40001for your retry logic, at every isolation level,READ COMMITTEDincluded, 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
UPDATEthat 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 withSQLSTATE 23505and 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,CHECKand 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
WHEREalong; 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
COPYstatements 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-COPYleaves every node agreeing with the answer theCOPYgot, none of the rows if it failed and all of them if it committed. ACOPYinsideBEGINcommits 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
COPYand a concurrentINSERTorCOPYon another node that take the same primary key or unique value never both commit: the second to commit fails with40001, as twoINSERTs would. - Bounded by disk, not memory. The staged rows wait on each replica's disk; a
COPYlarger than the execution memory budget commits whole. - Nothing left behind. Rows a
COPYstaged 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_tripsround trips), not after a whole commit wait. ACOPYwhose rows make no progress at all for[cluster.transactions] commit_answer_timeoutfails with40001and 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_totalcounts 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, aNOT NULLcolumn added with no default, aCHECK, aUNIQUEor primary key constraint and a foreign key judge every row of the table before the change is made, with PostgreSQL's error codes.NOT VALIDandVALIDATE CONSTRAINTwork 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
COMMITthat writes the table fails withSQLSTATE 40001; retry the transaction, as for any serialization failure. The change itself waits for transactions that hold the table as long aslock_timeoutallows, on every node, and fails withSQLSTATE 55P03when it runs out. - A column's default is read once. A column added with a default such as
now()orrandom()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 TABLEof 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 TABLEof a distributed table is refused, andDROP TABLEof 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 theLOCK TABLE, waits for it (or fails immediately withNOWAIT), following the same lock-mode conflict rules PostgreSQL uses (a plain read takesACCESS SHARE, a writeROW EXCLUSIVE,TRUNCATE/ALTER TABLE/DROP/CREATE INDEXtakeACCESS EXCLUSIVE, and so on).pg_advisory_lockand 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 UPDATEfirst, the one that commits first commits; the other fails at itsCOMMITwithSQLSTATE 40001(could not serialize access) and changes nothing. This holds at every isolation level,READ COMMITTEDincluded, so no update is ever lost. A row a transaction only selectedFOR UPDATEcounts 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. NOWAITrefuses a row whose writer is committing right now, at once, withSQLSTATE 55P03(could not obtain lock on row in relation ...). A row changed by a transaction that has not reached itsCOMMITyet is still returned, as the first rule describes. UnderSERIALIZABLE, 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 itsCOMMITwith40001when it writes anything.SKIP LOCKEDleaves 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 itsCOMMITyet is not left out: two workers can take the same row, and the second to commit gets40001.SKIP LOCKEDandNOWAITnever 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
COMMITwith40001, so no child is ever left without its parent; on a single node PostgreSQL would make one of them wait instead. lock_timeouthas no effect, and no row lock takes part in a deadlock, so40P01never 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
- Node discovery and Forming a cluster if you haven't stood up the cluster you're about to query.
- Read freshness and routing for the freshness knobs and the size ceiling behind the routes named above.
- Multi-zone and Multi-region for how node placement affects the network hops described above.
- Failover for what a query in flight experiences if a node it depends on goes away mid-statement.