Re: Costing for parallel scans with few/single row produced in the outer side - Mailing list pgsql-hackers
| From | Haibo Yan |
|---|---|
| Subject | Re: Costing for parallel scans with few/single row produced in the outer side |
| Date | |
| Msg-id | CABXr29FUZJsb1xHbzY_5=C2nuvbCUXWCNmBKzZXTgjkFqzdKAQ@mail.gmail.com Whole thread |
| In response to | Re: Costing for parallel scans with few/single row produced in the outer side (Matthias van de Meent <boekewurm@gmail.com>) |
| List | pgsql-hackers |
On Tue, Sep 29, 2026 at 2:05 PM Matthias van de Meent <boekewurm@gmail.com> wrote: > > On Tue, 29 Sept 2026, 16:50 Robert Haas, <robertmhaas@gmail.com> wrote: > > > > On Mon, Sep 28, 2026 at 6:06 PM Matthias van de Meent > > <boekewurm@gmail.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
pgsql-hackers by date: