| 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-09-30 05:36:56 |
| Message-ID: | CABXr29FUZJsb1xHbzY_5=C2nuvbCUXWCNmBKzZXTgjkFqzdKAQ@mail.gmail.com |
| Views: | Whole Thread | Raw Message | Download mbox | Resend email |
| Thread: | |
| Lists: | pgsql-hackers |
On Tue, Sep 29, 2026 at 2:05 PM Matthias van de Meent
<boekewurm(at)gmail(dot)com> wrote:
>
> On Tue, 29 Sept 2026, 16:50 Robert Haas, <robertmhaas(at)gmail(dot)com> wrote:
> >
> > On Mon, Sep 28, 2026 at 6:06 PM Matthias van de Meent
> > <boekewurm(at)gmail(dot)com> wrote:
> > > At Databricks, we noticed a weird plan in one of our TPCC benchmarks
> > > on pg18 (and reproduced on master), which was basically as follows:
> > >
> > > Gather
> > > NestedLoop
> > > NestedLoop
> > > Parallel Seq Scan (Filter: <primary key = constant>)
> > > BitmapScan
> > > BitmapScan
> > >
> > > At face value that plan looks OK, but when you look at it in detail
> > > it's quite a strange plan: The outermost plan is expected to only
> > > produce one row. Joins can generally only produce tuples once rows are
> > > produced on the outer side, so the paralellism is wasted: only a
> > > single worker will find a row, and thus only a single worker will
> > > execute the joins. Given that cap on parallelism, an index scan on
> > > the primary key would've been much cheaper.
> >
> > I went looking for other reports of this issue. The only really clear
> > example I found was
> > http://postgr.es/m/f16e6fd6-7b7e-4c09-b26e-d7979ef86d2e@gmail.com --
> > in that email, Mark Kirkwood complains of what looks like the exact
> > same problem.
>
> Indeed.
>
> > In
> > http://postgr.es/m/152840735359.22458.3303333403164396853@wrigleys.postgresql.org
> > there's an interesting case where the driving table returns only 23
> > rows, but needs to be joined to a large sequential scan. Note that the
> > *estimate* for the area table is just 1 row, so this is very close to
> > being a case that your patch would affect, but I think that it isn't,
> > quite.
>
> The plan indicated by the user seems to have a 1-row estimate in the
> non-parallel portion of the plan, so that shouldn't be affected by
> this patch, no.
>
> > http://postgr.es/m/872bffe7-82d0-86db-e3d6-2e20b1a72c4b@codata.eu
> > is a sort of opposite case: the estimates are high, the actual row
> > counts are low, and parallelism loses, but the issue there may have
> > more to do with worker startup and shutdown being expensive on that
> > machine than with the work distribution being uneven.
>
> This seems to be a significant overestimation indeed, but the costing
> for it won't change with this patch. The only queries whose planning
> could be negatively affected are queries that underestimate the output
> row count to a value below the useful worker count of the outermost
> relation, and whilst I do think those exist, I also think that those
> misestimations probably have bigger issues, because the cost of joins
> and other operations higher in the plan tree are likely to be severely
> underestimated for those too.
>
> > As far as the approach taken by the patch, I'm not sure that making
> > Path bigger for this is a good idea. It might be fine if we had lots
> > of reports of this being a problem, but it seems expensive as things
> > are. Still, that's not necessarily to say I think we should do
> > nothing.
>
> I'll adjust this if it proves necessary: We can safely change the `int
> *worker*` fields' types to int16, given the MAX_PARALLEL_WORKER_LIMIT
> of 1024. While it will need a bit more care to avoid overflows
> everywhere we sum up or otherwise process those fields, it's not a
> huge complication.
>
> > The approach makes me a little nervous in that it treats very
> > small number of rows as a very special case in need of very special
> > handling, but there's some argument to be made that this is actually
> > the case. I mean, small LIMIT values have similar problems, and we get
> > those wrong frequently, arguably because we don't treat that as a
> > sufficiently special case. Still, your patch takes the idea further:
> > the correction drops to zero as soon as #rows >= #workers, but lumpy
> > work distribution could still be an issue past that point (e.g. 5
> > rows, 4 workers). I'm not really sure what's best here.
>
> I agree that we probably should do more for such plans, but I don't
> know when and where to apply these corrections. It is trivial to
> understand (and nearly as easy to cost) that if we expect one row
> total from lower plan nodes, we shouldn't expect to process that one
> row in more than one worker. But figuring out which of the N workers
> will process the X planned tuples isn't as trivial to figure out, let
> alone cost.
>
> A better data-aware parallel costing model probably will need to be
> moved into cost_*() and/or baserel calculations, but I don't have
> enough of a theoretical or math background to patch together something
> that'd work for this.
>
> Note that parallel scans generally scan many more pages than it has
> workers, so that (assuming linear distribution of matched rows across
> pages) each worker will probably produce 1/Nth of the tuples. Join
> and matching-tuple skew will break this assumption, but I'm not sure
> we have the stats readily available to produce good estimates on this.
>
>
> Kind regards,
>
> Matthias van de Meent
> Databricks (https://www.databricks.com)
>
>
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:
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.
I don't think this is literally a double subtraction of the leader. The
helper appears to count the leader as a full occupancy slot, while
`get_parallel_divisor()` later gives the leader only the usual fractional
contribution. But those two interpretations do not line up at the
threshold, which produces the 3.1 -> 4.0 jump above.
I can send the minimized reproducers if useful.
Thanks,
Haibo
| From | Date | Subject | |
|---|---|---|---|
| Next Message | Virender Singla | 2026-09-30 05:39:40 | Re: [PATCH] Corruption Issue: Fix missing tts_tid in ExecForceStoreHeapTuple |
| Previous Message | Nisha Moond | 2026-09-30 05:34:49 | Re: Proposal: Conflict log history table for Logical Replication |