From 7815bb3ba93acfd0a7ad3899be726a886d6db4dd Mon Sep 17 00:00:00 2001 From: Nik Samokhvalov Date: Wed, 23 Sep 2026 18:34:32 -0700 Subject: [PATCH] Refresh autovacuum costs while parallel indexes finish Use an atomic completion count to keep the leader responsive while parallel index workers run. Preserve the normal parallel worker wait for final errors and WAL feedback. Add injection-point tests for reload and cost rebalance. --- src/backend/commands/vacuumparallel.c | 102 ++++++++++- .../t/001_parallel_autovacuum.pl | 163 ++++++++++++++++++ 2 files changed, 264 insertions(+), 1 deletion(-) diff --git a/src/backend/commands/vacuumparallel.c b/src/backend/commands/vacuumparallel.c index 767d162e57..9124caf5b9 100644 --- a/src/backend/commands/vacuumparallel.c +++ b/src/backend/commands/vacuumparallel.c @@ -43,11 +43,14 @@ #include "executor/instrument.h" #include "optimizer/paths.h" #include "pgstat.h" +#include "postmaster/interrupt.h" #include "storage/bufmgr.h" #include "storage/proc.h" #include "tcop/tcopprot.h" +#include "utils/injection_point.h" #include "utils/lsyscache.h" #include "utils/rel.h" +#include "utils/wait_event.h" /* * DSM keys for parallel vacuum. Unlike other parallel execution code, since @@ -60,6 +63,8 @@ #define PARALLEL_VACUUM_KEY_WAL_USAGE 4 #define PARALLEL_VACUUM_KEY_INDEX_STATS 5 +#define PARALLEL_VACUUM_COST_UPDATE_INTERVAL_MS 100 + /* * Struct for cost-based vacuum delay related parameters to share among an * autovacuum worker and its parallel vacuum workers. @@ -148,6 +153,9 @@ typedef struct PVShared /* Counter for vacuuming and cleanup */ pg_atomic_uint32 idx; + /* Number of indexes completed in the current phase */ + pg_atomic_uint32 completed_indexes; + /* DSA handle where the TidStore lives */ dsa_handle dead_items_dsa_handle; @@ -285,6 +293,8 @@ static int parallel_vacuum_compute_workers(Relation *indrels, int nindexes, int bool *will_parallel_vacuum); static void parallel_vacuum_process_all_indexes(ParallelVacuumState *pvs, int num_index_scans, bool vacuum, PVWorkerStats *wstats); +static void parallel_vacuum_update_leader_cost_params(void); +static void parallel_vacuum_wait_for_indexes(ParallelVacuumState *pvs); static void parallel_vacuum_process_safe_indexes(ParallelVacuumState *pvs); static void parallel_vacuum_process_unsafe_indexes(ParallelVacuumState *pvs); static void parallel_vacuum_process_one_index(ParallelVacuumState *pvs, Relation indrel, @@ -453,6 +463,7 @@ parallel_vacuum_init(Relation rel, Relation *indrels, int nindexes, pg_atomic_init_u32(&(shared->cost_balance), 0); pg_atomic_init_u32(&(shared->active_nworkers), 0); pg_atomic_init_u32(&(shared->idx), 0); + pg_atomic_init_u32(&(shared->completed_indexes), 0); shared->is_autovacuum = AmAutoVacuumWorkerProcess(); @@ -869,6 +880,7 @@ parallel_vacuum_process_all_indexes(ParallelVacuumState *pvs, int num_index_scan /* Reset the parallel index processing and progress counters */ pg_atomic_write_u32(&(pvs->shared->idx), 0); + pg_atomic_write_u32(&(pvs->shared->completed_indexes), 0); /* Setup the shared cost-based vacuum delay and launch workers */ if (nworkers > 0) @@ -929,6 +941,10 @@ parallel_vacuum_process_all_indexes(ParallelVacuumState *pvs, int num_index_scan /* Vacuum the indexes that can be processed by only leader process */ parallel_vacuum_process_unsafe_indexes(pvs); + if (pvs->shared->is_autovacuum && + pvs->pcxt->nworkers_launched > 0) + INJECTION_POINT("parallel-vacuum-leader-before-index", NULL); + /* * Join as a parallel worker. The leader vacuums alone processes all * parallel-safe indexes in the case where no workers are launched. @@ -941,7 +957,9 @@ parallel_vacuum_process_all_indexes(ParallelVacuumState *pvs, int num_index_scan */ if (nworkers > 0) { - /* Wait for all vacuum workers to finish */ + /* Wait for all vacuum workers to finish. */ + if (AmAutoVacuumWorkerProcess()) + parallel_vacuum_wait_for_indexes(pvs); WaitForParallelWorkersToFinish(pvs->pcxt); for (int i = 0; i < pvs->pcxt->nworkers_launched; i++) @@ -974,6 +992,82 @@ parallel_vacuum_process_all_indexes(ParallelVacuumState *pvs, int num_index_scan } } +/* + * Refresh the leader's cost parameters while it waits for parallel vacuum + * workers, and propagate any changes to them. + */ +static void +parallel_vacuum_update_leader_cost_params(void) +{ + Assert(AmAutoVacuumWorkerProcess()); + + if (ConfigReloadPending) + { + ConfigReloadPending = false; + ProcessConfigFile(PGC_SIGHUP); + VacuumUpdateCosts(); + } + else + AutoVacuumUpdateCostLimit(); + + parallel_vacuum_propagate_shared_delay_params(); + INJECTION_POINT("parallel-vacuum-leader-cost-updated", NULL); +} + +/* + * Keep autovacuum cost parameters current while workers vacuum indexes. + * Worker completion is counted atomically after publishing index results; + * the ordinary parallel-worker wait still handles final errors and cleanup. + */ +static void +parallel_vacuum_wait_for_indexes(ParallelVacuumState *pvs) +{ + ParallelContext *pcxt = pvs->pcxt; + + while (pg_atomic_read_membarrier_u32(&(pvs->shared->completed_indexes)) < + pvs->nindexes) + { + int nfinished = 0; + + CHECK_FOR_INTERRUPTS(); + parallel_vacuum_update_leader_cost_params(); + + for (int i = 0; i < pcxt->nworkers_launched; i++) + { + pid_t pid; + shm_mq *mq; + + if (pcxt->worker[i].error_mqh == NULL) + { + nfinished++; + continue; + } + + if (pcxt->worker[i].bgwhandle == NULL || + GetBackgroundWorkerPid(pcxt->worker[i].bgwhandle, &pid) != + BGWH_STOPPED) + continue; + + mq = shm_mq_get_queue(pcxt->worker[i].error_mqh); + if (shm_mq_get_sender(mq) == NULL) + ereport(ERROR, + (errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE), + errmsg("parallel worker failed to initialize"), + errhint("More details may be available in the server log."))); + } + + /* Let the final index-status check report any missing index. */ + if (nfinished == pcxt->nworkers_launched) + break; + + (void) WaitLatch(MyLatch, + WL_LATCH_SET | WL_TIMEOUT | WL_EXIT_ON_PM_DEATH, + PARALLEL_VACUUM_COST_UPDATE_INTERVAL_MS, + WAIT_EVENT_PARALLEL_FINISH); + ResetLatch(MyLatch); + } +} + /* * Index vacuum/cleanup routine used by the leader process and parallel * vacuum worker processes to vacuum the indexes in parallel. @@ -1097,6 +1191,11 @@ parallel_vacuum_process_one_index(ParallelVacuumState *pvs, Relation indrel, pvs->indname = pstrdup(RelationGetRelationName(indrel)); pvs->status = indstats->status; +#ifdef USE_INJECTION_POINTS + if (IsParallelWorker()) + INJECTION_POINT("parallel-vacuum-worker-before-index", NULL); +#endif + switch (indstats->status) { case PARALLEL_INDVAC_STATUS_NEED_BULKDELETE: @@ -1138,6 +1237,7 @@ parallel_vacuum_process_one_index(ParallelVacuumState *pvs, Relation indrel, * touches different indexes. */ indstats->status = PARALLEL_INDVAC_STATUS_COMPLETED; + pg_atomic_fetch_add_u32(&(pvs->shared->completed_indexes), 1); /* Reset error traceback information */ pvs->status = PARALLEL_INDVAC_STATUS_COMPLETED; diff --git a/src/test/modules/test_autovacuum/t/001_parallel_autovacuum.pl b/src/test/modules/test_autovacuum/t/001_parallel_autovacuum.pl index 33c86bbdc9..ad84e66758 100644 --- a/src/test/modules/test_autovacuum/t/001_parallel_autovacuum.pl +++ b/src/test/modules/test_autovacuum/t/001_parallel_autovacuum.pl @@ -252,6 +252,169 @@ $node->safe_psql('postgres', "SELECT injection_points_wakeup('autovacuum-worker-cost-balanced')"); $node->safe_psql('postgres', "SELECT injection_points_detach('autovacuum-worker-cost-balanced')"); +$node->wait_for_log( + qr/automatic vacuum of table "postgres\.public\.test_autovac"/, + $log_offset); +$node->poll_query_until( + 'postgres', q{ + SELECT count(*) = 0 FROM pg_stat_activity + WHERE backend_type = 'autovacuum worker' AND datname = 'regress_db2' +}) or die "second autovacuum worker did not finish"; + +# Test 4: +# Check whether a config reload is serviced while the autovacuum leader waits +# for a parallel worker to finish an index. Hold the worker after it claims +# an index, so the leader can process the remaining indexes and enter +# ParallelFinish. +my $postgresoid = $node->safe_psql('postgres', + "SELECT oid FROM pg_database WHERE datname = 'postgres'"); +my $testautovacid = + $node->safe_psql('postgres', "SELECT 'test_autovac'::regclass::oid"); + +$node->safe_psql( + 'postgres', qq{ + ALTER SYSTEM SET autovacuum_max_workers = 1; + ALTER SYSTEM SET autovacuum_vacuum_cost_limit = 700; + ALTER SYSTEM SET autovacuum_vacuum_cost_delay = 0; + SELECT pg_reload_conf(); +}); + +prepare_for_next_test($node, 4); +$log_offset = -s $node->logfile; + +$node->safe_psql( + 'postgres', q{ + SELECT injection_points_attach('parallel-vacuum-worker-before-index', 'wait'); + SELECT injection_points_attach('parallel-vacuum-leader-before-index', 'wait'); +}); +$node->safe_psql('postgres', + 'ALTER TABLE test_autovac SET (autovacuum_enabled = true)'); +$node->wait_for_event('autovacuum worker', + 'parallel-vacuum-leader-before-index'); +$node->wait_for_event('parallel worker', + 'parallel-vacuum-worker-before-index'); +$node->safe_psql( + 'postgres', q{ + SELECT injection_points_wakeup('parallel-vacuum-leader-before-index'); + SELECT injection_points_detach('parallel-vacuum-leader-before-index'); +}); +$node->poll_query_until( + 'postgres', q{ + SELECT count(*) > 0 + FROM pg_stat_activity worker + JOIN pg_stat_activity leader ON leader.pid = worker.leader_pid + WHERE worker.backend_type = 'parallel worker' + AND worker.wait_event = 'parallel-vacuum-worker-before-index' + AND leader.wait_event = 'ParallelFinish' +}) or die "autovacuum leader did not enter ParallelFinish"; + +$node->safe_psql( + 'postgres', qq{ + ALTER SYSTEM SET autovacuum_vacuum_cost_limit = 800; + ALTER SYSTEM SET autovacuum_vacuum_cost_delay = 8; + ALTER SYSTEM SET vacuum_cost_page_miss = 11; + ALTER SYSTEM SET vacuum_cost_page_dirty = 12; + ALTER SYSTEM SET vacuum_cost_page_hit = 13; + SELECT pg_reload_conf(); +}); + +$node->wait_for_log( + qr/Autovacuum VacuumUpdateCosts\(db=$postgresoid, rel=$testautovacid, dobalance=yes, cost_limit=800, cost_delay=8 /, + $log_offset); + +$node->safe_psql('postgres', + "SELECT injection_points_wakeup('parallel-vacuum-worker-before-index')"); +$node->safe_psql('postgres', + "SELECT injection_points_detach('parallel-vacuum-worker-before-index')"); +$node->wait_for_log( + qr/parallel autovacuum worker updated cost params: cost_limit=800, cost_delay=8, cost_page_miss=11, cost_page_dirty=12, cost_page_hit=13/, + $log_offset); +$node->wait_for_log( + qr/automatic vacuum of table "postgres\.public\.test_autovac"/, + $log_offset); +ok(1, "config reload is propagated while the leader waits for workers"); + +# Test 5: +# Check the same wait path for an un-signalled cost-limit rebalance. A second +# autovacuum worker joins the balance while the first leader and its parallel +# worker remain held. +$node->safe_psql( + 'postgres', qq{ + ALTER SYSTEM SET autovacuum_max_workers = 2; + ALTER SYSTEM SET autovacuum_vacuum_cost_limit = 600; + SELECT pg_reload_conf(); +}); + +prepare_for_next_test($node, 5); +$node->safe_psql('regress_db2', + 'ALTER TABLE filler SET (autovacuum_enabled = false)'); +$node->safe_psql('regress_db2', 'UPDATE filler SET id = id + 1'); + +$log_offset = -s $node->logfile; +$node->safe_psql( + 'postgres', q{ + SELECT injection_points_attach('parallel-vacuum-worker-before-index', 'wait'); + SELECT injection_points_attach('parallel-vacuum-leader-before-index', 'wait'); +}); +$node->safe_psql('postgres', + 'ALTER TABLE test_autovac SET (autovacuum_enabled = true)'); +$node->wait_for_event('autovacuum worker', + 'parallel-vacuum-leader-before-index'); +$node->wait_for_event('parallel worker', + 'parallel-vacuum-worker-before-index'); +$node->safe_psql( + 'postgres', q{ + SELECT injection_points_wakeup('parallel-vacuum-leader-before-index'); + SELECT injection_points_detach('parallel-vacuum-leader-before-index'); +}); +$node->poll_query_until( + 'postgres', q{ + SELECT count(*) > 0 + FROM pg_stat_activity worker + JOIN pg_stat_activity leader ON leader.pid = worker.leader_pid + WHERE worker.backend_type = 'parallel worker' + AND worker.wait_event = 'parallel-vacuum-worker-before-index' + AND leader.wait_event = 'ParallelFinish' +}) or die "autovacuum leader did not enter ParallelFinish"; + +$node->safe_psql('postgres', + "SELECT injection_points_attach('autovacuum-worker-cost-balanced', 'wait')" +); +$node->safe_psql('regress_db2', + 'ALTER TABLE filler SET (autovacuum_enabled = true)'); +$node->wait_for_log( + qr/VacuumUpdateCosts\(db=$db2oid, rel=$filleroid, dobalance=yes, cost_limit=300,/, + $log_offset); +$node->safe_psql('postgres', + "SELECT injection_points_attach('parallel-vacuum-leader-cost-updated', 'notice')" +); +$node->wait_for_log( + qr/notice triggered for injection point parallel-vacuum-leader-cost-updated/, + $log_offset); +$node->safe_psql('postgres', + "SELECT injection_points_detach('parallel-vacuum-leader-cost-updated')"); + +$node->safe_psql('postgres', + "SELECT injection_points_wakeup('parallel-vacuum-worker-before-index')"); +$node->safe_psql('postgres', + "SELECT injection_points_detach('parallel-vacuum-worker-before-index')"); +$node->wait_for_log( + qr/parallel autovacuum worker updated cost params: cost_limit=300,/, + $log_offset); + +$node->safe_psql('postgres', + "SELECT injection_points_wakeup('autovacuum-worker-cost-balanced')"); +$node->safe_psql('postgres', + "SELECT injection_points_detach('autovacuum-worker-cost-balanced')"); +$node->wait_for_log( + qr/automatic vacuum of table "postgres\.public\.test_autovac"/, + $log_offset); +$node->poll_query_until( + 'postgres', q{ + SELECT count(*) = 0 FROM pg_stat_activity + WHERE backend_type = 'autovacuum worker' AND datname = 'regress_db2' +}) or die "second autovacuum worker did not finish"; +ok(1, "cost rebalance is propagated while the leader waits for workers"); $node->stop; done_testing();