Re: Improving scalability of Parallel Bitmap Heap/Index Scan

From: John Naylor <johncnaylorls(at)gmail(dot)com>
To: David Geier <geidav(dot)pg(at)gmail(dot)com>
Cc: PostgreSQL Developers <pgsql-hackers(at)lists(dot)postgresql(dot)org>
Subject: Re: Improving scalability of Parallel Bitmap Heap/Index Scan
Date: 2026-10-07 12:08:34
Message-ID: CANWCAZYNnc8O1kmnsVWZOXsOp=3OAh-4nNsaHatVhq6oMUdgVA@mail.gmail.com
Views: Whole Thread | Raw Message | Download mbox | Resend email
Thread:
Lists: pgsql-hackers

On Tue, Oct 6, 2026 at 11:28 PM David Geier <geidav(dot)pg(at)gmail(dot)com> wrote:
> > Interesting. What's the unit of work for each worker? Can you describe
> > how and when repartitioning works? Let's say you have an "And" node
> > for two indexes.
>
> Above the BIS, the unit of work is a disjoint partition of the
> TIDBitmap, i.e. a non-overlapping set of HEAP blocks (and their TIDs in
> non-lossy mode). This allows every operation above the BIS — Bitmap And,
> Bitmap Or, and the Bitmap Heap Scan — to run locally in each worker
> without any synchronization.
>
> The unit of work inside the BIS is different: each worker scans a
> portion of the index and builds a private TIDBitmap that may contain
> TIDs of arbitrary HEAP blocks.
>
> Repartitioning happens at the end of the BIS, after the index has been
> read to completion by all participants. At that point, each participant
> has a private TIDBitmap covering arbitrary HEAP pages that it shares
> with the others through DSA. Then each participant create a new
> TIDBitmap containing only the pages assigned to it, by scanning the
> bitmaps of all other participants with the following condition to decide
> if the block belongs to it or not:

I have not done a full review of the patch, but the approach seems
promising. The devil's in the details, of course.

Okay, so worst case N workers will use N times more memory. (not
counting transient use during repartitioning)

One mitigation could be to use hash_mem_multiplier to expand work_mem,
and divide that by the number of workers. (I'm not sure why we don't
do that already; seems like a natural use for it). This will need
testing.

> Dividing by 256 keeps runs of 256 consecutive HEAP pages together in one

Does this need to be PAGES_PER_CHUNK?

> participant, which is important for TIDBitmap efficiency. Because every
> participant uses the same hash function and the same participant count,
> the resulting partitions are identical across workers. The exchange uses
> an unordered iterator, since order does not matter for this assignment.
>
> For your Bitmap And example:
>
> -> Bitmap Heap Scan
> -> Bitmap And
> -> Bitmap Index Scan(index_a)
> -> Bitmap Index Scan(index_b)
>
> Both leaf scans are parallel-aware and are run by every participant.

+ /*
+ * All leaves in a partial bitmap qual must be parallel scans so that
+ * each worker produces a disjoint partition of the same TIDBitmap.
+ */
+ if (!ipath->indexinfo->amcanparallel)
+ return NULL;

This disallows bitmap heap scans for non-btree indexes where they were
allowed before, doesn't it?

It seems it should be not too much additional work to allow one worker
to scan an index serially, and then publish it for repartitioning.

> Worker i first executes the scan on index_a, repartitions the result
> into partition A_i, then does the same for index_b to get partition B_i.
> Since both scans use the same partitioning scheme, A_i and B_i cover
> exactly the same set of HEAP pages. The Bitmap And node just intersects
> A_i and B_i locally with tbm_intersect(); no locks, no shared iteration
> state. The result is still partition i, which the Bitmap Heap Scan then
> consumes.
>
> So the only synchronization points are the two barriers inside each BIS
> (one after sharing the local bitmaps, one after everyone has finished
> reading the shared bitmaps). After that, the whole bitmap plan tree
> above the BIS is embarrassingly parallel.

Sounds good in principle (BTW, for that to work the workers need to
keep the same my_id across index scans, but that's computed inside
BitmapIndexScanPartition, so that seems broken unless I'm missing
something)

> >> There's also a forward-looking angle: ongoing work ([1], [2]) is
> >> converting TIDBitmap from a hash table to a radix tree. Once that lands,
> >> we could consider building a single cooperative bitmap rather than N
> >> private ones — in case this structure lends itself better to shared
> >> write access (which seems to be the case).
> >
> > It's not currently optimized for concurrency, so it would depend a lot
> > on the ordering of the input and could actually be worse especially
> > for e.g. a uuid index -- both cache misses and higher constant
> > overheads than a hash table. The advantage of the radix tree is that
> > 1) iteration is automatically ordered so no sorting step, and 2) the
>
> For that to pay off, the increased cost for inserting must be reasonably
> low. Is that the case?

Actually, no, I just reminded myself that in shared memory, the node
pointers use DSA pointers, and dereferencing them requires a function
call.

> > page table entries can be variable length, probably using less net
> > memory, and allowing the max number of TIDs per page to be decoupled
> > from the max number of tuples,
>
> Is this about wasting less memory for PagetableEntries that only contain
> very few TIDs?

It's not about the number of TIDs, it's the largest offset for a given
page. Real life tables don't have a single int, they may have say
30-50 columns -- you don't need 291 bits available to store any TID
from such a page. With less memory, there's less chance of regression
with workers getting a fraction of the memory. My patch has to invent
shared disjoint iteration under a lock during the heap scan phase, but
with your architecture, it wouldn't need that.

> > and 3) easier to be smart about lossification.
>
> Can you share what ideas exist at this point?

a) Since iteration is ordered, chunks are only created for a whole
range, not at random.
b) Once a whole range has been lossified, you can "close" that range
and any new page in that range will just go in the chunk.
c) You could decide a range is not worth lossifying

> The even better implementation, that would keep the advantages of the
> partitioning scheme while avoiding increased memory consumption and
> rereading all bitmaps by all participants, is to keep the partitions in
> shared memory from the staert and having each participant directly
> insert the TIDs into the right partition, e.g.
>
> BlockNumber blockno = ItemPointerGetBlockNumber(cur_tid);
> uint32 p = murmurhash32(blockno / 256) % nparticipants;
> tbm_add_tuple(tidbitmaps[p], cur_tid);
>
> If we use enough partitions, say 8 per participant, then we might get
> away with locking the entire tbm_add_tuple() operation, hoping that the
> TIDs distribute more or less uniformly across all partitions. Now when
> thinking about it, that might be the better approach, also for the
> TIDBitmap-based variant.

Parallel index scan with shared partitioned hash table has been tried
before without success, I believe.

--
John Naylor
Amazon Web Services

In response to

Browse pgsql-hackers by date

  From Date Subject
Next Message Matthias van de Meent 2026-10-07 12:09:29 Re: Let an ordering index scan hand its ORDER BY value to the target list
Previous Message Fujii Masao 2026-10-07 11:59:11 Re: [PATCH] psql: avoid CREATE command completion after GRANT/REVOKE CREATE