Re: Costing for parallel scans with few/single row produced in the outer side

From: Haibo Yan <tristan(dot)yim(at)gmail(dot)com>
To: Matthias van de Meent <boekewurm(at)gmail(dot)com>
Cc: Robert Haas <robertmhaas(at)gmail(dot)com>, PostgreSQL Hackers <pgsql-hackers(at)lists(dot)postgresql(dot)org>
Subject: Re: Costing for parallel scans with few/single row produced in the outer side
Date: 2026-10-04 02:05:53
Message-ID: CABXr29Gt=gv1GiCNh3qEYPmuUgtCnDNKydmd4YHBafPT2y9dUw@mail.gmail.com
Views: Whole Thread | Raw Message | Download mbox | Resend email
Thread:
Lists: pgsql-hackers

On Fri, Oct 2, 2026 at 2:10 PM Matthias van de Meent
<boekewurm(at)gmail(dot)com> wrote:
>
> On Wed, 30 Sept 2026 at 07:37, Haibo Yan <tristan(dot)yim(at)gmail(dot)com> wrote:
> >
> > Hi Matthias,
> >
> > I spent some time testing v1 and found two issues in the current
> > `effective_workers` propagation.
> >
> > The first one looks like a correctness bug.
> >
> > `subpath_adjusted_effective_workers()` can see `subpath->parent->rows == 0`
> > while building paths for `UPPERREL_PARTIAL_GROUP_AGG`. With leader
> > participation enabled this can make `effective_workers` negative.
> >
> > For example, with a four-worker Partial Aggregate followed by Sort/Group,
> > the calculation can end up with:
>
> Thanks, I've adjusted this in the latest version.
>
> > best_row_estimate = 0
> > effective_workers = ceil(0) - 1 = -1
> >
> > and `get_parallel_divisor()` then returns 0.3. With leader participation
> > disabled the corresponding value is 0, and the divisor becomes 0.
> >
> > That feeds into `compute_gather_rows()`, so I was able to get cases where
> > a partial path producing 100 rows per participant was reconstructed as a
> > Gather Merge producing only 30 rows with leader participation enabled,
> > or 1 row with it disabled.
> >
> > The second issue is a discontinuity around the low-cardinality threshold.
> >
> > With four planned workers and leader participation enabled I get:
> >
> > estimated rows effective_workers divisor
> > 2 1 1.7
> > 3 2 2.4
> > 4 3 3.1
> > 5 4 4.0
> >
> > In a parameterized Nested Loop reproducer, that makes the estimated
> > partial output decrease when the global row estimate increases:
> >
> > R=4: partial NL rows = 1290, total cost = 23659.35
> > R=5: partial NL rows = 1250, total cost = 23559.37
> >
> > The R=5 case has 25% more qualifying outer rows and 25% more final output,
> > but is costed lower.
>
> Yes, that's presumably because the parallel portion of the plan is
> spread across more workers, and the join portions are otherwise costed
> equivalently (although in this case it expects fewer rows per worker
> total).
>
> And, note that each worker up to 3 adds 0.7 to the divisor, whilst the
> 4th worker adds .9, and every worker after that contributes 1.0. This
> is just how the divisor is calculated, and changing that is not part
> of what I'm planning to do.
>
> --------------------
>
> Attached is v2, which changes a few things:
>
> * Changed types of the Path->*_workers fields to int16.
> This is primarily to avoid size changes, and safe because we don't
> support values larger than 1024.
> * Updated the calculations of subpath_adjusted_effective_workers to
> have better accounting in certain cases
> This fixes the underflow issues mentioned by Haibo Yan.
> * Added effective_worker to create_append_path.
> This was needed to avoid otherwise unrelated plan changes in the
> regression tests.
>
>
> Kind regards,
>
> Matthias van de Meent
> Databricks (https://www.databricks.com)

Hi Matthias,

Thanks for the v2 update. I tested it further.

The negative `effective_workers` case from v1 is fixed, but I think the
underlying issue is broader than the uninitialized upper-rel row count:
an upper relation's global cardinality is not generally an upper bound
on the partial tuple stream below Gather.

A simple example is a one-group partial aggregate. With four workers I
get:

leader on leader off
subpath->rows 1 1
subpath->parent->rows 0 0
incoming effective_workers 4 4
adjusted effective_workers 0 1
divisor 1 1
Gather Merge estimated rows 1 1
actual partial rows 5 4

The zero `parent->rows` is one way this becomes visible during upper-rel
construction, but even if that value were initialized to the correct
final cardinality of 1, it would still not be a valid cap here. One
global group can produce one partial aggregate state from every
participating process before finalization.

So in this case:

final/global groups = 1

does not imply:

partial tuples before Gather <= 1

The same distinction appears to apply to partial DISTINCT.

I think this means we need to distinguish the final/global cardinality
of an upper relation from the number of tuples/processes represented by
a partial path before Gather/finalization.

I also found two problems with the new Append propagation.

For two low-cardinality partial children, each child ends up with:

effective_workers = 0
divisor = 1
partial rows = 1000

The v2 Append calculation then adds one leader slot for each partial
child, so the parent gets:

effective_workers = 2
divisor = 2.4

The Append itself has 2000 partial rows, and `compute_gather_rows()`
therefore reconstructs:

2000 * 2.4 = 4800 rows

while the actual global result is 2000 rows.

There is only one leader process for the parallel plan, so adding one
leader slot per partial child seems to convert repeated references to
the same process into additional worker capacity.

I also tested the opposite case: one low-cardinality partial child plus
three large non-partial children under Parallel Append. The planner
gets:

parallel_workers = 4
effective_workers = 1
divisor = 1.7

but EXPLAIN ANALYZE shows substantial work on multiple workers, roughly:

partial child ~1k rows
non-partial child ~200k rows
non-partial child ~200k rows
non-partial child ~200k rows

I don't mean that the planner should predict that exact runtime worker
occupancy. The point is that a non-partial child's
`effective_workers = 0` does not mean that it contributes zero potential
worker capacity when Parallel Append schedules different child paths on
different workers.

There seems to be a more general compositional issue here. If two child
paths have worker sets S1 and S2, the parent occupancy depends on:

|S1 union S2|

while scalar child `effective_workers` values only describe bounds on:

|S1| and |S2|

For example, two children may both have EW=1, but they could either run
on the same worker or on two different workers. The child EW values are
identical in both cases, while the parent occupancy differs.

So an exact Append-level `effective_workers` does not seem derivable
from the child scalar values alone.

One other observation: I tried retaining v1's simpler Append behavior
(`effective_workers = parallel_workers`) while keeping the v2 fix for
the negative-EW case. The core regression suite was clean in that
configuration.

That makes me think the more complicated Append calculation is not
needed to eliminate the regression changes that motivated it; at least
the cases I tested appear to have originated from the upper-rel
propagation problem instead.

The R=4 -> R=5 behavior I mentioned earlier is also still present in v2:

R=4: EW=3, divisor=3.1, partial NL rows=1290, cost=50550.92
R=5: EW=4, divisor=4.0, partial NL rows=1250, cost=49764.42

I agree that the 3.1 -> 4.0 schedule itself comes from the existing
leader heuristic. What seems new is that `effective_workers` makes that
heuristic cardinality-dependent at a fixed requested worker count, so
increasing the global row count can reduce both the estimated
per-participant work and the total cost.

I think this is another symptom of the same semantic mismatch. The
helper treats the leader as a full occupancy slot when deriving
`effective_workers`, while `get_parallel_divisor()` interprets the
result as a worker count and then applies its own fractional leader
contribution.

I don't think the right next step is to patch the Append formula with
another scalar adjustment. The examples above suggest that
`effective_workers` is currently being used for at least three different
things: a planner-side bound/estimate of useful workers, leader
participation, and executor worker occupancy under Parallel Append.

Those are related, but they are not the same quantity. In particular, a
scalar per-child value cannot represent the union semantics of worker
occupancy across Append children.

So I think the first step should be to define the invariant for
`effective_workers`, and then decide which node types can propagate that
invariant safely. For Append, a conservative/simple behavior seems
preferable to mixing leader participation and child occupancy into one
scalar.

I can send the minimized SQL reproducers for the upper-rel and Append
cases if useful.

Thanks,
Haibo

In response to

Browse pgsql-hackers by date

  From Date Subject
Previous Message Andres Freund 2026-10-03 23:07:02 Re: Coverage with make coverage-html is broken on latest Debian using lcov v2