Skip to content

duckdb_pglake: parallelize postgres_scan on partitioned tables - #653

Open
sfc-gh-mslot wants to merge 1 commit into
mainfrom
marcoslot/partitioned-postgres-scan
Open

sfc-gh-mslot wants to merge 1 commit into
mainfrom
marcoslot/partitioned-postgres-scan

Conversation

@sfc-gh-mslot

Copy link
Copy Markdown
Collaborator

Problem

When scanning a partitioned PostgreSQL table, postgres_scan previously fell back to a single serial task with one thread and zero parallelism. Because partitioned parent tables have no physical heap storage (relkind = 'p'), their relation size is 0 pages. Furthermore, parallelizing a partitioned table directly by splitting the parent's ctid range is invalid: each physical partition has its own independent block numbers starting at block 0, so a ctid range filter on the parent evaluates against every partition simultaneously, concentrating rows in the first task and returning empty results for later tasks.

Solution

Add table-partition-scan.patch to duckdb-postgres to make postgres_scan and attached partitioned table scans partition-aware:

  1. Hierarchy discovery: Recursively query pg_inherits to resolve all underlying leaf partitions (relkind IN ('r', 'm', 'f')), handling arbitrary subpartitioning depth as well as partitions residing in separate schemas.
  2. Cardinality and thread estimation: Set the table's pages_approx to the sum of all leaf partition page counts (pg_relation_size(oid)), yielding accurate cardinality estimates and sizing max_threads across all partition tasks.
  3. Partition-level and ctid-level parallelism: Parallel tasks are scheduled across partitions. Partitions larger than pg_pages_per_task are subdivided into ctid-range chunks within the partition, while smaller or empty partitions are scanned as single tasks.
  4. Chunk completion fix: Fix ScanChunk to return non-empty output buffers when a task finishes with remaining tuples, preventing the next task from overwriting unconsumed results.
-- Scan a partitioned table in parallel across and within partitions
SELECT count(*) FROM postgres_scan('host=localhost dbname=mydb', 'public', 'measurements');

-- Pushdown filters and projection to underlying partitions
SELECT id, value FROM postgres_scan_pushdown('host=localhost dbname=mydb', 'public', 'measurements') WHERE id > 1000;

Test plan

  • Tested range-partitioned table scans returning all rows from leaf partitions.
  • Tested multi-level subpartitioning hierarchies with partitions in separate schemas.
  • Tested empty partitioned tables (no partitions attached) returning zero rows without error.
  • Tested parallel ctid chunking across and within partitions with pg_pages_per_task = 5.
  • Tested filter pushdown and projection pushdown on partitioned tables.
  • Tested binary and text protocols across multiple tasks without lost tuples.
  • Verified full test suite in pgduck_server/tests/pytests/test_postgres_scanner.py (30 passed).

When scanning a partitioned PostgreSQL table, postgres_scan previously saw
relpages = 0 because partitioned parent tables have no storage. As a result,
it fell back to a single task covering the entire table with 1 thread and no
parallelism. Attempting parallel ctid scans on a partitioned table directly
is also invalid because each physical partition maintains its own independent
ctid space starting at (0, 1).

Add table-partition-scan.patch to duckdb-postgres:
- Recursively traverses inheritance hierarchies via pg_inherits to discover all
  underlying leaf partitions and their physical relation sizes.
- Sets approx_num_pages and cardinality estimates from the sum of leaf partition pages.
- Sizes and schedules tasks across partitions and within partitions based on ctid.
- Correctly handles multi-level subpartitioning, cross-schema partitions, and empty tables.
- Ensures finished tasks with non-empty final chunks emit their tuples before proceeding.

Signed-off-by: Marco Slot <marco.slot@snowflake.com>

@sfc-gh-dachristensen sfc-gh-dachristensen left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is there a similar issue with table inheritance?

+ } else if (is_partitioned) {
+ idx_t total_tasks = 0;
+ for (auto &partition_entry : partitions) {
+ if (!use_ctid_scan || partition_entry.approx_num_pages == 0) {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

approx_num_pages is a proxy for being a branch partition?

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants