From bdfae4e59e424def7651c2a45c61ccfce70d9e89 Mon Sep 17 00:00:00 2001 From: David Rowley Date: Wed, 7 Oct 2026 12:21:22 +1300 Subject: [PATCH v1] Show Material node Maximum Storage for parallel workers 1eff8279d did this for the leader, but neglected to consider that the Material node may be below a Gather or Gather Merge and have parallel workers also participate in the execution of the Material node. Here we add EXPLAIN ANALYZE output to show Maximum Storage per worker and all the required shared memory related code to get those details back to the leader from each worker. This also adjusts the public tuplestore_get_stats() API function. Previously max_storage_type was populated to point to either a "Disk" or "Memory" string constant. That's been changed to populate a bool pointer with true if the tuplestore spilled to disk. Any external callers should follow the update made in show_storage_info() in explain.c. This change had to be made as we cannot store pointers to process-local memory inside shared memory. --- src/backend/commands/explain.c | 77 +++++++++++------ src/backend/executor/execParallel.c | 17 ++++ src/backend/executor/nodeMaterial.c | 110 +++++++++++++++++++++++++ src/backend/utils/sort/tuplestore.c | 9 +- src/include/executor/instrument_node.h | 21 +++++ src/include/executor/nodeMaterial.h | 7 ++ src/include/nodes/execnodes.h | 1 + src/include/utils/tuplestore.h | 2 +- src/tools/pgindent/typedefs.list | 2 + 9 files changed, 213 insertions(+), 33 deletions(-) diff --git a/src/backend/commands/explain.c b/src/backend/commands/explain.c index 96f2f0e2e74..4b1ad4227ac 100644 --- a/src/backend/commands/explain.c +++ b/src/backend/commands/explain.c @@ -119,7 +119,7 @@ static void show_window_def(WindowAggState *planstate, static void show_window_keys(StringInfo buf, PlanState *planstate, int nkeys, AttrNumber *keycols, List *ancestors, ExplainState *es); -static void show_storage_info(char *maxStorageType, int64 maxSpaceUsed, +static void show_storage_info(bool diskUsed, int64 maxSpaceUsed, ExplainState *es); static void show_tablesample(TableSampleClause *tsc, PlanState *planstate, List *ancestors, ExplainState *es); @@ -3006,13 +3006,14 @@ show_window_keys(StringInfo buf, PlanState *planstate, * Show information on storage method and maximum memory/disk space used. */ static void -show_storage_info(char *maxStorageType, int64 maxSpaceUsed, ExplainState *es) +show_storage_info(bool diskUsed, int64 maxSpaceUsed, ExplainState *es) { int64 maxSpaceUsedKB = BYTES_TO_KILOBYTES(maxSpaceUsed); + const char *storageType = diskUsed ? "Disk" : "Memory"; if (es->format != EXPLAIN_FORMAT_TEXT) { - ExplainPropertyText("Storage", maxStorageType, es); + ExplainPropertyText("Storage", storageType, es); ExplainPropertyInteger("Maximum Storage", "kB", maxSpaceUsedKB, es); } else @@ -3020,7 +3021,7 @@ show_storage_info(char *maxStorageType, int64 maxSpaceUsed, ExplainState *es) ExplainIndentText(es); appendStringInfo(es->str, "Storage: %s Maximum Storage: " INT64_FORMAT "kB\n", - maxStorageType, + storageType, maxSpaceUsedKB); } } @@ -3485,20 +3486,46 @@ show_hash_info(HashState *hashstate, ExplainState *es) static void show_material_info(MaterialState *mstate, ExplainState *es) { - char *maxStorageType; + bool usedDisk; int64 maxSpaceUsed; Tuplestorestate *tupstore = mstate->tuplestorestate; + if (!es->analyze) + return; + /* - * Nothing to show if ANALYZE option wasn't used or if execution didn't - * get as far as creating the tuplestore. + * Nothing to show for the leader if execution didn't get as far as + * creating the tuplestore. */ - if (!es->analyze || tupstore == NULL) + if (tupstore != NULL) + { + tuplestore_get_stats(tupstore, &usedDisk, &maxSpaceUsed); + show_storage_info(usedDisk, maxSpaceUsed, es); + } + + if (mstate->shared_info == NULL) return; - tuplestore_get_stats(tupstore, &maxStorageType, &maxSpaceUsed); - show_storage_info(maxStorageType, maxSpaceUsed, es); + /* Show details from parallel workers */ + for (int n = 0; n < mstate->shared_info->num_workers; n++) + { + MaterialInstrumentation *si; + + si = &mstate->shared_info->sinstrument[n]; + + /* Skip workers that didn't create a tuplestore */ + if (si->maxSpaceUsed == 0) + continue; + + if (es->workers_state) + ExplainOpenWorker(n, es); + + show_storage_info(si->usedDisk, si->maxSpaceUsed, es); + + if (es->workers_state) + ExplainCloseWorker(n, es); + } } /* @@ -3508,7 +3535,7 @@ show_material_info(MaterialState *mstate, ExplainState *es) static void show_windowagg_info(WindowAggState *winstate, ExplainState *es) { - char *maxStorageType; + bool usedDisk; int64 maxSpaceUsed; Tuplestorestate *tupstore = winstate->buffer; @@ -3520,8 +3547,8 @@ show_windowagg_info(WindowAggState *winstate, ExplainState *es) if (!es->analyze || tupstore == NULL) return; - tuplestore_get_stats(tupstore, &maxStorageType, &maxSpaceUsed); - show_storage_info(maxStorageType, maxSpaceUsed, es); + tuplestore_get_stats(tupstore, &usedDisk, &maxSpaceUsed); + show_storage_info(usedDisk, maxSpaceUsed, es); } /* @@ -3531,7 +3558,7 @@ show_windowagg_info(WindowAggState *winstate, ExplainState *es) static void show_ctescan_info(CteScanState *ctescanstate, ExplainState *es) { - char *maxStorageType; + bool usedDisk; int64 maxSpaceUsed; Tuplestorestate *tupstore = ctescanstate->leader->cte_table; @@ -3539,8 +3566,8 @@ show_ctescan_info(CteScanState *ctescanstate, ExplainState *es) if (!es->analyze || tupstore == NULL) return; - tuplestore_get_stats(tupstore, &maxStorageType, &maxSpaceUsed); - show_storage_info(maxStorageType, maxSpaceUsed, es); + tuplestore_get_stats(tupstore, &usedDisk, &maxSpaceUsed); + show_storage_info(usedDisk, maxSpaceUsed, es); } /* @@ -3550,7 +3577,7 @@ show_ctescan_info(CteScanState *ctescanstate, ExplainState *es) static void show_table_func_scan_info(TableFuncScanState *tscanstate, ExplainState *es) { - char *maxStorageType; + bool usedDisk; int64 maxSpaceUsed; Tuplestorestate *tupstore = tscanstate->tupstore; @@ -3558,8 +3585,8 @@ show_table_func_scan_info(TableFuncScanState *tscanstate, ExplainState *es) if (!es->analyze || tupstore == NULL) return; - tuplestore_get_stats(tupstore, &maxStorageType, &maxSpaceUsed); - show_storage_info(maxStorageType, maxSpaceUsed, es); + tuplestore_get_stats(tupstore, &usedDisk, &maxSpaceUsed); + show_storage_info(usedDisk, maxSpaceUsed, es); } /* @@ -3569,8 +3596,8 @@ show_table_func_scan_info(TableFuncScanState *tscanstate, ExplainState *es) static void show_recursive_union_info(RecursiveUnionState *rstate, ExplainState *es) { - char *maxStorageType, - *tempStorageType; + bool maxUsedDisk, + tempUsedDisk; int64 maxSpaceUsed, tempSpaceUsed; @@ -3582,16 +3609,16 @@ show_recursive_union_info(RecursiveUnionState *rstate, ExplainState *es) * from one of them which consumed more memory/disk than the other. The * storage size is sum of the two. */ - tuplestore_get_stats(rstate->working_table, &tempStorageType, + tuplestore_get_stats(rstate->working_table, &tempUsedDisk, &tempSpaceUsed); - tuplestore_get_stats(rstate->intermediate_table, &maxStorageType, + tuplestore_get_stats(rstate->intermediate_table, &maxUsedDisk, &maxSpaceUsed); if (tempSpaceUsed > maxSpaceUsed) - maxStorageType = tempStorageType; + maxUsedDisk = tempUsedDisk; maxSpaceUsed += tempSpaceUsed; - show_storage_info(maxStorageType, maxSpaceUsed, es); + show_storage_info(maxUsedDisk, maxSpaceUsed, es); } /* diff --git a/src/backend/executor/execParallel.c b/src/backend/executor/execParallel.c index 075bc5cd2c9..316dd3f67c6 100644 --- a/src/backend/executor/execParallel.c +++ b/src/backend/executor/execParallel.c @@ -36,6 +36,7 @@ #include "executor/nodeIncrementalSort.h" #include "executor/nodeIndexonlyscan.h" #include "executor/nodeIndexscan.h" +#include "executor/nodeMaterial.h" #include "executor/nodeMemoize.h" #include "executor/nodeSeqscan.h" #include "executor/nodeSort.h" @@ -338,6 +339,10 @@ ExecParallelEstimate(PlanState *planstate, ExecParallelEstimateContext *e) /* even when not parallel-aware, for EXPLAIN ANALYZE */ ExecMemoizeEstimate((MemoizeState *) planstate, e->pcxt); break; + case T_MaterialState: + /* even when not parallel-aware, for EXPLAIN ANALYZE */ + ExecMaterialEstimate((MaterialState *) planstate, e->pcxt); + break; default: break; } @@ -586,6 +591,10 @@ ExecParallelInitializeDSM(PlanState *planstate, /* even when not parallel-aware, for EXPLAIN ANALYZE */ ExecMemoizeInitializeDSM((MemoizeState *) planstate, d->pcxt); break; + case T_MaterialState: + /* even when not parallel-aware, for EXPLAIN ANALYZE */ + ExecMaterialInitializeDSM((MaterialState *) planstate, d->pcxt); + break; default: break; } @@ -1074,6 +1083,7 @@ ExecParallelReInitializeDSM(PlanState *planstate, case T_SortState: case T_IncrementalSortState: case T_MemoizeState: + case T_MaterialState: /* these nodes have DSM state, but no reinitialization is required */ break; @@ -1155,6 +1165,9 @@ ExecParallelRetrieveInstrumentation(PlanState *planstate, case T_MemoizeState: ExecMemoizeRetrieveInstrumentation((MemoizeState *) planstate); break; + case T_MaterialState: + ExecMaterialRetrieveInstrumentation((MaterialState *) planstate); + break; case T_BitmapHeapScanState: ExecBitmapHeapRetrieveInstrumentation((BitmapHeapScanState *) planstate); break; @@ -1484,6 +1497,10 @@ ExecParallelInitializeWorker(PlanState *planstate, ParallelWorkerContext *pwcxt) /* even when not parallel-aware, for EXPLAIN ANALYZE */ ExecMemoizeInitializeWorker((MemoizeState *) planstate, pwcxt); break; + case T_MaterialState: + /* even when not parallel-aware, for EXPLAIN ANALYZE */ + ExecMaterialInitializeWorker((MaterialState *) planstate, pwcxt); + break; default: break; } diff --git a/src/backend/executor/nodeMaterial.c b/src/backend/executor/nodeMaterial.c index 7f5aceaa3ca..72c8fe88e72 100644 --- a/src/backend/executor/nodeMaterial.c +++ b/src/backend/executor/nodeMaterial.c @@ -197,6 +197,7 @@ ExecInitMaterial(Material *node, EState *estate, int eflags) matstate->eof_underlying = false; matstate->tuplestorestate = NULL; + matstate->shared_info = NULL; /* * Miscellaneous initialization @@ -240,6 +241,29 @@ ExecInitMaterial(Material *node, EState *estate, int eflags) void ExecEndMaterial(MaterialState *node) { + /* + * When ending a parallel worker, copy the tuplestore usage details into + * shared memory so that they can be picked up by the main process to + * report in EXPLAIN ANALYZE. The Gather node above us may be rescanned, + * in which case a new set of workers will write to the same slots, so + * merge with any existing values rather than overwriting them. + */ + if (node->shared_info != NULL && IsParallelWorker() && + node->tuplestorestate != NULL) + { + MaterialInstrumentation *si; + bool usedDisk; + int64 maxSpaceUsed; + + Assert(ParallelWorkerNumber < node->shared_info->num_workers); + + tuplestore_get_stats(node->tuplestorestate, &usedDisk, &maxSpaceUsed); + + si = &node->shared_info->sinstrument[ParallelWorkerNumber]; + si->maxSpaceUsed = Max(si->maxSpaceUsed, maxSpaceUsed); + si->usedDisk |= usedDisk; + } + /* * Release tuplestore resources */ @@ -365,3 +389,89 @@ ExecReScanMaterial(MaterialState *node) node->eof_underlying = false; } } + +/* ---------------------------------------------------------------- + * Parallel Query Support + * ---------------------------------------------------------------- + */ + +/* ---------------------------------------------------------------- + * ExecMaterialEstimate + * + * Estimate space required to propagate material statistics. + * ---------------------------------------------------------------- + */ +void +ExecMaterialEstimate(MaterialState *node, ParallelContext *pcxt) +{ + Size size; + + /* don't need this if not instrumenting or no workers */ + if (!node->ss.ps.instrument || pcxt->nworkers == 0) + return; + + size = mul_size(pcxt->nworkers, sizeof(MaterialInstrumentation)); + size = add_size(size, offsetof(SharedMaterialInfo, sinstrument)); + shm_toc_estimate_chunk(&pcxt->estimator, size); + shm_toc_estimate_keys(&pcxt->estimator, 1); +} + +/* ---------------------------------------------------------------- + * ExecMaterialInitializeDSM + * + * Initialize DSM space for material statistics. + * ---------------------------------------------------------------- + */ +void +ExecMaterialInitializeDSM(MaterialState *node, ParallelContext *pcxt) +{ + Size size; + + /* don't need this if not instrumenting or no workers */ + if (!node->ss.ps.instrument || pcxt->nworkers == 0) + return; + + size = offsetof(SharedMaterialInfo, sinstrument) + + pcxt->nworkers * sizeof(MaterialInstrumentation); + node->shared_info = shm_toc_allocate(pcxt->toc, size); + /* ensure any unfilled slots will contain zeroes */ + memset(node->shared_info, 0, size); + node->shared_info->num_workers = pcxt->nworkers; + shm_toc_insert(pcxt->toc, node->ss.ps.plan->plan_node_id, + node->shared_info); +} + +/* ---------------------------------------------------------------- + * ExecMaterialInitializeWorker + * + * Attach worker to DSM space for material statistics. + * ---------------------------------------------------------------- + */ +void +ExecMaterialInitializeWorker(MaterialState *node, ParallelWorkerContext *pwcxt) +{ + node->shared_info = + shm_toc_lookup(pwcxt->toc, node->ss.ps.plan->plan_node_id, true); +} + +/* ---------------------------------------------------------------- + * ExecMaterialRetrieveInstrumentation + * + * Transfer material statistics from DSM to private memory. + * ---------------------------------------------------------------- + */ +void +ExecMaterialRetrieveInstrumentation(MaterialState *node) +{ + Size size; + SharedMaterialInfo *si; + + if (node->shared_info == NULL) + return; + + size = offsetof(SharedMaterialInfo, sinstrument) + + node->shared_info->num_workers * sizeof(MaterialInstrumentation); + si = palloc(size); + memcpy(si, node->shared_info, size); + node->shared_info = si; +} diff --git a/src/backend/utils/sort/tuplestore.c b/src/backend/utils/sort/tuplestore.c index 35e438ceaa6..610ddf00532 100644 --- a/src/backend/utils/sort/tuplestore.c +++ b/src/backend/utils/sort/tuplestore.c @@ -1575,16 +1575,11 @@ tuplestore_updatemax(Tuplestorestate *state) * tuplestore_trim() or tuplestore_clear(). */ void -tuplestore_get_stats(Tuplestorestate *state, char **max_storage_type, - int64 *max_space) +tuplestore_get_stats(Tuplestorestate *state, bool *used_disk, int64 *max_space) { tuplestore_updatemax(state); - if (state->usedDisk) - *max_storage_type = "Disk"; - else - *max_storage_type = "Memory"; - + *used_disk = state->usedDisk; *max_space = state->maxSpace; } diff --git a/src/include/executor/instrument_node.h b/src/include/executor/instrument_node.h index 41a7d33f19c..4759a628afe 100644 --- a/src/include/executor/instrument_node.h +++ b/src/include/executor/instrument_node.h @@ -173,6 +173,27 @@ typedef struct SharedMemoizeInfo } SharedMemoizeInfo; +/* --------------------- + * Instrumentation information for Material + * --------------------- + */ +typedef struct MaterialInstrumentation +{ + int64 maxSpaceUsed; /* peak memory/disk usage in bytes, or 0 if no + * tuplestore was created */ + bool usedDisk; /* did the tuplestore spill to disk? */ +} MaterialInstrumentation; + +/* + * Shared memory container for per-worker material information + */ +typedef struct SharedMaterialInfo +{ + int num_workers; + MaterialInstrumentation sinstrument[FLEXIBLE_ARRAY_MEMBER]; +} SharedMaterialInfo; + + /* --------------------- * Instrumentation information for Sorts. * --------------------- diff --git a/src/include/executor/nodeMaterial.h b/src/include/executor/nodeMaterial.h index 19b254db4c8..e60b837923b 100644 --- a/src/include/executor/nodeMaterial.h +++ b/src/include/executor/nodeMaterial.h @@ -14,6 +14,7 @@ #ifndef NODEMATERIAL_H #define NODEMATERIAL_H +#include "access/parallel.h" #include "nodes/execnodes.h" extern MaterialState *ExecInitMaterial(Material *node, EState *estate, int eflags); @@ -21,5 +22,11 @@ extern void ExecEndMaterial(MaterialState *node); extern void ExecMaterialMarkPos(MaterialState *node); extern void ExecMaterialRestrPos(MaterialState *node); extern void ExecReScanMaterial(MaterialState *node); +extern void ExecMaterialEstimate(MaterialState *node, ParallelContext *pcxt); +extern void ExecMaterialInitializeDSM(MaterialState *node, + ParallelContext *pcxt); +extern void ExecMaterialInitializeWorker(MaterialState *node, + ParallelWorkerContext *pwcxt); +extern void ExecMaterialRetrieveInstrumentation(MaterialState *node); #endif /* NODEMATERIAL_H */ diff --git a/src/include/nodes/execnodes.h b/src/include/nodes/execnodes.h index 0998c9b2b9d..e3cea52e7df 100644 --- a/src/include/nodes/execnodes.h +++ b/src/include/nodes/execnodes.h @@ -2259,6 +2259,7 @@ typedef struct MaterialState int eflags; /* capability flags to pass to tuplestore */ bool eof_underlying; /* reached end of underlying plan? */ Tuplestorestate *tuplestorestate; + SharedMaterialInfo *shared_info; /* statistics for parallel workers */ } MaterialState; struct MemoizeEntry; diff --git a/src/include/utils/tuplestore.h b/src/include/utils/tuplestore.h index aaacd281874..12dc51b81d2 100644 --- a/src/include/utils/tuplestore.h +++ b/src/include/utils/tuplestore.h @@ -66,7 +66,7 @@ extern void tuplestore_copy_read_pointer(Tuplestorestate *state, extern void tuplestore_trim(Tuplestorestate *state); -extern void tuplestore_get_stats(Tuplestorestate *state, char **max_storage_type, +extern void tuplestore_get_stats(Tuplestorestate *state, bool *used_disk, int64 *max_space); extern bool tuplestore_in_memory(Tuplestorestate *state); diff --git a/src/tools/pgindent/typedefs.list b/src/tools/pgindent/typedefs.list index d001f56efae..3390a5799fe 100644 --- a/src/tools/pgindent/typedefs.list +++ b/src/tools/pgindent/typedefs.list @@ -1726,6 +1726,7 @@ MVNDistinctItem ManyTestResource ManyTestResourceKind Material +MaterialInstrumentation MaterialPath MaterialState MdPathStr @@ -2861,6 +2862,7 @@ SharedInvalSmgrMsg SharedInvalSnapshotMsg SharedInvalidationMessage SharedJitInstrumentation +SharedMaterialInfo SharedMemoizeInfo SharedRecordTableEntry SharedRecordTableKey -- 2.53.0