Re: Improving scalability of Parallel Bitmap Heap/Index Scan

From: David Geier <geidav(dot)pg(at)gmail(dot)com>
To: John Naylor <johncnaylorls(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-06 16:27:58
Message-ID: aad4ac2f-6fec-4d2b-b463-0ce59bce06e8@gmail.com
Views: Whole Thread | Raw Message | Download mbox | Resend email
Thread:
Lists: pgsql-hackers

>> Each parallel BIS worker builds its own private TIDBitmap from the index
>> blocks it reads. Once all BIS participants finish, the bitmaps are
>> repartitioned by hash(pageno / 256) % nparticipants. This way every
>> worker ends up with roughly 1/N-th of the pages / TIDs, while keeping
>> runs of 256 local to one participant which is important to keep
>> TIDBitmap for working efficiently. From that point on, each participant
>> only processes its private partition. No shared iteration state exists
>> and hence no synchronization is needed above the BIS.
>
> 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:

murmurhash32(blockno / 256) % nparticipants == my_id

Dividing by 256 keeps runs of 256 consecutive HEAP pages together in one
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.
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.

>> 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?

> 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?

> and 3) easier to be smart about lossification.

Can you share what ideas exist at this point?

> I have made progress this summer, but there are still a lot of finicky
> details to get right.

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.

If we do it this way, integrating the TID radix tree would actually be
straightforward.

--
David Geier

In response to

Browse pgsql-hackers by date

  From Date Subject
Next Message Greg Burd 2026-10-06 16:29:12 Re: Let an ordering index scan hand its ORDER BY value to the target list
Previous Message Sami Imseih 2026-10-06 16:20:16 Re: pgstat: allow a stats kind to use its own dedicated dsa/dshash