From 2bad2aa8654ae9bf2232f5eb151dac6051a265da Mon Sep 17 00:00:00 2001 From: Matthias van de Meent Date: Mon, 28 Sep 2026 20:56:39 +0200 Subject: [PATCH v1] optimizer: Better costing for low-rowcount parallel plans If outer nodes of parallel plans produce fewer rows than there are workers, then some of those workers will necessarily be unable to contribute to the execution. These workers should therefore not be included in the cost model, which is what this patch contributes. --- src/backend/optimizer/path/costsize.c | 17 ++-- src/backend/optimizer/path/indxpath.c | 1 + src/backend/optimizer/util/pathnode.c | 119 ++++++++++++++++++++++++++ src/include/nodes/pathnodes.h | 9 ++ 4 files changed, 140 insertions(+), 6 deletions(-) diff --git a/src/backend/optimizer/path/costsize.c b/src/backend/optimizer/path/costsize.c index 7bbddb8bee4..de2d788afd0 100644 --- a/src/backend/optimizer/path/costsize.c +++ b/src/backend/optimizer/path/costsize.c @@ -748,6 +748,8 @@ cost_index(IndexPath *path, PlannerInfo *root, double loop_count, if (partial_path) { + int workers; + /* * For index only scans compute workers based on number of index pages * fetched; the number of heap pages we fetch might be so small as to @@ -762,10 +764,13 @@ cost_index(IndexPath *path, PlannerInfo *root, double loop_count, * sequential as for parallel scans the pages are accessed in random * order. */ - path->path.parallel_workers = compute_parallel_worker(baserel, - rand_heap_pages, - index_pages, - max_parallel_workers_per_gather); + workers = compute_parallel_worker(baserel, + rand_heap_pages, + index_pages, + max_parallel_workers_per_gather); + + path->path.parallel_workers = workers; + path->path.effective_workers = workers; /* * Fall out if workers can't be assigned for parallel scan, because in @@ -6731,7 +6736,7 @@ page_size(double tuples, int width) static double get_parallel_divisor(Path *path) { - double parallel_divisor = path->parallel_workers; + double parallel_divisor = path->effective_workers; /* * Early experience with parallel query suggests that when there is only @@ -6748,7 +6753,7 @@ get_parallel_divisor(Path *path) { double leader_contribution; - leader_contribution = 1.0 - (0.3 * path->parallel_workers); + leader_contribution = 1.0 - (0.3 * path->effective_workers); if (leader_contribution > 0) parallel_divisor += leader_contribution; } diff --git a/src/backend/optimizer/path/indxpath.c b/src/backend/optimizer/path/indxpath.c index 3f5d4fa3182..f9ffd054fff 100644 --- a/src/backend/optimizer/path/indxpath.c +++ b/src/backend/optimizer/path/indxpath.c @@ -2037,6 +2037,7 @@ bitmap_scan_cost_est(PlannerInfo *root, RelOptInfo *rel, Path *ipath) * Parallel bitmap heap path will be considered at later stage. */ bpath.path.parallel_workers = 0; + bpath.path.effective_workers = 0; /* Now we can do cost_bitmap_heap_scan */ cost_bitmap_heap_scan(&bpath.path, root, rel, diff --git a/src/backend/optimizer/util/pathnode.c b/src/backend/optimizer/util/pathnode.c index 67c47fd7c6c..6b9ae8d4cab 100644 --- a/src/backend/optimizer/util/pathnode.c +++ b/src/backend/optimizer/util/pathnode.c @@ -53,6 +53,7 @@ static List *reparameterize_pathlist_by_child(PlannerInfo *root, RelOptInfo *child_rel); static bool pathlist_is_reparameterizable_by_child(List *pathlist, RelOptInfo *child_rel); +static int subpath_adjusted_effective_workers(const Path *subpath); /***************************************************************************** @@ -237,6 +238,65 @@ compare_path_costs_fuzzily(Path *path1, Path *path2, double fuzz_factor) #undef CONSIDER_PATH_STARTUP_COST } +/* + * subpath_adjusted_effective_workers + * calculate the effective number of workers participating in a node + * + * Joins' inner plans generally only execute if the plan node produced tuples. + * In parallel plans, this means that if outer plans produce fewer tuples + * than there are workers, amortization of inner plans' costs needs to take + * this reduced parallelism into account. This reduced parallelism is + * accounted for in Path.effective_workers, and then applied in + * costsize.c's get_parallel_divisor(). + * + * Note that subpath can benefit from higher degrees of parallelism than + * the Path that'll have it as subpath, because (e.g.) removing N-1 tuples + * from the table with N tuples will be faster with more workers, even if + * it'll produce just one tuple; only the node above it will have to deal + * with the reduced concurrency. + */ +static int +subpath_adjusted_effective_workers(const Path *subpath) +{ + int leader_factor = 0; + int result = 0; + double best_row_estimate; + + /* can't go below 0 workers */ + if (subpath->parallel_workers == 0) + return 0; + + if (parallel_leader_participation) + leader_factor = 1; + + /* most basic estimate for row counts */ + best_row_estimate = subpath->rows * + (double) (subpath->effective_workers + leader_factor); + + /* check stats of the plan, ex parallel adjustments upward */ + if (subpath->param_info) + best_row_estimate = Min(subpath->param_info->ppi_rows, + best_row_estimate); + + /* same but for the RelOptInfo's data */ + if (subpath->parent) + best_row_estimate = Min(subpath->parent->rows, + best_row_estimate); + + result = subpath->effective_workers; + + /* + * We'll use at most ceil(row_estimate) processes. If the leader + * participates, we'll have to subtract it: we don't count its + * contribution to parallelism higher up the plan tree for these + * nodes. + */ + if (best_row_estimate < result + leader_factor) + result = (int) ceil(best_row_estimate) - leader_factor; + + return result; +} + /* * set_cheapest * Find the minimum-cost paths from among a relation's paths, @@ -1036,6 +1096,7 @@ create_seqscan_path(PlannerInfo *root, RelOptInfo *rel, pathnode->parallel_aware = (parallel_workers > 0); pathnode->parallel_safe = rel->consider_parallel; pathnode->parallel_workers = parallel_workers; + pathnode->effective_workers = parallel_workers; pathnode->pathkeys = NIL; /* seqscan has unordered result */ cost_seqscan(pathnode, root, rel, pathnode->param_info); @@ -1060,6 +1121,7 @@ create_samplescan_path(PlannerInfo *root, RelOptInfo *rel, Relids required_outer pathnode->parallel_aware = false; pathnode->parallel_safe = rel->consider_parallel; pathnode->parallel_workers = 0; + pathnode->effective_workers = 0; pathnode->pathkeys = NIL; /* samplescan has unordered result */ cost_samplescan(pathnode, root, rel, pathnode->param_info); @@ -1112,6 +1174,7 @@ create_index_path(PlannerInfo *root, pathnode->path.parallel_aware = false; pathnode->path.parallel_safe = rel->consider_parallel; pathnode->path.parallel_workers = 0; + pathnode->path.effective_workers = 0; pathnode->path.pathkeys = pathkeys; pathnode->indexinfo = index; @@ -1163,6 +1226,7 @@ create_bitmap_heap_path(PlannerInfo *root, pathnode->path.parallel_aware = (parallel_degree > 0); pathnode->path.parallel_safe = rel->consider_parallel; pathnode->path.parallel_workers = parallel_degree; + pathnode->path.effective_workers = parallel_degree; pathnode->path.pathkeys = NIL; /* always unordered */ pathnode->bitmapqual = bitmapqual; @@ -1215,6 +1279,7 @@ create_bitmap_and_path(PlannerInfo *root, pathnode->path.parallel_aware = false; pathnode->path.parallel_safe = rel->consider_parallel; pathnode->path.parallel_workers = 0; + pathnode->path.effective_workers = 0; pathnode->path.pathkeys = NIL; /* always unordered */ @@ -1267,6 +1332,7 @@ create_bitmap_or_path(PlannerInfo *root, pathnode->path.parallel_aware = false; pathnode->path.parallel_safe = rel->consider_parallel; pathnode->path.parallel_workers = 0; + pathnode->path.effective_workers = 0; pathnode->path.pathkeys = NIL; /* always unordered */ @@ -1296,6 +1362,7 @@ create_tidscan_path(PlannerInfo *root, RelOptInfo *rel, List *tidquals, pathnode->path.parallel_aware = false; pathnode->path.parallel_safe = rel->consider_parallel; pathnode->path.parallel_workers = 0; + pathnode->path.effective_workers = 0; pathnode->path.pathkeys = NIL; /* always unordered */ pathnode->tidquals = tidquals; @@ -1326,6 +1393,7 @@ create_tidrangescan_path(PlannerInfo *root, RelOptInfo *rel, pathnode->path.parallel_aware = (parallel_workers > 0); pathnode->path.parallel_safe = rel->consider_parallel; pathnode->path.parallel_workers = parallel_workers; + pathnode->path.effective_workers = parallel_workers; pathnode->path.pathkeys = NIL; /* always unordered */ pathnode->tidrangequals = tidrangequals; @@ -1388,6 +1456,7 @@ create_append_path(PlannerInfo *root, pathnode->path.parallel_aware = parallel_aware; pathnode->path.parallel_safe = rel->consider_parallel; pathnode->path.parallel_workers = parallel_workers; + pathnode->path.effective_workers = parallel_workers; pathnode->path.pathkeys = pathkeys; /* @@ -1549,6 +1618,7 @@ create_merge_append_path(PlannerInfo *root, pathnode->path.parallel_aware = false; pathnode->path.parallel_safe = rel->consider_parallel; pathnode->path.parallel_workers = 0; + pathnode->path.effective_workers = 0; pathnode->path.pathkeys = pathkeys; pathnode->subpaths = subpaths; @@ -1675,6 +1745,7 @@ create_group_result_path(PlannerInfo *root, RelOptInfo *rel, pathnode->path.parallel_aware = false; pathnode->path.parallel_safe = rel->consider_parallel; pathnode->path.parallel_workers = 0; + pathnode->path.effective_workers = 0; pathnode->path.pathkeys = NIL; pathnode->quals = havingqual; @@ -1725,6 +1796,8 @@ create_material_path(RelOptInfo *rel, Path *subpath, bool enabled) pathnode->path.parallel_safe = rel->consider_parallel && subpath->parallel_safe; pathnode->path.parallel_workers = subpath->parallel_workers; + pathnode->path.effective_workers = + subpath_adjusted_effective_workers(subpath); pathnode->path.pathkeys = subpath->pathkeys; pathnode->subpath = subpath; @@ -1761,6 +1834,8 @@ create_memoize_path(PlannerInfo *root, RelOptInfo *rel, Path *subpath, pathnode->path.parallel_safe = rel->consider_parallel && subpath->parallel_safe; pathnode->path.parallel_workers = subpath->parallel_workers; + pathnode->path.effective_workers = + subpath_adjusted_effective_workers(subpath); pathnode->path.pathkeys = subpath->pathkeys; pathnode->subpath = subpath; @@ -1879,6 +1954,7 @@ create_gather_path(PlannerInfo *root, RelOptInfo *rel, Path *subpath, pathnode->path.parallel_aware = false; pathnode->path.parallel_safe = false; pathnode->path.parallel_workers = 0; + pathnode->path.effective_workers = 0; pathnode->path.pathkeys = NIL; /* Gather has unordered result */ pathnode->subpath = subpath; @@ -1923,6 +1999,8 @@ create_subqueryscan_path(PlannerInfo *root, RelOptInfo *rel, Path *subpath, pathnode->path.parallel_safe = rel->consider_parallel && subpath->parallel_safe; pathnode->path.parallel_workers = subpath->parallel_workers; + pathnode->path.effective_workers = + subpath_adjusted_effective_workers(subpath); pathnode->path.pathkeys = pathkeys; pathnode->subpath = subpath; @@ -1951,6 +2029,7 @@ create_functionscan_path(PlannerInfo *root, RelOptInfo *rel, pathnode->parallel_aware = false; pathnode->parallel_safe = rel->consider_parallel; pathnode->parallel_workers = 0; + pathnode->effective_workers = 0; pathnode->pathkeys = pathkeys; cost_functionscan(pathnode, root, rel, pathnode->param_info); @@ -1977,6 +2056,7 @@ create_tablefuncscan_path(PlannerInfo *root, RelOptInfo *rel, pathnode->parallel_aware = false; pathnode->parallel_safe = rel->consider_parallel; pathnode->parallel_workers = 0; + pathnode->effective_workers = 0; pathnode->pathkeys = NIL; /* result is always unordered */ cost_tablefuncscan(pathnode, root, rel, pathnode->param_info); @@ -2003,6 +2083,7 @@ create_valuesscan_path(PlannerInfo *root, RelOptInfo *rel, pathnode->parallel_aware = false; pathnode->parallel_safe = rel->consider_parallel; pathnode->parallel_workers = 0; + pathnode->effective_workers = 0; pathnode->pathkeys = NIL; /* result is always unordered */ cost_valuesscan(pathnode, root, rel, pathnode->param_info); @@ -2029,6 +2110,7 @@ create_ctescan_path(PlannerInfo *root, RelOptInfo *rel, pathnode->parallel_aware = false; pathnode->parallel_safe = rel->consider_parallel; pathnode->parallel_workers = 0; + pathnode->effective_workers = 0; pathnode->pathkeys = pathkeys; cost_ctescan(pathnode, root, rel, pathnode->param_info); @@ -2055,6 +2137,7 @@ create_namedtuplestorescan_path(PlannerInfo *root, RelOptInfo *rel, pathnode->parallel_aware = false; pathnode->parallel_safe = rel->consider_parallel; pathnode->parallel_workers = 0; + pathnode->effective_workers = 0; pathnode->pathkeys = NIL; /* result is always unordered */ cost_namedtuplestorescan(pathnode, root, rel, pathnode->param_info); @@ -2081,6 +2164,7 @@ create_resultscan_path(PlannerInfo *root, RelOptInfo *rel, pathnode->parallel_aware = false; pathnode->parallel_safe = rel->consider_parallel; pathnode->parallel_workers = 0; + pathnode->effective_workers = 0; pathnode->pathkeys = NIL; /* result is always unordered */ cost_resultscan(pathnode, root, rel, pathnode->param_info); @@ -2107,6 +2191,7 @@ create_worktablescan_path(PlannerInfo *root, RelOptInfo *rel, pathnode->parallel_aware = false; pathnode->parallel_safe = rel->consider_parallel; pathnode->parallel_workers = 0; + pathnode->effective_workers = 0; pathnode->pathkeys = NIL; /* result is always unordered */ /* Cost is the same as for a regular CTE scan */ @@ -2150,6 +2235,7 @@ create_foreignscan_path(PlannerInfo *root, RelOptInfo *rel, pathnode->path.parallel_aware = false; pathnode->path.parallel_safe = rel->consider_parallel; pathnode->path.parallel_workers = 0; + pathnode->path.effective_workers = 0; pathnode->path.rows = rows; pathnode->path.disabled_nodes = disabled_nodes; pathnode->path.startup_cost = startup_cost; @@ -2204,6 +2290,7 @@ create_foreign_join_path(PlannerInfo *root, RelOptInfo *rel, pathnode->path.parallel_aware = false; pathnode->path.parallel_safe = rel->consider_parallel; pathnode->path.parallel_workers = 0; + pathnode->path.effective_workers = 0; pathnode->path.rows = rows; pathnode->path.disabled_nodes = disabled_nodes; pathnode->path.startup_cost = startup_cost; @@ -2253,6 +2340,7 @@ create_foreign_upper_path(PlannerInfo *root, RelOptInfo *rel, pathnode->path.parallel_aware = false; pathnode->path.parallel_safe = rel->consider_parallel; pathnode->path.parallel_workers = 0; + pathnode->path.effective_workers = 0; pathnode->path.rows = rows; pathnode->path.disabled_nodes = disabled_nodes; pathnode->path.startup_cost = startup_cost; @@ -2419,6 +2507,8 @@ create_nestloop_path(PlannerInfo *root, outer_path->parallel_safe && inner_path->parallel_safe; /* This is a foolish way to estimate parallel_workers, but for now... */ pathnode->jpath.path.parallel_workers = outer_path->parallel_workers; + pathnode->jpath.path.effective_workers = + subpath_adjusted_effective_workers(outer_path); pathnode->jpath.path.pathkeys = pathkeys; pathnode->jpath.jointype = jointype; pathnode->jpath.inner_unique = extra->inner_unique; @@ -2485,6 +2575,8 @@ create_mergejoin_path(PlannerInfo *root, outer_path->parallel_safe && inner_path->parallel_safe; /* This is a foolish way to estimate parallel_workers, but for now... */ pathnode->jpath.path.parallel_workers = outer_path->parallel_workers; + pathnode->jpath.path.effective_workers = + subpath_adjusted_effective_workers(outer_path); pathnode->jpath.path.pathkeys = pathkeys; pathnode->jpath.jointype = jointype; pathnode->jpath.inner_unique = extra->inner_unique; @@ -2551,6 +2643,8 @@ create_hashjoin_path(PlannerInfo *root, outer_path->parallel_safe && inner_path->parallel_safe; /* This is a foolish way to estimate parallel_workers, but for now... */ pathnode->jpath.path.parallel_workers = outer_path->parallel_workers; + pathnode->jpath.path.effective_workers = + subpath_adjusted_effective_workers(outer_path); /* * A hashjoin never has pathkeys, since its output ordering is @@ -2619,6 +2713,8 @@ create_projection_path(PlannerInfo *root, subpath->parallel_safe && is_parallel_safe(root, (Node *) target->exprs); pathnode->path.parallel_workers = subpath->parallel_workers; + pathnode->path.effective_workers = + subpath_adjusted_effective_workers(subpath); /* Projection does not change the sort order */ pathnode->path.pathkeys = subpath->pathkeys; @@ -2803,6 +2899,8 @@ create_set_projection_path(PlannerInfo *root, subpath->parallel_safe && is_parallel_safe(root, (Node *) target->exprs); pathnode->path.parallel_workers = subpath->parallel_workers; + pathnode->path.effective_workers = + subpath_adjusted_effective_workers(subpath); /* Projection does not change the sort order XXX? */ pathnode->path.pathkeys = subpath->pathkeys; @@ -2873,6 +2971,8 @@ create_incremental_sort_path(PlannerInfo *root, pathnode->path.parallel_safe = rel->consider_parallel && subpath->parallel_safe; pathnode->path.parallel_workers = subpath->parallel_workers; + pathnode->path.effective_workers = + subpath_adjusted_effective_workers(subpath); pathnode->path.pathkeys = pathkeys; pathnode->subpath = subpath; @@ -2921,6 +3021,8 @@ create_sort_path(PlannerInfo *root, pathnode->path.parallel_safe = rel->consider_parallel && subpath->parallel_safe; pathnode->path.parallel_workers = subpath->parallel_workers; + pathnode->path.effective_workers = + subpath_adjusted_effective_workers(subpath); pathnode->path.pathkeys = pathkeys; pathnode->subpath = subpath; @@ -2967,6 +3069,8 @@ create_group_path(PlannerInfo *root, pathnode->path.parallel_safe = rel->consider_parallel && subpath->parallel_safe; pathnode->path.parallel_workers = subpath->parallel_workers; + pathnode->path.effective_workers = + subpath_adjusted_effective_workers(subpath); /* Group doesn't change sort ordering */ pathnode->path.pathkeys = subpath->pathkeys; @@ -3022,6 +3126,8 @@ create_unique_path(PlannerInfo *root, pathnode->path.parallel_safe = rel->consider_parallel && subpath->parallel_safe; pathnode->path.parallel_workers = subpath->parallel_workers; + pathnode->path.effective_workers = + subpath_adjusted_effective_workers(subpath); /* Unique doesn't change the input ordering */ pathnode->path.pathkeys = subpath->pathkeys; @@ -3088,6 +3194,8 @@ create_agg_path(PlannerInfo *root, pathnode->path.parallel_safe = rel->consider_parallel && subpath->parallel_safe; pathnode->path.parallel_workers = subpath->parallel_workers; + pathnode->path.effective_workers = + subpath_adjusted_effective_workers(subpath); if (aggstrategy == AGG_SORTED) { @@ -3172,6 +3280,8 @@ create_groupingsets_path(PlannerInfo *root, pathnode->path.parallel_safe = rel->consider_parallel && subpath->parallel_safe; pathnode->path.parallel_workers = subpath->parallel_workers; + pathnode->path.effective_workers = + subpath_adjusted_effective_workers(subpath); pathnode->subpath = subpath; /* @@ -3332,6 +3442,7 @@ create_minmaxagg_path(PlannerInfo *root, pathnode->path.parallel_aware = false; pathnode->path.parallel_safe = true; /* might change below */ pathnode->path.parallel_workers = 0; + pathnode->path.effective_workers = 0; /* Result is one unordered row */ pathnode->path.rows = 1; pathnode->path.pathkeys = NIL; @@ -3427,6 +3538,8 @@ create_windowagg_path(PlannerInfo *root, pathnode->path.parallel_safe = rel->consider_parallel && subpath->parallel_safe; pathnode->path.parallel_workers = subpath->parallel_workers; + pathnode->path.effective_workers = + subpath_adjusted_effective_workers(subpath); /* WindowAgg preserves the input sort order */ pathnode->path.pathkeys = subpath->pathkeys; @@ -3626,6 +3739,8 @@ create_recursiveunion_path(PlannerInfo *root, leftpath->parallel_safe && rightpath->parallel_safe; /* Foolish, but we'll do it like joins for now: */ pathnode->path.parallel_workers = leftpath->parallel_workers; + pathnode->path.effective_workers = + subpath_adjusted_effective_workers(leftpath); /* RecursiveUnion result is always unsorted */ pathnode->path.pathkeys = NIL; @@ -3664,6 +3779,7 @@ create_lockrows_path(PlannerInfo *root, RelOptInfo *rel, pathnode->path.parallel_aware = false; pathnode->path.parallel_safe = false; pathnode->path.parallel_workers = 0; + pathnode->path.effective_workers = 0; pathnode->path.rows = subpath->rows; /* @@ -3743,6 +3859,7 @@ create_modifytable_path(PlannerInfo *root, RelOptInfo *rel, pathnode->path.parallel_aware = false; pathnode->path.parallel_safe = false; pathnode->path.parallel_workers = 0; + pathnode->path.effective_workers = 0; pathnode->path.pathkeys = NIL; /* @@ -3830,6 +3947,8 @@ create_limit_path(PlannerInfo *root, RelOptInfo *rel, pathnode->path.parallel_safe = rel->consider_parallel && subpath->parallel_safe; pathnode->path.parallel_workers = subpath->parallel_workers; + pathnode->path.effective_workers = + subpath_adjusted_effective_workers(subpath); pathnode->path.rows = subpath->rows; pathnode->path.disabled_nodes = subpath->disabled_nodes; pathnode->path.startup_cost = subpath->startup_cost; diff --git a/src/include/nodes/pathnodes.h b/src/include/nodes/pathnodes.h index 1c6d1fe3d04..7fdf274383b 100644 --- a/src/include/nodes/pathnodes.h +++ b/src/include/nodes/pathnodes.h @@ -2000,6 +2000,15 @@ typedef struct Path bool parallel_safe; /* desired # of workers; 0 = not parallel */ int parallel_workers; + /* + * Effective # of workers that produce tuples + * + * Use this to cost partial plan nodes. Parallel scans may produce + * few tuples, and the cost of join nodes should not be amortized across + * workers if just one worker will produce a tuple that participates in + * the higher partial nodes. + */ + int effective_workers; /* estimated size/costs for path (see costsize.c for more info) */ Cardinality rows; /* estimated number of result tuples */ -- 2.54.0 (Apple Git-157)