InletDownload

PostgreSQL EXPLAIN

Merge Join in PostgreSQL EXPLAIN

A Merge Join reads two inputs that are both sorted on the join key and walks them side by side, like merging two sorted lists. It shines when indexes already provide the order; when it needs big Sort nodes underneath, it’s often slower than a hash join.

Updated 9 October 2026

What it does

A Merge Join needs both inputs sorted on the join columns. It reads the first row of each, compares the keys, and advances whichever side is behind. When the keys match it outputs the joined rows. Each input is read once, from start to finish (or until it can stop), and no hash table is built.

The order comes either from an index scan that returns rows in key order or from a Sort node placed under the join. The joined rows come out in join-key order, which can save a sort later for ORDER BY on the same key.

Two details explain odd-looking numbers:

  • It can stop early. Once one input runs out, there can be no more matches, so it stops reading the other. The plan then shows far fewer actual rows on that side than estimated.
  • Duplicate keys mean re-reading. If several rows on the outer side share a key, the join goes back and re-reads the matching inner rows for each of them. Many duplicates make the inner side do more work than its row count suggests.

Variants: Merge Left Join, Merge Full Join, Merge Semi Join, Merge Anti Join and so on. Merge joins, like hash joins, need an equality condition.

When the planner picks it

  • Both inputs are large and already sorted, typically through indexes on both join columns.
  • The query wants the result ordered by the join key (ORDER BY o.id when joining on o.id), so a merge join’s output order saves a sort.
  • The hashed side of a hash join wouldn’t fit in memory, and sorting both sides is estimated cheaper.
  • One side has a limited range of keys and the other is an index: the join can stop after reading only the overlapping part of the index.

PostgreSQL 18 can also feed a merge join through an Incremental Sort when an input is already sorted on a prefix of the join key.

Reading its numbers

Merge Join  (cost=13043.44..25098.19 rows=60793 width=22) (actual time=12.369..25.646 rows=60000.00 loops=1)
  Merge Cond: (o.id = l.order_id)
  ->  Index Scan using orders_pkey on orders o  (cost=0.42..34425.43 rows=1000000 width=14) (actual time=0.043..4.208 rows=20001.00 loops=1)
  ->  Sort  (cost=13042.98..13194.96 rows=60793 width=16) (actual time=12.320..14.294 rows=60000.00 loops=1)
        Sort Key: l.order_id
        Sort Method: quicksort  Memory: 3412kB
  • Merge Cond: the equality the inputs are merged on.
  • Join Filter / Rows Removed by Join Filter: any further join conditions, checked on matched pairs.
  • An input with fewer actual rows than estimated: here the index scan on orders was estimated at 1,000,000 rows and read 20,001. That’s the early stop, not a bad estimate: the other side only had order ids up to 20,000. The PostgreSQL docs call this out as a measurement artefact.
  • A Sort child: its Sort Method tells you whether the sort fitted in work_mem (quicksort) or used temporary files (external sort / external merge with Disk:).
  • Startup time (actual time=12.369..): the first row can’t come out until any Sort below has finished.

When it’s a problem

  • Big Sorts underneath. If both inputs need sorting and they’re large, the sorts are most of the cost, and they may spill to disk. A hash join usually wins then; if the planner still chose a merge join, compare with SET enable_mergejoin = off in your session (as a test only). An index on the join column of the sorted side removes that sort.
  • The index scan visits many pages. An index provides order, but if the table isn’t stored in that order each row may be on a different page. In the second example below, merging 900,000 order lines through their primary key cost 895,000 buffer accesses because the table’s rows aren’t stored in order_id order.
  • Many duplicate keys on both sides make the inner side re-read the same rows repeatedly. Joining on a more selective key, or aggregating one side first, helps.
  • Sorts that spill. Raise work_mem for the query (SET LOCAL work_mem = '64MB' in a transaction) if the sort is the slow part and you can’t avoid it.

Example

PostgreSQL 18.6, default settings:

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;

-- three lines for each of the first 300,000 orders
CREATE TABLE order_lines (
  order_id   bigint NOT NULL,
  line_no    int NOT NULL,
  product_id int NOT NULL,
  qty        int NOT NULL,
  PRIMARY KEY (order_id, line_no)
);
INSERT INTO order_lines
SELECT o, l, 1 + ((o * 31 + l * 17) % 2000) * 50, 1 + (o + l) % 5
FROM generate_series(1, 300000) AS o, generate_series(1, 3) AS l;
VACUUM ANALYZE orders, order_lines;

