Why PostgreSQL Uses Fewer Workers for Large, Complex Queries

PostgreSQL can split one single query processing into multiple CPU cores instead of doing all the tasks alone in one core. For example, a sequential scan on a large table can be split to several workers, and these results can be combined together by a Gather node. This is how a parallel query works.When you are running a query that traverses through a billion rows and suppose the query contains about six joins, your system has 32 cores, but when you run EXPLAIN, you might see something like:

Gather  (cost=1000.00..184521.33 rows=1 width=8)
  Workers Planned: 2

Have you ever wondered why just two workers run this complex query? PostgreSQL's worker count does not have almost anything to do with how complex the running query is, and the actual reason is a function with a while loop in it that most people do not know.

The function that decides the workers

Every parallel plans worker count originates in compute_parallel_worker(), in src/backend/optimizer/path/allpaths.c.

heap_parallel_threshold = Max(min_parallel_table_scan_size, 1);
while (heap_pages >= (BlockNumber) (heap_parallel_threshold * 3))
{
heap_parallel_workers++;
heap_parallel_threshold *= 3;
if (heap_parallel_threshold > INT_MAX / 3)
break; /* avoid overflow */
}

Here in this code you will see *3; this will reduce the worker count. Each additional worker costs three times as much data as the previous one, which means the relation has to be triple in size to get one more process or worker here.

The thresholds it produces goes approximatly like this:


Table sizeWorkers requested
8 MB1
24 MB2
72 MB3
216 MB4
648 MB5
1.9 GB6
5.7 GB7
17 GB8
51 GB9
154 GB10
461 GB11

Growing the table will give you more parallelism; this is the intended behaviour, not a bug. This is to be kept in mind before you try to get more workers from settings or configuration parameters.

The same loop runs one more time for the index pages, against min_parallel_index_scan_size (the default 512 kB). When both are applied, the planner will take the minimum of the two:

if (parallel_workers > 0)
    parallel_workers = Min(parallel_workers, index_parallel_workers);

So the parallel index scan will happen by which one is smaller, either parallel_workers or index_parallel_workers.

But complexity is not what the loop measures.

The compute_parallel_workers function is called from actually 4 places in the query planning workflow:

the sequential scan path (allpaths.c), the index scan path (costsize.c), the bitmap heap scan path (allpaths.c), and the TID scan path (tidpath.c). There is a fifth caller in planner.c, but that is plan_create_index_workers(), which is used for parallel create index queries.

So how does a join decide its worker count? You can see this code part in /src/backend/optimizer/util/pathnode.c

/* This is a foolish way to estimate parallel_workers, but for now... */
pathnode->jpath.path.parallel_workers = outer_path->parallel_workers;

This line of code is repeated in pathnode.c for nested loop, merge join, and hash join. This one denotes that the parallelism of your entire query is decided by one scan at the bottom of the plan that is the outermost running relation, and nothing else above it can raise this count.

Join a 900 GB fact table with one 20 MB dimension table, and you might get 10 workers. Let the planner change the join order so that the smaller table drives, and the plan shrinks to 1 worker count. That means what the worker calculation is the join order, not the complexity of the query. In the last line of the function compute_parallel_worker(), you can see:

parallel_workers = Min(parallel_workers, max_workers);

For query planning, max_workers is max_parallel_workers_per_gather, which defaults to 2. All that careful logical reasoning about table size is, on a default install, immediately clamped to 2.

The very first branch of compute_parallel_worker() is if we have set parallel workers for a table with an ALTER TABLE command:

if (rel->rel_parallel_workers != -1)
    parallel_workers = rel->rel_parallel_workers;

ALTER TABLE facts SET (parallel_workers = 8) skips the flow entirely. Only the Min() with the GUC still applies. For a known hot fact table, this is the easiest tool available to manage workers, and its existence is an admission that the log of 3 heuristics is a guess.

We would normally think of a powerful 32-core server as throwing everything it has at a complex, billion-row query. But like what we see in the source code in compute_parallel_worker(), PostgreSQL favours a conservative, table-size-based heuristic over dynamic complexity analysis.

WhatsApp