InletDownload

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 BY on 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 BY done 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 Planned and Workers 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=3 and rows=10000.00: three processes, each sorting about 10,000 rows (an average).
  • Sort Method per process: quicksort (in memory), top-N heapsort (in memory, only the top rows kept), or external merge with Disk: (it ran out of work_mem and used temporary files).
  • Each process has its own work_mem budget, so a parallel sort can use up to (workers + 1) × work_mem in 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. Raise work_mem for 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, a LIMIT query 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 single Index Scan Backward that 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.

Related

Sources