From 6e7d8a2bf027ffbbda1287264b3e88dc64917dea Mon Sep 17 00:00:00 2001 From: Nik Samokhvalov Date: Mon, 21 Sep 2026 08:57:49 -0700 Subject: [PATCH] Refresh autovacuum costs while waiting for parallel workers --- src/backend/access/transam/parallel.c | 37 +++- src/backend/commands/vacuumparallel.c | 47 ++++- src/include/access/parallel.h | 4 + .../t/001_parallel_autovacuum.pl | 163 ++++++++++++++++++ 4 files changed, 246 insertions(+), 5 deletions(-) diff --git a/src/backend/access/transam/parallel.c b/src/backend/access/transam/parallel.c index 89e9d224ee..b222588e66 100644 --- a/src/backend/access/transam/parallel.c +++ b/src/backend/access/transam/parallel.c @@ -801,14 +801,17 @@ WaitForParallelWorkersToAttach(ParallelContext *pcxt) * Also, we want to update our notion of XactLastRecEnd based on worker * feedback. */ -void -WaitForParallelWorkersToFinish(ParallelContext *pcxt) +static void +WaitForParallelWorkersToFinishInternal(ParallelContext *pcxt, + ParallelWorkerWaitCallback callback, + void *arg, long wait_interval) { for (;;) { bool anyone_alive = false; int nfinished = 0; int i; + uint32 wait_events = WL_LATCH_SET | WL_EXIT_ON_PM_DEATH; /* * This will process any parallel messages that are pending, which may @@ -816,6 +819,8 @@ WaitForParallelWorkersToFinish(ParallelContext *pcxt) * error propagated from a worker. */ CHECK_FOR_INTERRUPTS(); + if (callback != NULL) + callback(arg); for (i = 0; i < pcxt->nworkers_launched; ++i) { @@ -892,7 +897,11 @@ WaitForParallelWorkersToFinish(ParallelContext *pcxt) } } - (void) WaitLatch(MyLatch, WL_LATCH_SET | WL_EXIT_ON_PM_DEATH, -1, + if (wait_interval >= 0) + wait_events |= WL_TIMEOUT; + + (void) WaitLatch(MyLatch, wait_events, + wait_interval, WAIT_EVENT_PARALLEL_FINISH); ResetLatch(MyLatch); } @@ -907,6 +916,28 @@ WaitForParallelWorkersToFinish(ParallelContext *pcxt) } } +void +WaitForParallelWorkersToFinish(ParallelContext *pcxt) +{ + WaitForParallelWorkersToFinishInternal(pcxt, NULL, NULL, -1); +} + +/* + * As above, but run callback before checking worker state and after each + * wakeup. While workers remain active, wake at least once every + * wait_interval milliseconds. + */ +void +WaitForParallelWorkersToFinishWithCallback(ParallelContext *pcxt, + ParallelWorkerWaitCallback callback, + void *arg, long wait_interval) +{ + Assert(callback != NULL); + Assert(wait_interval > 0); + + WaitForParallelWorkersToFinishInternal(pcxt, callback, arg, wait_interval); +} + /* * Wait for all workers to exit. * diff --git a/src/backend/commands/vacuumparallel.c b/src/backend/commands/vacuumparallel.c index 767d162e57..b5d30f2152 100644 --- a/src/backend/commands/vacuumparallel.c +++ b/src/backend/commands/vacuumparallel.c @@ -43,9 +43,11 @@ #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" @@ -60,6 +62,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. @@ -285,6 +289,7 @@ 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 *arg); 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, @@ -929,6 +934,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,8 +950,14 @@ parallel_vacuum_process_all_indexes(ParallelVacuumState *pvs, int num_index_scan */ if (nworkers > 0) { - /* Wait for all vacuum workers to finish */ - WaitForParallelWorkersToFinish(pvs->pcxt); + /* Wait for all vacuum workers to finish. */ + if (AmAutoVacuumWorkerProcess()) + WaitForParallelWorkersToFinishWithCallback(pvs->pcxt, + parallel_vacuum_update_leader_cost_params, + NULL, + PARALLEL_VACUUM_COST_UPDATE_INTERVAL_MS); + else + WaitForParallelWorkersToFinish(pvs->pcxt); for (int i = 0; i < pvs->pcxt->nworkers_launched; i++) InstrAccumParallelQuery(&pvs->buffer_usage[i], &pvs->wal_usage[i]); @@ -974,6 +989,29 @@ 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 *arg) +{ + Assert(AmAutoVacuumWorkerProcess()); + Assert(arg == NULL); + + if (ConfigReloadPending) + { + ConfigReloadPending = false; + ProcessConfigFile(PGC_SIGHUP); + VacuumUpdateCosts(); + } + else + AutoVacuumUpdateCostLimit(); + + parallel_vacuum_propagate_shared_delay_params(); + INJECTION_POINT("parallel-vacuum-leader-cost-updated", NULL); +} + /* * Index vacuum/cleanup routine used by the leader process and parallel * vacuum worker processes to vacuum the indexes in parallel. @@ -1097,6 +1135,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: diff --git a/src/include/access/parallel.h b/src/include/access/parallel.h index 60f857675e..4ca5a80b1a 100644 --- a/src/include/access/parallel.h +++ b/src/include/access/parallel.h @@ -23,6 +23,7 @@ #include "storage/shm_toc.h" typedef void (*parallel_worker_main_type) (dsm_segment *seg, shm_toc *toc); +typedef void (*ParallelWorkerWaitCallback) (void *arg); typedef struct ParallelWorkerInfo { @@ -69,6 +70,9 @@ extern void ReinitializeParallelWorkers(ParallelContext *pcxt, int nworkers_to_l extern void LaunchParallelWorkers(ParallelContext *pcxt); extern void WaitForParallelWorkersToAttach(ParallelContext *pcxt); extern void WaitForParallelWorkersToFinish(ParallelContext *pcxt); +extern void WaitForParallelWorkersToFinishWithCallback(ParallelContext *pcxt, + ParallelWorkerWaitCallback callback, + void *arg, long wait_interval); extern void DestroyParallelContext(ParallelContext *pcxt); extern bool ParallelContextActive(void); 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(); base-commit: b368bdd230181c60e085267fe65d42fd372d3a24 -- 2.50.1 (Apple Git-155)