Hi all,
I have been running recovery of Postgres clusters and small AWS
instances and
with a long history to replay, this quickly eats up the I/O budget for the
device. For example, with a gp3 with an IOPS limit of 3000 you quickly
eat it up
and get long stalls to recover the cluster. This is because the WAL
replay reads
the WAL in 8 KiB block rather than honoring io_combine_limit, which
typically
could allow reading 256 KiB in one go.
Since Postgres already have the option io_combine_limit to combine multiple
block reads into a single read, it is straightforward to honor it during
recovery using a simple cache (see attached patch). Since WAL replay is
mostly
from beginning to end, it is a simple way to reduce IOPS usage.
With this patch, and the default io_combine_limit of 16, I managed to
reduce the
number of I/O reads significantly:
| | reads | read_bytes | avg bytes/read |
|------------|--------|-------------|----------------|
| WAL before | 31,283 | 256,270,336 | 8,192 |
| WAL after | 1,955 | 256,131,072 | 131,013 |
Some back-of-the-envelope calculations give that with this patch we
should not be
able to saturate the device interface with the maximum read bandwidth
for gp3 is
approx 125 MiB/s, but the maximum read speed if you want to not hit the
limit
when you read 8 KiB blocks is 23-24 MiB/s.
I thought I'd share the patch and the results and hear if people think
this is a
useful improvement. There is already some work on pre-fetching, but this
is such
a simple patch that I think it could be worth to add to the current version.
(I also have some Perl scripts to set up crash recovery and some bpftrace
scripts that checks the reads, if anybody is interested, but I need to clean
them up a little.)
Best wishes,
Mats Kindahl, Multigres Engineer, Supabase
From 5ad1746c6850cdd81b3b90f7cbc9c7c55effa198 Mon Sep 17 00:00:00 2001
From: Mats Kindahl <[email protected]>
Date: Tue, 21 Jul 2026 09:49:58 +0200
Subject: Combine WAL segment reads during recovery into larger pread() calls
XLogPageRead() unconditionally read exactly one XLOG_BLCKSZ page per
pread() call regardless of how much more of the segment was already
known to be safely readable.
This means up to ~32x more I/O operations than necessary for sequential
WAL replay (crash recovery, archive recovery, standby WAL replay), since
a single pread() request could instead cover up to io_combine_limit
pages.
Add a read-ahead cache private to XLogPageRead() but keep the existing
usage of returning a single XLOG_BLCKSZ page for each call. On a miss,
read as many pages as are safe in one pread() into the cache (up to
io_combine_limit pages). On a hit, memcpy the requested page out of the
cache.
Also introduces a new XlogFileClose() helper that also invalidates the
cache and use that instead of explicitly closing the file each time.
---
src/backend/access/transam/xlogrecovery.c | 222 +++++++++++++++-------
1 file changed, 157 insertions(+), 65 deletions(-)
diff --git a/src/backend/access/transam/xlogrecovery.c b/src/backend/access/transam/xlogrecovery.c
index 5f3b065b894..e2c3f4bfd9a 100644
--- a/src/backend/access/transam/xlogrecovery.c
+++ b/src/backend/access/transam/xlogrecovery.c
@@ -52,6 +52,7 @@
#include "replication/slot.h"
#include "replication/slotsync.h"
#include "replication/walreceiver.h"
+#include "storage/bufmgr.h"
#include "storage/fd.h"
#include "storage/ipc.h"
#include "storage/latch.h"
@@ -63,6 +64,7 @@
#include "utils/fmgrprotos.h"
#include "utils/guc.h"
#include "utils/guc_hooks.h"
+#include "utils/memutils.h"
#include "utils/pgstat_internal.h"
#include "utils/pg_lsn.h"
#include "utils/ps_status.h"
@@ -238,6 +240,24 @@ static uint32 readOff = 0;
static uint32 readLen = 0;
static XLogSource readSource = XLOG_FROM_ANY;
+/*
+ * Read-ahead cache for WAL segment reads, allocated once by InitWalRecovery()
+ * and filled by XLogPageRead() with up to io_combine_limit XLOG_BLCKSZ pages in
+ * a single pread().
+ *
+ * This is private to XLogPageRead() and does not change the page read contract:
+ * callers still get back one XLOG_BLCKSZ page per call in the caller-supplied
+ * buffer.
+ *
+ * If readAheadLen == 0 means the cache is empty or invalid and must be reset to
+ * 0 at every point that already closes or otherwise invalidates readFile (see
+ * XLogReadAheadInvalidate()).
+ */
+static char *readAheadBuf = NULL;
+static XLogSegNo readAheadSegNo = 0;
+static uint32 readAheadOff = 0;
+static uint32 readAheadLen = 0;
+
/*
* Keeps track of which source we're currently reading from. This is
* different from readSource in that this is always set, even when we don't
@@ -372,6 +392,8 @@ static XLogRecord *ReadRecord(XLogPrefetcher *xlogprefetcher,
static int XLogPageRead(XLogReaderState *xlogreader, XLogRecPtr targetPagePtr,
int reqLen, XLogRecPtr targetRecPtr, char *readBuf);
+static void XLogFileClose(void);
+static void XLogReadAheadInvalidate(void);
static XLogPageReadResult WaitForWALToBecomeAvailable(XLogRecPtr RecPtr,
bool randAccess,
bool fetching_ckpt,
@@ -518,6 +540,13 @@ InitWalRecovery(ControlFileData *ControlFile, bool *wasShutdown_ptr,
*/
XLogReaderSetDecodeBuffer(xlogreader, NULL, wal_decode_buffer_size);
+ /*
+ * Allocate the read-ahead cache used by XLogPageRead() to combine several
+ * XLOG_BLCKSZ page reads into one larger pread().
+ */
+ Assert(readAheadBuf == NULL);
+ readAheadBuf = palloc(MAX_IO_COMBINE_LIMIT * XLOG_BLCKSZ);
+
/* Create a WAL prefetcher. */
xlogprefetcher = XLogPrefetcherAllocate(xlogreader);
@@ -1511,11 +1540,7 @@ FinishWalRecovery(void)
* If the ending log segment is still open, close it (to avoid
* problems on Windows with trying to rename or delete an open file).
*/
- if (readFile >= 0)
- {
- close(readFile);
- readFile = -1;
- }
+ XLogFileClose();
}
/*
@@ -1577,11 +1602,7 @@ ShutdownWalRecovery(void)
XLogPrefetcherComputeStats(xlogprefetcher);
/* Shut down xlogreader */
- if (readFile >= 0)
- {
- close(readFile);
- readFile = -1;
- }
+ XLogFileClose();
pfree(xlogreader->private_data);
XLogReaderFree(xlogreader);
XLogPrefetcherFree(xlogprefetcher);
@@ -3154,11 +3175,7 @@ ReadRecord(XLogPrefetcher *xlogprefetcher, int emode,
missingContrecPtr = xlogreader->missingContrecPtr;
}
- if (readFile >= 0)
- {
- close(readFile);
- readFile = -1;
- }
+ XLogFileClose();
/*
* We only end up here without a message when XLogPageRead()
@@ -3250,6 +3267,39 @@ ReadRecord(XLogPrefetcher *xlogprefetcher, int emode,
}
}
+/*
+ * Invalidate the WAL read-ahead cache.
+ *
+ * Must be called at every point that closes or otherwise invalidates readFile,
+ * but can be called at other times too and will force a re-read of the next
+ * page.
+ */
+static void
+XLogReadAheadInvalidate(void)
+{
+ readAheadLen = 0;
+}
+
+/*
+ * Close the current WAL file and invalidate the read-ahead cache.
+ *
+ * It is strictly speaking not necessary to set readSource and readLen to
+ * their initial values, but the original code did that, so we keep it for
+ * now.
+ */
+static void
+XLogFileClose(void)
+{
+ if (readFile >= 0)
+ {
+ close(readFile);
+ readFile = -1;
+ readSource = XLOG_FROM_ANY;
+ readLen = 0;
+ XLogReadAheadInvalidate();
+ }
+}
+
/*
* Read the XLOG page containing targetPagePtr into readBuf (if not read
* already). Returns number of bytes read, if the page is read successfully,
@@ -3287,7 +3337,6 @@ XLogPageRead(XLogReaderState *xlogreader, XLogRecPtr targetPagePtr, int reqLen,
int emode = private->emode;
uint32 targetPageOff;
XLogSegNo targetSegNo PG_USED_FOR_ASSERTS_ONLY;
- ssize_t r;
instr_time io_start;
Assert(AmStartupProcess() || !IsUnderPostmaster);
@@ -3316,9 +3365,7 @@ XLogPageRead(XLogReaderState *xlogreader, XLogRecPtr targetPagePtr, int reqLen,
}
}
- close(readFile);
- readFile = -1;
- readSource = XLOG_FROM_ANY;
+ XLogFileClose();
}
XLByteToSeg(targetPagePtr, readSegNo, wal_segment_size);
@@ -3346,11 +3393,7 @@ retry:
case XLREAD_WOULDBLOCK:
return XLREAD_WOULDBLOCK;
case XLREAD_FAIL:
- if (readFile >= 0)
- close(readFile);
- readFile = -1;
- readLen = 0;
- readSource = XLOG_FROM_ANY;
+ XLogFileClose();
return XLREAD_FAIL;
case XLREAD_SUCCESS:
break;
@@ -3383,45 +3426,101 @@ retry:
/* Read the requested page */
readOff = targetPageOff;
- /* Measure I/O timing when reading segment */
- io_start = pgstat_prepare_io_time(track_wal_io_timing);
-
- pgstat_report_wait_start(WAIT_EVENT_WAL_READ);
- r = pg_pread(readFile, readBuf, XLOG_BLCKSZ, (pgoff_t) readOff);
- if (r != XLOG_BLCKSZ)
+ /*
+ * If the requested page is not in the read-ahead cache from a previous,
+ * larger read, the cache is invalid, or all blocks from the cache is
+ * read; fetch a large block from disk (up to io_combine_limit pages).
+ *
+ * The cache is invalidated (see XLogReadAheadInvalidate()) at every point
+ * that could make readSegNo mean a different underlying file, so a hit
+ * here is always safe to trust.
+ */
+ if (readAheadLen <= 0 ||
+ readAheadSegNo != readSegNo ||
+ readOff < readAheadOff ||
+ readOff + XLOG_BLCKSZ > readAheadOff + readAheadLen)
{
- char fname[MAXFNAMELEN];
- int save_errno = errno;
+ uint32 wanted;
+ ssize_t r;
- pgstat_report_wait_end();
+ Assert(readAheadBuf != NULL);
+
+ /* Read at most to the end of the WAL segment and not beyond */
+ wanted = (uint32) wal_segment_size - targetPageOff;
+
+ /*
+ * If we stream from the primary, we cannot read more than we have in
+ * the buffer, so cap the read at the number of bytes in the buffer.
+ */
+ if (readSource == XLOG_FROM_STREAM)
+ wanted = Min(wanted, readLen);
- /* Count I/O stats only for successful short reads */
- if (r > 0)
- pgstat_count_io_op_time(IOOBJECT_WAL, IOCONTEXT_NORMAL, IOOP_READ,
- io_start, 1, r);
+ /*
+ * Never read more blocks than what io_combine_limit says we should
+ * read
+ */
+ wanted = Min(wanted, io_combine_limit * XLOG_BLCKSZ);
+
+ /* Read at least one page */
+ wanted = Max(wanted, XLOG_BLCKSZ);
+
+ /* Measure I/O timing when reading segment */
+ io_start = pgstat_prepare_io_time(track_wal_io_timing);
+
+ pgstat_report_wait_start(WAIT_EVENT_WAL_READ);
+ r = pg_pread(readFile, readAheadBuf, wanted, (pgoff_t) readOff);
+ pgstat_report_wait_end();
- XLogFileName(fname, curFileTLI, readSegNo, wal_segment_size);
- if (r < 0)
+ if (r < XLOG_BLCKSZ)
{
- errno = save_errno;
- ereport(emode_for_corrupt_record(emode, targetPagePtr + reqLen),
- (errcode_for_file_access(),
- errmsg("could not read from WAL segment %s, LSN %X/%08X, offset %u: %m",
- fname, LSN_FORMAT_ARGS(targetPagePtr),
- readOff)));
+ char fname[MAXFNAMELEN];
+ int save_errno = errno;
+
+ /* Count I/O stats only for successful short reads */
+ if (r > 0)
+ pgstat_count_io_op_time(IOOBJECT_WAL, IOCONTEXT_NORMAL, IOOP_READ,
+ io_start, 1, r);
+
+ XLogFileName(fname, curFileTLI, readSegNo, wal_segment_size);
+ if (r < 0)
+ {
+ errno = save_errno;
+ ereport(emode_for_corrupt_record(emode, targetPagePtr + reqLen),
+ (errcode_for_file_access(),
+ errmsg("could not read from WAL segment %s, LSN %X/%08X, offset %u: %m",
+ fname, LSN_FORMAT_ARGS(targetPagePtr),
+ readOff)));
+ }
+ else
+ ereport(emode_for_corrupt_record(emode, targetPagePtr + reqLen),
+ (errcode(ERRCODE_DATA_CORRUPTED),
+ errmsg("could not read from WAL segment %s, LSN %X/%08X, offset %u: read %zd of %zu",
+ fname, LSN_FORMAT_ARGS(targetPagePtr),
+ readOff, r, (Size) wanted)));
+
+ /*
+ * Need to invalidate the cache so that next read will read in a
+ * complete page. The current page is not complete, so trying to
+ * use it in a later call might cause a corrupted page to be used.
+ */
+ XLogReadAheadInvalidate();
+ goto next_record_is_invalid;
}
- else
- ereport(emode_for_corrupt_record(emode, targetPagePtr + reqLen),
- (errcode(ERRCODE_DATA_CORRUPTED),
- errmsg("could not read from WAL segment %s, LSN %X/%08X, offset %u: read %zd of %zu",
- fname, LSN_FORMAT_ARGS(targetPagePtr),
- readOff, r, (Size) XLOG_BLCKSZ)));
- goto next_record_is_invalid;
+
+ pgstat_count_io_op_time(IOOBJECT_WAL, IOCONTEXT_NORMAL, IOOP_READ,
+ io_start, 1, r);
+
+ readAheadSegNo = readSegNo;
+ readAheadOff = readOff;
+ readAheadLen = r;
}
- pgstat_report_wait_end();
- pgstat_count_io_op_time(IOOBJECT_WAL, IOCONTEXT_NORMAL, IOOP_READ,
- io_start, 1, r);
+ /*
+ * We copy out the page from the cache since readBuf is allocated
+ * elsewhere. It was considered pointing into the cache, but that seems
+ * like an unnecessary risk and more work than it's worth.
+ */
+ memcpy(readBuf, readAheadBuf + (readOff - readAheadOff), XLOG_BLCKSZ);
Assert(targetSegNo == readSegNo);
Assert(targetPageOff == readOff);
@@ -3491,11 +3590,7 @@ next_record_is_invalid:
lastSourceFailed = true;
- if (readFile >= 0)
- close(readFile);
- readFile = -1;
- readLen = 0;
- readSource = XLOG_FROM_ANY;
+ XLogFileClose();
/* In standby-mode, keep trying */
if (StandbyMode)
@@ -3773,11 +3868,8 @@ WaitForWALToBecomeAvailable(XLogRecPtr RecPtr, bool randAccess,
Assert(!WalRcvStreaming());
/* Close any old file we might have open. */
- if (readFile >= 0)
- {
- close(readFile);
- readFile = -1;
- }
+ XLogFileClose();
+
/* Reset curFileTLI if random fetch. */
if (randAccess)
curFileTLI = 0;
--
2.43.0