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

From: Matthias van de Meent <boekewurm(at)gmail(dot)com>
To: Robert Haas <robertmhaas(at)gmail(dot)com>
Cc: 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-29 21:05:25
Message-ID: CAEze2Wj8jMehea76X5F1a4eq8aQs-i1_+_HGcD-mcTDwfzT6NA@mail.gmail.com
Views: Whole Thread | Raw Message | Download mbox | Resend email
Thread:
Lists: pgsql-hackers

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)

In response to

Browse pgsql-hackers by date

  From Date Subject
Previous Message Alexander Lakhin 2026-09-29 21:00:00 Re: Bug in logical decoding with DDL and subtransactions