From d339891d8455fb7264a5d06e76ac4171249c2cd2 Mon Sep 17 00:00:00 2001 From: David Geier Date: Thu, 17 Sep 2026 13:18:46 +0200 Subject: [PATCH v2] Support for parallel bitmap operations --- src/backend/access/gin/ginget.c | 10 +- src/backend/access/gin/ginscan.c | 2 +- src/backend/access/heap/heapam.c | 2 +- src/backend/access/index/indexam.c | 40 ++ src/backend/access/nbtree/nbtree.c | 4 + src/backend/executor/execParallel.c | 16 +- src/backend/executor/nodeBitmapHeapscan.c | 230 +------ src/backend/executor/nodeBitmapIndexscan.c | 362 ++++++++++- src/backend/executor/nodeBitmapOr.c | 7 +- src/backend/nodes/tidbitmap.c | 571 ++++-------------- src/backend/optimizer/path/allpaths.c | 104 +++- src/backend/optimizer/path/indxpath.c | 101 ++-- src/backend/optimizer/plan/createplan.c | 31 +- src/backend/optimizer/util/pathnode.c | 40 +- src/include/access/genam.h | 6 + src/include/access/gin_private.h | 2 +- src/include/access/relscan.h | 2 +- src/include/executor/nodeBitmapHeapscan.h | 14 +- src/include/executor/nodeBitmapIndexscan.h | 1 + src/include/nodes/execnodes.h | 13 +- src/include/nodes/plannodes.h | 3 - src/include/nodes/tidbitmap.h | 72 +-- src/test/regress/expected/bitmapops.out | 98 +++ src/test/regress/expected/memoize.out | 4 +- src/test/regress/expected/select_parallel.out | 6 +- src/test/regress/sql/bitmapops.sql | 49 ++ src/test/regress/sql/memoize.sql | 2 + src/tools/pgindent/typedefs.list | 8 +- 28 files changed, 931 insertions(+), 869 deletions(-) diff --git a/src/backend/access/gin/ginget.c b/src/backend/access/gin/ginget.c index 2bcb32ca3d0..543f9a9bb38 100644 --- a/src/backend/access/gin/ginget.c +++ b/src/backend/access/gin/ginget.c @@ -374,7 +374,7 @@ restartScanEntry: if (entry->matchBitmap) { if (entry->matchIterator) - tbm_end_private_iterate(entry->matchIterator); + tbm_end_ordered_iterate(&entry->matchIterator); entry->matchIterator = NULL; tbm_free(entry->matchBitmap); entry->matchBitmap = NULL; @@ -387,7 +387,7 @@ restartScanEntry: if (entry->matchBitmap && !tbm_is_empty(entry->matchBitmap)) { entry->matchIterator = - tbm_begin_private_iterate(entry->matchBitmap); + tbm_begin_ordered_iterate(entry->matchBitmap); entry->isFinished = false; } } @@ -828,7 +828,7 @@ entryGetItem(GinState *ginstate, GinScanEntry entry, { /* * If we've exhausted all items on this block, move to next block - * in the bitmap. tbm_private_iterate() sets matchResult.blockno + * in the bitmap. tbm_ordered_iterate() sets matchResult.blockno * to InvalidBlockNumber when the bitmap is exhausted. */ while ((!BlockNumberIsValid(entry->matchResult.blockno)) || @@ -838,11 +838,11 @@ entryGetItem(GinState *ginstate, GinScanEntry entry, (ItemPointerIsLossyPage(&advancePast) && entry->matchResult.blockno == advancePastBlk)) { - if (!tbm_private_iterate(entry->matchIterator, &entry->matchResult)) + if (!tbm_ordered_iterate(entry->matchIterator, &entry->matchResult)) { Assert(!BlockNumberIsValid(entry->matchResult.blockno)); ItemPointerSetInvalid(&entry->curItem); - tbm_end_private_iterate(entry->matchIterator); + tbm_end_ordered_iterate(&entry->matchIterator); entry->matchIterator = NULL; entry->isFinished = true; break; diff --git a/src/backend/access/gin/ginscan.c b/src/backend/access/gin/ginscan.c index 1564b409bdb..a13b784f2e5 100644 --- a/src/backend/access/gin/ginscan.c +++ b/src/backend/access/gin/ginscan.c @@ -250,7 +250,7 @@ ginFreeScanKeys(GinScanOpaque so) if (entry->list) pfree(entry->list); if (entry->matchIterator) - tbm_end_private_iterate(entry->matchIterator); + tbm_end_ordered_iterate(&entry->matchIterator); if (entry->matchBitmap) tbm_free(entry->matchBitmap); } diff --git a/src/backend/access/heap/heapam.c b/src/backend/access/heap/heapam.c index 4207f0e0e08..4638159dcb9 100644 --- a/src/backend/access/heap/heapam.c +++ b/src/backend/access/heap/heapam.c @@ -330,7 +330,7 @@ bitmapheap_stream_read_next(ReadStream *pgsr, void *private_data, CHECK_FOR_INTERRUPTS(); /* no more entries in the bitmap */ - if (!tbm_iterate(&sscan->st.rs_tbmiterator, tbmres)) + if (!tbm_ordered_iterate(sscan->st.rs_tbmiterator, tbmres)) return InvalidBlockNumber; /* diff --git a/src/backend/access/index/indexam.c b/src/backend/access/index/indexam.c index db811056991..ef7dc37ec61 100644 --- a/src/backend/access/index/indexam.c +++ b/src/backend/access/index/indexam.c @@ -304,6 +304,46 @@ index_beginscan_bitmap(Relation indexRelation, NULL, instrument, false, false, SO_NONE); } +/* + * index_beginscan_bitmap_parallel - start a parallel scan of an index with amgetbitmap. + * + * Unlike index_beginscan_parallel, no heap relation is supplied: + * bitmap index scans never access the heap directly. + */ +IndexScanDesc +index_beginscan_bitmap_parallel(Relation indexRelation, + Snapshot snapshot, + IndexScanInstrumentation *instrument, + int nkeys, + ParallelIndexScanDesc pscan) +{ + Snapshot restoredsnap; + + Assert(snapshot != InvalidSnapshot); + Assert(IsMVCCLikeSnapshot(snapshot)); + + /* + * Use the snapshot serialized in the parallel scan descriptor, just like + * index_beginscan_parallel does. Register it so that index_endscan can + * unregister it from the same resource owner. + */ + restoredsnap = RestoreSnapshot(pscan->ps_snapshot_data); + RegisterSnapshot(restoredsnap); + + return index_beginscan_internal(indexRelation, NULL, nkeys, 0, restoredsnap, + pscan, instrument, false, true, SO_NONE); +} + +/* + * index_am_supports_parallel_scan - does the index access method support + * parallel scans? + */ +bool +index_am_supports_parallel_scan(Relation indexRelation) +{ + return indexRelation->rd_indam->amcanparallel; +} + /* * index_beginscan_internal --- common code for index_beginscan variants * diff --git a/src/backend/access/nbtree/nbtree.c b/src/backend/access/nbtree/nbtree.c index 0abdd7b49f5..1a77df5aa79 100644 --- a/src/backend/access/nbtree/nbtree.c +++ b/src/backend/access/nbtree/nbtree.c @@ -329,6 +329,10 @@ btgetbitmap(IndexScanDesc scan, TIDBitmap *tbm) /* Now see if we need another primitive index scan */ } while (so->numArrayKeys && _bt_start_prim_scan(scan)); + /* Mark parallel scan finished if we're part of one */ + if (scan->parallel_scan != NULL) + _bt_parallel_done(scan); + return ntids; } diff --git a/src/backend/executor/execParallel.c b/src/backend/executor/execParallel.c index 6b85508a696..b59c199cb53 100644 --- a/src/backend/executor/execParallel.c +++ b/src/backend/executor/execParallel.c @@ -306,9 +306,6 @@ ExecParallelEstimate(PlanState *planstate, ExecParallelEstimateContext *e) e->pcxt); break; case T_BitmapHeapScanState: - if (planstate->plan->parallel_aware) - ExecBitmapHeapEstimate((BitmapHeapScanState *) planstate, - e->pcxt); /* even when not parallel-aware, for EXPLAIN ANALYZE */ ExecBitmapHeapInstrumentEstimate((BitmapHeapScanState *) planstate, e->pcxt); @@ -554,9 +551,6 @@ ExecParallelInitializeDSM(PlanState *planstate, d->pcxt); break; case T_BitmapHeapScanState: - if (planstate->plan->parallel_aware) - ExecBitmapHeapInitializeDSM((BitmapHeapScanState *) planstate, - d->pcxt); /* even when not parallel-aware, for EXPLAIN ANALYZE */ ExecBitmapHeapInstrumentInitDSM((BitmapHeapScanState *) planstate, d->pcxt); @@ -1060,9 +1054,6 @@ ExecParallelReInitializeDSM(PlanState *planstate, pcxt); break; case T_BitmapHeapScanState: - if (planstate->plan->parallel_aware) - ExecBitmapHeapReInitializeDSM((BitmapHeapScanState *) planstate, - pcxt); break; case T_HashJoinState: if (planstate->plan->parallel_aware) @@ -1070,6 +1061,10 @@ ExecParallelReInitializeDSM(PlanState *planstate, pcxt); break; case T_BitmapIndexScanState: + if (planstate->plan->parallel_aware) + ExecBitmapIndexScanReInitializeDSM((BitmapIndexScanState *) planstate, + pcxt); + break; case T_HashState: case T_SortState: case T_IncrementalSortState: @@ -1451,9 +1446,6 @@ ExecParallelInitializeWorker(PlanState *planstate, ParallelWorkerContext *pwcxt) pwcxt); break; case T_BitmapHeapScanState: - if (planstate->plan->parallel_aware) - ExecBitmapHeapInitializeWorker((BitmapHeapScanState *) planstate, - pwcxt); /* even when not parallel-aware, for EXPLAIN ANALYZE */ ExecBitmapHeapInstrumentInitWorker((BitmapHeapScanState *) planstate, pwcxt); diff --git a/src/backend/executor/nodeBitmapHeapscan.c b/src/backend/executor/nodeBitmapHeapscan.c index a9e83f30687..958dd57c589 100644 --- a/src/backend/executor/nodeBitmapHeapscan.c +++ b/src/backend/executor/nodeBitmapHeapscan.c @@ -44,97 +44,33 @@ #include "miscadmin.h" #include "pgstat.h" #include "storage/bufmgr.h" -#include "storage/condition_variable.h" -#include "utils/dsa.h" #include "utils/rel.h" #include "utils/spccache.h" -#include "utils/wait_event.h" static void BitmapTableScanSetup(BitmapHeapScanState *node); static TupleTableSlot *BitmapHeapNext(BitmapHeapScanState *node); -static inline void BitmapDoneInitializingSharedState(ParallelBitmapHeapState *pstate); -static bool BitmapShouldInitializeSharedState(ParallelBitmapHeapState *pstate); - - -/* ---------------- - * SharedBitmapState information - * - * BM_INITIAL TIDBitmap creation is not yet started, so first worker - * to see this state will set the state to BM_INPROGRESS - * and that process will be responsible for creating - * TIDBitmap. - * BM_INPROGRESS TIDBitmap creation is in progress; workers need to - * sleep until it's finished. - * BM_FINISHED TIDBitmap creation is done, so now all workers can - * proceed to iterate over TIDBitmap. - * ---------------- - */ -typedef enum -{ - BM_INITIAL, - BM_INPROGRESS, - BM_FINISHED, -} SharedBitmapState; - -/* ---------------- - * ParallelBitmapHeapState information - * tbmiterator iterator for scanning current pages - * state current state of the TIDBitmap - * cv conditional wait variable - * ---------------- - */ -typedef struct ParallelBitmapHeapState -{ - dsa_pointer tbmiterator; - pg_atomic_uint32 state; - ConditionVariable cv; -} ParallelBitmapHeapState; /* - * Do the underlying index scan, build the bitmap, set up the parallel state - * needed for parallel workers to iterate through the bitmap, and set up the - * underlying table scan descriptor. + * Do the underlying index scan, build the bitmap, and set up the underlying + * table scan descriptor. With parallel-aware BitmapHeapScans, the underlying + * BitmapIndexScan(s) have already partitioned their output per worker. */ static void BitmapTableScanSetup(BitmapHeapScanState *node) { - TBMIterator tbmiterator = {0}; - ParallelBitmapHeapState *pstate = node->pstate; - dsa_area *dsa = node->ss.ps.state->es_query_dsa; + TBMOrderedIterator *tbmiterator; - if (!pstate) - { - node->tbm = (TIDBitmap *) MultiExecProcNode(outerPlanState(node)); + node->tbm = (TIDBitmap *) MultiExecProcNode(outerPlanState(node)); - if (!node->tbm || !IsA(node->tbm, TIDBitmap)) - elog(ERROR, "unrecognized result from subplan"); - } - else if (BitmapShouldInitializeSharedState(pstate)) - { - /* - * The leader will immediately come out of the function, but others - * will be blocked until leader populates the TBM and wakes them up. - */ - node->tbm = (TIDBitmap *) MultiExecProcNode(outerPlanState(node)); - if (!node->tbm || !IsA(node->tbm, TIDBitmap)) - elog(ERROR, "unrecognized result from subplan"); - - /* - * Prepare to iterate over the TBM. This will return the dsa_pointer - * of the iterator state which will be used by multiple processes to - * iterate jointly. - */ - pstate->tbmiterator = tbm_prepare_shared_iterate(node->tbm); - - /* We have initialized the shared state so wake up others. */ - BitmapDoneInitializingSharedState(pstate); - } + if (!node->tbm || !IsA(node->tbm, TIDBitmap)) + elog(ERROR, "unrecognized result from subplan"); - tbmiterator = tbm_begin_iterate(node->tbm, dsa, - pstate ? - pstate->tbmiterator : - InvalidDsaPointer); + /* + * The bitmap we receive is already a per-worker partition, so a private + * iterator is sufficient. + */ + tbmiterator = tbm_begin_ordered_iterate(node->tbm); /* * If this is the first scan of the underlying table, create the table @@ -217,19 +153,6 @@ BitmapHeapNext(BitmapHeapScanState *node) return ExecClearTuple(slot); } -/* - * BitmapDoneInitializingSharedState - Shared state is initialized - * - * By this time the leader has already populated the TBM and initialized the - * shared state so wake up other processes. - */ -static inline void -BitmapDoneInitializingSharedState(ParallelBitmapHeapState *pstate) -{ - pg_atomic_write_membarrier_u32(&pstate->state, BM_FINISHED); - ConditionVariableBroadcast(&pstate->cv); -} - /* * BitmapHeapRecheck -- access method routine to recheck a tuple in EvalPlanQual */ @@ -279,8 +202,8 @@ ExecReScanBitmapHeapScan(BitmapHeapScanState *node) * End iteration on iterators saved in scan descriptor if they have * not already been cleaned up. */ - if (!tbm_exhausted(&scan->st.rs_tbmiterator)) - tbm_end_iterate(&scan->st.rs_tbmiterator); + if (scan->st.rs_tbmiterator != NULL) + tbm_end_ordered_iterate(&scan->st.rs_tbmiterator); /* rescan to release any page pin */ table_rescan(node->ss.ss_currentScanDesc, NULL); @@ -359,8 +282,8 @@ ExecEndBitmapHeapScan(BitmapHeapScanState *node) * End iteration on iterators saved in scan descriptor if they have * not already been cleaned up. */ - if (!tbm_exhausted(&scanDesc->st.rs_tbmiterator)) - tbm_end_iterate(&scanDesc->st.rs_tbmiterator); + if (scanDesc->st.rs_tbmiterator != NULL) + tbm_end_ordered_iterate(&scanDesc->st.rs_tbmiterator); /* * close table scan @@ -410,7 +333,6 @@ ExecInitBitmapHeapScan(BitmapHeapScan *node, EState *estate, int eflags) memset(&scanstate->stats, 0, sizeof(BitmapHeapScanInstrumentation)); scanstate->initialized = false; - scanstate->pstate = NULL; scanstate->recheck = true; /* @@ -460,126 +382,6 @@ ExecInitBitmapHeapScan(BitmapHeapScan *node, EState *estate, int eflags) return scanstate; } -/*---------------- - * BitmapShouldInitializeSharedState - * - * The first process to come here and see the state to the BM_INITIAL - * will become the leader for the parallel bitmap scan and will be - * responsible for populating the TIDBitmap. The other processes will - * be blocked by the condition variable until the leader wakes them up. - * --------------- - */ -static bool -BitmapShouldInitializeSharedState(ParallelBitmapHeapState *pstate) -{ - uint32 state; - - while (1) - { - state = BM_INITIAL; - pg_atomic_compare_exchange_u32(&pstate->state, &state, BM_INPROGRESS); - - /* Exit if bitmap is done, or if we're the leader. */ - if (state != BM_INPROGRESS) - break; - - /* Wait for the leader to wake us up. */ - ConditionVariableSleep(&pstate->cv, WAIT_EVENT_PARALLEL_BITMAP_SCAN); - } - - ConditionVariableCancelSleep(); - - return (state == BM_INITIAL); -} - -/* ---------------------------------------------------------------- - * ExecBitmapHeapEstimate - * - * Compute the amount of space we'll need in the parallel - * query DSM, and inform pcxt->estimator about our needs. - * ---------------------------------------------------------------- - */ -void -ExecBitmapHeapEstimate(BitmapHeapScanState *node, - ParallelContext *pcxt) -{ - shm_toc_estimate_chunk(&pcxt->estimator, - MAXALIGN(sizeof(ParallelBitmapHeapState))); - shm_toc_estimate_keys(&pcxt->estimator, 1); -} - -/* ---------------------------------------------------------------- - * ExecBitmapHeapInitializeDSM - * - * Set up a parallel bitmap heap scan descriptor. - * ---------------------------------------------------------------- - */ -void -ExecBitmapHeapInitializeDSM(BitmapHeapScanState *node, - ParallelContext *pcxt) -{ - ParallelBitmapHeapState *pstate; - dsa_area *dsa = node->ss.ps.state->es_query_dsa; - - /* If there's no DSA, there are no workers; initialize nothing. */ - if (dsa == NULL) - return; - - pstate = (ParallelBitmapHeapState *) - shm_toc_allocate(pcxt->toc, - MAXALIGN(sizeof(ParallelBitmapHeapState))); - - pstate->tbmiterator = 0; - - pg_atomic_init_u32(&pstate->state, BM_INITIAL); - - ConditionVariableInit(&pstate->cv); - - shm_toc_insert(pcxt->toc, node->ss.ps.plan->plan_node_id, pstate); - node->pstate = pstate; -} - -/* ---------------------------------------------------------------- - * ExecBitmapHeapReInitializeDSM - * - * Reset shared state before beginning a fresh scan. - * ---------------------------------------------------------------- - */ -void -ExecBitmapHeapReInitializeDSM(BitmapHeapScanState *node, - ParallelContext *pcxt) -{ - ParallelBitmapHeapState *pstate = node->pstate; - dsa_area *dsa = node->ss.ps.state->es_query_dsa; - - /* If there's no DSA, there are no workers; do nothing. */ - if (dsa == NULL) - return; - - pg_atomic_write_u32(&pstate->state, BM_INITIAL); - - if (DsaPointerIsValid(pstate->tbmiterator)) - tbm_free_shared_area(dsa, pstate->tbmiterator); - - pstate->tbmiterator = InvalidDsaPointer; -} - -/* ---------------------------------------------------------------- - * ExecBitmapHeapInitializeWorker - * - * Copy relevant information from TOC into planstate. - * ---------------------------------------------------------------- - */ -void -ExecBitmapHeapInitializeWorker(BitmapHeapScanState *node, - ParallelWorkerContext *pwcxt) -{ - Assert(node->ss.ps.state->es_query_dsa != NULL); - - node->pstate = (ParallelBitmapHeapState *) - shm_toc_lookup(pwcxt->toc, node->ss.ps.plan->plan_node_id, false); -} - /* * Compute the amount of space we'll need for the shared instrumentation and * inform pcxt->estimator. diff --git a/src/backend/executor/nodeBitmapIndexscan.c b/src/backend/executor/nodeBitmapIndexscan.c index 90b010f9b71..51528b5c2cd 100644 --- a/src/backend/executor/nodeBitmapIndexscan.c +++ b/src/backend/executor/nodeBitmapIndexscan.c @@ -22,12 +22,33 @@ #include "postgres.h" #include "access/genam.h" +#include "access/relscan.h" +#include "common/hashfn.h" #include "executor/executor.h" #include "executor/instrument.h" #include "executor/nodeBitmapIndexscan.h" #include "executor/nodeIndexscan.h" #include "miscadmin.h" #include "nodes/tidbitmap.h" +#include "storage/barrier.h" +#include "utils/wait_event.h" + +/* + * Per-BitmapIndexScan shared state used to share locally-built bitmaps among + * workers and coordinate the partition/free phases. + */ +typedef struct SharedBitmapIndexState +{ + int max_participants; /* leader + max number of workers */ + Barrier barrier; + dsa_pointer worker_tbmiter[FLEXIBLE_ARRAY_MEMBER]; +} SharedBitmapIndexState; + +#define PARALLEL_KEY_BITMAP_INDEX_OFFSET UINT64CONST(0xD100000000000000) + +static Size BitmapIndexScanSharedStateSize(int nworkers); +static TIDBitmap *BitmapIndexScanPartition(BitmapIndexScanState *node, + TIDBitmap *tbm); /* ---------------------------------------------------------------- @@ -64,6 +85,24 @@ MultiExecBitmapIndexScan(BitmapIndexScanState *node) */ scandesc = node->biss_ScanDesc; + /* + * Serial execution of a parallel-aware plan (no workers launched) needs a + * scan descriptor that we did not create during DSM setup. + */ + if (scandesc == NULL) + { + scandesc = index_beginscan_bitmap(node->biss_RelationDesc, + node->ss.ps.state->es_snapshot, + node->biss_Instrument, + node->biss_NumScanKeys); + node->biss_ScanDesc = scandesc; + + if (node->biss_NumRuntimeKeys == 0 || node->biss_RuntimeKeysReady) + index_rescan(scandesc, + node->biss_ScanKeys, node->biss_NumScanKeys, + NULL, 0); + } + /* * If we have runtime keys and they've not already been set up, do it now. * Array keys are also treated as runtime keys; note that if ExecReScan @@ -79,6 +118,43 @@ MultiExecBitmapIndexScan(BitmapIndexScanState *node) else doscan = true; + /* + * If we're running as part of a parallel query, we need to attach to the + * barrier before scanning the index. This ensures that workers that arrive + * late don't waste time scanning the index only to discard their bitmap. + */ + if (node->biss_ParallelState != NULL) + { + SharedBitmapIndexState *sstate = node->biss_ParallelState; + int phase; + + phase = BarrierAttach(&sstate->barrier); + if (phase != 0) + { + /* + * We attached after the build phase was already done. We cannot + * contribute a partition, so just detach and return an empty bitmap. + */ + TIDBitmap *empty_bitmap; + + BarrierDetach(&sstate->barrier); + empty_bitmap = tbm_create(get_hash_memory_limit(), NULL); + if (node->ss.ps.instrument) + InstrStopNode(node->ss.ps.instrument, 0); + return (Node *) empty_bitmap; + } + + /* + * Index access methods that don't support parallel scans cannot divide + * the index scan work among workers. Elect exactly one participant to + * build the full bitmap; the others will leave their bitmap empty and + * still receive a partition during the sharing phase below. + */ + if (!index_am_supports_parallel_scan(node->biss_RelationDesc) && + !BarrierArriveAndWait(&sstate->barrier, WAIT_EVENT_PARALLEL_BITMAP_SCAN)) + doscan = false; + } + /* * Prepare the result bitmap. Normally we just create a new one to pass * back; however, our parent node is allowed to store a pre-made one into @@ -94,7 +170,7 @@ MultiExecBitmapIndexScan(BitmapIndexScanState *node) { /* XXX should we use less than work_mem for this? */ tbm = tbm_create(work_mem * (Size) 1024, - ((BitmapIndexScan *) node->ss.ps.plan)->isshared ? + node->ss.ps.plan->parallel_aware ? node->ss.ps.state->es_query_dsa : NULL); } @@ -115,6 +191,13 @@ MultiExecBitmapIndexScan(BitmapIndexScanState *node) NULL, 0); } + /* + * If we're running as part of a parallel query, partition the full bitmap + * so that each worker gets a disjoint subset of the heap blocks. + */ + if (node->biss_ParallelState != NULL) + tbm = BitmapIndexScanPartition(node, tbm); + /* must provide our own instrumentation support */ if (node->ss.ps.instrument) InstrStopNode(node->ss.ps.instrument, nTuples); @@ -157,13 +240,13 @@ ExecReScanBitmapIndexScan(BitmapIndexScanState *node) if (node->biss_NumArrayKeys != 0) node->biss_RuntimeKeysReady = ExecIndexEvalArrayKeys(econtext, - node->biss_ArrayKeys, - node->biss_NumArrayKeys); + node->biss_ArrayKeys, + node->biss_NumArrayKeys); else node->biss_RuntimeKeysReady = true; /* reset index scan */ - if (node->biss_RuntimeKeysReady) + if (node->biss_RuntimeKeysReady && node->biss_ScanDesc) index_rescan(node->biss_ScanDesc, node->biss_ScanKeys, node->biss_NumScanKeys, NULL, 0); @@ -227,6 +310,7 @@ ExecInitBitmapIndexScan(BitmapIndexScan *node, EState *estate, int eflags) { BitmapIndexScanState *indexstate; LOCKMODE lockmode; + bool parallel_aware = node->scan.plan.parallel_aware; /* check for unsupported flags */ Assert(!(eflags & (EXEC_FLAG_BACKWARD | EXEC_FLAG_MARK))); @@ -246,9 +330,15 @@ ExecInitBitmapIndexScan(BitmapIndexScan *node, EState *estate, int eflags) * We do not open or lock the base relation here. We assume that an * ancestor BitmapHeapScan node is holding AccessShareLock (or better) on * the heap relation throughout the execution of the plan tree. + * + * For a parallel-aware scan, however, we need access to the heap relation + * to initialize the parallel index scan descriptor. */ - - indexstate->ss.ss_currentRelation = NULL; + if (parallel_aware && !(eflags & EXEC_FLAG_EXPLAIN_ONLY)) + indexstate->ss.ss_currentRelation = + ExecOpenScanRelation(estate, node->scan.scanrelid, eflags); + else + indexstate->ss.ss_currentRelation = NULL; indexstate->ss.ss_currentScanDesc = NULL; /* @@ -325,23 +415,27 @@ ExecInitBitmapIndexScan(BitmapIndexScan *node, EState *estate, int eflags) } /* - * Initialize scan descriptor. + * Initialize scan descriptor. For parallel-aware scans this is delayed + * until the DSM is set up. */ - indexstate->biss_ScanDesc = - index_beginscan_bitmap(indexstate->biss_RelationDesc, - estate->es_snapshot, - indexstate->biss_Instrument, - indexstate->biss_NumScanKeys); + if (!parallel_aware) + { + indexstate->biss_ScanDesc = + index_beginscan_bitmap(indexstate->biss_RelationDesc, + estate->es_snapshot, + indexstate->biss_Instrument, + indexstate->biss_NumScanKeys); - /* - * If no run-time keys to calculate, go ahead and pass the scankeys to the - * index AM. - */ - if (indexstate->biss_NumRuntimeKeys == 0 && - indexstate->biss_NumArrayKeys == 0) - index_rescan(indexstate->biss_ScanDesc, - indexstate->biss_ScanKeys, indexstate->biss_NumScanKeys, - NULL, 0); + /* + * If no run-time keys to calculate, go ahead and pass the scankeys to the + * index AM. + */ + if (indexstate->biss_NumRuntimeKeys == 0 && + indexstate->biss_NumArrayKeys == 0) + index_rescan(indexstate->biss_ScanDesc, + indexstate->biss_ScanKeys, indexstate->biss_NumScanKeys, + NULL, 0); + } /* * all done. @@ -359,13 +453,30 @@ ExecInitBitmapIndexScan(BitmapIndexScan *node, EState *estate, int eflags) void ExecBitmapIndexScanEstimate(BitmapIndexScanState *node, ParallelContext *pcxt) { + EState *estate = node->ss.ps.state; Size size; + if (pcxt->nworkers == 0) + return; + /* - * Parallel bitmap index scans are not supported, but we still need to - * store the scan's instrumentation in DSM during parallel query + * Parallel-aware scans need space for the parallel scan descriptor and for + * the per-worker bitmap sharing/partitioning state. */ - if (!node->ss.ps.instrument || pcxt->nworkers == 0) + if (node->ss.ps.plan->parallel_aware) + { + node->biss_PscanLen = + index_parallelscan_estimate(node->biss_RelationDesc, + node->biss_NumScanKeys, + 0, + estate->es_snapshot); + shm_toc_estimate_chunk(&pcxt->estimator, node->biss_PscanLen); + shm_toc_estimate_chunk(&pcxt->estimator, + BitmapIndexScanSharedStateSize(pcxt->nworkers)); + shm_toc_estimate_keys(&pcxt->estimator, 2); + } + + if (!node->ss.ps.instrument) return; size = offsetof(SharedIndexScanInstrumentation, winstrument) + @@ -377,17 +488,61 @@ ExecBitmapIndexScanEstimate(BitmapIndexScanState *node, ParallelContext *pcxt) /* ---------------------------------------------------------------- * ExecBitmapIndexScanInitializeDSM * - * Set up bitmap index scan shared instrumentation. + * Set up shared state for a parallel-aware bitmap index scan. * ---------------------------------------------------------------- */ void ExecBitmapIndexScanInitializeDSM(BitmapIndexScanState *node, ParallelContext *pcxt) { + EState *estate = node->ss.ps.state; + ParallelIndexScanDesc piscan; + SharedBitmapIndexState *sstate; Size size; - /* don't need this if not instrumenting or no workers */ - if (!node->ss.ps.instrument || pcxt->nworkers == 0) + if (pcxt->nworkers == 0) + return; + + if (node->ss.ps.plan->parallel_aware) + { + piscan = shm_toc_allocate(pcxt->toc, node->biss_PscanLen); + index_parallelscan_initialize(node->ss.ss_currentRelation, + node->biss_RelationDesc, + estate->es_snapshot, + piscan); + shm_toc_insert(pcxt->toc, + node->ss.ps.plan->plan_node_id, + piscan); + + size = BitmapIndexScanSharedStateSize(pcxt->nworkers); + sstate = shm_toc_allocate(pcxt->toc, size); + memset(sstate, 0, size); + sstate->max_participants = pcxt->nworkers + 1; + BarrierInit(&sstate->barrier, 0); + shm_toc_insert(pcxt->toc, + node->ss.ps.plan->plan_node_id + + PARALLEL_KEY_BITMAP_INDEX_OFFSET, + sstate); + node->biss_ParallelState = sstate; + + node->biss_ScanDesc = + index_beginscan_bitmap_parallel(node->biss_RelationDesc, + estate->es_snapshot, + node->biss_Instrument, + node->biss_NumScanKeys, + piscan); + + /* + * If no run-time keys to calculate or they are ready, go ahead and pass + * the scankeys to the index AM. + */ + if (node->biss_NumRuntimeKeys == 0 || node->biss_RuntimeKeysReady) + index_rescan(node->biss_ScanDesc, + node->biss_ScanKeys, node->biss_NumScanKeys, + NULL, 0); + } + + if (!node->ss.ps.instrument) return; size = offsetof(SharedIndexScanInstrumentation, winstrument) + @@ -405,6 +560,34 @@ ExecBitmapIndexScanInitializeDSM(BitmapIndexScanState *node, node->biss_SharedInfo->num_workers = pcxt->nworkers; } +/* ---------------------------------------------------------------- + * ExecBitmapIndexScanReInitializeDSM + * + * Reset shared state before beginning a fresh scan. + * ---------------------------------------------------------------- + */ +void +ExecBitmapIndexScanReInitializeDSM(BitmapIndexScanState *node, + ParallelContext *pcxt) +{ + SharedBitmapIndexState *sstate = node->biss_ParallelState; + Size size; + + Assert(node->ss.ps.plan->parallel_aware); + + if (node->biss_ScanDesc) + index_parallelrescan(node->biss_ScanDesc); + + if (sstate == NULL) + return; + + /* Clear the stored per-worker iterators from the previous scan. */ + size = sstate->max_participants * sizeof(dsa_pointer); + memset(sstate->worker_tbmiter, 0, size); + + BarrierInit(&sstate->barrier, 0); +} + /* ---------------------------------------------------------------- * ExecBitmapIndexScanInitializeWorker * @@ -415,7 +598,38 @@ void ExecBitmapIndexScanInitializeWorker(BitmapIndexScanState *node, ParallelWorkerContext *pwcxt) { - /* don't need this if not instrumenting */ + EState *estate = node->ss.ps.state; + ParallelIndexScanDesc piscan; + SharedBitmapIndexState *sstate; + + if (node->ss.ps.plan->parallel_aware) + { + piscan = shm_toc_lookup(pwcxt->toc, + node->ss.ps.plan->plan_node_id, + false); + sstate = shm_toc_lookup(pwcxt->toc, + node->ss.ps.plan->plan_node_id + + PARALLEL_KEY_BITMAP_INDEX_OFFSET, + false); + node->biss_ParallelState = sstate; + + node->biss_ScanDesc = + index_beginscan_bitmap_parallel(node->biss_RelationDesc, + estate->es_snapshot, + node->biss_Instrument, + node->biss_NumScanKeys, + piscan); + + /* + * If no run-time keys to calculate or they are ready, go ahead and pass + * the scankeys to the index AM. + */ + if (node->biss_NumRuntimeKeys == 0 || node->biss_RuntimeKeysReady) + index_rescan(node->biss_ScanDesc, + node->biss_ScanKeys, node->biss_NumScanKeys, + NULL, 0); + } + if (!node->ss.ps.instrument) return; @@ -447,3 +661,95 @@ ExecBitmapIndexScanRetrieveInstrumentation(BitmapIndexScanState *node) node->biss_SharedInfo = palloc(size); memcpy(node->biss_SharedInfo, SharedInfo, size); } + +/* + * Compute the size of the SharedBitmapIndexState for a given number of + * requested workers. + */ +static Size +BitmapIndexScanSharedStateSize(int nworkers) +{ + return add_size(offsetof(SharedBitmapIndexState, worker_tbmiter), + mul_size(nworkers + 1, sizeof(dsa_pointer))); +} + + + +/* + * Given a full per-worker TIDBitmap built by a parallel-aware BitmapIndexScan, + * share it with all other workers and return a new TIDBitmap containing only + * the heap blocks assigned to this worker by a hash of the block number. + * + * The original TIDBitmap is freed here, after all workers have finished using + * the shared iterator state. + */ +static TIDBitmap * +BitmapIndexScanPartition(BitmapIndexScanState *node, TIDBitmap *tbm) +{ + /* + * Use a stable participant ID across multiple Bitmap Index Scans in the + * same query. This ensures consistent partitioning for Bitmap And / Or + * operations. + */ + int my_id = (IsParallelWorker() ? ParallelWorkerNumber + 1 : 0); + SharedBitmapIndexState *sstate = node->biss_ParallelState; + dsa_area *dsa = node->ss.ps.state->es_query_dsa; + int total; + int i; + dsa_pointer dp = InvalidDsaPointer; + TIDBitmap *partition; + + Assert(dsa != NULL); + Assert(my_id < sstate->max_participants); + + /* Share our local bitmap, unless it is empty. */ + if (!tbm_is_empty(tbm)) + dp = tbm_prepare_shared_unordered_iterate(tbm); + sstate->worker_tbmiter[my_id] = dp; + + /* Wait until every participant has stored its bitmap. */ + BarrierArriveAndWait(&sstate->barrier, WAIT_EVENT_PARALLEL_BITMAP_SCAN); + + total = BarrierParticipants(&sstate->barrier); + Assert(total > 0); + + /* Build the per-worker partition. */ + partition = tbm_create(work_mem * (Size) 1024, NULL); + + for (i = 0; i < sstate->max_participants; i++) + { + dsa_pointer slot = sstate->worker_tbmiter[i]; + TBMUnorderedIterator *iter; + TBMIterateResult tbmres; + + if (!DsaPointerIsValid(slot)) + continue; + + /* + * Scan the whole shared bitmap with a private cursor; we must not + * use the joint shared cursor, or the participants would divide the + * bitmap's pages among themselves instead of each scanning them all. + * + * Use the unsorted iterator since order doesn't matter for partition + * assignment - this avoids the complex merge sort logic needed for + * sorted iteration. + */ + iter = tbm_begin_shared_unordered_iterate(dsa, slot); + + while (tbm_shared_unordered_iterate(iter, &tbmres)) + if (murmurhash32(tbmres.blockno / 256) % (uint32) total == (uint32) my_id) + tbm_copy_page(partition, &tbmres); + + tbm_end_shared_unordered_iterate(&iter); + } + + /* Wait until everyone is done reading the shared bitmaps. */ + BarrierArriveAndWait(&sstate->barrier, WAIT_EVENT_PARALLEL_BITMAP_SCAN); + + /* Free the original bitmap and the shared iterator state we created. */ + tbm_free(tbm); + sstate->worker_tbmiter[my_id] = InvalidDsaPointer; + BarrierDetach(&sstate->barrier); + + return partition; +} diff --git a/src/backend/executor/nodeBitmapOr.c b/src/backend/executor/nodeBitmapOr.c index e84f97eebd4..77aea499c64 100644 --- a/src/backend/executor/nodeBitmapOr.c +++ b/src/backend/executor/nodeBitmapOr.c @@ -141,14 +141,13 @@ MultiExecBitmapOr(BitmapOrState *node) * tbm_union step for each child: just pass down the current result * bitmap and let the child OR directly into it. */ - if (IsA(subnode, BitmapIndexScanState)) + if (IsA(subnode, BitmapIndexScanState) && + node->ps.state->es_query_dsa == NULL) { if (result == NULL) /* first subplan */ { /* XXX should we use less than work_mem for this? */ - result = tbm_create(work_mem * (Size) 1024, - ((BitmapOr *) node->ps.plan)->isshared ? - node->ps.state->es_query_dsa : NULL); + result = tbm_create(work_mem * (Size) 1024, NULL); } ((BitmapIndexScanState *) subnode)->biss_result = result; diff --git a/src/backend/nodes/tidbitmap.c b/src/backend/nodes/tidbitmap.c index f1f925cb13b..a342e96226b 100644 --- a/src/backend/nodes/tidbitmap.c +++ b/src/backend/nodes/tidbitmap.c @@ -51,7 +51,7 @@ /* * When we have to switch over to lossy storage, we use a data structure * with one bit per page, where all pages having the same number DIV - * PAGES_PER_CHUNK are aggregated into one chunk. When a chunk is present + * TBM_MAX_PAGES_PER_CHUNK are aggregated into one chunk. When a chunk is present * and has the bit set for a given page, there must not be a per-page entry * for that page in the page table. * @@ -59,11 +59,11 @@ * table, using identical data structures. (This is because the memory * management for hashtables doesn't easily/efficiently allow space to be * transferred easily from one hashtable to another.) Therefore it's best - * if PAGES_PER_CHUNK is the same as TBM_MAX_TUPLES_PER_PAGE, or at least not - * too different. But we also want PAGES_PER_CHUNK to be a power of 2 to + * if TBM_MAX_PAGES_PER_CHUNK is the same as TBM_MAX_TUPLES_PER_PAGE, or at least not + * too different. But we also want TBM_MAX_PAGES_PER_CHUNK to be a power of 2 to * avoid expensive integer remainder operations. So, define it like this: */ -#define PAGES_PER_CHUNK (BLCKSZ / 32) +#define TBM_MAX_PAGES_PER_CHUNK (BLCKSZ / 32) /* We use BITS_PER_BITMAPWORD and typedef bitmapword from nodes/bitmapset.h */ @@ -73,13 +73,13 @@ /* number of active words for an exact page: */ #define WORDS_PER_PAGE ((TBM_MAX_TUPLES_PER_PAGE - 1) / BITS_PER_BITMAPWORD + 1) /* number of active words for a lossy chunk: */ -#define WORDS_PER_CHUNK ((PAGES_PER_CHUNK - 1) / BITS_PER_BITMAPWORD + 1) +#define WORDS_PER_CHUNK ((TBM_MAX_PAGES_PER_CHUNK - 1) / BITS_PER_BITMAPWORD + 1) /* * The hashtable entries are represented by this data structure. For * an exact page, blockno is the page number and bit k of the bitmap * represents tuple offset k+1. For a lossy chunk, blockno is the first - * page in the chunk (this must be a multiple of PAGES_PER_CHUNK) and + * page in the chunk (this must be a multiple of TBM_MAX_PAGES_PER_CHUNK) and * bit k represents page blockno+k. Note that it is not possible to * have exact storage for the first page of a chunk if we are using * lossy storage for any page in the chunk's range, since the same @@ -98,15 +98,6 @@ typedef struct PagetableEntry bitmapword words[Max(WORDS_PER_PAGE, WORDS_PER_CHUNK)]; } PagetableEntry; -/* - * Holds array of pagetable entries. - */ -typedef struct PTEntryArray -{ - pg_atomic_uint32 refcount; /* no. of iterator attached */ - PagetableEntry ptentry[FLEXIBLE_ARRAY_MEMBER]; -} PTEntryArray; - /* * We want to avoid the overhead of creating the hashtable, which is * comparatively large, when not necessary. Particularly when we are using a @@ -156,18 +147,36 @@ struct TIDBitmap PagetableEntry **schunks; /* sorted lossy-chunk list, or NULL */ dsa_pointer dsapagetable; /* dsa_pointer to the element array */ dsa_pointer dsapagetableold; /* dsa_pointer to the old element array */ - dsa_pointer ptpages; /* dsa_pointer to the page array */ - dsa_pointer ptchunks; /* dsa_pointer to the chunk array */ dsa_area *dsa; /* reference to per-query dsa area */ }; +/* define hashtable mapping block numbers to PagetableEntry's */ +#define SH_USE_NONDEFAULT_ALLOCATOR +#define SH_PREFIX pagetable +#define SH_ELEMENT_TYPE PagetableEntry +#define SH_KEY_TYPE BlockNumber +#define SH_KEY blockno +#define SH_HASH_KEY(tb, key) murmurhash32(key) +#define SH_EQUAL(tb, a, b) a == b +#define SH_SCOPE static inline +#define SH_DEFINE +#define SH_DECLARE +#include "lib/simplehash.h" + +#define ST_SORT qsort_pagetable +#define ST_ELEMENT_TYPE_VOID +#define ST_COMPARE(a, b) pg_cmp_u32((*(const PagetableEntry *const *)a)->blockno, (*(const PagetableEntry *const *)b)->blockno) +#define ST_SCOPE static +#define ST_DEFINE +#include "lib/sort_template.h" + /* * When iterating over a backend-local bitmap in sorted order, a - * TBMPrivateIterator is used to track our progress. There can be several + * TBMOrderedIterator is used to track our progress. There can be several * iterators scanning the same bitmap concurrently. Note that the bitmap * becomes read-only as soon as any iterator is created. */ -struct TBMPrivateIterator +struct TBMOrderedIterator { TIDBitmap *tbm; /* TIDBitmap we're iterating over */ int spageptr; /* next spages index */ @@ -181,38 +190,19 @@ struct TBMPrivateIterator */ typedef struct TBMSharedIteratorState { - int nentries; /* number of entries in pagetable */ - int maxentries; /* limit on same to meet maxbytes */ - int npages; /* number of exact entries in pagetable */ - int nchunks; /* number of lossy entries in pagetable */ - dsa_pointer pagetable; /* dsa pointers to head of pagetable data */ - dsa_pointer spages; /* dsa pointer to page array */ - dsa_pointer schunks; /* dsa pointer to chunk array */ - LWLock lock; /* lock to protect below members */ - int spageptr; /* next spages index */ - int schunkptr; /* next schunks index */ - int schunkbit; /* next bit to check in current schunk */ + uint64 size; + uint64 sizemask; + dsa_pointer pagetable; } TBMSharedIteratorState; /* - * pagetable iteration array. - */ -typedef struct PTIterationArray -{ - pg_atomic_uint32 refcount; /* no. of iterator attached */ - int index[FLEXIBLE_ARRAY_MEMBER]; /* index array */ -} PTIterationArray; - -/* - * same as TBMPrivateIterator, but it is used for joint iteration, therefore + * same as TBMOrderedIterator, but it is used for joint iteration, therefore * this also holds a reference to the shared state. */ -struct TBMSharedIterator +struct TBMUnorderedIterator { - TBMSharedIteratorState *state; /* shared state */ - PTEntryArray *ptbase; /* pagetable element array */ - PTIterationArray *ptpages; /* sorted exact page index list */ - PTIterationArray *ptchunks; /* sorted lossy page index list */ + struct pagetable_hash pagetable; + pagetable_iterator iter; }; /* Local function prototypes */ @@ -225,23 +215,6 @@ static PagetableEntry *tbm_get_pageentry(TIDBitmap *tbm, BlockNumber pageno); static bool tbm_page_is_lossy(const TIDBitmap *tbm, BlockNumber pageno); static void tbm_mark_page_lossy(TIDBitmap *tbm, BlockNumber pageno); static void tbm_lossify(TIDBitmap *tbm); -static int tbm_comparator(const void *left, const void *right); -static int tbm_shared_comparator(const void *left, const void *right, - void *arg); - -/* define hashtable mapping block numbers to PagetableEntry's */ -#define SH_USE_NONDEFAULT_ALLOCATOR -#define SH_PREFIX pagetable -#define SH_ELEMENT_TYPE PagetableEntry -#define SH_KEY_TYPE BlockNumber -#define SH_KEY blockno -#define SH_HASH_KEY(tb, key) murmurhash32(key) -#define SH_EQUAL(tb, a, b) a == b -#define SH_SCOPE static inline -#define SH_DEFINE -#define SH_DECLARE -#include "lib/simplehash.h" - /* * tbm_create - create an initially-empty bitmap @@ -268,8 +241,6 @@ tbm_create(Size maxbytes, dsa_area *dsa) tbm->dsa = dsa; tbm->dsapagetable = InvalidDsaPointer; tbm->dsapagetableold = InvalidDsaPointer; - tbm->ptpages = InvalidDsaPointer; - tbm->ptchunks = InvalidDsaPointer; return tbm; } @@ -317,46 +288,11 @@ tbm_free(TIDBitmap *tbm) pfree(tbm->spages); if (tbm->schunks) pfree(tbm->schunks); + if (DsaPointerIsValid(tbm->dsapagetable)) + dsa_free(tbm->dsa, tbm->dsapagetable); pfree(tbm); } -/* - * tbm_free_shared_area - free shared state - * - * Free shared iterator state, Also free shared pagetable and iterator arrays - * memory if they are not referred by any of the shared iterator i.e recount - * is becomes 0. - */ -void -tbm_free_shared_area(dsa_area *dsa, dsa_pointer dp) -{ - TBMSharedIteratorState *istate = dsa_get_address(dsa, dp); - PTEntryArray *ptbase; - PTIterationArray *ptpages; - PTIterationArray *ptchunks; - - if (DsaPointerIsValid(istate->pagetable)) - { - ptbase = dsa_get_address(dsa, istate->pagetable); - if (pg_atomic_sub_fetch_u32(&ptbase->refcount, 1) == 0) - dsa_free(dsa, istate->pagetable); - } - if (DsaPointerIsValid(istate->spages)) - { - ptpages = dsa_get_address(dsa, istate->spages); - if (pg_atomic_sub_fetch_u32(&ptpages->refcount, 1) == 0) - dsa_free(dsa, istate->spages); - } - if (DsaPointerIsValid(istate->schunks)) - { - ptchunks = dsa_get_address(dsa, istate->schunks); - if (pg_atomic_sub_fetch_u32(&ptchunks->refcount, 1) == 0) - dsa_free(dsa, istate->schunks); - } - - dsa_free(dsa, dp); -} - /* * tbm_add_tuples - add some tuple IDs to a TIDBitmap * @@ -438,6 +374,23 @@ tbm_add_page(TIDBitmap *tbm, BlockNumber pageno) tbm_lossify(tbm); } +/* + * tbm_copy_page - Copy a single page worth of TIDs from a TBMIterateResult into a TIDBitmap. + * + * This is optimized for the partitioning case where we know: + * - The destination bitmap is initially empty + * - Each page is inserted exactly once + * - No hash collisions or duplicate entries need to be handled + * + * Therefore, we can reuse tbm_union_page which already handles all the edge cases. + */ +void +tbm_copy_page(TIDBitmap *dest, TBMIterateResult *src) +{ + PagetableEntry *src_page = (PagetableEntry *) src->internal_page; + tbm_union_page(dest, src_page); +} + /* * tbm_union - set union * @@ -660,30 +613,30 @@ tbm_is_empty(const TIDBitmap *tbm) } /* - * tbm_begin_private_iterate - prepare to iterate through a TIDBitmap + * tbm_begin_ordered_iterate - prepare to iterate through a TIDBitmap * - * The TBMPrivateIterator struct is created in the caller's memory context. - * For a clean shutdown of the iteration, call tbm_end_private_iterate; but + * The TBMOrderedIterator struct is created in the caller's memory context. + * For a clean shutdown of the iteration, call tbm_end_ordered_iterate; but * it's okay to just allow the memory context to be released, too. It is - * caller's responsibility not to touch the TBMPrivateIterator anymore once + * caller's responsibility not to touch the TBMOrderedIterator anymore once * the TIDBitmap is freed. * * NB: after this is called, it is no longer allowed to modify the contents * of the bitmap. However, you can call this multiple times to scan the * contents repeatedly, including parallel scans. */ -TBMPrivateIterator * -tbm_begin_private_iterate(TIDBitmap *tbm) +TBMOrderedIterator * +tbm_begin_ordered_iterate(TIDBitmap *tbm) { - TBMPrivateIterator *iterator; + TBMOrderedIterator *iterator; Assert(tbm->iterating != TBM_ITERATING_SHARED); /* - * Create the TBMPrivateIterator struct, with enough trailing space to + * Create the TBMOrderedIterator struct, with enough trailing space to * serve the needs of the TBMIterateResult sub-struct. */ - iterator = palloc_object(TBMPrivateIterator); + iterator = palloc_object(TBMOrderedIterator); iterator->tbm = tbm; /* @@ -727,11 +680,9 @@ tbm_begin_private_iterate(TIDBitmap *tbm) Assert(npages == tbm->npages); Assert(nchunks == tbm->nchunks); if (npages > 1) - qsort(tbm->spages, npages, sizeof(PagetableEntry *), - tbm_comparator); + qsort_pagetable(tbm->spages, npages, sizeof(PagetableEntry *)); if (nchunks > 1) - qsort(tbm->schunks, nchunks, sizeof(PagetableEntry *), - tbm_comparator); + qsort_pagetable(tbm->schunks, nchunks, sizeof(PagetableEntry *)); } tbm->iterating = TBM_ITERATING_PRIVATE; @@ -740,151 +691,51 @@ tbm_begin_private_iterate(TIDBitmap *tbm) } /* - * tbm_prepare_shared_iterate - prepare shared iteration state for a TIDBitmap. + * tbm_prepare_shared_unordered_iterate - prepare shared iteration state for a TIDBitmap. * * The necessary shared state will be allocated from the DSA passed to * tbm_create, so that multiple processes can attach to it and iterate jointly. * * This will convert the pagetable hash into page and chunk array of the index - * into pagetable array. + * into pagetable array. Note that this function does not sort the arrays, + * as the current use case (partition assignment) does not require sorted order. */ dsa_pointer -tbm_prepare_shared_iterate(TIDBitmap *tbm) +tbm_prepare_shared_unordered_iterate(TIDBitmap *tbm) { dsa_pointer dp; TBMSharedIteratorState *istate; - PTEntryArray *ptbase = NULL; - PTIterationArray *ptpages = NULL; - PTIterationArray *ptchunks = NULL; Assert(tbm->dsa != NULL); Assert(tbm->iterating != TBM_ITERATING_PRIVATE); + Assert(tbm->iterating == TBM_NOT_ITERATING); - /* - * Allocate TBMSharedIteratorState from DSA to hold the shared members and - * lock, this will also be used by multiple worker for shared iterate. - */ dp = dsa_allocate0(tbm->dsa, sizeof(TBMSharedIteratorState)); istate = dsa_get_address(tbm->dsa, dp); - /* - * If we're not already iterating, create and fill the sorted page lists. - * (If we are, the sorted page lists are already stored in the TIDBitmap, - * and we can just reuse them.) - */ - if (tbm->iterating == TBM_NOT_ITERATING) + if (tbm->status == TBM_ONE_PAGE) { - pagetable_iterator i; - PagetableEntry *page; - int idx; - int npages; - int nchunks; - /* - * Allocate the page and chunk array memory from the DSA to share - * across multiple processes. + * In one page mode allocate the space for one pagetable entry, + * initialize it, and directly store its index (i.e. 0) in the + * page array. */ - if (tbm->npages) - { - tbm->ptpages = dsa_allocate(tbm->dsa, sizeof(PTIterationArray) + - tbm->npages * sizeof(int)); - ptpages = dsa_get_address(tbm->dsa, tbm->ptpages); - pg_atomic_init_u32(&ptpages->refcount, 0); - } - if (tbm->nchunks) - { - tbm->ptchunks = dsa_allocate(tbm->dsa, sizeof(PTIterationArray) + - tbm->nchunks * sizeof(int)); - ptchunks = dsa_get_address(tbm->dsa, tbm->ptchunks); - pg_atomic_init_u32(&ptchunks->refcount, 0); - } + void *ptbase; - /* - * If TBM status is TBM_HASH then iterate over the pagetable and - * convert it to page and chunk arrays. But if it's in the - * TBM_ONE_PAGE mode then directly allocate the space for one entry - * from the DSA. - */ - npages = nchunks = 0; - if (tbm->status == TBM_HASH) - { - ptbase = dsa_get_address(tbm->dsa, tbm->dsapagetable); - - pagetable_start_iterate(tbm->pagetable, &i); - while ((page = pagetable_iterate(tbm->pagetable, &i)) != NULL) - { - idx = page - ptbase->ptentry; - if (page->ischunk) - ptchunks->index[nchunks++] = idx; - else - ptpages->index[npages++] = idx; - } - - Assert(npages == tbm->npages); - Assert(nchunks == tbm->nchunks); - } - else if (tbm->status == TBM_ONE_PAGE) - { - /* - * In one page mode allocate the space for one pagetable entry, - * initialize it, and directly store its index (i.e. 0) in the - * page array. - */ - tbm->dsapagetable = dsa_allocate(tbm->dsa, sizeof(PTEntryArray) + - sizeof(PagetableEntry)); - ptbase = dsa_get_address(tbm->dsa, tbm->dsapagetable); - memcpy(ptbase->ptentry, &tbm->entry1, sizeof(PagetableEntry)); - ptpages->index[0] = 0; - } - - if (ptbase != NULL) - pg_atomic_init_u32(&ptbase->refcount, 0); - if (npages > 1) - qsort_arg(ptpages->index, npages, sizeof(int), - tbm_shared_comparator, ptbase->ptentry); - if (nchunks > 1) - qsort_arg(ptchunks->index, nchunks, sizeof(int), - tbm_shared_comparator, ptbase->ptentry); + tbm->dsapagetable = dsa_allocate(tbm->dsa, sizeof(PagetableEntry)); + ptbase = dsa_get_address(tbm->dsa, tbm->dsapagetable); + memcpy(ptbase, &tbm->entry1, sizeof(PagetableEntry)); + istate->size = 1; + istate->sizemask = 1; + } + else if (tbm->status == TBM_HASH) + { + istate->size = tbm->pagetable->size; + istate->sizemask = tbm->pagetable->sizemask; } - /* - * Store the TBM members in the shared state so that we can share them - * across multiple processes. - */ - istate->nentries = tbm->nentries; - istate->maxentries = tbm->maxentries; - istate->npages = tbm->npages; - istate->nchunks = tbm->nchunks; istate->pagetable = tbm->dsapagetable; - istate->spages = tbm->ptpages; - istate->schunks = tbm->ptchunks; - - ptbase = dsa_get_address(tbm->dsa, tbm->dsapagetable); - ptpages = dsa_get_address(tbm->dsa, tbm->ptpages); - ptchunks = dsa_get_address(tbm->dsa, tbm->ptchunks); - - /* - * For every shared iterator referring to pagetable and iterator array, - * increase the refcount by 1 so that while freeing the shared iterator we - * don't free pagetable and iterator array until its refcount becomes 0. - */ - if (ptbase != NULL) - pg_atomic_add_fetch_u32(&ptbase->refcount, 1); - if (ptpages != NULL) - pg_atomic_add_fetch_u32(&ptpages->refcount, 1); - if (ptchunks != NULL) - pg_atomic_add_fetch_u32(&ptchunks->refcount, 1); - - /* Initialize the iterator lock */ - LWLockInitialize(&istate->lock, LWTRANCHE_SHARED_TIDBITMAP); - - /* Initialize the shared iterator state */ - istate->schunkbit = 0; - istate->schunkptr = 0; - istate->spageptr = 0; - tbm->iterating = TBM_ITERATING_SHARED; - return dp; } @@ -936,7 +787,7 @@ tbm_advance_schunkbit(PagetableEntry *chunk, int *schunkbitp) { int schunkbit = *schunkbitp; - while (schunkbit < PAGES_PER_CHUNK) + while (schunkbit < TBM_MAX_PAGES_PER_CHUNK) { int wordnum = WORDNUM(schunkbit); int bitnum = BITNUM(schunkbit); @@ -950,7 +801,7 @@ tbm_advance_schunkbit(PagetableEntry *chunk, int *schunkbitp) } /* - * tbm_private_iterate - scan through next page of a TIDBitmap + * tbm_ordered_iterate - scan through next page of a TIDBitmap * * Caller must pass in a TBMIterateResult to be filled. * @@ -971,7 +822,7 @@ tbm_advance_schunkbit(PagetableEntry *chunk, int *schunkbitp) * always set true when lossy is true.) */ bool -tbm_private_iterate(TBMPrivateIterator *iterator, TBMIterateResult *tbmres) +tbm_ordered_iterate(TBMOrderedIterator *iterator, TBMIterateResult *tbmres) { TIDBitmap *tbm = iterator->tbm; @@ -987,7 +838,7 @@ tbm_private_iterate(TBMPrivateIterator *iterator, TBMIterateResult *tbmres) int schunkbit = iterator->schunkbit; tbm_advance_schunkbit(chunk, &schunkbit); - if (schunkbit < PAGES_PER_CHUNK) + if (schunkbit < TBM_MAX_PAGES_PER_CHUNK) { iterator->schunkbit = schunkbit; break; @@ -1044,121 +895,58 @@ tbm_private_iterate(TBMPrivateIterator *iterator, TBMIterateResult *tbmres) } /* - * tbm_shared_iterate - scan through next page of a TIDBitmap + * tbm_shared_unordered_iterate - scan through next page of a shared TIDBitmap * - * As above, but this will iterate using an iterator which is shared - * across multiple processes. We need to acquire the iterator LWLock, - * before accessing the shared members. + * This iterates through the pagetable entries in their natural (unsorted) order, + * which is much simpler than the sorted case since we don't need to merge + * pages and chunks in numerical order. This is useful for operations like + * partition assignment where only the block number matters, not the order. + * + * Like tbm_shared_iterate_private, the cursor is kept in the backend-private + * TBMUnorderedIterator and no lock is taken. */ bool -tbm_shared_iterate(TBMSharedIterator *iterator, TBMIterateResult *tbmres) +tbm_shared_unordered_iterate(TBMUnorderedIterator *iterator, TBMIterateResult *tbmres) { - TBMSharedIteratorState *istate = iterator->state; - PagetableEntry *ptbase = NULL; - int *idxpages = NULL; - int *idxchunks = NULL; - - if (iterator->ptbase != NULL) - ptbase = iterator->ptbase->ptentry; - if (iterator->ptpages != NULL) - idxpages = iterator->ptpages->index; - if (iterator->ptchunks != NULL) - idxchunks = iterator->ptchunks->index; - - /* Acquire the LWLock before accessing the shared members */ - LWLockAcquire(&istate->lock, LW_EXCLUSIVE); + PagetableEntry *page = pagetable_iterate(&iterator->pagetable, &iterator->iter); - /* - * If lossy chunk pages remain, make sure we've advanced schunkptr/ - * schunkbit to the next set bit. - */ - while (istate->schunkptr < istate->nchunks) + if (page != NULL) { - PagetableEntry *chunk = &ptbase[idxchunks[istate->schunkptr]]; - int schunkbit = istate->schunkbit; - - tbm_advance_schunkbit(chunk, &schunkbit); - if (schunkbit < PAGES_PER_CHUNK) - { - istate->schunkbit = schunkbit; - break; - } - /* advance to next chunk */ - istate->schunkptr++; - istate->schunkbit = 0; - } - - /* - * If both chunk and per-page data remain, must output the numerically - * earlier page. - */ - if (istate->schunkptr < istate->nchunks) - { - PagetableEntry *chunk = &ptbase[idxchunks[istate->schunkptr]]; - BlockNumber chunk_blockno; - - chunk_blockno = chunk->blockno + istate->schunkbit; - - if (istate->spageptr >= istate->npages || - chunk_blockno < ptbase[idxpages[istate->spageptr]].blockno) - { - /* Return a lossy page indicator from the chunk */ - tbmres->blockno = chunk_blockno; - tbmres->lossy = true; - tbmres->recheck = true; - tbmres->internal_page = NULL; - istate->schunkbit++; - - LWLockRelease(&istate->lock); - return true; - } - } - - if (istate->spageptr < istate->npages) - { - PagetableEntry *page = &ptbase[idxpages[istate->spageptr]]; - tbmres->internal_page = page; tbmres->blockno = page->blockno; - tbmres->lossy = false; + tbmres->lossy = page->ischunk; tbmres->recheck = page->recheck; - istate->spageptr++; - - LWLockRelease(&istate->lock); - return true; } - LWLockRelease(&istate->lock); - - /* Nothing more in the bitmap */ - tbmres->blockno = InvalidBlockNumber; return false; } /* - * tbm_end_private_iterate - finish an iteration over a TIDBitmap + * tbm_end_ordered_iterate - finish an iteration over a TIDBitmap * * Currently this is just a pfree, but it might do more someday. (For * instance, it could be useful to count open iterators and allow the * bitmap to return to read/write status when there are no more iterators.) */ void -tbm_end_private_iterate(TBMPrivateIterator *iterator) +tbm_end_ordered_iterate(TBMOrderedIterator **iterator) { - pfree(iterator); + pfree(*iterator); + *iterator = NULL; } /* - * tbm_end_shared_iterate - finish a shared iteration over a TIDBitmap + * tbm_end_shared_unordered_iterate - finish a shared iteration over a TIDBitmap * * This doesn't free any of the shared state associated with the iterator, * just our backend-private state. */ void -tbm_end_shared_iterate(TBMSharedIterator *iterator) +tbm_end_shared_unordered_iterate(TBMUnorderedIterator **iterator) { - pfree(iterator); + pfree(*iterator); + *iterator = NULL; } /* @@ -1258,7 +1046,7 @@ tbm_page_is_lossy(const TIDBitmap *tbm, BlockNumber pageno) return false; Assert(tbm->status == TBM_HASH); - bitno = pageno % PAGES_PER_CHUNK; + bitno = pageno % TBM_MAX_PAGES_PER_CHUNK; chunk_pageno = pageno - bitno; page = pagetable_lookup(tbm->pagetable, chunk_pageno); @@ -1294,7 +1082,7 @@ tbm_mark_page_lossy(TIDBitmap *tbm, BlockNumber pageno) if (tbm->status != TBM_HASH) tbm_create_pagetable(tbm); - bitno = pageno % PAGES_PER_CHUNK; + bitno = pageno % TBM_MAX_PAGES_PER_CHUNK; chunk_pageno = pageno - bitno; /* @@ -1380,7 +1168,7 @@ tbm_lossify(TIDBitmap *tbm) * If the page would become a chunk header, we won't save anything by * converting it to lossy, so skip it. */ - if ((page->blockno % PAGES_PER_CHUNK) == 0) + if ((page->blockno % TBM_MAX_PAGES_PER_CHUNK) == 0) continue; /* This does the dirty work ... */ @@ -1419,38 +1207,7 @@ tbm_lossify(TIDBitmap *tbm) } /* - * qsort comparator to handle PagetableEntry pointers. - */ -static int -tbm_comparator(const void *left, const void *right) -{ - BlockNumber l = (*((PagetableEntry *const *) left))->blockno; - BlockNumber r = (*((PagetableEntry *const *) right))->blockno; - - return pg_cmp_u32(l, r); -} - -/* - * As above, but this will get index into PagetableEntry array. Therefore, - * it needs to get actual PagetableEntry using the index before comparing the - * blockno. - */ -static int -tbm_shared_comparator(const void *left, const void *right, void *arg) -{ - PagetableEntry *base = (PagetableEntry *) arg; - PagetableEntry *lpage = &base[*(const int *) left]; - PagetableEntry *rpage = &base[*(const int *) right]; - - if (lpage->blockno < rpage->blockno) - return -1; - else if (lpage->blockno > rpage->blockno) - return 1; - return 0; -} - -/* - * tbm_attach_shared_iterate + * tbm_begin_shared_unordered_iterate * * Allocate a backend-private iterator and attach the shared iterator state * to it so that multiple processed can iterate jointly. @@ -1458,28 +1215,16 @@ tbm_shared_comparator(const void *left, const void *right, void *arg) * We also converts the DSA pointers to local pointers and store them into * our private iterator. */ -TBMSharedIterator * -tbm_attach_shared_iterate(dsa_area *dsa, dsa_pointer dp) +TBMUnorderedIterator * +tbm_begin_shared_unordered_iterate(dsa_area *dsa, dsa_pointer dp) { - TBMSharedIterator *iterator; - TBMSharedIteratorState *istate; - - /* - * Create the TBMSharedIterator struct, with enough trailing space to - * serve the needs of the TBMIterateResult sub-struct. - */ - iterator = palloc0_object(TBMSharedIterator); - - istate = (TBMSharedIteratorState *) dsa_get_address(dsa, dp); - - iterator->state = istate; - - iterator->ptbase = dsa_get_address(dsa, istate->pagetable); + TBMSharedIteratorState *istate = (TBMSharedIteratorState *)dsa_get_address(dsa, dp); - if (istate->npages) - iterator->ptpages = dsa_get_address(dsa, istate->spages); - if (istate->nchunks) - iterator->ptchunks = dsa_get_address(dsa, istate->schunks); + TBMUnorderedIterator *iterator = palloc_object(TBMUnorderedIterator); + iterator->pagetable.size = istate->size; + iterator->pagetable.sizemask = istate->sizemask; + iterator->pagetable.data = dsa_get_address(dsa, istate->pagetable); + pagetable_start_iterate(&iterator->pagetable, &iterator->iter); return iterator; } @@ -1494,7 +1239,6 @@ static inline void * pagetable_allocate(pagetable_hash *pagetable, Size size) { TIDBitmap *tbm = (TIDBitmap *) pagetable->private_data; - PTEntryArray *ptbase; if (tbm->dsa == NULL) return MemoryContextAllocExtended(pagetable->ctx, size, @@ -1505,12 +1249,8 @@ pagetable_allocate(pagetable_hash *pagetable, Size size) * new memory so that pagetable_free can free the old entry. */ tbm->dsapagetableold = tbm->dsapagetable; - tbm->dsapagetable = dsa_allocate_extended(tbm->dsa, - sizeof(PTEntryArray) + size, - DSA_ALLOC_HUGE | DSA_ALLOC_ZERO); - ptbase = dsa_get_address(tbm->dsa, tbm->dsapagetable); - - return ptbase->ptentry; + tbm->dsapagetable = dsa_allocate_extended(tbm->dsa, size, DSA_ALLOC_HUGE | DSA_ALLOC_ZERO); + return dsa_get_address(tbm->dsa, tbm->dsapagetable); } /* @@ -1556,68 +1296,3 @@ tbm_calculate_entries(Size maxbytes) return (int) nbuckets; } - -/* - * Create a shared or private bitmap iterator and start iteration. - * - * `tbm` is only used to create the private iterator and dsa and dsp are only - * used to create the shared iterator. - * - * Before invoking tbm_begin_iterate() to create a shared iterator, one - * process must already have invoked tbm_prepare_shared_iterate() to create - * and set up the TBMSharedIteratorState. - */ -TBMIterator -tbm_begin_iterate(TIDBitmap *tbm, dsa_area *dsa, dsa_pointer dsp) -{ - TBMIterator iterator = {0}; - - /* Allocate a private iterator and attach the shared state to it */ - if (DsaPointerIsValid(dsp)) - { - iterator.shared = true; - iterator.i.shared_iterator = tbm_attach_shared_iterate(dsa, dsp); - } - else - { - iterator.shared = false; - iterator.i.private_iterator = tbm_begin_private_iterate(tbm); - } - - return iterator; -} - -/* - * Clean up shared or private bitmap iterator. - */ -void -tbm_end_iterate(TBMIterator *iterator) -{ - Assert(iterator && !tbm_exhausted(iterator)); - - if (iterator->shared) - tbm_end_shared_iterate(iterator->i.shared_iterator); - else - tbm_end_private_iterate(iterator->i.private_iterator); - - *iterator = (TBMIterator) - { - 0 - }; -} - -/* - * Populate the next TBMIterateResult using the shared or private bitmap - * iterator. Returns false when there is nothing more to scan. - */ -bool -tbm_iterate(TBMIterator *iterator, TBMIterateResult *tbmres) -{ - Assert(iterator); - Assert(tbmres); - - if (iterator->shared) - return tbm_shared_iterate(iterator->i.shared_iterator, tbmres); - else - return tbm_private_iterate(iterator->i.private_iterator, tbmres); -} diff --git a/src/backend/optimizer/path/allpaths.c b/src/backend/optimizer/path/allpaths.c index de8f29c2b57..13d1bc4e723 100644 --- a/src/backend/optimizer/path/allpaths.c +++ b/src/backend/optimizer/path/allpaths.c @@ -4919,6 +4919,90 @@ remove_unused_subquery_outputs(Query *subquery, RelOptInfo *rel, } } +/* + * build_parallel_bitmapqual + * Build a parallel-aware version of a bitmap qualification tree by + * replacing each leaf IndexPath with a partial (parallel-aware) IndexPath. + * This is used to create the bitmap qual for a partial BitmapHeapPath so + * that each worker can build its own bitmap in parallel. + * + * If a leaf IndexPath cannot be turned into a parallel-aware path (e.g. + * because no workers can be assigned), this function returns NULL and the + * caller must not create a partial BitmapHeapPath from this bitmap qual. + */ +static Path * +build_parallel_bitmapqual(PlannerInfo *root, RelOptInfo *rel, Path *path) +{ + if (IsA(path, IndexPath)) + { + IndexPath *ipath = (IndexPath *) path; + Relids required_outer = PATH_REQ_OUTER((Path *) ipath); + + /* + * Partial bitmap heap paths are not parameterized, so loop_count is 1. + */ + ipath = (IndexPath *) create_index_path(root, + ipath->indexinfo, + ipath->indexclauses, + NIL, /* no ordering for bitmap */ + NIL, + NIL, /* bitmap scans unordered */ + ForwardScanDirection, + false, /* never index-only */ + required_outer, + 1.0, + true); /* partial path */ + + /* + * If the path is not parallel-aware (no workers could be assigned), + * the partial bitmap heap path must not be created. + */ + if (!ipath->path.parallel_aware) + { + pfree(ipath); + return NULL; + } + + return (Path *) ipath; + } + else if (IsA(path, BitmapAndPath)) + { + BitmapAndPath *apath = (BitmapAndPath *) path; + List *newquals = NIL; + ListCell *lc; + + foreach(lc, apath->bitmapquals) + { + Path *sub = build_parallel_bitmapqual(root, rel, + (Path *) lfirst(lc)); + + if (sub == NULL) + return NULL; + newquals = lappend(newquals, sub); + } + return (Path *) create_bitmap_and_path(root, rel, newquals); + } + else if (IsA(path, BitmapOrPath)) + { + BitmapOrPath *opath = (BitmapOrPath *) path; + List *newquals = NIL; + ListCell *lc; + + foreach(lc, opath->bitmapquals) + { + Path *sub = build_parallel_bitmapqual(root, rel, + (Path *) lfirst(lc)); + + if (sub == NULL) + return NULL; + newquals = lappend(newquals, sub); + } + return (Path *) create_bitmap_or_path(root, rel, newquals); + } + else + elog(ERROR, "unrecognized node type: %d", nodeTag(path)); +} + /* * create_partial_bitmap_paths * Build partial bitmap heap path for the relation @@ -4929,6 +5013,11 @@ create_partial_bitmap_paths(PlannerInfo *root, RelOptInfo *rel, { int parallel_workers; double pages_fetched; + Path *parbitmapqual; + + /* A partial path cannot use a parameterized bitmap qual */ + if (bitmapqual->param_info != NULL) + return; /* Compute heap pages for bitmap heap scan */ pages_fetched = compute_bitmap_pages(root, rel, bitmapqual, 1.0, @@ -4940,8 +5029,21 @@ create_partial_bitmap_paths(PlannerInfo *root, RelOptInfo *rel, if (parallel_workers <= 0) return; + /* + * Build a parallel-aware version of the bitmapqual so that each worker + * can run the bitmap index scan(s) in parallel. The original non-parallel + * bitmapqual is used by the regular BitmapHeapPath. + */ + parbitmapqual = build_parallel_bitmapqual(root, rel, bitmapqual); + + if (parbitmapqual == NULL) + return; + add_partial_path(rel, (Path *) create_bitmap_heap_path(root, rel, - bitmapqual, rel->lateral_relids, 1.0, parallel_workers)); + parbitmapqual, + rel->lateral_relids, + 1.0, + parallel_workers)); } /* diff --git a/src/backend/optimizer/path/indxpath.c b/src/backend/optimizer/path/indxpath.c index 3f5d4fa3182..74dbdffa802 100644 --- a/src/backend/optimizer/path/indxpath.c +++ b/src/backend/optimizer/path/indxpath.c @@ -103,11 +103,13 @@ static bool eclass_already_used(EquivalenceClass *parent_ec, Relids oldrelids, static void get_index_paths(PlannerInfo *root, RelOptInfo *rel, IndexOptInfo *index, IndexClauseSet *clauses, List **bitindexpaths); -static List *build_index_paths(PlannerInfo *root, RelOptInfo *rel, +static void build_index_paths(PlannerInfo *root, RelOptInfo *rel, IndexOptInfo *index, IndexClauseSet *clauses, bool useful_predicate, ScanTypeControl scantype, - bool *skip_nonnative_saop); + bool *skip_nonnative_saop, + List **paths, + List **partial_paths); static List *build_paths_for_OR(PlannerInfo *root, RelOptInfo *rel, List *clauses, List *other_clauses); static List *generate_bitmap_or_paths(PlannerInfo *root, RelOptInfo *rel, @@ -716,19 +718,23 @@ get_index_paths(PlannerInfo *root, RelOptInfo *rel, IndexOptInfo *index, IndexClauseSet *clauses, List **bitindexpaths) { - List *indexpaths; bool skip_nonnative_saop = false; ListCell *lc; + List *paths = NIL; + List *partial_paths = NIL; + List *saoppaths = NIL; /* * Build simple index paths using the clauses. Allow ScalarArrayOpExpr * clauses only if the index AM supports them natively. */ - indexpaths = build_index_paths(root, rel, - index, clauses, - index->predOK, - ST_ANYSCAN, - &skip_nonnative_saop); + build_index_paths(root, rel, + index, clauses, + index->predOK, + ST_ANYSCAN, + &skip_nonnative_saop, + &paths, + &partial_paths); /* * Submit all the ones that can form plain IndexScan plans to add_path. (A @@ -742,12 +748,12 @@ get_index_paths(PlannerInfo *root, RelOptInfo *rel, * only interested in paths that have some selectivity; we should discard * anything that was generated solely for ordering purposes. */ - foreach(lc, indexpaths) + foreach(lc, paths) { IndexPath *ipath = (IndexPath *) lfirst(lc); if (index->amhasgettuple) - add_path(rel, (Path *) ipath); + add_path(rel, (Path *) ipath); if (index->amhasgetbitmap && (ipath->path.pathkeys == NIL || @@ -755,6 +761,20 @@ get_index_paths(PlannerInfo *root, RelOptInfo *rel, *bitindexpaths = lappend(*bitindexpaths, ipath); } + /* + * Partial paths generated above are intended only for plain parallel + * index scans; do not reuse them as bitmap-qual children. Bitmap-qual + * parallelism is represented at the BitmapHeapPath level, or by building + * a separate parallel-aware bitmapqual tree when needed. + */ + foreach(lc, partial_paths) + { + IndexPath *ipath = (IndexPath *) lfirst(lc); + + if (index->amhasgettuple) + add_partial_path(rel, (Path *) ipath); + } + /* * If there were ScalarArrayOpExpr clauses that the index can't handle * natively, generate bitmap scan paths relying on executor-managed @@ -762,12 +782,14 @@ get_index_paths(PlannerInfo *root, RelOptInfo *rel, */ if (skip_nonnative_saop) { - indexpaths = build_index_paths(root, rel, - index, clauses, - false, - ST_BITMAPSCAN, - NULL); - *bitindexpaths = list_concat(*bitindexpaths, indexpaths); + build_index_paths(root, rel, + index, clauses, + false, + ST_BITMAPSCAN, + NULL, + &saoppaths, + NULL); + *bitindexpaths = list_concat(*bitindexpaths, saoppaths); } } @@ -805,14 +827,15 @@ get_index_paths(PlannerInfo *root, RelOptInfo *rel, * 'scantype' indicates whether we need plain or bitmap scan support * 'skip_nonnative_saop' indicates whether to accept SAOP if index AM doesn't */ -static List * +void build_index_paths(PlannerInfo *root, RelOptInfo *rel, IndexOptInfo *index, IndexClauseSet *clauses, bool useful_predicate, ScanTypeControl scantype, - bool *skip_nonnative_saop) + bool *skip_nonnative_saop, + List **paths, + List **partial_paths) { - List *result = NIL; IndexPath *ipath; List *index_clauses; Relids outer_relids; @@ -835,11 +858,11 @@ build_index_paths(PlannerInfo *root, RelOptInfo *rel, { case ST_INDEXSCAN: if (!index->amhasgettuple) - return NIL; + return; break; case ST_BITMAPSCAN: if (!index->amhasgetbitmap) - return NIL; + return; break; case ST_ANYSCAN: /* either or both are OK */ @@ -898,7 +921,7 @@ build_index_paths(PlannerInfo *root, RelOptInfo *rel, * clauses.) */ if (index_clauses == NIL && !index->amoptionalkey) - return NIL; + return; } /* We do not want the index's rel itself listed in outer_relids */ @@ -975,15 +998,16 @@ build_index_paths(PlannerInfo *root, RelOptInfo *rel, outer_relids, loop_count, false); - result = lappend(result, ipath); + *paths = lappend(*paths, ipath); /* * If appropriate, consider parallel index scan. We don't allow * parallel index scan for bitmap index scans. */ - if (index->amcanparallel && + if (index->amcanparallel && index->amhasgettuple && rel->consider_parallel && outer_relids == NULL && - scantype != ST_BITMAPSCAN) + scantype != ST_BITMAPSCAN && + partial_paths != NULL) { ipath = create_index_path(root, index, index_clauses, @@ -1001,7 +1025,7 @@ build_index_paths(PlannerInfo *root, RelOptInfo *rel, * parallel workers, just free it. */ if (ipath->path.parallel_workers > 0) - add_partial_path(rel, (Path *) ipath); + *partial_paths = lappend(*partial_paths, ipath); else pfree(ipath); } @@ -1028,12 +1052,13 @@ build_index_paths(PlannerInfo *root, RelOptInfo *rel, outer_relids, loop_count, false); - result = lappend(result, ipath); + *paths = lappend(*paths, ipath); /* If appropriate, consider parallel index scan */ - if (index->amcanparallel && + if (index->amcanparallel && index->amhasgettuple && rel->consider_parallel && outer_relids == NULL && - scantype != ST_BITMAPSCAN) + scantype != ST_BITMAPSCAN && + partial_paths != NULL) { ipath = create_index_path(root, index, index_clauses, @@ -1051,14 +1076,12 @@ build_index_paths(PlannerInfo *root, RelOptInfo *rel, * using parallel workers, just free it. */ if (ipath->path.parallel_workers > 0) - add_partial_path(rel, (Path *) ipath); + *partial_paths = lappend(*partial_paths, ipath); else pfree(ipath); } } } - - return result; } /* @@ -1099,7 +1122,7 @@ build_paths_for_OR(PlannerInfo *root, RelOptInfo *rel, { IndexOptInfo *index = (IndexOptInfo *) lfirst(lc); IndexClauseSet clauseset; - List *indexpaths; + List *indexpaths = NIL; bool useful_predicate; /* Ignore index if it doesn't support bitmap scans */ @@ -1160,11 +1183,13 @@ build_paths_for_OR(PlannerInfo *root, RelOptInfo *rel, /* * Construct paths if possible. */ - indexpaths = build_index_paths(root, rel, - index, &clauseset, - useful_predicate, - ST_BITMAPSCAN, - NULL); + build_index_paths(root, rel, + index, &clauseset, + useful_predicate, + ST_BITMAPSCAN, + NULL, + &indexpaths, + NULL); result = list_concat(result, indexpaths); } diff --git a/src/backend/optimizer/plan/createplan.c b/src/backend/optimizer/plan/createplan.c index 2f8a08ad78e..6b42f355be3 100644 --- a/src/backend/optimizer/plan/createplan.c +++ b/src/backend/optimizer/plan/createplan.c @@ -128,7 +128,6 @@ static BitmapHeapScan *create_bitmap_scan_plan(PlannerInfo *root, List *tlist, List *scan_clauses); static Plan *create_bitmap_subplan(PlannerInfo *root, Path *bitmapqual, List **qual, List **indexqual, List **indexECs); -static void bitmap_subplan_mark_shared(Plan *plan); static TidScan *create_tidscan_plan(PlannerInfo *root, TidPath *best_path, List *tlist, List *scan_clauses); static TidRangeScan *create_tidrangescan_plan(PlannerInfo *root, @@ -3067,9 +3066,6 @@ create_bitmap_scan_plan(PlannerInfo *root, &bitmapqualorig, &indexquals, &indexECs); - if (best_path->path.parallel_aware) - bitmap_subplan_mark_shared(bitmapqualplan); - /* * The qpqual list must contain all restrictions not automatically handled * by the index, other than pseudoconstant clauses which will be handled @@ -3217,7 +3213,7 @@ create_bitmap_subplan(PlannerInfo *root, Path *bitmapqual, plan->plan_rows = clamp_row_est(apath->bitmapselectivity * apath->path.parent->tuples); plan->plan_width = 0; /* meaningless */ - plan->parallel_aware = false; + plan->parallel_aware = apath->path.parallel_aware; plan->parallel_safe = apath->path.parallel_safe; *qual = subquals; *indexqual = subindexquals; @@ -3281,7 +3277,7 @@ create_bitmap_subplan(PlannerInfo *root, Path *bitmapqual, plan->plan_rows = clamp_row_est(opath->bitmapselectivity * opath->path.parent->tuples); plan->plan_width = 0; /* meaningless */ - plan->parallel_aware = false; + plan->parallel_aware = opath->path.parallel_aware; plan->parallel_safe = opath->path.parallel_safe; } @@ -3328,7 +3324,7 @@ create_bitmap_subplan(PlannerInfo *root, Path *bitmapqual, plan->plan_rows = clamp_row_est(ipath->indexselectivity * ipath->path.parent->tuples); plan->plan_width = 0; /* meaningless */ - plan->parallel_aware = false; + plan->parallel_aware = ipath->path.parallel_aware; plan->parallel_safe = ipath->path.parallel_safe; /* Extract original index clauses, actual index quals, relevant ECs */ subquals = NIL; @@ -5508,27 +5504,6 @@ label_incrementalsort_with_costsize(PlannerInfo *root, IncrementalSort *plan, plan->sort.plan.parallel_safe = lefttree->parallel_safe; } -/* - * bitmap_subplan_mark_shared - * Set isshared flag in bitmap subplan so that it will be created in - * shared memory. - */ -static void -bitmap_subplan_mark_shared(Plan *plan) -{ - if (IsA(plan, BitmapAnd)) - bitmap_subplan_mark_shared(linitial(((BitmapAnd *) plan)->bitmapplans)); - else if (IsA(plan, BitmapOr)) - { - ((BitmapOr *) plan)->isshared = true; - bitmap_subplan_mark_shared(linitial(((BitmapOr *) plan)->bitmapplans)); - } - else if (IsA(plan, BitmapIndexScan)) - ((BitmapIndexScan *) plan)->isshared = true; - else - elog(ERROR, "unrecognized node type: %d", nodeTag(plan)); -} - /***************************************************************************** * * PLAN NODE BUILDING ROUTINES diff --git a/src/backend/optimizer/util/pathnode.c b/src/backend/optimizer/util/pathnode.c index 67c47fd7c6c..6813a01323d 100644 --- a/src/backend/optimizer/util/pathnode.c +++ b/src/backend/optimizer/util/pathnode.c @@ -1193,8 +1193,9 @@ create_bitmap_and_path(PlannerInfo *root, /* * Identify the required outer rels as the union of what the child paths - * depend on. (Alternatively, we could insist that the caller pass this - * in, but it's more convenient and reliable to compute it here.) + * depend on, and propagate parallel-aware from children. If any child is + * parallel-aware, this AND must be too, because a parallel-aware child + * can only appear inside a parallel-aware plan. */ foreach(lc, bitmapquals) { @@ -1202,19 +1203,22 @@ create_bitmap_and_path(PlannerInfo *root, required_outer = bms_add_members(required_outer, PATH_REQ_OUTER(bitmapqual)); + if (bitmapqual->parallel_aware) + { + pathnode->path.parallel_aware = true; + if (bitmapqual->parallel_workers > pathnode->path.parallel_workers) + pathnode->path.parallel_workers = bitmapqual->parallel_workers; + } } pathnode->path.param_info = get_baserel_parampathinfo(root, rel, required_outer); /* - * Currently, a BitmapHeapPath, BitmapAndPath, or BitmapOrPath will be - * parallel-safe if and only if rel->consider_parallel is set. So, we can - * set the flag for this path based only on the relation-level flag, - * without actually iterating over the list of children. + * A BitmapHeapPath, BitmapAndPath, or BitmapOrPath will be parallel-safe if + * and only if rel->consider_parallel is set. If we found any + * parallel-aware children above, this path must be parallel-aware as well. */ - pathnode->path.parallel_aware = false; pathnode->path.parallel_safe = rel->consider_parallel; - pathnode->path.parallel_workers = 0; pathnode->path.pathkeys = NIL; /* always unordered */ @@ -1245,8 +1249,9 @@ create_bitmap_or_path(PlannerInfo *root, /* * Identify the required outer rels as the union of what the child paths - * depend on. (Alternatively, we could insist that the caller pass this - * in, but it's more convenient and reliable to compute it here.) + * depend on, and propagate parallel-aware from children. If any child is + * parallel-aware, this OR must be too, because a parallel-aware child can + * only appear inside a parallel-aware plan. */ foreach(lc, bitmapquals) { @@ -1254,19 +1259,22 @@ create_bitmap_or_path(PlannerInfo *root, required_outer = bms_add_members(required_outer, PATH_REQ_OUTER(bitmapqual)); + if (bitmapqual->parallel_aware) + { + pathnode->path.parallel_aware = true; + if (bitmapqual->parallel_workers > pathnode->path.parallel_workers) + pathnode->path.parallel_workers = bitmapqual->parallel_workers; + } } pathnode->path.param_info = get_baserel_parampathinfo(root, rel, required_outer); /* - * Currently, a BitmapHeapPath, BitmapAndPath, or BitmapOrPath will be - * parallel-safe if and only if rel->consider_parallel is set. So, we can - * set the flag for this path based only on the relation-level flag, - * without actually iterating over the list of children. + * A BitmapHeapPath, BitmapAndPath, or BitmapOrPath will be parallel-safe if + * and only if rel->consider_parallel is set. If we found any + * parallel-aware children above, this path must be parallel-aware as well. */ - pathnode->path.parallel_aware = false; pathnode->path.parallel_safe = rel->consider_parallel; - pathnode->path.parallel_workers = 0; pathnode->path.pathkeys = NIL; /* always unordered */ diff --git a/src/include/access/genam.h b/src/include/access/genam.h index 1bfdc1ef6b1..31e1a28b256 100644 --- a/src/include/access/genam.h +++ b/src/include/access/genam.h @@ -166,6 +166,12 @@ extern IndexScanDesc index_beginscan_bitmap(Relation indexRelation, Snapshot snapshot, IndexScanInstrumentation *instrument, int nkeys); +extern IndexScanDesc index_beginscan_bitmap_parallel(Relation indexRelation, + Snapshot snapshot, + IndexScanInstrumentation *instrument, + int nkeys, + ParallelIndexScanDesc pscan); +extern bool index_am_supports_parallel_scan(Relation indexRelation); extern void index_rescan(IndexScanDesc scan, ScanKey keys, int nkeys, ScanKey orderbys, int norderbys); diff --git a/src/include/access/gin_private.h b/src/include/access/gin_private.h index 3c5fd6ba817..b144574c429 100644 --- a/src/include/access/gin_private.h +++ b/src/include/access/gin_private.h @@ -349,7 +349,7 @@ typedef struct GinScanEntryData /* for a partial-match or full-scan query, we accumulate all TIDs here */ TIDBitmap *matchBitmap; - TBMPrivateIterator *matchIterator; + TBMOrderedIterator *matchIterator; /* * If blockno is InvalidBlockNumber, all of the other fields in the diff --git a/src/include/access/relscan.h b/src/include/access/relscan.h index d7614104730..08e50b75be6 100644 --- a/src/include/access/relscan.h +++ b/src/include/access/relscan.h @@ -45,7 +45,7 @@ typedef struct TableScanDescData union { /* Iterator for Bitmap Table Scans */ - TBMIterator rs_tbmiterator; + TBMOrderedIterator *rs_tbmiterator; /* * Range of ItemPointers for table_scan_getnextslot_tidrange() to diff --git a/src/include/executor/nodeBitmapHeapscan.h b/src/include/executor/nodeBitmapHeapscan.h index 5c0b6592a1f..f4349b6b4fe 100644 --- a/src/include/executor/nodeBitmapHeapscan.h +++ b/src/include/executor/nodeBitmapHeapscan.h @@ -20,20 +20,12 @@ extern BitmapHeapScanState *ExecInitBitmapHeapScan(BitmapHeapScan *node, EState *estate, int eflags); extern void ExecEndBitmapHeapScan(BitmapHeapScanState *node); extern void ExecReScanBitmapHeapScan(BitmapHeapScanState *node); -extern void ExecBitmapHeapEstimate(BitmapHeapScanState *node, - ParallelContext *pcxt); -extern void ExecBitmapHeapInitializeDSM(BitmapHeapScanState *node, - ParallelContext *pcxt); -extern void ExecBitmapHeapReInitializeDSM(BitmapHeapScanState *node, - ParallelContext *pcxt); -extern void ExecBitmapHeapInitializeWorker(BitmapHeapScanState *node, - ParallelWorkerContext *pwcxt); extern void ExecBitmapHeapInstrumentEstimate(BitmapHeapScanState *node, - ParallelContext *pcxt); + ParallelContext *pcxt); extern void ExecBitmapHeapInstrumentInitDSM(BitmapHeapScanState *node, - ParallelContext *pcxt); + ParallelContext *pcxt); extern void ExecBitmapHeapInstrumentInitWorker(BitmapHeapScanState *node, - ParallelWorkerContext *pwcxt); + ParallelWorkerContext *pwcxt); extern void ExecBitmapHeapRetrieveInstrumentation(BitmapHeapScanState *node); #endif /* NODEBITMAPHEAPSCAN_H */ diff --git a/src/include/executor/nodeBitmapIndexscan.h b/src/include/executor/nodeBitmapIndexscan.h index 5a43bda943d..da8e59a05ae 100644 --- a/src/include/executor/nodeBitmapIndexscan.h +++ b/src/include/executor/nodeBitmapIndexscan.h @@ -23,6 +23,7 @@ extern void ExecEndBitmapIndexScan(BitmapIndexScanState *node); extern void ExecReScanBitmapIndexScan(BitmapIndexScanState *node); extern void ExecBitmapIndexScanEstimate(BitmapIndexScanState *node, ParallelContext *pcxt); extern void ExecBitmapIndexScanInitializeDSM(BitmapIndexScanState *node, ParallelContext *pcxt); +extern void ExecBitmapIndexScanReInitializeDSM(BitmapIndexScanState *node, ParallelContext *pcxt); extern void ExecBitmapIndexScanInitializeWorker(BitmapIndexScanState *node, ParallelWorkerContext *pwcxt); extern void ExecBitmapIndexScanRetrieveInstrumentation(BitmapIndexScanState *node); diff --git a/src/include/nodes/execnodes.h b/src/include/nodes/execnodes.h index 91bb0bd2e13..969936e1bdb 100644 --- a/src/include/nodes/execnodes.h +++ b/src/include/nodes/execnodes.h @@ -1809,6 +1809,10 @@ typedef struct IndexOnlyScanState * SharedInfo parallel worker instrumentation (no leader entry) * ---------------- */ + +/* this struct is defined in nodeBitmapIndexscan.c */ +typedef struct SharedBitmapIndexState SharedBitmapIndexState; + typedef struct BitmapIndexScanState { ScanState ss; /* its first field is NodeTag */ @@ -1825,6 +1829,8 @@ typedef struct BitmapIndexScanState struct IndexScanDescData *biss_ScanDesc; IndexScanInstrumentation *biss_Instrument; SharedIndexScanInstrumentation *biss_SharedInfo; + Size biss_PscanLen; + SharedBitmapIndexState *biss_ParallelState; } BitmapIndexScanState; @@ -1835,15 +1841,11 @@ typedef struct BitmapIndexScanState * tbm bitmap obtained from child index scan(s) * stats execution statistics * initialized is node is ready to iterate - * pstate shared state for parallel bitmap scan - * sinstrument statistics for parallel workers + * * sinstrument statistics for parallel workers * recheck do current page's tuples need recheck * ---------------- */ -/* this struct is defined in nodeBitmapHeapscan.c */ -typedef struct ParallelBitmapHeapState ParallelBitmapHeapState; - typedef struct BitmapHeapScanState { ScanState ss; /* its first field is NodeTag */ @@ -1851,7 +1853,6 @@ typedef struct BitmapHeapScanState TIDBitmap *tbm; BitmapHeapScanInstrumentation stats; bool initialized; - ParallelBitmapHeapState *pstate; SharedBitmapHeapInstrumentation *sinstrument; bool recheck; } BitmapHeapScanState; diff --git a/src/include/nodes/plannodes.h b/src/include/nodes/plannodes.h index 09a1ec73180..38f4e4ccd93 100644 --- a/src/include/nodes/plannodes.h +++ b/src/include/nodes/plannodes.h @@ -522,7 +522,6 @@ typedef struct BitmapAnd typedef struct BitmapOr { Plan plan; - bool isshared; List *bitmapplans; } BitmapOr; @@ -688,8 +687,6 @@ typedef struct BitmapIndexScan Scan scan; /* OID of index to scan */ Oid indexid; - /* Create shared bitmap if set */ - bool isshared; /* list of index quals (OpExprs) */ List *indexqual; /* the same in original form */ diff --git a/src/include/nodes/tidbitmap.h b/src/include/nodes/tidbitmap.h index 4a09ab41cdf..cdf4a017f35 100644 --- a/src/include/nodes/tidbitmap.h +++ b/src/include/nodes/tidbitmap.h @@ -33,31 +33,34 @@ */ #define TBM_MAX_TUPLES_PER_PAGE MaxHeapTuplesPerPage +/* + * When we have to switch over to lossy storage, we use a data structure + * with one bit per page, where all pages having the same number DIV + * TBM_MAX_PAGES_PER_CHUNK are aggregated into one chunk. When a chunk is present + * and has the bit set for a given page, there must not be a per-page entry + * for that page in the page table. + * + * We actually store both exact pages and lossy chunks in the same hash + * table, using identical data structures. (This is because the memory + * management for hashtables doesn't easily/efficiently allow space to be + * transferred easily from one hashtable to another.) Therefore it's best + * if TBM_MAX_PAGES_PER_CHUNK is the same as TBM_MAX_TUPLES_PER_PAGE, or at least not + * too different. But we also want TBM_MAX_PAGES_PER_CHUNK to be a power of 2 to + * avoid expensive integer remainder operations. So, define it like this: + */ +#define TBM_MAX_PAGES_PER_CHUNK (BLCKSZ / 32) + /* * Actual bitmap representation is private to tidbitmap.c. Callers can * do IsA(x, TIDBitmap) on it, but nothing else. */ typedef struct TIDBitmap TIDBitmap; -/* Likewise, TBMPrivateIterator is private */ -typedef struct TBMPrivateIterator TBMPrivateIterator; -typedef struct TBMSharedIterator TBMSharedIterator; +/* Likewise, iterators are private */ +typedef struct TBMOrderedIterator TBMOrderedIterator; +typedef struct TBMUnorderedIterator TBMUnorderedIterator; -/* - * Callers with both private and shared implementations can use this unified - * API. - */ -typedef struct TBMIterator -{ - bool shared; - union - { - TBMPrivateIterator *private_iterator; - TBMSharedIterator *shared_iterator; - } i; -} TBMIterator; - -/* Result structure for tbm_iterate */ +/* Result structure for iteration */ typedef struct TBMIterateResult { BlockNumber blockno; /* block number containing tuples */ @@ -82,12 +85,12 @@ typedef struct TBMIterateResult extern TIDBitmap *tbm_create(Size maxbytes, dsa_area *dsa); extern void tbm_free(TIDBitmap *tbm); -extern void tbm_free_shared_area(dsa_area *dsa, dsa_pointer dp); extern void tbm_add_tuples(TIDBitmap *tbm, const ItemPointerData *tids, int ntids, bool recheck); extern void tbm_add_page(TIDBitmap *tbm, BlockNumber pageno); +extern void tbm_copy_page(TIDBitmap *dest, TBMIterateResult *src); extern void tbm_union(TIDBitmap *a, const TIDBitmap *b); extern void tbm_intersect(TIDBitmap *a, const TIDBitmap *b); @@ -98,30 +101,15 @@ extern int tbm_extract_page_tuple(TBMIterateResult *iteritem, extern bool tbm_is_empty(const TIDBitmap *tbm); -extern TBMPrivateIterator *tbm_begin_private_iterate(TIDBitmap *tbm); -extern dsa_pointer tbm_prepare_shared_iterate(TIDBitmap *tbm); -extern bool tbm_private_iterate(TBMPrivateIterator *iterator, TBMIterateResult *tbmres); -extern bool tbm_shared_iterate(TBMSharedIterator *iterator, TBMIterateResult *tbmres); -extern void tbm_end_private_iterate(TBMPrivateIterator *iterator); -extern void tbm_end_shared_iterate(TBMSharedIterator *iterator); -extern TBMSharedIterator *tbm_attach_shared_iterate(dsa_area *dsa, - dsa_pointer dp); -extern int tbm_calculate_entries(Size maxbytes); - -extern TBMIterator tbm_begin_iterate(TIDBitmap *tbm, - dsa_area *dsa, dsa_pointer dsp); -extern void tbm_end_iterate(TBMIterator *iterator); +extern TBMOrderedIterator *tbm_begin_ordered_iterate(TIDBitmap *tbm); +extern bool tbm_ordered_iterate(TBMOrderedIterator *iterator, TBMIterateResult *tbmres); +extern void tbm_end_ordered_iterate(TBMOrderedIterator **iterator); -extern bool tbm_iterate(TBMIterator *iterator, TBMIterateResult *tbmres); +extern dsa_pointer tbm_prepare_shared_unordered_iterate(TIDBitmap *tbm); +extern TBMUnorderedIterator *tbm_begin_shared_unordered_iterate(dsa_area *dsa, dsa_pointer dp); +extern bool tbm_shared_unordered_iterate(TBMUnorderedIterator *iterator, TBMIterateResult *tbmres); +extern void tbm_end_shared_unordered_iterate(TBMUnorderedIterator **iterator); -static inline bool -tbm_exhausted(TBMIterator *iterator) -{ - /* - * It doesn't matter if we check the private or shared iterator here. If - * tbm_end_iterate() was called, they will be NULL - */ - return !iterator->i.private_iterator; -} +extern int tbm_calculate_entries(Size maxbytes); #endif /* TIDBITMAP_H */ diff --git a/src/test/regress/expected/bitmapops.out b/src/test/regress/expected/bitmapops.out index 64068e0469c..f9470af52eb 100644 --- a/src/test/regress/expected/bitmapops.out +++ b/src/test/regress/expected/bitmapops.out @@ -44,5 +44,103 @@ SELECT count(*) FROM bmscantest WHERE a = 1 OR b = 1; 2485 (1 row) +-- Test parallel bitmap-and and bitmap-or. +ALTER TABLE bmscantest SET (parallel_workers = 2); +set min_parallel_table_scan_size = 0; +set min_parallel_index_scan_size = 0; +set parallel_setup_cost = 0; +set parallel_tuple_cost = 0; +set max_parallel_workers_per_gather = 2; +set cpu_tuple_cost = 1; +EXPLAIN (COSTS OFF) +SELECT count(*) FROM bmscantest WHERE a = 1 AND b = 1; + QUERY PLAN +------------------------------------------------------------------ + Aggregate + -> Gather + Workers Planned: 2 + -> Parallel Bitmap Heap Scan on bmscantest + Recheck Cond: ((b = 1) AND (a = 1)) + -> Parallel BitmapAnd + -> Parallel Bitmap Index Scan on i_bmtest_b + Index Cond: (b = 1) + -> Parallel Bitmap Index Scan on i_bmtest_a + Index Cond: (a = 1) +(10 rows) + +SELECT count(*) FROM bmscantest WHERE a = 1 AND b = 1; + count +------- + 23 +(1 row) + +EXPLAIN (COSTS OFF) +SELECT count(*) FROM bmscantest WHERE a = 1 OR b = 1; + QUERY PLAN +------------------------------------------------------------------------ + Finalize Aggregate + -> Gather + Workers Planned: 2 + -> Partial Aggregate + -> Parallel Bitmap Heap Scan on bmscantest + Recheck Cond: ((a = 1) OR (b = 1)) + -> Parallel BitmapOr + -> Parallel Bitmap Index Scan on i_bmtest_a + Index Cond: (a = 1) + -> Parallel Bitmap Index Scan on i_bmtest_b + Index Cond: (b = 1) +(11 rows) + +SELECT count(*) FROM bmscantest WHERE a = 1 OR b = 1; + count +------- + 2485 +(1 row) + +reset min_parallel_table_scan_size; +reset min_parallel_index_scan_size; +reset parallel_setup_cost; +reset parallel_tuple_cost; +reset max_parallel_workers_per_gather; +reset cpu_tuple_cost; -- clean up DROP TABLE bmscantest; +-- Test parallel bitmap heap scan with a non-parallel-aware index (GIN). +-- The BitmapIndexScan node is parallel-aware, but only one worker can +-- actually scan the GIN index; all workers still repartition the bitmap. +CREATE TABLE bmgintest (a int, arr int[]) WITH (autovacuum_enabled = false); +INSERT INTO bmgintest SELECT i, ARRAY[i % 100] FROM generate_series(1,100000) i; +CREATE INDEX i_bmgintest ON bmgintest USING gin(arr); +ANALYZE bmgintest; +ALTER TABLE bmgintest SET (parallel_workers = 2); +set min_parallel_table_scan_size = 0; +set min_parallel_index_scan_size = 0; +set parallel_setup_cost = 0; +set parallel_tuple_cost = 0; +set max_parallel_workers_per_gather = 2; +EXPLAIN (COSTS OFF) +SELECT count(*) FROM bmgintest WHERE arr @> ARRAY[5]; + QUERY PLAN +------------------------------------------------------------------- + Finalize Aggregate + -> Gather + Workers Planned: 2 + -> Partial Aggregate + -> Parallel Bitmap Heap Scan on bmgintest + Recheck Cond: (arr @> '{5}'::integer[]) + -> Parallel Bitmap Index Scan on i_bmgintest + Index Cond: (arr @> '{5}'::integer[]) +(8 rows) + +SELECT count(*) FROM bmgintest WHERE arr @> ARRAY[5]; + count +------- + 1000 +(1 row) + +reset min_parallel_table_scan_size; +reset min_parallel_index_scan_size; +reset parallel_setup_cost; +reset parallel_tuple_cost; +reset max_parallel_workers_per_gather; +DROP TABLE bmgintest; diff --git a/src/test/regress/expected/memoize.out b/src/test/regress/expected/memoize.out index 2d24f4480c4..b62c5a766fb 100644 --- a/src/test/regress/expected/memoize.out +++ b/src/test/regress/expected/memoize.out @@ -463,6 +463,7 @@ RESET enable_bitmapscan; RESET enable_hashjoin; -- Test parallel plans with Memoize SET min_parallel_table_scan_size TO 0; +SET min_parallel_index_scan_size TO 0; SET parallel_setup_cost TO 0; SET parallel_tuple_cost TO 0; SET max_parallel_workers_per_gather TO 2; @@ -480,7 +481,7 @@ WHERE t1.unique1 < 1000; -> Nested Loop -> Parallel Bitmap Heap Scan on tenk1 t1 Recheck Cond: (unique1 < 1000) - -> Bitmap Index Scan on tenk1_unique1 + -> Parallel Bitmap Index Scan on tenk1_unique1 Index Cond: (unique1 < 1000) -> Memoize Cache Key: t1.twenty @@ -501,6 +502,7 @@ WHERE t1.unique1 < 1000; RESET max_parallel_workers_per_gather; RESET parallel_tuple_cost; RESET parallel_setup_cost; +RESET min_parallel_index_scan_size; RESET min_parallel_table_scan_size; -- Ensure memoize works for ANTI joins CREATE TABLE tab_anti (a int, b boolean); diff --git a/src/test/regress/expected/select_parallel.out b/src/test/regress/expected/select_parallel.out index e1344215644..fc19433b90c 100644 --- a/src/test/regress/expected/select_parallel.out +++ b/src/test/regress/expected/select_parallel.out @@ -545,8 +545,8 @@ END $$; set work_mem='64kB'; --set small work mem to force lossy pages explain (costs off) select count(*) from tenk1, tenk2 where tenk1.hundred > 1 and tenk2.thousand=0; - QUERY PLAN ------------------------------------------------------------- + QUERY PLAN +--------------------------------------------------------------------- Aggregate -> Nested Loop -> Gather @@ -558,7 +558,7 @@ explain (costs off) Workers Planned: 4 -> Parallel Bitmap Heap Scan on tenk1 Recheck Cond: (hundred > 1) - -> Bitmap Index Scan on tenk1_hundred + -> Parallel Bitmap Index Scan on tenk1_hundred Index Cond: (hundred > 1) (13 rows) diff --git a/src/test/regress/sql/bitmapops.sql b/src/test/regress/sql/bitmapops.sql index 1b175f6ff96..6aeba2429fe 100644 --- a/src/test/regress/sql/bitmapops.sql +++ b/src/test/regress/sql/bitmapops.sql @@ -42,6 +42,55 @@ SELECT count(*) FROM bmscantest WHERE a = 1 AND b = 1; -- Test bitmap-or. SELECT count(*) FROM bmscantest WHERE a = 1 OR b = 1; +-- Test parallel bitmap-and and bitmap-or. +ALTER TABLE bmscantest SET (parallel_workers = 2); +set min_parallel_table_scan_size = 0; +set min_parallel_index_scan_size = 0; +set parallel_setup_cost = 0; +set parallel_tuple_cost = 0; +set max_parallel_workers_per_gather = 2; +set cpu_tuple_cost = 1; + +EXPLAIN (COSTS OFF) +SELECT count(*) FROM bmscantest WHERE a = 1 AND b = 1; +SELECT count(*) FROM bmscantest WHERE a = 1 AND b = 1; + +EXPLAIN (COSTS OFF) +SELECT count(*) FROM bmscantest WHERE a = 1 OR b = 1; +SELECT count(*) FROM bmscantest WHERE a = 1 OR b = 1; + +reset min_parallel_table_scan_size; +reset min_parallel_index_scan_size; +reset parallel_setup_cost; +reset parallel_tuple_cost; +reset max_parallel_workers_per_gather; +reset cpu_tuple_cost; -- clean up DROP TABLE bmscantest; + +-- Test parallel bitmap heap scan with a non-parallel-aware index (GIN). +-- The BitmapIndexScan node is parallel-aware, but only one worker can +-- actually scan the GIN index; all workers still repartition the bitmap. +CREATE TABLE bmgintest (a int, arr int[]) WITH (autovacuum_enabled = false); +INSERT INTO bmgintest SELECT i, ARRAY[i % 100] FROM generate_series(1,100000) i; +CREATE INDEX i_bmgintest ON bmgintest USING gin(arr); +ANALYZE bmgintest; +ALTER TABLE bmgintest SET (parallel_workers = 2); +set min_parallel_table_scan_size = 0; +set min_parallel_index_scan_size = 0; +set parallel_setup_cost = 0; +set parallel_tuple_cost = 0; +set max_parallel_workers_per_gather = 2; + +EXPLAIN (COSTS OFF) +SELECT count(*) FROM bmgintest WHERE arr @> ARRAY[5]; +SELECT count(*) FROM bmgintest WHERE arr @> ARRAY[5]; + +reset min_parallel_table_scan_size; +reset min_parallel_index_scan_size; +reset parallel_setup_cost; +reset parallel_tuple_cost; +reset max_parallel_workers_per_gather; + +DROP TABLE bmgintest; diff --git a/src/test/regress/sql/memoize.sql b/src/test/regress/sql/memoize.sql index c02a0d51af4..18499247e01 100644 --- a/src/test/regress/sql/memoize.sql +++ b/src/test/regress/sql/memoize.sql @@ -228,6 +228,7 @@ RESET enable_hashjoin; -- Test parallel plans with Memoize SET min_parallel_table_scan_size TO 0; +SET min_parallel_index_scan_size TO 0; SET parallel_setup_cost TO 0; SET parallel_tuple_cost TO 0; SET max_parallel_workers_per_gather TO 2; @@ -246,6 +247,7 @@ WHERE t1.unique1 < 1000; RESET max_parallel_workers_per_gather; RESET parallel_tuple_cost; RESET parallel_setup_cost; +RESET min_parallel_index_scan_size; RESET min_parallel_table_scan_size; -- Ensure memoize works for ANTI joins diff --git a/src/tools/pgindent/typedefs.list b/src/tools/pgindent/typedefs.list index 656f1f60862..851cea53111 100644 --- a/src/tools/pgindent/typedefs.list +++ b/src/tools/pgindent/typedefs.list @@ -288,6 +288,7 @@ BitmapHeapScanInstrumentation BitmapHeapScanState BitmapIndexScan BitmapIndexScanState +SharedBitmapIndexState BitmapOr BitmapOrPath BitmapOrState @@ -2102,7 +2103,6 @@ PSQL_ECHO PSQL_ECHO_HIDDEN PSQL_ERROR_ROLLBACK PSQL_SEND_MODE -PTEntryArray PTIterationArray PTOKEN_PRIVILEGES PTOKEN_USER @@ -2131,7 +2131,6 @@ ParallelAppendState ParallelApplyWorkerEntry ParallelApplyWorkerInfo ParallelApplyWorkerShared -ParallelBitmapHeapState ParallelBlockTableScanDesc ParallelBlockTableScanWorker ParallelBlockTableScanWorkerData @@ -3068,9 +3067,8 @@ TAPtype TAR_MEMBER TBMIterateResult TBMIteratingState -TBMIterator -TBMPrivateIterator -TBMSharedIterator +TBMOrderedIterator +TBMUnorderedIterator TBMSharedIteratorState TBMStatus TBlockState -- 2.55.0