max_parallel_workers_per_gather = 2 gets you three processes. The third is the backend your client is connected to, the leader, which launches the two workers and then, instead of waiting for them, runs its own copy of the same plan. parallel_leader_participation decides whether it does that. It is a boolean, the default is on, the context is user (any session can change it, and ALTER ROLE ... SET works), and it arrived in PostgreSQL 11.
You have already seen the leader at work, in the line of EXPLAIN (ANALYZE, VERBOSE) output that confuses everyone the first time (costs, timing and buffers switched off here to keep it short):
1 Finalize Aggregate (actual rows=1.00 loops=1)
2 Output: count(*)
3 -> Gather (actual rows=3.00 loops=1)
4 Output: (PARTIAL count(*))
5 Workers Planned: 2
6 Workers Launched: 2
7 -> Partial Aggregate (actual rows=1.00 loops=3)
8 Output: PARTIAL count(*)
9 Worker 0: actual rows=1.00 loops=1
10 Worker 1: actual rows=1.00 loops=1
11 -> Parallel Seq Scan on public.t (actual rows=3333333.33 loops=3)
12 Output: id, k, pad
13 Worker 0: actual rows=3778598.00 loops=1
14 Worker 1: actual rows=2529480.00 loops=1
Two workers, loops=3. Row counts under a Gather are per-loop averages, and VERBOSE prints a line for each worker but none for the leader; you find its share by subtraction. Here it scanned 3,691,922 of ten million rows, 37% of the table. Set the parameter to off and the same plan reports loops=2, the workers take about five million rows each, and the leader spends the query showing IPC:ExecuteGather in pg_stat_activity.
A worker with a second job
The leader has a duty the workers do not: emptying their tuple queues. Each worker writes its rows into a 64kB queue and stops when the queue is full, so a leader that wanders off stalls everybody. The Gather loop is written accordingly. It reads from the queues first, and only when every queue is empty does it run its own copy of the plan until that produces a row, and then it checks the queues again. A plan whose workers send up one row each (the partial aggregate above) leaves the leader free to be a full participant. A plan that pours rows through the Gather keeps it at the queues more of the time. (Gather Merge is stricter about it, since the leader is one of the inputs to the merge and has to produce its next row whenever the sort order calls for it.)
The planner’s model of this is one fixed guess, applied to every plan: the leader spends 30% of its time servicing each worker. That makes the leader worth 1 − 0.3 × workers of a worker: 0.7 with one worker, 0.4 with two, 0.1 with three, and nothing at all from four up. The leader on my test machine had not read the formula. With one worker it scanned about half the table, where the planner expected 41%. With two it scanned between 30% and 45% from run to run, where the planner expected 17%. With nine million rows streaming up through the Gather, it still scanned between 23% and 31% of them itself. (That machine has two cores, so at two workers there were three processes sharing them.)
Turning the parameter off takes that fraction out of the planner’s arithmetic as well as taking the leader out of the executor. At two workers, the estimated rows per process on the scan above go from 4,167,052 to 5,000,463 and the total cost of the plan from 146,592 to 157,010, so parallel plans lose to serial ones a little more often. At four workers or more the costs do not move, because the planner had already written the leader off.
At one worker, a parallel scan costs what the serial scan costs with parallel_setup_cost added, and loses. With max_parallel_workers_per_gather = 1 and the leader off, every scan, sort and aggregate I tried went serial. The exception was not an improvement. Parallel Hash is allowed a shared hash table of work_mem × hash_mem_multiplier × (workers + 1) before it has to split into batches, and the + 1 is there whether the leader shows up or not, so at work_mem = '64MB' the planner chose a Parallel Hash Join for the memory. One worker ran all of it while the leader watched, and it finished no sooner than the serial plan did.
What off is for
Testing, mostly. Robert Haas committed Thomas Munro’s patch in November 2017 with the note that it was sometimes useful for testing and that “it’s possible that could work out better even in production,” which is about as lukewarm as a commit message gets. PostgreSQL’s own regression and isolation tests turn it off thirteen times. It is also the right setting when you are measuring how a query scales with worker count, since otherwise your one-worker data point is two processes.
The production argument, per the documentation, is that workers stall less often when the leader is not busy elsewhere. A stalled worker shows IPC:MessageQueueSend in pg_stat_activity. I went looking for the effect and did not find it. Streaming those nine million rows, the workers were in MessageQueueSend in half to two-thirds of my samples with the leader on, and in half to two-thirds of them with it off, while a leader with nothing else to do sat idle in ExecuteGather for about 45% of the query. The query ran 3% to 11% slower with the leader off, in every pair of runs. Seeing that wait event on a worker does not tell you the leader was distracted.
Kaarel Moppel has better numbers than two cores can produce: a full-scan aggregate on 16-vCPU machines, with max_parallel_workers_per_gather at 2, 4, 8 and 16. With the data in memory, turning the leader off made the query about 19% slower. With a table twice the size of RAM, it was 0.6% faster unpartitioned and 6% faster partitioned, and 11% was the best result on any of his six machines.
The cost the documentation warns about is that with the leader off, nothing comes back until a worker has started. That is true and small: on my test machine the first row arrived in about 0.15 ms with the leader on and in 2 to 3 ms with it off. Nor can off strand a query. If no workers can be had (see max_parallel_workers for how that happens), the leader runs the plan itself whatever the setting says.
One limit on all of the above: the parameter governs Gather and Gather Merge and nothing else. The leaders of parallel index builds and parallel VACUUM always take part; max_parallel_maintenance_workers is the knob there.
Leave it on. At the default of two workers, off sends home a third of the processes working on your query, and I could not measure anything it bought in return. If you have one large analytic query running at more than eight workers against a table that does not fit in memory and want to know, put SET LOCAL parallel_leader_participation = off in its transaction, time it both ways, and keep the winner. It does not belong in postgresql.conf. If you find it there, someone was benchmarking and did not clean up. (EXPLAIN (SETTINGS) lists it whenever it is not on, which is the quick way to find out.)