(Our test orders table also had a foreign key and two secondary indexes; they aren’t used here.)

The lines of the first 20,000 orders, in order:

EXPLAIN (ANALYZE, BUFFERS)
SELECT o.id, o.total, l.product_id, l.qty
FROM orders o JOIN order_lines l ON l.order_id = o.id
WHERE l.order_id <= 20000
ORDER BY o.id;
Merge Join  (cost=13043.44..25098.19 rows=60793 width=22) (actual time=12.369..25.646 rows=60000.00 loops=1)
  Merge Cond: (o.id = l.order_id)
  Buffers: shared hit=752 read=167
  ->  Index Scan using orders_pkey on orders o  (cost=0.42..34425.43 rows=1000000 width=14) (actual time=0.043..4.208 rows=20001.00 loops=1)
        Index Searches: 1
        Buffers: shared hit=57 read=167
  ->  Sort  (cost=13042.98..13194.96 rows=60793 width=16) (actual time=12.320..14.294 rows=60000.00 loops=1)
        Sort Key: l.order_id
        Sort Method: quicksort  Memory: 3412kB
        Buffers: shared hit=695
        ->  Bitmap Heap Scan on order_lines l  (cost=1719.57..8212.48 rows=60793 width=16) (actual time=1.483..5.791 rows=60000.00 loops=1)
              Recheck Cond: (order_id <= 20000)
              Heap Blocks: exact=386
              Buffers: shared hit=695
              ->  Bitmap Index Scan on order_lines_pkey  (cost=0.00..1704.37 rows=60793 width=0) (actual time=1.449..1.449 rows=60000.00 loops=1)
                    Index Cond: (order_id <= 20000)
                    Index Searches: 1
                    Buffers: shared hit=309
Planning:
  Buffers: shared hit=22 read=1
Planning Time: 0.173 ms
Execution Time: 27.719 ms

The order lines were fetched and sorted (60,000 rows, 3.4 MB in memory). The orders came pre-sorted from the primary key, and the join stopped reading them at order 20,001. The output is already in o.id order, so there’s no Sort at the top for the ORDER BY.

The same query on PostgreSQL 14.24 chose the same plan, but its Sort didn’t fit in the default 4 MB:

->  Sort  (cost=12796.84..12944.00 rows=58863 width=16) (actual time=78.068..81.557 rows=60000 loops=1)
      Sort Key: l.order_id
      Sort Method: external sort  Disk: 1768kB

PostgreSQL 15 reduced the memory in-memory sorts need, which is why 18 sorted the same rows in memory.

All lines of all orders, with no WHERE:

EXPLAIN (ANALYZE, BUFFERS)
SELECT o.id, o.total, l.product_id, l.qty
FROM orders o JOIN order_lines l ON l.order_id = o.id
ORDER BY o.id;
Merge Join  (cost=7.64..76187.54 rows=900000 width=22) (actual time=0.008..502.617 rows=900000.00 loops=1)
  Merge Cond: (o.id = l.order_id)
  Buffers: shared hit=895188 read=12735
  ->  Index Scan using orders_pkey on orders o  (cost=0.42..34425.43 rows=1000000 width=14) (actual time=0.005..79.902 rows=300001.00 loops=1)
        Index Searches: 1
        Buffers: shared hit=224 read=3099
  ->  Index Scan using order_lines_pkey on order_lines l  (cost=0.42..53794.22 rows=900000 width=16) (actual time=0.002..321.530 rows=900000.00 loops=1)
        Index Searches: 1
        Buffers: shared hit=894964 read=9636
Planning Time: 0.117 ms
Execution Time: 528.780 ms

No sorts at all: both sides come from their primary keys. But look at the buffers on order_lines: 894,964 page accesses for a 5,733-page table. The rows were inserted line number by line number, not order by order, so walking the index in order_id order jumps to a different page for almost every row.

In Inlet

Inlet draws EXPLAIN ANALYZE as a tree and highlights the slowest step and badly misestimated row counts. An early-stopping input like the 20,001 rows above shows a big gap between estimate and actual; that one is harmless, which is why it pays to read the node before acting on it.

Related

Sources