From 7a1aebfeac9b98c314a839d6ac412bfc6591278a Mon Sep 17 00:00:00 2001 From: Peter Geoghegan Date: Sun, 6 Sep 2026 13:20:05 -0400 Subject: [PATCH v37 09/11] WIP: Add read stream request deduplication. Add a stream flag under which the stream remembers the block numbers it queued most recently (a ring of 32) and does not read or pin a repeat of one of them in advance. A repeat entry records the block's place in the sequence, ordered relative to the queue of buffers by a gap count, and the stream pins the block afresh when that place is reached. Repeat entries do not count against the look-ahead distance. --- src/include/storage/read_stream.h | 14 + src/backend/access/heap/heapam_indexscan.c | 2 +- src/backend/storage/aio/read_stream.c | 338 +++++++++++++++++- .../modules/test_aio/t/004_read_stream.pl | 55 +++ src/test/modules/test_aio/test_aio--1.0.sql | 2 +- src/test/modules/test_aio/test_aio.c | 24 +- src/tools/pgindent/typedefs.list | 1 + 7 files changed, 424 insertions(+), 12 deletions(-) diff --git a/src/include/storage/read_stream.h b/src/include/storage/read_stream.h index e2dcf1be5..7f6049e9f 100644 --- a/src/include/storage/read_stream.h +++ b/src/include/storage/read_stream.h @@ -63,6 +63,20 @@ */ #define READ_STREAM_USE_BATCHING 0x08 +/* --- + * Opt-in to deduplicating recently streamed blocks. + * + * The callback may return a block number that it returned a short time + * before. With this flag, the stream remembers the block numbers it queued + * most recently and does not read or pin them in advance a second time. + * + * Only for streams created with a per_buffer_data_size of 0. A repeated block + * is pinned when its turn comes, and nothing holds per-buffer data for it + * until then. + * --- + */ +#define READ_STREAM_DEDUP_RECENT 0x10 + struct ReadStream; typedef struct ReadStream ReadStream; struct IOStats; diff --git a/src/backend/access/heap/heapam_indexscan.c b/src/backend/access/heap/heapam_indexscan.c index b21323b5e..3489a7497 100644 --- a/src/backend/access/heap/heapam_indexscan.c +++ b/src/backend/access/heap/heapam_indexscan.c @@ -1300,7 +1300,7 @@ heapam_index_start_prefetching(IndexScanDesc scan, IndexScanHeapData *hscan) Assert(!scan->batchringbuf.prefetchPos.valid); hscan->xs_read_stream = - read_stream_begin_relation(READ_STREAM_DEFAULT, NULL, + read_stream_begin_relation(READ_STREAM_DEDUP_RECENT, NULL, scan->heapRelation, MAIN_FORKNUM, heapam_index_prefetch_next_block, scan, 0); diff --git a/src/backend/storage/aio/read_stream.c b/src/backend/storage/aio/read_stream.c index 5ab82ffb8..ab7e3864e 100644 --- a/src/backend/storage/aio/read_stream.c +++ b/src/backend/storage/aio/read_stream.c @@ -62,6 +62,15 @@ * the range 42..44 requires an I/O wait before its buffers are returned, as * does block 60. * + * A stream created with READ_STREAM_DEDUP_RECENT also remembers the block + * numbers it queued most recently. When the callback returns one of them + * again, no buffer is queued for it; instead a repeat entry records its place + * in the sequence, and the block is pinned afresh when that place is reached. + * Repeat entries hold no pins and no I/Os, and the pending read is free to + * keep growing past them. Such a stream must be created with a + * per_buffer_data_size of 0, as a repeat entry has nowhere to keep per-buffer + * data between the callback and the hand-out. + * * * Portions Copyright (c) 2024-2026, PostgreSQL Global Development Group * Portions Copyright (c) 1994, Regents of the University of California @@ -80,6 +89,7 @@ #include "storage/smgr.h" #include "storage/read_stream.h" #include "utils/memdebug.h" +#include "utils/memutils.h" #include "utils/rel.h" #include "utils/spccache.h" @@ -89,6 +99,17 @@ typedef struct InProgressIO ReadBuffersOperation op; } InProgressIO; +typedef struct ReadStreamRepeat +{ + BlockNumber blocknum; + uint32 seq; +} ReadStreamRepeat; + +#define DEDUP_RING_SIZE 32 +#define DEDUP_REPEATS_PER_BUFFER 8 +#define DEDUP_MAX_REPEATS 2048 +#define DEDUP_INITIAL_REPEATS 64 + /* * State for managing a stream of reads. */ @@ -122,6 +143,7 @@ struct ReadStream int read_buffers_flags; bool sync_mode; /* using io_method=sync */ bool batch_mode; /* READ_STREAM_USE_BATCHING */ + bool dedup; /* READ_STREAM_DEDUP_RECENT */ bool advice_enabled; bool temporary; @@ -160,6 +182,32 @@ struct ReadStream bool fast_path; + /* + * READ_STREAM_DEDUP_RECENT queued blocks. + * + * recent_blocknums is a ring of the block numbers most recently queued + * for reading, unused slots holding InvalidBlockNumber. repeats is a + * circular queue of repeat entries, ordered relative to the queue of + * buffers by queued_buffers, a count of the buffers queued for reading so + * far: a repeat entry records the count at the time it was queued, and is + * handed out, before the oldest buffer, once the buffers handed out so + * far (queued_buffers less those still pinned or pending) reach it. + * + * repeats is allocated when the first repeat entry is queued and grows + * geometrically up to repeat_cap, so that a stream whose callback never + * returns a recently queued block pays nothing for it. max_repeats is + * the capacity it has now, and is 0 until it is first allocated. + */ + BlockNumber *recent_blocknums; + int16 recent_index; /* next ring slot to overwrite */ + ReadStreamRepeat *repeats; + int16 max_repeats; + int16 repeat_cap; + int16 nrepeats; + int16 oldest_repeat_index; + int16 next_repeat_index; + uint32 queued_buffers; + /* Circular queue of buffers. */ int16 oldest_buffer_index; /* Next pinned buffer to return */ int16 next_buffer_index; /* Index of next buffer to pin */ @@ -176,6 +224,147 @@ get_per_buffer_data(ReadStream *stream, int16 buffer_index) stream->per_buffer_data_size * buffer_index; } +/* + * Is blocknum one of the block numbers most recently queued for reading? + */ +static inline bool +read_stream_recent_block(ReadStream *stream, BlockNumber blocknum) +{ + int found = 0; + + Assert(stream->dedup); + + /* + * Accumulate matches rather than returning as soon as we find one, so + * that the compiler can vectorize + */ + for (int i = 0; i < DEDUP_RING_SIZE; i++) + found |= (stream->recent_blocknums[i] == blocknum); + + return found != 0; +} + +/* + * Remember that blocknum has just been queued for reading, replacing the + * oldest entry in the ring of recent block numbers. + */ +static inline void +read_stream_remember_block(ReadStream *stream, BlockNumber blocknum) +{ + StaticAssertStmt((DEDUP_RING_SIZE & (DEDUP_RING_SIZE - 1)) == 0, + "DEDUP_RING_SIZE must be a power of 2"); + + Assert(stream->dedup); + + stream->recent_blocknums[stream->recent_index] = blocknum; + stream->recent_index = (stream->recent_index + 1) & (DEDUP_RING_SIZE - 1); + stream->queued_buffers++; +} + +/* + * Make room for one more repeat entry, allocating the circular queue if this + * is the first one, and otherwise doubling it, up to repeat_cap. + */ +static pg_noinline void +read_stream_grow_repeats(ReadStream *stream) +{ + ReadStreamRepeat *repeats; + int16 nrepeats = stream->nrepeats; + int16 new_max; + int16 src; + + Assert(stream->dedup); + Assert(stream->nrepeats == stream->max_repeats); + Assert(stream->max_repeats < stream->repeat_cap); + + if (stream->max_repeats == 0) + new_max = Min(DEDUP_INITIAL_REPEATS, stream->repeat_cap); + else + new_max = Min(stream->max_repeats * 2, stream->repeat_cap); + + /* + * One slot more than the capacity, so that the look-ahead at the end of + * read_stream_next_repeat() cannot queue a repeat over the entry just + * handed out. The queue lives as long as the stream does. + */ + repeats = MemoryContextAlloc(GetMemoryChunkContext(stream), + sizeof(ReadStreamRepeat) * (new_max + 1)); + + /* Unwrap the live entries into the front of the new queue */ + for (int16 i = 0; i < nrepeats; i++) + { + src = stream->oldest_repeat_index + i; + if (src > stream->max_repeats) + src -= stream->max_repeats + 1; + repeats[i] = stream->repeats[src]; + } + + if (stream->repeats != NULL) + pfree(stream->repeats); + stream->repeats = repeats; + stream->max_repeats = new_max; + stream->oldest_repeat_index = 0; + stream->next_repeat_index = nrepeats; +} + +/* + * Queue a repeat entry for blocknum, which the callback returned while it was + * still in the ring of recent block numbers. The entry's place in the + * sequence is the number of buffers queued before it. + */ +static inline void +read_stream_queue_repeat(ReadStream *stream, BlockNumber blocknum) +{ + ReadStreamRepeat *repeat; + + Assert(stream->dedup); + + if (unlikely(stream->nrepeats == stream->max_repeats)) + read_stream_grow_repeats(stream); + Assert(stream->nrepeats < stream->max_repeats); + + repeat = &stream->repeats[stream->next_repeat_index]; + repeat->blocknum = blocknum; + repeat->seq = stream->queued_buffers; + + if (++stream->next_repeat_index == stream->max_repeats + 1) + stream->next_repeat_index = 0; + stream->nrepeats++; +} + +/* + * Has the oldest repeat entry's turn come? It has once every buffer queued + * before it has been handed out, i.e. once the buffers handed out so far + * (those queued, less those still pinned or pending) reach its seq. + */ +static inline bool +read_stream_repeat_due(ReadStream *stream) +{ + Assert(stream->dedup); + Assert(stream->nrepeats > 0); + + return stream->repeats[stream->oldest_repeat_index].seq == + stream->queued_buffers - stream->pinned_buffers - + stream->pending_read_nblocks; +} + +/* + * Forget all repeat entries (they hold no pins) and all recent block numbers. + */ +static void +read_stream_dedup_reset(ReadStream *stream) +{ + Assert(stream->dedup); + + for (int i = 0; i < DEDUP_RING_SIZE; i++) + stream->recent_blocknums[i] = InvalidBlockNumber; + stream->recent_index = 0; + stream->nrepeats = 0; + stream->oldest_repeat_index = 0; + stream->next_repeat_index = 0; + stream->queued_buffers = 0; +} + /* * General-use ReadStreamBlockNumberCB for block range scans. Loops over the * blocks [current_blocknum, last_exclusive). @@ -562,6 +751,11 @@ read_stream_should_look_ahead(ReadStream *stream) if (stream->ios_in_progress >= stream->max_ios) return false; + /* never queue more repeat entries than we can make space for */ + if (stream->dedup && stream->nrepeats == stream->max_repeats && + stream->max_repeats == stream->repeat_cap) + return false; + /* * Allow looking further ahead if we are in the process of building a * larger IO, the IO is not yet big enough, and we don't yet have IO in @@ -698,11 +892,24 @@ read_stream_look_ahead(ReadStream *stream) break; } + /* + * A recently queued block is not read or pinned again in advance. Its + * repeat entry will be pinned when its turn comes, and the pending + * read is free to keep growing past it. + */ + if (stream->dedup && read_stream_recent_block(stream, blocknum)) + { + read_stream_queue_repeat(stream, blocknum); + continue; + } + /* Can we merge it with the pending read? */ if (stream->pending_read_nblocks > 0 && stream->pending_read_blocknum + stream->pending_read_nblocks == blocknum) { stream->pending_read_nblocks++; + if (stream->dedup) + read_stream_remember_block(stream, blocknum); continue; } @@ -723,6 +930,8 @@ read_stream_look_ahead(ReadStream *stream) /* This is the start of a new pending read. */ stream->pending_read_blocknum = blocknum; stream->pending_read_nblocks = 1; + if (stream->dedup) + read_stream_remember_block(stream, blocknum); } /* @@ -736,17 +945,67 @@ read_stream_look_ahead(ReadStream *stream) read_stream_start_pending_read(stream); /* - * There should always be something pinned when we leave this function, - * whether started by this call or not, unless we've hit the end of the - * stream. In the worst case we can always make progress one buffer at a - * time. + * There should always be something pinned or a repeat entry queued when + * we leave this function, whether started by this call or not, unless + * we've hit the end of the stream. In the worst case we can always make + * progress one buffer at a time. */ - Assert(stream->pinned_buffers > 0 || stream->readahead_distance <= 0); + Assert(stream->pinned_buffers > 0 || stream->nrepeats > 0 || + stream->readahead_distance <= 0); if (stream->batch_mode) pgaio_exit_batchmode(); } +/* + * Hand out the oldest repeat entry, whose turn has come: pin its block afresh. + * The block was queued recently, so this is normally a buffer mapping hit. + * If the buffer was evicted in the meantime, we read the block synchronously. + * + * A repeat hand-out leaves the look-ahead distances alone: the entry had no + * I/O of its own, and the read of the block it repeats has already been + * accounted for. + */ +static pg_noinline Buffer +read_stream_next_repeat(ReadStream *stream) +{ + ReadStreamRepeat *repeat = &stream->repeats[stream->oldest_repeat_index]; + ReadBuffersOperation operation; + Buffer buffer = InvalidBuffer; + int flags = stream->read_buffers_flags | READ_BUFFERS_SYNCHRONOUSLY; + + Assert(stream->dedup); + Assert(read_stream_repeat_due(stream)); + + /* Read from the same relation, with the same strategy, as the stream */ + operation.rel = stream->ios[0].op.rel; + operation.smgr = stream->ios[0].op.smgr; + operation.persistence = stream->ios[0].op.persistence; + operation.forknum = stream->ios[0].op.forknum; + operation.strategy = stream->ios[0].op.strategy; + + if (unlikely(StartReadBuffer(&operation, &buffer, repeat->blocknum, flags))) + { + WaitReadBuffers(&operation); + + /* update I/O stats */ + read_stream_count_io(stream, 1, stream->ios_in_progress + 1); + read_stream_count_wait(stream); + } + Assert(BufferGetBlockNumber(buffer) == repeat->blocknum); + + read_stream_count_prefetch(stream); + + if (++stream->oldest_repeat_index == stream->max_repeats + 1) + stream->oldest_repeat_index = 0; + stream->nrepeats--; + + /* Prepare for the next call; there is space for one more repeat entry. */ + read_stream_look_ahead(stream); + + return buffer; +} + /* * Create a new read stream object that can be used to perform the equivalent * of a series of ReadBuffer() calls for one fork of one relation. @@ -771,6 +1030,7 @@ read_stream_begin_impl(int flags, size_t size; int16 queue_size; int16 queue_overflow; + int16 repeat_cap; int max_ios; int strategy_pin_limit; uint32 max_pinned_buffers; @@ -866,6 +1126,22 @@ read_stream_begin_impl(int flags, */ queue_size = max_pinned_buffers + 1; + /* + * Streams that deduplicate recent blocks must not use per-buffer data at + * all, because a repeat entry has nowhere to keep it. Only the ring of + * recent block numbers is carved out of the stream's own allocation; the + * queue of repeat entries is allocated on demand, and repeat_cap is how + * large it may grow. + */ + if (flags & READ_STREAM_DEDUP_RECENT) + { + Assert(per_buffer_data_size == 0); + repeat_cap = Min(max_pinned_buffers * DEDUP_REPEATS_PER_BUFFER, + DEDUP_MAX_REPEATS); + } + else + repeat_cap = 0; + /* * Allocate the object, the buffers, the ios and per_buffer_data space in * one big chunk. Though we have queue_size buffers, we want to be able @@ -878,6 +1154,11 @@ read_stream_begin_impl(int flags, size += sizeof(InProgressIO) * Max(1, max_ios); size += per_buffer_data_size * queue_size; size += MAXIMUM_ALIGNOF * 2; + if (repeat_cap > 0) + { + size += sizeof(BlockNumber) * DEDUP_RING_SIZE; + size += MAXIMUM_ALIGNOF; + } stream = (ReadStream *) palloc(size); memset(stream, 0, offsetof(ReadStream, buffers)); stream->ios = (InProgressIO *) @@ -885,6 +1166,15 @@ read_stream_begin_impl(int flags, if (per_buffer_data_size > 0) stream->per_buffer_data = (void *) MAXALIGN(&stream->ios[Max(1, max_ios)]); + if (repeat_cap > 0) + { + Assert(stream->per_buffer_data == NULL); + stream->recent_blocknums = (BlockNumber *) + MAXALIGN(&stream->ios[Max(1, max_ios)]); + stream->dedup = true; + stream->repeat_cap = repeat_cap; + read_stream_dedup_reset(stream); + } stream->sync_mode = io_method == IOMETHOD_SYNC; stream->batch_mode = flags & READ_STREAM_USE_BATCHING; @@ -1053,6 +1343,7 @@ read_stream_next_buffer(ReadStream *stream, void **per_buffer_data) Assert(stream->combine_distance == 1); Assert(stream->pending_read_nblocks == 0); Assert(stream->per_buffer_data_size == 0); + Assert(stream->nrepeats == 0); Assert(stream->initialized_buffers > stream->oldest_buffer_index); /* We're going to return the buffer we pinned last time. */ @@ -1116,6 +1407,16 @@ read_stream_next_buffer(ReadStream *stream, void **per_buffer_data) stream->ios[0].buffer_index = oldest_buffer_index; stream->seq_blocknum = next_blocknum + 1; + /* + * Blocks pinned in the fast path aren't entered into the ring of + * recent blocks, to keep the fast path fast; a repeat of one of + * them just after leaving the fast path is pinned again like + * before. This block's read is in progress, though, so remember + * it. + */ + if (stream->dedup) + read_stream_remember_block(stream, next_blocknum); + /* * XXX: It might be worth triggering additional read-ahead here, * to avoid having to effectively do another synchronous IO for @@ -1143,7 +1444,7 @@ read_stream_next_buffer(ReadStream *stream, void **per_buffer_data) } #endif - if (unlikely(stream->pinned_buffers == 0)) + if (unlikely(stream->pinned_buffers == 0 && stream->nrepeats == 0)) { Assert(stream->oldest_buffer_index == stream->next_buffer_index); @@ -1160,13 +1461,27 @@ read_stream_next_buffer(ReadStream *stream, void **per_buffer_data) read_stream_look_ahead(stream); /* End of stream reached? */ - if (stream->pinned_buffers == 0) + if (stream->pinned_buffers == 0 && stream->nrepeats == 0) { Assert(stream->readahead_distance == 0); return InvalidBuffer; } } + /* + * Either a buffer is pinned, or a repeat entry is queued whose turn has + * come. With nothing pinned nothing is pending either, as look-ahead can + * always pin one buffer, so every queued repeat entry is due. + */ + Assert(stream->pinned_buffers > 0 || read_stream_repeat_due(stream)); + + /* + * A repeat entry whose turn has come goes out before the oldest pinned + * buffer. Only READ_STREAM_DEDUP_RECENT streams have any. + */ + if (stream->nrepeats > 0 && read_stream_repeat_due(stream)) + return read_stream_next_repeat(stream); + /* Grab the oldest pinned buffer and associated per-buffer data. */ Assert(stream->pinned_buffers > 0); oldest_buffer_index = stream->oldest_buffer_index; @@ -1355,7 +1670,8 @@ read_stream_next_buffer(ReadStream *stream, void **per_buffer_data) stream->readahead_distance == 1 && stream->combine_distance == 1 && stream->pending_read_nblocks == 0 && - stream->per_buffer_data_size == 0) + stream->per_buffer_data_size == 0 && + stream->nrepeats == 0) { /* * The fast path spins on one buffer entry repeatedly instead of @@ -1462,6 +1778,10 @@ read_stream_reset(ReadStream *stream) stream->buffered_blocknum = InvalidBlockNumber; stream->fast_path = false; + /* Forget repeat entries and recent blocks. */ + if (stream->dedup) + read_stream_dedup_reset(stream); + /* Unpin anything that wasn't consumed. */ while ((buffer = read_stream_next_buffer(stream, NULL)) != InvalidBuffer) ReleaseBuffer(buffer); @@ -1502,5 +1822,7 @@ void read_stream_end(ReadStream *stream) { read_stream_reset(stream); + if (stream->repeats != NULL) + pfree(stream->repeats); pfree(stream); } diff --git a/src/test/modules/test_aio/t/004_read_stream.pl b/src/test/modules/test_aio/t/004_read_stream.pl index 32311c07a..dab956294 100644 --- a/src/test/modules/test_aio/t/004_read_stream.pl +++ b/src/test/modules/test_aio/t/004_read_stream.pl @@ -89,6 +89,53 @@ sub test_repeated_blocks ARRAY[0, 1, 2, 3, 4, 5, 6, 5, 4, 3, 2, 1, 0]);/); ok(1, "$io_method: stream accessing same block"); + # The same sequences with READ_STREAM_DEDUP_RECENT, where a repeat of + # a block that is still in flight or was consumed a moment ago is pinned + # afresh when its turn comes, instead of being read again. + $psql->query_safe(qq/SELECT evict_rel('largeish');/); + $psql->query_safe( + qq/SELECT * FROM read_stream_for_blocks('largeish', + ARRAY[0, 2, 2, 4, 4], true);/); + ok(1, "$io_method: dedup stream missing the same block repeatedly"); + + # Repeats interleaved with new blocks, some of which are still in flight + $psql->query_safe(qq/SELECT evict_rel('largeish');/); + $psql->query_safe( + qq/SELECT * FROM read_stream_for_blocks('largeish', + ARRAY[0, 1, 0, 2, 1, 3, 2, 4, 3, 5, 4, 6, 5, 0, 6, 1], true);/); + ok(1, "$io_method: dedup stream missing jittered blocks"); + + # More repeats than the stream has repeat entries (two pinned buffers + # allowed, so 16 repeat entries), so that look-ahead has to stop and + # resume as they are consumed + $psql->query_safe(qq/SET effective_io_concurrency = 1;/); + $psql->query_safe(qq/SELECT evict_rel('largeish');/); + $psql->query_safe( + qq/SELECT * FROM read_stream_for_blocks('largeish', + (SELECT array_agg(i % 2) FROM generate_series(0, 99) i), true);/); + ok(1, "$io_method: dedup stream with more repeats than repeat entries"); + $psql->query_safe(qq/RESET effective_io_concurrency;/); + + # The same kinds of sequences with I/O combining, so that a pending read + # of several blocks grows past repeat entries queued in the middle of it + $psql->query_safe(qq/SET io_combine_limit = 16;/); + $psql->query_safe(qq/SELECT evict_rel('largeish');/); + $psql->query_safe( + qq/SELECT * FROM read_stream_for_blocks('largeish', + ARRAY[0, 1, 2, 0, 3, 4, 1, 5, 6, 2, 7, 8, 5, 9, 3], true);/); + ok(1, "$io_method: dedup stream combining reads past repeats"); + $psql->query_safe(qq/SET io_combine_limit = 1;/); + + # A repeat entry holds no pin, so the block it repeats can be evicted + # before its turn comes. evict_before evicts the block of the request it + # names, just before that request is pulled from the stream, leaving the + # repeat entry to read the block again for itself. + $psql->query_safe(qq/SELECT evict_rel('largeish');/); + $psql->query_safe( + qq/SELECT * FROM read_stream_for_blocks('largeish', + ARRAY[0, 7, 7], true, evict_before=>2);/); + ok(1, "$io_method: dedup stream repeating an evicted block"); + # Test repeated blocks with a temp table, using invalidate_rel_block() # to evict individual local buffers. $psql->query_safe( @@ -111,6 +158,14 @@ sub test_repeated_blocks ARRAY[0, 2, 2, 4, 4]);/); ok(1, "$io_method: temp stream hitting the same block repeatedly"); + # And with READ_STREAM_DEDUP_RECENT, missing and hitting local buffers + $psql->query_safe(qq/SELECT invalidate_rel_block('largeish_temp', 0);/); + $psql->query_safe(qq/SELECT invalidate_rel_block('largeish_temp', 2);/); + $psql->query_safe(qq/SELECT invalidate_rel_block('largeish_temp', 4);/); + $psql->query_safe( + qq/SELECT * FROM read_stream_for_blocks('largeish_temp', + ARRAY[0, 2, 0, 4, 2, 4, 0], true);/); + ok(1, "$io_method: temp dedup stream missing recent blocks"); $psql->quit(); } diff --git a/src/test/modules/test_aio/test_aio--1.0.sql b/src/test/modules/test_aio/test_aio--1.0.sql index ff3b200cf..dc6a18bca 100644 --- a/src/test/modules/test_aio/test_aio--1.0.sql +++ b/src/test/modules/test_aio/test_aio--1.0.sql @@ -60,7 +60,7 @@ AS 'MODULE_PATHNAME' LANGUAGE C; /* * Read stream related functions */ -CREATE FUNCTION read_stream_for_blocks(rel regclass, blocks int4[], OUT blockoff int4, OUT blocknum int4, OUT buf int4) +CREATE FUNCTION read_stream_for_blocks(rel regclass, blocks int4[], dedup bool DEFAULT false, evict_before int4 DEFAULT -1, OUT blockoff int4, OUT blocknum int4, OUT buf int4) RETURNS SETOF record STRICT AS 'MODULE_PATHNAME' LANGUAGE C; diff --git a/src/test/modules/test_aio/test_aio.c b/src/test/modules/test_aio/test_aio.c index cebcf18d2..208dba57b 100644 --- a/src/test/modules/test_aio/test_aio.c +++ b/src/test/modules/test_aio/test_aio.c @@ -872,10 +872,13 @@ read_stream_for_blocks(PG_FUNCTION_ARGS) { Oid relid = PG_GETARG_OID(0); ArrayType *blocksarray = PG_GETARG_ARRAYTYPE_P(1); + bool dedup = PG_GETARG_BOOL(2); + int evict_before = PG_GETARG_INT32(3); ReturnSetInfo *rsinfo = (ReturnSetInfo *) fcinfo->resultinfo; Relation rel; BlocksReadStreamData stream_data; ReadStream *stream; + int flags = READ_STREAM_FULL; InitMaterializedSRF(fcinfo, 0); @@ -895,7 +898,10 @@ read_stream_for_blocks(PG_FUNCTION_ARGS) rel = relation_open(relid, AccessShareLock); - stream = read_stream_begin_relation(READ_STREAM_FULL, + if (dedup) + flags |= READ_STREAM_DEDUP_RECENT; + + stream = read_stream_begin_relation(flags, NULL, rel, MAIN_FORKNUM, @@ -905,12 +911,26 @@ read_stream_for_blocks(PG_FUNCTION_ARGS) for (int i = 0; i < stream_data.nblocks; i++) { - Buffer buf = read_stream_next_buffer(stream, NULL); + Buffer buf; Datum values[3] = {0}; bool nulls[3] = {0}; + /* + * Evict the block that this request asks for, so that a stream which + * deduplicates recent blocks has to read it again when the turn of + * its repeat entry comes. That only works for a block the stream has + * already handed out, as an evicted buffer must not be pinned. + */ + if (i == evict_before) + invalidate_one_block(rel, MAIN_FORKNUM, stream_data.blocks[i]); + + buf = read_stream_next_buffer(stream, NULL); + if (!BufferIsValid(buf)) elog(ERROR, "read_stream_next_buffer() call %d is unexpectedly invalid", i); + if (BufferGetBlockNumber(buf) != stream_data.blocks[i]) + elog(ERROR, "read_stream_next_buffer() call %d returned block %u, expected block %u", + i, BufferGetBlockNumber(buf), stream_data.blocks[i]); values[0] = Int32GetDatum(i); values[1] = UInt32GetDatum(stream_data.blocks[i]); diff --git a/src/tools/pgindent/typedefs.list b/src/tools/pgindent/typedefs.list index fe82a9f74..f165e2ba0 100644 --- a/src/tools/pgindent/typedefs.list +++ b/src/tools/pgindent/typedefs.list @@ -2548,6 +2548,7 @@ ReadNodeContext ReadReplicationSlotCmd ReadStream ReadStreamBlockNumberCB +ReadStreamRepeat ReassignOwnedStmt RecheckForeignScan_function RecordCacheArrayEntry -- 2.55.0