Glossary: Planner
work_mem is per node, not per query
Also called: sort memory, hash memory, hash_mem_multiplier.
Definition, revised in place. Last updated .
work_mem is how much memory a single executor node may use before it finishes its work on disk instead. It is a unit, not a budget. One statement may contain several nodes that each take the full amount, a parallel plan takes it once per worker per node, and hash-based nodes take a multiple of it set by hash_mem_multiplier. Nothing divides the setting between the nodes of a plan, and nothing caps the total.
Multiplying the unit
The arithmetic people do when sizing a server is connections multiplied by work_mem, and it understates the answer in two directions at once. A single connection running a query with two sorts and a hash join has three nodes, each entitled to the setting, and the hash join is entitled to a multiple of it. If that query goes parallel with two workers, most of those allocations happen once per worker.
Nobody should compute the theoretical maximum from that, because it describes a plan shape no workload actually runs. What it does mean is that the setting is a knob on the shape of individual plans rather than a memory ceiling, and that raising it globally to fix one report changes the memory profile of every statement on the server. Setting it for the session, the role or the one statement that needs it is the version of that decision with a bounded blast radius.
Two sorts, one setting
A merge join needs both inputs sorted, so it contains two sort nodes and the plan reports each one separately.
CREATE TABLE reading AS
SELECT g AS seq_key, ((g::bigint * 7919) % 500000)::int AS scattered_key, repeat('.', 60) AS pad
FROM generate_series(1, 500000) g;
SET max_parallel_workers_per_gather = 0;
SET enable_hashjoin = off;
SET enable_nestloop = off;
SET work_mem = '48MB';
EXPLAIN (ANALYZE, COSTS OFF, TIMING OFF, SUMMARY OFF, BUFFERS OFF)
SELECT count(*) FROM reading a JOIN reading b ON a.scattered_key = b.seq_key;
SELECT 500000
SET
SET
SET
SET
QUERY PLAN
----------------------------------------------------------------------
Aggregate (actual rows=1 loops=1)
-> Merge Join (actual rows=499999 loops=1)
Merge Cond: (a.scattered_key = b.seq_key)
-> Sort (actual rows=500000 loops=1)
Sort Key: a.scattered_key
Sort Method: quicksort Memory: 35726kB
-> Seq Scan on reading a (actual rows=500000 loops=1)
-> Sort (actual rows=500000 loops=1)
Sort Key: b.seq_key
Sort Method: quicksort Memory: 35726kB
-> Seq Scan on reading b (actual rows=500000 loops=1)
(11 rows)
That is 14.24, where each sort took about thirty-five megabytes and the pair took about seventy, against a setting of forty-eight. Neither node spilled, because neither node individually exceeded the limit. The same statement on each of the four later versions this ran against sorts the same rows in about twelve megabytes each, which is worth knowing before carrying a tuned value across a major upgrade.
Setting it somewhere narrower than the server
The useful question is which statements actually need more, and the answer comes from the plans rather than from the setting: a node reporting an external merge wanted more than it had. When a node that used to fit starts spilling with no change in data volume, the grant is usually not the problem and the row estimate is, which is where query plan regression picks the story up. The disk side of the same event is temporary files.
Multiplying that out for a whole server, including the hash multiplier, the parallel workers and the autovacuum workers running beside them, is what the memory budget calculator does; it also checks whether effective_cache_size is promising the planner a cache that has room to exist.