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>, 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-08 07:52:27
Message-ID: f45e438a-f19d-4326-a143-ce49c19a29c8@googlemail.com
Views: Whole Thread | Raw Message | Download mbox | Resend email
Thread:
Lists: pgsql-hackers

John, thanks for taking a closer look at the patch! Your input is very
helpful.

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

Yes, but only temporarily and only in the BIS. The per-worker memory
consumption also stays the same. It's just that there are more workers
now that each can encounter the same TIDs / blocks. And once the
repartitioning step is complete, the old TIDBitmaps are freed and the
memory consumption is as before because now every participant holds a
TIDBitmap containing only 1/nth of the total set of TIDs / blocks.

> One mitigation could be to use hash_mem_multiplier to expand work_mem,
> and divide that by the number of workers.

I thought work_mem is per-node, per-worker. Given that the per-node,
per-worker memory consumption doesn't change, things should be fine?

> (I'm not sure why we don't
> do that already; seems like a natural use for it). This will need
> testing.

Using hash_mem_multiplier in TIDBitmap makes sense, given that TIDBitmap
is hash-based. I'll add that.

Additionally, we could improve on the increased memory consumption by
having the participants feed the found TIDs into lock-free ring buffers,
one per partition. The participants would alternate between reading from
the index to feed the ring buffers and pulling data out from their
assigned ring buffer to store it in their local partition TIDBitmap.

That's for sure less work than making TIDBitmap parallel-safe for insert
(which includes lossification), should be on-par performance-wise and
avoids the increase in memory consumption.

>> Dividing by 256 keeps runs of 256 consecutive HEAP pages together in one
>
> Does this need to be PAGES_PER_CHUNK?

Yes. I didn't use that because it's defined in tidbitmap.c. I'll move it
to tidbitmap.h and rename it to TBM_MAX_PAGES_PER_CHUNK, right next to
TBM_MAX_TUPLES_PER_PAGE.

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

Will check and fix if needed. pg_regress tests are also not passing at
the moment. Will get them to work as well.

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

Indeed. Will do.
>> 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)

Good catch. I haven't tested (much) yet with Bitmap And/Or. I'll add
parallel bitmap queries to the pg_regress tests to cover these cases.
>>> 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.

That makes sense: for what I understood, in parallel VACUUM the TidStore
is filled by a single leader but then read by multiple per-index workers.

Then using the partitioned architecture is actually making it easier to
replace TIDBitmap with the radix tree.

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

Got it.

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

Makes sense.

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

Good ideas. Regarding (c): TIDBitmap could also be smarter about that
today already. Currently, it eagerly lossifies any PagetableEntry which
wouldn't turn into a chunk (is not a multiple of PAGES_PER_CHUNK).
Instead we should lossify based on the number of TIDs. I'll code that up
and add it as a separate patch.

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

As I shared in my other mail: this approach doesn't fly because it
requires simplehash.h to work in shared memory. I think the better
approach is what I described above: lock-free, per-partition ring
buffers. I'll give this a try.

> Parallel index scan with shared partitioned hash table has been tried
> before without success, I believe.
Do you have any pointers to mailing list discussions or similar?

I'll share an updated patch the next days.

--
David Geier

In response to

Responses

Browse pgsql-hackers by date

  From Date Subject
Next Message Hayato Kuroda (Fujitsu) 2026-10-08 07:53:50 RE: Incorrect CONTEXT reported for errors from parallel apply worker in logical replication
Previous Message Michael Paquier 2026-10-08 07:51:28 Re: pgstat: allow a stats kind to use its own dedicated dsa/dshash