From 616383b2f5e8db79841c2ead89eec99bfbbb337d Mon Sep 17 00:00:00 2001 From: Matthias van de Meent Date: Mon, 28 Sep 2026 20:56:39 +0200 Subject: [PATCH v2 1/2] optimizer: Improve 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/allpaths.c | 85 ++++++++++-- src/backend/optimizer/path/costsize.c | 22 ++- src/backend/optimizer/path/indxpath.c | 1 + src/backend/optimizer/path/joinrels.c | 2 +- src/backend/optimizer/plan/planner.c | 1 + src/backend/optimizer/prep/prepunion.c | 37 +++++- src/backend/optimizer/util/pathnode.c | 177 +++++++++++++++++++++++-- src/include/nodes/pathnodes.h | 11 +- src/include/optimizer/pathnode.h | 10 +- src/include/optimizer/paths.h | 4 +- 10 files changed, 306 insertions(+), 44 deletions(-) diff --git a/src/backend/optimizer/path/allpaths.c b/src/backend/optimizer/path/allpaths.c index de8f29c2b57..a1ff219ee42 100644 --- a/src/backend/optimizer/path/allpaths.c +++ b/src/backend/optimizer/path/allpaths.c @@ -28,6 +28,7 @@ #include "nodes/makefuncs.h" #include "nodes/nodeFuncs.h" #include "nodes/supportnodes.h" +#include "postmaster/bgworker_internals.h" #ifdef OPTIMIZER_DEBUG #include "nodes/print.h" #endif @@ -865,7 +866,7 @@ set_plain_rel_pathlist(PlannerInfo *root, RelOptInfo *rel, RangeTblEntry *rte) static void create_plain_partial_paths(PlannerInfo *root, RelOptInfo *rel) { - int parallel_workers; + int16 parallel_workers; parallel_workers = compute_parallel_worker(rel, rel->pages, -1, max_parallel_workers_per_gather); @@ -1616,13 +1617,14 @@ add_paths_to_append_rel(PlannerInfo *root, RelOptInfo *rel, */ if (unparameterized_valid) add_path(rel, (Path *) create_append_path(root, rel, unparameterized, - NIL, NULL, 0, false, + NIL, NULL, 0, 0, false, -1)); /* build an AppendPath for the cheap startup paths, if valid */ if (startup_valid) add_path(rel, (Path *) create_append_path(root, rel, startup, - NIL, NULL, 0, false, -1)); + NIL, NULL, 0, 0, false, + -1)); /* * Consider an append of unordered, unparameterized partial paths. Make @@ -1633,6 +1635,7 @@ add_paths_to_append_rel(PlannerInfo *root, RelOptInfo *rel, AppendPath *appendpath; ListCell *lc; int parallel_workers = 0; + int effective_workers = 0; /* Find the highest number of workers requested for any subpath. */ foreach(lc, partial_only.partial_subpaths) @@ -1640,7 +1643,16 @@ add_paths_to_append_rel(PlannerInfo *root, RelOptInfo *rel, Path *path = lfirst(lc); parallel_workers = Max(parallel_workers, path->parallel_workers); + effective_workers += path->effective_workers; } + + /* + * If the parallel leader participates, include that in the effective + * worker count. + */ + if (parallel_leader_participation) + effective_workers += list_length(partial_only.partial_subpaths); + Assert(parallel_workers > 0); /* @@ -1651,19 +1663,38 @@ add_paths_to_append_rel(PlannerInfo *root, RelOptInfo *rel, * want to end up with a radically different answer for a table with N * partitions vs. an unpartitioned table with the same data, so the * use of some kind of log-scaling here seems to make some sense. + * + * If any workers are added here, we also adjust the effective_worker + * count accordingly. */ if (enable_parallel_append) { + int init_parallel_workers = parallel_workers; + parallel_workers = Max(parallel_workers, pg_leftmost_one_pos32(list_length(live_childrels)) + 1); parallel_workers = Min(parallel_workers, max_parallel_workers_per_gather); + + /* + * Distribute added workers to the effective worker count, too; + * non-parallel subpaths will contribute across those workers. + */ + if (init_parallel_workers < parallel_workers) + effective_workers += parallel_workers - init_parallel_workers; } - Assert(parallel_workers > 0); + + effective_workers = Min(effective_workers, parallel_workers); + + Assert(MAX_PARALLEL_WORKER_LIMIT >= parallel_workers && + parallel_workers > 0); + Assert(parallel_workers >= effective_workers && + effective_workers >= 0); /* Generate a partial append path. */ appendpath = create_append_path(root, rel, partial_only, - NIL, NULL, parallel_workers, + NIL, NULL, + parallel_workers, effective_workers, enable_parallel_append, -1); @@ -1687,7 +1718,9 @@ add_paths_to_append_rel(PlannerInfo *root, RelOptInfo *rel, { AppendPath *appendpath; ListCell *lc; + int init_parallel_workers; int parallel_workers = 0; + int effective_workers = 0; /* * Find the highest number of workers requested for any partial @@ -1698,8 +1731,14 @@ add_paths_to_append_rel(PlannerInfo *root, RelOptInfo *rel, Path *path = lfirst(lc); parallel_workers = Max(parallel_workers, path->parallel_workers); + effective_workers += path->effective_workers; } + if (parallel_leader_participation) + effective_workers += list_length(parallel_append.partial_subpaths); + + init_parallel_workers = parallel_workers; + /* * Same formula here as above. It's even more important in this * instance because the non-partial paths won't contribute anything to @@ -1709,10 +1748,20 @@ add_paths_to_append_rel(PlannerInfo *root, RelOptInfo *rel, pg_leftmost_one_pos32(list_length(live_childrels)) + 1); parallel_workers = Min(parallel_workers, max_parallel_workers_per_gather); - Assert(parallel_workers > 0); + + if (parallel_workers > init_parallel_workers) + effective_workers += parallel_workers - init_parallel_workers; + + effective_workers = Min(parallel_workers, effective_workers); + + Assert(MAX_PARALLEL_WORKER_LIMIT >= parallel_workers && + parallel_workers > 0); + Assert(parallel_workers >= effective_workers && + effective_workers >= 0); appendpath = create_append_path(root, rel, parallel_append, - NIL, NULL, parallel_workers, true, + NIL, NULL, parallel_workers, + effective_workers, true, partial_rows); add_partial_path(rel, (Path *) appendpath); } @@ -1774,8 +1823,8 @@ add_paths_to_append_rel(PlannerInfo *root, RelOptInfo *rel, if (parameterized_valid) add_path(rel, (Path *) create_append_path(root, rel, parameterized, - NIL, required_outer, 0, false, - -1)); + NIL, required_outer, 0, 0, + false, -1)); } /* @@ -1802,8 +1851,8 @@ add_paths_to_append_rel(PlannerInfo *root, RelOptInfo *rel, append.partial_subpaths = list_make1(path); appendpath = create_append_path(root, rel, append, NIL, NULL, - path->parallel_workers, true, - partial_rows); + path->parallel_workers, path->effective_workers, + true, partial_rows); add_partial_path(rel, (Path *) appendpath); } } @@ -2097,6 +2146,7 @@ generate_orderedappend_paths(PlannerInfo *root, RelOptInfo *rel, pathkeys, NULL, 0, + 0, false, -1)); if (startup_neq_total) @@ -2106,6 +2156,7 @@ generate_orderedappend_paths(PlannerInfo *root, RelOptInfo *rel, pathkeys, NULL, 0, + 0, false, -1)); @@ -2116,6 +2167,7 @@ generate_orderedappend_paths(PlannerInfo *root, RelOptInfo *rel, pathkeys, NULL, 0, + 0, false, -1)); } @@ -2372,7 +2424,7 @@ set_dummy_rel_pathlist(RelOptInfo *rel) /* Set up the dummy path */ add_path(rel, (Path *) create_append_path(NULL, rel, in, NIL, rel->lateral_relids, - 0, false, -1)); + 0, 0, false, -1)); /* * We set the cheapest-path fields immediately, just in case they were @@ -4959,7 +5011,7 @@ create_partial_bitmap_paths(PlannerInfo *root, RelOptInfo *rel, * "max_workers" is caller's limit on the number of workers. This typically * comes from a GUC. */ -int +int16 compute_parallel_worker(RelOptInfo *rel, double heap_pages, double index_pages, int max_workers) { @@ -4968,6 +5020,8 @@ compute_parallel_worker(RelOptInfo *rel, double heap_pages, double index_pages, /* * If the user has set the parallel_workers reloption, use that; otherwise * select a default number of workers. + * + * We do try to avoid exceeding the parallel worker limit */ if (rel->rel_parallel_workers != -1) parallel_workers = rel->rel_parallel_workers; @@ -5035,7 +5089,10 @@ compute_parallel_worker(RelOptInfo *rel, double heap_pages, double index_pages, /* In no case use more than caller supplied maximum number of workers */ parallel_workers = Min(parallel_workers, max_workers); - return parallel_workers; + Assert(parallel_workers <= MAX_PARALLEL_WORKER_LIMIT); + Assert(parallel_workers >= 0); + + return (int16) parallel_workers; } /* diff --git a/src/backend/optimizer/path/costsize.c b/src/backend/optimizer/path/costsize.c index 7bbddb8bee4..c7e8cb57863 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 @@ -2462,7 +2467,7 @@ cost_append(AppendPath *apath, PlannerInfo *root) */ if (i == 0) apath->path.startup_cost = subpath->startup_cost; - else if (i < apath->path.parallel_workers) + else if (i < (int) apath->path.parallel_workers) apath->path.startup_cost = Min(apath->path.startup_cost, subpath->startup_cost); @@ -6731,7 +6736,10 @@ 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; + + Assert(0 <= path->effective_workers && + path->effective_workers <= path->parallel_workers); /* * Early experience with parallel query suggests that when there is only @@ -6748,7 +6756,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/path/joinrels.c b/src/backend/optimizer/path/joinrels.c index 10fb3e39d28..4124c448d3a 100644 --- a/src/backend/optimizer/path/joinrels.c +++ b/src/backend/optimizer/path/joinrels.c @@ -1558,7 +1558,7 @@ mark_dummy_rel(RelOptInfo *rel) /* Set up the dummy path */ add_path(rel, (Path *) create_append_path(NULL, rel, in, NIL, rel->lateral_relids, - 0, false, -1)); + 0, 0, false, -1)); /* Set or update cheapest_total_path and related fields */ set_cheapest(rel); diff --git a/src/backend/optimizer/plan/planner.c b/src/backend/optimizer/plan/planner.c index 9458efab480..518ac13cb4d 100644 --- a/src/backend/optimizer/plan/planner.c +++ b/src/backend/optimizer/plan/planner.c @@ -4302,6 +4302,7 @@ create_degenerate_grouping_paths(PlannerInfo *root, RelOptInfo *input_rel, NIL, NULL, 0, + 0, false, -1); } diff --git a/src/backend/optimizer/prep/prepunion.c b/src/backend/optimizer/prep/prepunion.c index b136f12ff3b..c3478ecece1 100644 --- a/src/backend/optimizer/prep/prepunion.c +++ b/src/backend/optimizer/prep/prepunion.c @@ -31,6 +31,7 @@ #include "nodes/makefuncs.h" #include "nodes/nodeFuncs.h" #include "optimizer/cost.h" +#include "optimizer/optimizer.h" #include "optimizer/pathnode.h" #include "optimizer/paths.h" #include "optimizer/planner.h" @@ -38,6 +39,7 @@ #include "optimizer/tlist.h" #include "parser/parse_coerce.h" #include "port/pg_bitutils.h" +#include "postmaster/bgworker_internals.h" #include "utils/selfuncs.h" @@ -860,7 +862,7 @@ generate_union_paths(SetOperationStmt *op, PlannerInfo *root, * union child. */ apath = (Path *) create_append_path(root, result_rel, cheapest, - NIL, NULL, 0, false, -1); + NIL, NULL, 0, 0, false, -1); /* * Initialize the result row estimate to the total input size. This is @@ -877,6 +879,7 @@ generate_union_paths(SetOperationStmt *op, PlannerInfo *root, { Path *papath; int parallel_workers = 0; + int effective_workers = 0; /* Find the highest number of workers requested for any subpath. */ foreach(lc, partial.partial_subpaths) @@ -885,8 +888,17 @@ generate_union_paths(SetOperationStmt *op, PlannerInfo *root, parallel_workers = Max(parallel_workers, subpath->parallel_workers); + effective_workers += subpath->effective_workers; } - Assert(parallel_workers > 0); + + /* + * If the parallel leader participates, include that in the effective + * worker count. + */ + if (parallel_leader_participation) + effective_workers += list_length(partial.partial_subpaths); + + Assert(parallel_workers > 0 && effective_workers >= 0); /* * If the use of parallel append is permitted, always request at least @@ -894,19 +906,34 @@ generate_union_paths(SetOperationStmt *op, PlannerInfo *root, * extra workers in this case because they will be spread out across * the children. The precise formula is just a guess; see * add_paths_to_append_rel. + * + * If any workers are added here, we also adjust the effective_worker + * count accordingly. */ if (enable_parallel_append) { + int init_parallel_workers = parallel_workers; + parallel_workers = Max(parallel_workers, pg_leftmost_one_pos32(list_length(partial.partial_subpaths)) + 1); parallel_workers = Min(parallel_workers, max_parallel_workers_per_gather); + + if (init_parallel_workers < parallel_workers) + effective_workers += parallel_workers - init_parallel_workers; } - Assert(parallel_workers > 0); + + effective_workers = Min(parallel_workers, effective_workers); + + Assert(MAX_PARALLEL_WORKER_LIMIT >= parallel_workers && + parallel_workers > 0); + Assert(parallel_workers >= effective_workers && + effective_workers >= 0); papath = (Path *) create_append_path(root, result_rel, partial, - NIL, NULL, parallel_workers, + NIL, NULL, + parallel_workers, effective_workers, enable_parallel_append, -1); gpath = (Path *) create_gather_path(root, result_rel, papath, @@ -1223,7 +1250,7 @@ generate_nonunion_paths(SetOperationStmt *op, PlannerInfo *root, */ apath = (Path *) create_append_path(root, result_rel, append, NIL, NULL, 0, - false, -1); + 0, false, -1); add_path(result_rel, apath); diff --git a/src/backend/optimizer/util/pathnode.c b/src/backend/optimizer/util/pathnode.c index 67c47fd7c6c..34d7b22b648 100644 --- a/src/backend/optimizer/util/pathnode.c +++ b/src/backend/optimizer/util/pathnode.c @@ -14,6 +14,8 @@ */ #include "postgres.h" +#include + #include "access/htup_details.h" #include "executor/nodeSetOp.h" #include "foreign/fdwapi.h" @@ -28,6 +30,7 @@ #include "optimizer/planmain.h" #include "optimizer/tlist.h" #include "parser/parsetree.h" +#include "postmaster/bgworker_internals.h" #include "utils/memutils.h" #include "utils/selfuncs.h" @@ -53,6 +56,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 int16 subpath_adjusted_effective_workers(const Path *subpath); /***************************************************************************** @@ -237,6 +241,96 @@ 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 int16 +subpath_adjusted_effective_workers(const Path *subpath) +{ + int leader_factor = 0; + int result; + double best_row_estimate; + + /* non-parallel subpaths don't have any effective worker count */ + if (subpath->parallel_workers == 0) + return 0; + + /* + * If every effective worker of the subpath produces more than one + * row, then parallelism shouldn't be reduced. + */ + if (subpath->rows > 1) + return subpath->effective_workers; + + if (parallel_leader_participation) + leader_factor = 1; + + /* + * If we're already down to one effective process, don't further + * expect parallelism to reduce. + */ + if (subpath->effective_workers + leader_factor <= 1) + return subpath->effective_workers; + + /* If there is just, */ + if (Max(0, subpath->rows) == 0.0) + return subpath->effective_workers; + + /* most basic estimate for row counts */ + best_row_estimate = Max(0, subpath->rows) * + (double) (subpath->effective_workers + leader_factor); + + result = subpath->effective_workers; + + Assert(subpath->effective_workers + leader_factor >= 1); + + /* Adjust down for parent rel rowcount */ + if (subpath->parent->rows >= 0) + best_row_estimate = Min(subpath->parent->rows, best_row_estimate); + + /* Adjust down for parameterization, where relevant */ + if (subpath->param_info && subpath->param_info->ppi_rows >= 0) + best_row_estimate = Min(subpath->param_info->ppi_rows, + best_row_estimate); + + /* + * 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 < (double) (result + leader_factor)) + { + result = (int16) (ceil(best_row_estimate) - leader_factor); + + /* + * If the row estimate is zero and a leader is present, the + * worker count could underflow. + */ + if (result < 1 - leader_factor) + result = 1 - leader_factor; + } + + Assert(result + leader_factor >= 1); + Assert(result <= MAX_PARALLEL_WORKER_LIMIT); + + return (int16) result; +} + /* * set_cheapest * Find the minimum-cost paths from among a relation's paths, @@ -1024,7 +1118,7 @@ add_partial_path_precheck(RelOptInfo *parent_rel, int disabled_nodes, */ Path * create_seqscan_path(PlannerInfo *root, RelOptInfo *rel, - Relids required_outer, int parallel_workers) + Relids required_outer, int16 parallel_workers) { Path *pathnode = makeNode(Path); @@ -1036,6 +1130,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 +1155,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 +1208,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; @@ -1151,7 +1248,7 @@ create_bitmap_heap_path(PlannerInfo *root, Path *bitmapqual, Relids required_outer, double loop_count, - int parallel_degree) + int16 parallel_degree) { BitmapHeapPath *pathnode = makeNode(BitmapHeapPath); @@ -1163,6 +1260,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 +1313,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 +1366,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 +1396,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; @@ -1314,7 +1415,7 @@ create_tidscan_path(PlannerInfo *root, RelOptInfo *rel, List *tidquals, TidRangePath * create_tidrangescan_path(PlannerInfo *root, RelOptInfo *rel, List *tidrangequals, Relids required_outer, - int parallel_workers) + int16 parallel_workers) { TidRangePath *pathnode = makeNode(TidRangePath); @@ -1326,6 +1427,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; @@ -1353,13 +1455,14 @@ create_append_path(PlannerInfo *root, RelOptInfo *rel, AppendPathInput input, List *pathkeys, Relids required_outer, - int parallel_workers, bool parallel_aware, - double rows) + int16 parallel_workers, int16 effective_workers, + bool parallel_aware, double rows) { AppendPath *pathnode = makeNode(AppendPath); ListCell *l; - Assert(!parallel_aware || parallel_workers > 0); + Assert(!parallel_aware || + (parallel_workers > 0 && effective_workers >= 0)); pathnode->child_append_relid_sets = input.child_append_relid_sets; pathnode->path.pathtype = T_Append; @@ -1388,6 +1491,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 = effective_workers; pathnode->path.pathkeys = pathkeys; /* @@ -1549,6 +1653,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 +1780,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 +1831,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 +1869,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 +1989,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 +2034,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 +2064,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 +2091,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 +2118,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 +2145,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 +2172,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 +2199,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 +2226,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 +2270,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 +2325,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 +2375,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 +2542,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 +2610,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 +2678,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 +2748,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 +2934,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 +3006,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 +3056,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 +3104,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 +3161,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 +3229,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 +3315,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 +3477,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 +3573,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; @@ -3496,8 +3644,14 @@ create_setop_path(PlannerInfo *root, pathnode->path.parallel_aware = false; pathnode->path.parallel_safe = rel->consider_parallel && leftpath->parallel_safe && rightpath->parallel_safe; + + /* avoid exceeding parallel worker count limits */ pathnode->path.parallel_workers = - leftpath->parallel_workers + rightpath->parallel_workers; + Min(leftpath->parallel_workers + rightpath->parallel_workers, + max_parallel_workers); + pathnode->path.effective_workers = + Min(leftpath->effective_workers + rightpath->effective_workers, + max_parallel_workers); /* SetOp preserves the input sort order if in sort mode */ pathnode->path.pathkeys = (strategy == SETOP_SORTED) ? leftpath->pathkeys : NIL; @@ -3626,6 +3780,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 +3820,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 +3900,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 +3988,7 @@ 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->effective_workers; pathnode->path.rows = subpath->rows; pathnode->path.disabled_nodes = subpath->disabled_nodes; pathnode->path.startup_cost = subpath->startup_cost; @@ -4039,8 +4198,8 @@ reparameterize_path(PlannerInfo *root, Path *path, create_append_path(root, rel, new_append, apath->path.pathkeys, required_outer, apath->path.parallel_workers, - apath->path.parallel_aware, - -1); + apath->path.effective_workers, + apath->path.parallel_aware, -1); } case T_Material: { diff --git a/src/include/nodes/pathnodes.h b/src/include/nodes/pathnodes.h index 1c6d1fe3d04..a9f13293174 100644 --- a/src/include/nodes/pathnodes.h +++ b/src/include/nodes/pathnodes.h @@ -1999,7 +1999,16 @@ typedef struct Path /* OK to use as part of parallel plan? */ bool parallel_safe; /* desired # of workers; 0 = not parallel */ - int parallel_workers; + int16 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. + */ + int16 effective_workers; /* estimated size/costs for path (see costsize.c for more info) */ Cardinality rows; /* estimated number of result tuples */ diff --git a/src/include/optimizer/pathnode.h b/src/include/optimizer/pathnode.h index da2d9b384b5..3ffe8b38bb5 100644 --- a/src/include/optimizer/pathnode.h +++ b/src/include/optimizer/pathnode.h @@ -65,7 +65,7 @@ extern bool add_partial_path_precheck(RelOptInfo *parent_rel, Cost total_cost, List *pathkeys); extern Path *create_seqscan_path(PlannerInfo *root, RelOptInfo *rel, - Relids required_outer, int parallel_workers); + Relids required_outer, int16 parallel_workers); extern Path *create_samplescan_path(PlannerInfo *root, RelOptInfo *rel, Relids required_outer); extern IndexPath *create_index_path(PlannerInfo *root, @@ -84,7 +84,7 @@ extern BitmapHeapPath *create_bitmap_heap_path(PlannerInfo *root, Path *bitmapqual, Relids required_outer, double loop_count, - int parallel_degree); + int16 parallel_degree); extern BitmapAndPath *create_bitmap_and_path(PlannerInfo *root, RelOptInfo *rel, List *bitmapquals); @@ -97,13 +97,13 @@ extern TidRangePath *create_tidrangescan_path(PlannerInfo *root, RelOptInfo *rel, List *tidrangequals, Relids required_outer, - int parallel_workers); + int16 parallel_workers); extern AppendPath *create_append_path(PlannerInfo *root, RelOptInfo *rel, AppendPathInput input, List *pathkeys, Relids required_outer, - int parallel_workers, bool parallel_aware, - double rows); + int16 parallel_workers, int16 effective_workers, + bool parallel_aware, double rows); extern MergeAppendPath *create_merge_append_path(PlannerInfo *root, RelOptInfo *rel, List *subpaths, diff --git a/src/include/optimizer/paths.h b/src/include/optimizer/paths.h index d3853d1c076..bf88e5b407d 100644 --- a/src/include/optimizer/paths.h +++ b/src/include/optimizer/paths.h @@ -68,8 +68,8 @@ extern void generate_useful_gather_paths(PlannerInfo *root, RelOptInfo *rel, bool override_rows); extern void generate_grouped_paths(PlannerInfo *root, RelOptInfo *grouped_rel, RelOptInfo *rel); -extern int compute_parallel_worker(RelOptInfo *rel, double heap_pages, - double index_pages, int max_workers); +extern int16 compute_parallel_worker(RelOptInfo *rel, double heap_pages, + double index_pages, int max_workers); extern void create_partial_bitmap_paths(PlannerInfo *root, RelOptInfo *rel, Path *bitmapqual); extern void generate_partitionwise_join_paths(PlannerInfo *root, -- 2.54.0 (Apple Git-157)