From 582533d27934914e99c6a5367d7431fe2081c01f Mon Sep 17 00:00:00 2001 From: Jakub Wartak Date: Tue, 18 Aug 2026 10:36:55 +0200 Subject: [PATCH v06_18082026] Use atomics for speed instead of spinlocks in xlogpipeline --- src/backend/access/transam/xlogpipeline.c | 119 ++++++++++------------ src/backend/access/transam/xlogrecovery.c | 18 +--- src/include/access/xlogpipeline.h | 29 +++--- 3 files changed, 70 insertions(+), 96 deletions(-) diff --git a/src/backend/access/transam/xlogpipeline.c b/src/backend/access/transam/xlogpipeline.c index 26acb84edf1..f81df19c420 100644 --- a/src/backend/access/transam/xlogpipeline.c +++ b/src/backend/access/transam/xlogpipeline.c @@ -88,6 +88,12 @@ static dsm_segment *consumer_dsm_seg = NULL; static shm_mq *consumer_mq = NULL; static shm_mq_handle *consumer_mq_handle = NULL; +/* Private walpipeline counters, updated with atomic stores */ +static uint64 producer_records_sent = 0; +static uint64 producer_bytes_sent = 0; +static uint64 consumer_records_received = 0; +static uint64 consumer_bytes_received = 0; + /* * Local buffer containing msg header that will be sent together with the * decoded data, to the msg queue @@ -243,6 +249,17 @@ WalPipelineShmemInit(void *arg) memset(WalPipelineShm, 0, sizeof(WalPipelineShmCtl)); SpinLockInit(&WalPipelineShm->mutex); + + pg_atomic_init_u32(&WalPipelineShm->shutdown_requested, 0); + pg_atomic_init_u32(&WalPipelineShm->producerWaiting, 0); + + pg_atomic_init_u64(&WalPipelineShm->producer_lsn, 0); + pg_atomic_init_u64(&WalPipelineShm->records_sent, 0); + pg_atomic_init_u64(&WalPipelineShm->bytes_sent, 0); + pg_atomic_init_u64(&WalPipelineShm->consumer_lsn, 0); + pg_atomic_init_u64(&WalPipelineShm->applied_lsn, 0); + pg_atomic_init_u64(&WalPipelineShm->records_received, 0); + pg_atomic_init_u64(&WalPipelineShm->bytes_received, 0); } /* @@ -368,9 +385,7 @@ WalPipeline_RequestShutdown(void) if (!WalPipelineShm) return; - SpinLockAcquire(&WalPipelineShm->mutex); - WalPipelineShm->shutdown_requested = true; - SpinLockRelease(&WalPipelineShm->mutex); + pg_atomic_write_u32(&WalPipelineShm->shutdown_requested, 1); } /* @@ -528,14 +543,8 @@ WalPipeline_ProducerMain(Datum main_arg) /* Main decoding loop */ while (true) { - bool shutdown_requested; - - /* Check if consumer requested to stop decoding */ - SpinLockAcquire(&WalPipelineShm->mutex); - shutdown_requested = WalPipelineShm->shutdown_requested; - SpinLockRelease(&WalPipelineShm->mutex); - - if (shutdown_requested) + /* Check if the consumer requested to stop the decoding process (us) */ + if (pg_atomic_read_u32(&WalPipelineShm->shutdown_requested)) { elog(DEBUG1, "[walpipeline] producer: got shutdown request from the consumer."); break; @@ -560,10 +569,8 @@ WalPipeline_ProducerMain(Datum main_arg) break; } - /* Update our position for monitoring */ - SpinLockAcquire(&WalPipelineShm->mutex); - WalPipelineShm->producer_lsn = xlogreader->EndRecPtr; - SpinLockRelease(&WalPipelineShm->mutex); + /* Announce our position */ + pg_atomic_write_u64(&WalPipelineShm->producer_lsn, xlogreader->EndRecPtr); CHECK_FOR_INTERRUPTS(); } @@ -579,10 +586,8 @@ WalPipeline_ProducerMain(Datum main_arg) /* compute and set rusage delta */ pipeline_save_rusage(&ru0); - SpinLockAcquire(&WalPipelineShm->mutex); - records_sent = WalPipelineShm->records_sent; - records_received = WalPipelineShm->records_received; - SpinLockRelease(&WalPipelineShm->mutex); + records_sent = pg_atomic_read_u64(&WalPipelineShm->records_sent); + records_received = pg_atomic_read_u64(&WalPipelineShm->records_received); elog(LOG, "[walpipeline] producer: exiting: sent=" UINT64_FORMAT " received=" UINT64_FORMAT, records_sent, records_received); @@ -659,7 +664,9 @@ WalPipeline_SendRecord(XLogReaderState *record) /* Serialize the decoded data */ msglen = serialize_wal_record(record, iov); - res = shm_mq_sendv(producer_mq_handle, iov, 2, false, true); + /* TODO: force_flush=false avoids kill() flood, however from it may be necessary? */ + //res = shm_mq_sendv(producer_mq_handle, iov, 2, false, true); + res = shm_mq_sendv(producer_mq_handle, iov, 2, false, false); /* * Reset the offsets to exact ptrs because decoded record data still @@ -669,10 +676,11 @@ WalPipeline_SendRecord(XLogReaderState *record) if (res == SHM_MQ_SUCCESS) { - SpinLockAcquire(&WalPipelineShm->mutex); - WalPipelineShm->records_sent++; - WalPipelineShm->bytes_sent += msglen; - SpinLockRelease(&WalPipelineShm->mutex); + /* Publish stats lock-free; only the producer writes these */ + producer_records_sent++; + producer_bytes_sent += msglen; + pg_atomic_write_u64(&WalPipelineShm->records_sent, producer_records_sent); + pg_atomic_write_u64(&WalPipelineShm->bytes_sent, producer_bytes_sent); return true; } @@ -738,12 +746,12 @@ WalPipeline_ReceiveRecord(XLogReaderState *startup_reader) case WAL_MSG_RECORD: record = deserialize_wal_record((char *) data, nbytes, startup_reader); - /* Update statistics */ - SpinLockAcquire(&WalPipelineShm->mutex); - WalPipelineShm->records_received++; - WalPipelineShm->bytes_received += nbytes; - WalPipelineShm->consumer_lsn = hdr->endRecPtr; - SpinLockRelease(&WalPipelineShm->mutex); + /* Publish stats */ + consumer_records_received++; + consumer_bytes_received += nbytes; + pg_atomic_write_u64(&WalPipelineShm->records_received, consumer_records_received); + pg_atomic_write_u64(&WalPipelineShm->bytes_received, consumer_bytes_received); + pg_atomic_write_u64(&WalPipelineShm->consumer_lsn, hdr->endRecPtr); return record; @@ -795,14 +803,17 @@ bool WalPipeline_IsActive(void) { bool active; + uint32_t local_shutdown; if (!WalPipelineShm) return false; SpinLockAcquire(&WalPipelineShm->mutex); - active = WalPipelineShm->initialized && !WalPipelineShm->shutdown_requested; + active = WalPipelineShm->initialized; SpinLockRelease(&WalPipelineShm->mutex); + local_shutdown = pg_atomic_read_u32(&WalPipelineShm->shutdown_requested); + active = active && !local_shutdown; return active; } @@ -838,13 +849,7 @@ WalPipeline_WaitForConsumerShutdownRequest(void) while (true) { - bool shutdown_requested; - - SpinLockAcquire(&WalPipelineShm->mutex); - shutdown_requested = WalPipelineShm->shutdown_requested; - SpinLockRelease(&WalPipelineShm->mutex); - - if (shutdown_requested) + if (pg_atomic_read_u32(&WalPipelineShm->shutdown_requested)) break; if (++iters >= MAX_SHUTDOWN_WAIT_ITERS) @@ -875,10 +880,8 @@ WalPipeline_WaitForConsumerCatchup(void) for (;;) { - SpinLockAcquire(&WalPipelineShm->mutex); - producer_lsn = WalPipelineShm->producer_lsn; - consumer_lsn = WalPipelineShm->applied_lsn; - SpinLockRelease(&WalPipelineShm->mutex); + producer_lsn = pg_atomic_read_u64(&WalPipelineShm->producer_lsn); + consumer_lsn = pg_atomic_read_u64(&WalPipelineShm->applied_lsn); if (producer_lsn == consumer_lsn) return; @@ -898,18 +901,14 @@ void WalPipeline_GetStats(uint64 *records_sent, uint64 *records_received, XLogRecPtr *producer_lsn, XLogRecPtr *consumer_lsn) { - SpinLockAcquire(&WalPipelineShm->mutex); - if (records_sent) - *records_sent = WalPipelineShm->records_sent; + *records_sent = pg_atomic_read_u64(&WalPipelineShm->records_sent); if (records_received) - *records_received = WalPipelineShm->records_received; + *records_received = pg_atomic_read_u64(&WalPipelineShm->records_received); if (producer_lsn) - *producer_lsn = WalPipelineShm->producer_lsn; + *producer_lsn = pg_atomic_read_u64(&WalPipelineShm->producer_lsn); if (consumer_lsn) - *consumer_lsn = WalPipelineShm->consumer_lsn; - - SpinLockRelease(&WalPipelineShm->mutex); + *consumer_lsn = pg_atomic_read_u64(&WalPipelineShm->consumer_lsn); } @@ -1020,21 +1019,13 @@ bool AmWalPipeline(void) void SetProducerStartWaiting(void) { if (wal_pipeline_enabled && AmWalPipeline()) - { - SpinLockAcquire(&WalPipelineShm->mutex); - WalPipelineShm->producerWaiting = true; - SpinLockRelease(&WalPipelineShm->mutex); - } + pg_atomic_write_u32(&WalPipelineShm->producerWaiting, 1); } void SetProducerDoneWaiting(void) { if (wal_pipeline_enabled && AmWalPipeline()) - { - SpinLockAcquire(&WalPipelineShm->mutex); - WalPipelineShm->producerWaiting = false; - SpinLockRelease(&WalPipelineShm->mutex); - } + pg_atomic_write_u32(&WalPipelineShm->producerWaiting, 0); } /* @@ -1185,19 +1176,13 @@ void ProcessPipelineBgwInterrupts(void) { - bool shutdown_requested; - if (got_SIGHUP) { got_SIGHUP = false; PipelineRereadConfig(); } - SpinLockAcquire(&WalPipelineShm->mutex); - shutdown_requested = WalPipelineShm->shutdown_requested; - SpinLockRelease(&WalPipelineShm->mutex); - - if (shutdown_requested) + if (pg_atomic_read_u32(&WalPipelineShm->shutdown_requested)) proc_exit(0); CHECK_FOR_INTERRUPTS(); diff --git a/src/backend/access/transam/xlogrecovery.c b/src/backend/access/transam/xlogrecovery.c index 92d61b56da9..6a3a78c9333 100644 --- a/src/backend/access/transam/xlogrecovery.c +++ b/src/backend/access/transam/xlogrecovery.c @@ -2253,12 +2253,8 @@ ApplyWalRecord(XLogReaderState *xlogreader, XLogRecord *record, TimeLineID *repl if (wal_pipeline_enabled) { - bool idle = false; - - SpinLockAcquire(&WalPipelineShm->mutex); - - /* Keep track of consumer location */ - WalPipelineShm->applied_lsn = xlogreader->EndRecPtr; + /* Update consumer location */ + pg_atomic_write_u64(&WalPipelineShm->applied_lsn, xlogreader->EndRecPtr); /* * Check if producer is started waiting for more wal and current queue @@ -2267,13 +2263,9 @@ ApplyWalRecord(XLogReaderState *xlogreader, XLogRecord *record, TimeLineID *repl * this maintanance task in the pipeline worker so better to call it * here via flagging. */ - if ((WalPipelineShm->records_received == WalPipelineShm->records_sent) - && WalPipelineShm->producerWaiting) - idle = true; - - SpinLockRelease(&WalPipelineShm->mutex); - - if (idle) + if (pg_atomic_read_u32(&WalPipelineShm->producerWaiting) && + (pg_atomic_read_u64(&WalPipelineShm->records_received) == + pg_atomic_read_u64(&WalPipelineShm->records_sent))) KnownAssignedTransactionIdsIdleMaintenance(); } } diff --git a/src/include/access/xlogpipeline.h b/src/include/access/xlogpipeline.h index bdf2808ea59..d1e6f1b72c0 100644 --- a/src/include/access/xlogpipeline.h +++ b/src/include/access/xlogpipeline.h @@ -25,6 +25,7 @@ #include "access/xlogreader.h" #include "access/xlogrecovery.h" #include "access/xlogutils.h" +#include "port/atomics.h" #include "storage/dsm.h" #include "storage/shm_mq.h" #include "storage/spin.h" @@ -115,34 +116,30 @@ typedef struct WalPipelineParams */ typedef struct WalPipelineShmCtl { - /* Lifecycle management */ slock_t mutex; bool initialized; - bool shutdown_requested; - bool producerWaiting; - - /* Producer state */ pid_t producer_pid; - XLogRecPtr producer_lsn; /* Last LSN read by producer */ - - /* Consumer state */ pid_t consumer_pid; - XLogRecPtr consumer_lsn; /* Last LSN recieved by consumer */ - XLogRecPtr applied_lsn; /* Last LSN applied by consumer */ /* Queue handles */ dsm_handle dsm_seg_handle; shm_mq_handle *producer_mq_handle; shm_mq_handle *consumer_mq_handle; - /* Statistics */ - uint64 records_sent; - uint64 records_received; - uint64 bytes_sent; - uint64 bytes_received; - /* cpu usage delta by the producer */ PGRUsageDelta producer_rusage; + + pg_atomic_uint32 shutdown_requested pg_attribute_aligned(PG_CACHE_LINE_SIZE); + pg_atomic_uint32 producerWaiting pg_attribute_aligned(PG_CACHE_LINE_SIZE); + + pg_atomic_uint64 producer_lsn pg_attribute_aligned(PG_CACHE_LINE_SIZE); + pg_atomic_uint64 consumer_lsn pg_attribute_aligned(PG_CACHE_LINE_SIZE); + pg_atomic_uint64 records_sent pg_attribute_aligned(PG_CACHE_LINE_SIZE); + pg_atomic_uint64 bytes_sent pg_attribute_aligned(PG_CACHE_LINE_SIZE); + pg_atomic_uint64 applied_lsn pg_attribute_aligned(PG_CACHE_LINE_SIZE); + pg_atomic_uint64 records_received pg_attribute_aligned(PG_CACHE_LINE_SIZE); + pg_atomic_uint64 bytes_received pg_attribute_aligned(PG_CACHE_LINE_SIZE); + } WalPipelineShmCtl; /* consumer may have to compute prefetecher stats */ -- 2.43.0