PostgreSQL EXPLAIN
Gather Merge in PostgreSQL EXPLAIN
Gather Merge collects rows from parallel workers like Gather does, but each worker hands over its rows already sorted and the leader merges them, so the output stays in order. You’ll usually find a Sort in each worker beneath it.
Updated 9 October 2026
What it does
Gather Merge is the ordered version of Gather. It starts the parallel workers and
runs the plan below it in each of them and in the leader. The difference is what it expects back: every
process produces its rows already sorted, and Gather Merge merges those sorted streams into one, the way
you’d merge sorted piles of cards by always taking the smallest top card.
The merge is cheap. The expensive part is usually what’s under it: a Sort in each process, or an index scan that returns rows in order.
When the planner picks it
When a parallel plan needs its output in order:
ORDER BYon a large table:Gather Merge→Sort→Parallel Seq Scan.ORDER BY … LIMIT n: each process keeps only its own top n (a top-N heapsort), so very little crosses to the leader.GROUP BYdone as a sorted aggregate:Finalize GroupAggregate→Gather Merge→Sort→Partial HashAggregate. The leader needs the partial groups in order to combine them.- Merge joins and window functions above a parallel part.
Reading its numbers
Gather Merge shows the same fields as Gather:
Workers PlannedandWorkers Launched: asked for and actually started. Fewer started means the server’s worker pool (max_parallel_workers,max_worker_processes) was busy.rows: the total from all processes, in order.
The node beneath it usually tells you more. A parallel Sort prints the leader’s sort method first, then one line per worker:
-> Sort (cost=14376.24..14406.94 rows=12278 width=18) (actual time=45.677..46.070 rows=10000.00 loops=3)
Sort Key: total DESC
Sort Method: quicksort Memory: 786kB
Worker 0: Sort Method: quicksort Memory: 788kB
Worker 1: Sort Method: quicksort Memory: 752kB
loops=3androws=10000.00: three processes, each sorting about 10,000 rows (an average).Sort Methodper process:quicksort(in memory),top-N heapsort(in memory, only the top rows kept), orexternal mergewithDisk:(it ran out ofwork_memand used temporary files).- Each process has its own
work_membudget, so a parallel sort can use up to (workers + 1) ×work_memin total.
When it’s a problem
- The sorts spill to disk.
Sort Method: external merge Disk: …in the workers means each one wrote temporary files. Raisework_memfor the session or the one query (SET LOCAL work_mem = '64MB'inside a transaction) rather than server-wide, because every sort in every session can use that much. - An index could provide the order. If an index matches the filter and the
ORDER BY, aLIMITquery can read the first rows straight from it, with no sort and no parallel workers. For the example below,CREATE INDEX ON orders (status, total)turned the plan into a singleIndex Scan Backwardthat read 23 pages and finished in 0.14 ms instead of 17.7 ms. The trade-off is the usual one for indexes: more disk, slightly slower writes. - The output is huge. All rows still pass through the leader one at a time. For a large sorted result, the merge in the leader can become the bottleneck; consider whether the client really needs every row.
- Fewer workers launched than planned. Same cause and fixes as Gather.
Example
PostgreSQL 18.6, default settings. orders has 1,000,000 rows, 30,000 of them refunded (setup below).
Refunded orders, largest first:
EXPLAIN (ANALYZE, BUFFERS)
SELECT id, customer_id, total FROM orders WHERE status = 'refunded' ORDER BY total DESC;
Gather Merge (cost=15376.27..18808.18 rows=29467 width=18) (actual time=47.351..53.918 rows=30000.00 loops=1)
Workers Planned: 2
Workers Launched: 2
Buffers: shared hit=74 read=8334
-> Sort (cost=14376.24..14406.94 rows=12278 width=18) (actual time=45.677..46.070 rows=10000.00 loops=3)
Sort Key: total DESC
Sort Method: quicksort Memory: 786kB
Buffers: shared hit=74 read=8334
Worker 0: Sort Method: quicksort Memory: 788kB
Worker 1: Sort Method: quicksort Memory: 752kB
-> Parallel Seq Scan on orders (cost=0.00..13542.33 rows=12278 width=18) (actual time=0.305..42.218 rows=10000.00 loops=3)
Filter: (status = 'refunded'::text)
Rows Removed by Filter: 323333
Buffers: shared read=8334
Planning:
Buffers: shared hit=143 read=3
Planning Time: 0.371 ms
Execution Time: 54.848 ms
Each process scanned a third of the table, sorted its 10,000 rows in memory, and Gather Merge merged the three sorted streams into 30,000 rows in order.
With LIMIT 20, each process keeps only its top 20 in a small heap, and Gather Merge stops after 20
rows:
EXPLAIN (ANALYZE, BUFFERS)
SELECT id, customer_id, total FROM orders WHERE status = 'refunded' ORDER BY total DESC LIMIT 20;
Limit (cost=14869.07..14871.40 rows=20 width=18) (actual time=16.535..17.707 rows=20.00 loops=1)
Buffers: shared hit=203 read=8205
-> Gather Merge (cost=14869.07..18300.99 rows=29467 width=18) (actual time=16.534..17.704 rows=20.00 loops=1)
Workers Planned: 2
Workers Launched: 2
Buffers: shared hit=203 read=8205
-> Sort (cost=13869.05..13899.74 rows=12278 width=18) (actual time=15.264..15.265 rows=17.00 loops=3)
Sort Key: total DESC
Sort Method: top-N heapsort Memory: 27kB
Buffers: shared hit=203 read=8205
Worker 0: Sort Method: top-N heapsort Memory: 27kB
Worker 1: Sort Method: top-N heapsort Memory: 27kB
-> Parallel Seq Scan on orders (cost=0.00..13542.33 rows=12278 width=18) (actual time=0.094..14.057 rows=10000.00 loops=3)
Filter: (status = 'refunded'::text)
Rows Removed by Filter: 323333
Buffers: shared hit=129 read=8205
Planning Time: 0.077 ms
Execution Time: 17.726 ms
rows=17.00 loops=3 on the Sort: on average each process handed over 17 rows before the merge had its
20.
The setup:
CREATE SCHEMA seo_explain;
SET search_path = seo_explain;
CREATE TABLE orders (
id bigint PRIMARY KEY,
customer_id int NOT NULL,
status text NOT NULL,
total numeric(10,2) NOT NULL,
created_at timestamptz NOT NULL
);
INSERT INTO orders
SELECT i,
1 + (i::bigint * 7919) % 50000,
CASE WHEN i % 100 < 90 THEN 'shipped'
WHEN i % 100 < 97 THEN 'pending'
ELSE 'refunded' END,
round(((i * 37) % 100000) / 100.0, 2),
timestamptz '2024-01-01' + i * interval '1 minute'
FROM generate_series(1, 1000000) AS i;
VACUUM ANALYZE orders;
(Our test table also had a foreign key to a customers table and indexes on customer_id and
created_at; none of them are used by these queries.)
In Inlet
Inlet draws EXPLAIN ANALYZE as a tree and highlights the slowest step and badly misestimated row
counts, so you can see at once whether the time is in the merge, the sorts or the scan.