From 5b090985cd3d3a79a06f04136f08602f667d65f6 Mon Sep 17 00:00:00 2001 From: "Joshua D. Drake" Date: Mon, 14 Sep 2026 19:02:46 +0000 Subject: [PATCH 1/3] fix: share table-AM parallel scan groups across workers phs_nallocated was a first-wins mutex, so launched workers sat idle. Claim it as a group index, the same way the custom scan shares work. Co-authored-by: Cursor --- CHANGELOG.md | 13 +++ src/columnar_reader.c | 35 +++---- src/columnar_tableam.c | 6 +- test/check_ledger.tsv | 10 ++ test/check_ledger_budget.txt | 2 +- test/parallel_am_scan.sh | 108 +++++++++++++++++++++ test/pytest/TESTS.md | 24 +++++ test/pytest/expected_tests.txt | 3 +- test/pytest/test_compare_to_bash.py | 3 +- test/pytest/test_parallel_am_scan.py | 139 +++++++++++++++++++++++++++ test/run_all_versions.sh | 1 + 11 files changed, 321 insertions(+), 23 deletions(-) create mode 100755 test/parallel_am_scan.sh create mode 100644 test/pytest/test_parallel_am_scan.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 2ef4489a..bd6af60d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -578,6 +578,18 @@ true until the next version shipped. This is the third arm in this file to be repaired for counting a string across a whole file. The `deltuples` comment 15 lines above records the first, fixed by scoping; these two were left as whole-file counts and did the same thing again. +- A table-AM parallel scan was a single claimer. + + `pgcolumnar_read_start` treated `phs_nallocated` as a first-wins flag: the + first participant loaded every row group and the others marked themselves + exhausted. Workers launched, then sat idle while one backend (usually the + leader) read the table. The custom-scan path already claims distinct groups + from a shared counter; the AM path now uses `phs_nallocated` the same way, + as a group index, not a mutex. + + Measured with the custom scan off, two workers, and leader participation + off: both workers produced rows (19000 and 31000 of 50000). Restoring + first-wins returns one worker to 0. - `compare_to_bash.py`'s corpus arm called a WRAPPED name fabricated. A name too long for one line is written as adjacent literals, and Python joins them at parse time, @@ -2302,6 +2314,7 @@ true until the next version shipped. ### Fixed + - The standing parity arm graded a hand-written list, and nothing enforced it (#432, #1046). diff --git a/src/columnar_reader.c b/src/columnar_reader.c index b980d020..c6303dd2 100644 --- a/src/columnar_reader.c +++ b/src/columnar_reader.c @@ -986,9 +986,12 @@ PgColumnarRuntimeGroupsRemoved(PgColumnarReadState *readState) /* * pgcolumnar_read_start - * Lazily load the stripe list on the first fetch. For a parallel scan a - * single worker claims the whole scan and the others see it exhausted, - * which is a correct (if not parallel-accelerated) behaviour. + * Lazily load the stripe list on the first fetch. Every parallel + * participant loads the group list; work is claimed per group in + * pgcolumnar_next_group_index from phs_nallocated, the same way the + * custom scan claims from its DSM counter. The old first-wins use of + * that counter left one backend (usually the leader) to read every + * group and the launched workers idle. */ static void pgcolumnar_read_start(PgColumnarReadState *readState) @@ -998,16 +1001,6 @@ pgcolumnar_read_start(PgColumnarReadState *readState) readState->started = true; - if (readState->parallelScan != NULL) - { - ParallelBlockTableScanDesc bpscan = - (ParallelBlockTableScanDesc) readState->parallelScan; - uint64 claim = pg_atomic_fetch_add_u64(&bpscan->phs_nallocated, 1); - - if (claim != 0) - readState->exhausted = true; - } - if (!readState->exhausted) { /* @@ -3156,20 +3149,28 @@ PgColumnarReadFoldColumn(PgColumnarReadState *readState, int attidx, * The next native row group to scan, or -1 when none remain. The native * counterpart of pgcolumnar_next_stripe_index: a parallel custom scan claims * it from the shared atomic so each worker reads distinct row groups (gap - * 23, D6e); a serial scan walks rowGroupIndex. + * 23, D6e); a table-AM parallel scan claims from phs_nallocated the same + * way; a serial scan walks rowGroupIndex. */ static int64 pgcolumnar_next_group_index(PgColumnarReadState *readState) { int ngroups = list_length(readState->rowGroupList); - uint32 gi; + uint64 gi; if (readState->parallelCounter != NULL) gi = pg_atomic_fetch_add_u32(readState->parallelCounter, 1); + else if (readState->parallelScan != NULL) + { + ParallelBlockTableScanDesc bpscan = + (ParallelBlockTableScanDesc) readState->parallelScan; + + gi = pg_atomic_fetch_add_u64(&bpscan->phs_nallocated, 1); + } else - gi = (uint32) readState->rowGroupIndex++; + gi = (uint64) readState->rowGroupIndex++; - return (gi < (uint32) ngroups) ? (int64) gi : -1; + return (gi < (uint64) ngroups) ? (int64) gi : -1; } void diff --git a/src/columnar_tableam.c b/src/columnar_tableam.c index 77de1a3c..592c8924 100644 --- a/src/columnar_tableam.c +++ b/src/columnar_tableam.c @@ -749,7 +749,7 @@ pgcolumnar_scan_getnextslot(TableScanDesc sscan, ScanDirection direction, } /* ------------------------------------------------------------------------- - * parallel scan: single-worker claim (see pgcolumnar_reader.c) + * parallel scan: shared group claim via phs_nallocated (see pgcolumnar_reader.c) * ------------------------------------------------------------------------- */ static Size @@ -1983,8 +1983,8 @@ pgcolumnar_index_build_range_scan(Relation table_rel, Relation index_rel, /* * Obtain the reader. A parallel index build passes the TableScanDesc it * opened with table_beginscan_parallel; that scan already holds a reader - * bound to the shared parallel scan, whose single-participant claim (see - * pgcolumnar_read_start) makes exactly one participant read the whole table. + * bound to the shared parallel scan, whose per-group claim (see + * pgcolumnar_next_group_index) hands each participant distinct row groups. * We must read through that reader, not a private one: a private full-table * reader in every participant would index every row once per participant, * producing duplicate (key, TID) entries. When no scan is supplied (a serial diff --git a/test/check_ledger.tsv b/test/check_ledger.tsv index 2d7c4bdd..112d1bc7 100644 --- a/test/check_ledger.tsv +++ b/test/check_ledger.tsv @@ -1228,3 +1228,13 @@ native_join_vector_agg native_join_vector_agg unique join fold answer equals GUC native_join_vector_agg native_join_vector_agg unique join fold answer equals heap 15;16;17;18;19 never - native_join_vector_agg native_join_vector_agg unique join uses core Agg when GUC off 15;16;17;18;19 never - native_join_vector_agg native_join_vector_agg unique join uses vectorized agg when GUC on 15;16;17;18;19 2026-09-12 drop JOINREL fold +parallel_am_scan parallel_am_scan a parallel table-AM scan returns the same row count as serial 18 never - +parallel_am_scan parallel_am_scan premise: ANALYZE printed a rows= line per launched worker 18 never - +parallel_am_scan parallel_am_scan premise: EXPLAIN ANALYZE launched two workers 18 never - +parallel_am_scan parallel_am_scan premise: the parallel plan has Gather 18 never - +parallel_am_scan parallel_am_scan premise: the parallel plan is still a Seq Scan, not a custom scan 18 never - +parallel_am_scan parallel_am_scan premise: the parallel plan uses two workers 18 never - +parallel_am_scan parallel_am_scan premise: the serial plan is not a columnar custom scan 18 never - +parallel_am_scan parallel_am_scan premise: the table holds every inserted row 18 never - +parallel_am_scan parallel_am_scan premise: with the custom scan off the serial plan is a Seq Scan 18 never - +parallel_am_scan parallel_am_scan workers share the table-AM scan, it is not a single claimer 18 never - diff --git a/test/check_ledger_budget.txt b/test/check_ledger_budget.txt index b2ec91dc..c48cc743 100644 --- a/test/check_ledger_budget.txt +++ b/test/check_ledger_budget.txt @@ -58,4 +58,4 @@ suites_not_covered 249 # that is not this one. Neither survives. Re-derived by COUNTING on the merged tree, # which is the only resolution this number has: # awk -F'\t' '$5=="never"' test/check_ledger.tsv | wc -l -checks_never_observed_red 1222 +checks_never_observed_red 1232 diff --git a/test/parallel_am_scan.sh b/test/parallel_am_scan.sh new file mode 100755 index 00000000..7e38cf87 --- /dev/null +++ b/test/parallel_am_scan.sh @@ -0,0 +1,108 @@ +#!/usr/bin/env bash +# +# pgColumnar: a table-AM parallel scan must share work across workers. +# +# With the custom scan off, Parallel Seq Scan goes through the AM. The AM +# used to treat phs_nallocated as a first-wins flag: one backend claimed +# the whole scan and the others marked themselves exhausted. Workers +# launched, one backend read. +# +# The custom-scan path already claims distinct row groups from a shared +# counter. This suite pins the AM path to the same property, via EXPLAIN +# ANALYZE worker rows -- not internals. Leader participation is off so +# the two launched workers are the claimers under test, not the leader. +# Many small row groups keep both workers busy before either finishes +# the table. +# +# Independent of test/pytest/test_parallel_am_scan.py: same public seam, own +# fixture, own observations. +# +# Usage: test/parallel_am_scan.sh [PG_CONFIG] +# Written fresh for pgColumnar. + +set -uo pipefail +. "$(dirname "${BASH_SOURCE[0]}")/lib.sh" +pgc_setup "${1:-/usr/local/pg17/bin/pg_config}" + +N=50000 +psql_run "CREATE TABLE pam (id int, k int, payload text) USING pgcolumnar;" +psql_run "SELECT pgcolumnar.set_options('pam', chunk_group_row_limit => 100, stripe_row_limit => 1000);" +psql_run "INSERT INTO pam SELECT g, g%23, md5(g::text) FROM generate_series(1,$N) g;" +psql_run "ALTER TABLE pam SET (parallel_workers = 2);" +psql_run "ANALYZE pam;" + +setg() { q "ALTER DATABASE $PGC_DB SET $1 = $2;" >/dev/null; } +setg pgcolumnar.enable_custom_scan off +setg parallel_setup_cost 0 +setg parallel_tuple_cost 0 +setg min_parallel_table_scan_size 0 +setg jit off +setg parallel_leader_participation off + +explain_text() { + # $1 = max_parallel_workers_per_gather + # $2 = ANALYZE or empty + env PATH="$PGC_BINDIR:$PATH" psql -h 127.0.0.1 -p "$PGC_PORT" -U postgres \ + -d "$PGC_DB" -Atq \ + -c "SET max_parallel_workers_per_gather = $1;" \ + -c "EXPLAIN (COSTS OFF, VERBOSE $2) SELECT id FROM pam;" +} + +serial_plan="$(explain_text 0 "")" +par_plan="$(explain_text 2 "")" +par_ana="$(explain_text 2 ", ANALYZE, TIMING OFF, SUMMARY OFF")" + +echo "-- serial plan --" +echo "$serial_plan" +echo "-- parallel plan --" +echo "$par_plan" +echo "-- parallel analyze --" +echo "$par_ana" + +# Per-worker actual rows from ANALYZE text. A worker that produced nothing +# still prints rows=0, so a missing line is not a zero -- it is no measurement. +worker_rows() { + echo "$1" | grep -oE 'Worker [0-9]+:.*rows=[0-9]+' \ + | grep -oE 'rows=[0-9]+' | grep -oE '[0-9]+' +} + +rows_list="$(worker_rows "$par_ana")" +n_lines="$(echo "$rows_list" | grep -c . || true)" +n_busy="$(echo "$rows_list" | awk '$1>0{n++} END{print n+0}')" +echo "-- worker rows: $(echo "$rows_list" | tr "\n" " ") busy=$n_busy lines=$n_lines" + +check "premise: the table holds every inserted row" \ + "$(q "SELECT count(*) FROM pam")" "$N" + +check "premise: with the custom scan off the serial plan is a Seq Scan" \ + "$(echo "$serial_plan" | grep -c 'Seq Scan')" "1" + +check "premise: the serial plan is not a columnar custom scan" \ + "$(echo "$serial_plan" | grep -c 'Custom Scan')" "0" + +check "premise: the parallel plan has Gather" \ + "$(echo "$par_plan" | grep -c 'Gather')" "1" + +check "premise: the parallel plan uses two workers" \ + "$(echo "$par_plan" | grep -oE 'Workers Planned: [0-9]+' | head -1 | grep -oE '[0-9]+')" "2" + +check "premise: the parallel plan is still a Seq Scan, not a custom scan" \ + "$(echo "$par_ana" | grep -c 'Seq Scan')" "1" + +check "premise: EXPLAIN ANALYZE launched two workers" \ + "$(echo "$par_ana" | grep -oE 'Workers Launched: [0-9]+' | head -1 | grep -oE '[0-9]+')" "2" + +check "premise: ANALYZE printed a rows= line per launched worker" \ + "$n_lines" "2" + +serial_cnt="$(q "SET max_parallel_workers_per_gather = 0; SELECT count(*) FROM pam;" | grep -v '^SET$' | tail -1)" +par_cnt="$(q "SET max_parallel_workers_per_gather = 2; SELECT count(*) FROM pam;" | grep -v '^SET$' | tail -1)" +check "a parallel table-AM scan returns the same row count as serial" \ + "$par_cnt" "$serial_cnt" + +# THE DEFECT: one backend's rows=N and every other worker's rows=0. +# Sharing means both launched workers produced rows. +check "workers share the table-AM scan, it is not a single claimer" \ + "$n_busy" "2" + +pgc_summary diff --git a/test/pytest/TESTS.md b/test/pytest/TESTS.md index 779a0735..c7d75dc6 100644 --- a/test/pytest/TESTS.md +++ b/test/pytest/TESTS.md @@ -89,6 +89,7 @@ behaviour, the source of that number is named. - [41. test_projections.py: a second copy of some columns, kept honest](#41-test_projectionspy-a-second-copy-of-some-columns-kept-honest) - [42. test_compression_reaches_the_cascade.py: the codec setting decides encodings too](#42-test_compression_reaches_the_cascadepy-the-codec-setting-decides-encodings-too) - [43. test_pgxn_metadata.py: the published distribution metadata, which nothing read](#43-test_pgxn_metadatapy-the-published-distribution-metadata-which-nothing-read) +- [44. test_parallel_am_scan.py: a table-AM parallel scan must share work](#44-test_parallel_am_scanpy-a-table-am-parallel-scan-must-share-work) ## 1. How to read a test in here @@ -4338,3 +4339,26 @@ Removal proof, run on both harnesses: restore `META.json` as it shipped and the substantive arms redden on each side while every premise stays green. The premises hold because the file still parses and still names *a* script -- it names the wrong one, which is exactly the distinction the arms draw. + +## 44. test_parallel_am_scan.py: a table-AM parallel scan must share work + +The port of `test/parallel_am_scan.sh`. With the custom scan off, Parallel Seq +Scan goes through the table AM. `phs_nallocated` was a first-wins flag: one +backend claimed the whole scan and every launched worker reported 0 rows. +The custom-scan path already claims distinct row groups; this pair pins the +AM path to the same property. + +Public seam: `EXPLAIN ANALYZE` worker rows on a Parallel Seq Scan. Leader +participation is off so the two launched workers are the claimers under +test. The shell twin uses its own table (`pam`, 50000 rows, groups of 100); +this file uses `ampar`, 80000 rows, groups of 200. Assertion names match. + +### Every arm + +| test | what it holds | +| --- | --- | +| `test_parallel_am_scan` | the serial plan is a Seq Scan, not a custom scan; the parallel plan is a Seq Scan under Gather with two workers launched; a parallel AM scan returns the same count as serial; both launched workers produced rows | + +The load-bearing assertion is `workers share the table-AM scan, it is not a +single claimer`. It is unreachable while `phs_nallocated` is first-wins, and +reachable only when each worker claims its own row groups. diff --git a/test/pytest/expected_tests.txt b/test/pytest/expected_tests.txt index 0d9e1bd5..d1c7a8a9 100644 --- a/test/pytest/expected_tests.txt +++ b/test/pytest/expected_tests.txt @@ -237,4 +237,5 @@ guard_tests 346 # `guard_tests` was re-derived in the same run and did NOT move -- 342 -- which is # the expected answer for a file that needs a cluster, and checking it was the point # rather than assuming it. -cluster_tests 410 +# cluster_tests re-derived by collection on the rebased tree, never by adding a delta measured on another tree. +cluster_tests 411 diff --git a/test/pytest/test_compare_to_bash.py b/test/pytest/test_compare_to_bash.py index e417f61d..ae82f58f 100644 --- a/test/pytest/test_compare_to_bash.py +++ b/test/pytest/test_compare_to_bash.py @@ -86,7 +86,8 @@ # The shape is `SHELL_REFERENCES`' in `test_harness_deps.py`, asserted in both # directions for the same reason: a one-way list rots into a permanent exemption. COMPLETE = ["differential", "hilbert_cluster", "hilbert_locality", - "native_ownership", "native_projection", "projection_privilege", + "native_ownership", "native_projection", "parallel_am_scan", + "projection_privilege", "projections", "sorted_pathkeys", "stats_privilege", "zonemap_boundaries"] diff --git a/test/pytest/test_parallel_am_scan.py b/test/pytest/test_parallel_am_scan.py new file mode 100644 index 00000000..c7dc8688 --- /dev/null +++ b/test/pytest/test_parallel_am_scan.py @@ -0,0 +1,139 @@ +"""A table-AM parallel scan must share work across workers. + +With the custom scan off, Parallel Seq Scan goes through the AM. The AM +used to treat phs_nallocated as a first-wins flag: one backend claimed +the whole scan and the others marked themselves exhausted. Workers +launched, one backend read. + +Independent of test/parallel_am_scan.sh: same public seam (EXPLAIN ANALYZE +of a Parallel Seq Scan), own fixture, own observations. Assertion names +match the shell suite so the two can be compared by name, not by importing +each other. Leader participation is off so the two launched workers are +the claimers under test. Many small row groups keep both workers busy +before either finishes the table. +""" + + +def _nodes(plan): + stack = [plan[0]["Plan"]] + while stack: + node = stack.pop(0) + yield node + stack.extend(node.get("Plans") or ()) + + +def _first(plan, node_type): + for node in _nodes(plan): + if node.get("Node Type") == node_type: + return node + return None + + +def _worker_rows(plan): + rows = [] + scan = _first(plan, "Seq Scan") + if scan is None: + return rows + for worker in scan.get("Workers") or (): + if "Actual Rows" in worker: + rows.append(worker["Actual Rows"]) + return rows + + +def _plan(conn, workers, analyze=False): + opts = "ANALYZE, VERBOSE, TIMING OFF, SUMMARY OFF, " if analyze else "VERBOSE, " + with conn.cursor() as cur: + cur.execute(f"SET max_parallel_workers_per_gather = {workers}") + cur.execute( + f"EXPLAIN ({opts}FORMAT JSON, COSTS OFF) SELECT id FROM ampar" + ) + return cur.fetchone()[0] + + +def test_parallel_am_scan(pgc_conn, expect): + n = 80000 + with pgc_conn.cursor() as cur: + cur.execute( + "CREATE TABLE ampar (id int, k int, payload text) USING pgcolumnar" + ) + cur.execute( + "SELECT pgcolumnar.set_options('ampar', " + "chunk_group_row_limit => 200, stripe_row_limit => 1000)" + ) + cur.execute( + f"INSERT INTO ampar SELECT g, g, md5(g::text) " + f"FROM generate_series(1, {n}) g" + ) + cur.execute("ALTER TABLE ampar SET (parallel_workers = 2)") + cur.execute("ANALYZE ampar") + cur.execute("SELECT count(*) FROM ampar") + expect.num(cur.fetchone()[0], n, "premise: the table holds every inserted row") + + cur.execute("SET pgcolumnar.enable_custom_scan = off") + cur.execute("SET parallel_setup_cost = 0") + cur.execute("SET parallel_tuple_cost = 0") + cur.execute("SET min_parallel_table_scan_size = 0") + cur.execute("SET jit = off") + cur.execute("SET parallel_leader_participation = off") + + serial = _plan(pgc_conn, 0) + parallel = _plan(pgc_conn, 2) + analyzed = _plan(pgc_conn, 2, analyze=True) + + expect.text( + "Seq Scan" if _first(serial, "Seq Scan") else "none", + "Seq Scan", + "premise: with the custom scan off the serial plan is a Seq Scan", + ) + expect.text( + "none" if _first(serial, "Custom Scan") is None else "Custom Scan", + "none", + "premise: the serial plan is not a columnar custom scan", + ) + expect.text( + "Gather" if _first(parallel, "Gather") is not None else "none", + "Gather", + "premise: the parallel plan has Gather", + ) + expect.num( + (_first(parallel, "Gather") or {}).get("Workers Planned"), + 2, + "premise: the parallel plan uses two workers", + ) + expect.text( + "Seq Scan" if _first(analyzed, "Seq Scan") else "none", + "Seq Scan", + "premise: the parallel plan is still a Seq Scan, not a custom scan", + ) + expect.num( + (_first(analyzed, "Gather") or {}).get("Workers Launched"), + 2, + "premise: EXPLAIN ANALYZE launched two workers", + ) + + with pgc_conn.cursor() as cur: + cur.execute("SET max_parallel_workers_per_gather = 0") + cur.execute("SELECT count(*) FROM ampar") + serial_cnt = cur.fetchone()[0] + cur.execute("SET max_parallel_workers_per_gather = 2") + cur.execute("SELECT count(*) FROM ampar") + par_cnt = cur.fetchone()[0] + expect.num( + par_cnt, + serial_cnt, + "a parallel table-AM scan returns the same row count as serial", + ) + + worker_rows = _worker_rows(analyzed) + expect.num( + len(worker_rows), + 2, + "premise: ANALYZE printed a rows= line per launched worker", + ) + n_busy = sum(1 for r in worker_rows if r and r > 0) + print(f"-- worker rows {worker_rows} busy={n_busy}") + expect.num( + n_busy, + 2, + "workers share the table-AM scan, it is not a single claimer", + ) diff --git a/test/run_all_versions.sh b/test/run_all_versions.sh index 3f297e3c..7556d322 100755 --- a/test/run_all_versions.sh +++ b/test/run_all_versions.sh @@ -223,6 +223,7 @@ SUITES=( objstore_tls_read objstore_userinfo parallel + parallel_am_scan parallel_copy parallel_copy_dedup parallel_degree From e6faa3b572e8b0d6fb8c65befa1d935b296382a8 Mon Sep 17 00:00:00 2001 From: "Joshua D. Drake" Date: Wed, 16 Sep 2026 20:35:09 +0000 Subject: [PATCH 2/3] test: merge PG15-18 logs so parallel_am_scan names those majors CI suites (PG 17) refused these checks because a PG18-only seed left majors=18. The suite was run on PGDG 15.19, 16.15, 17.11 and Ubuntu 18.6 and those logs were merged. PG19 is not installed here. Co-authored-by: Cursor --- test/check_ledger.tsv | 20 ++++++++++---------- 1 file changed, 10 insertions(+), 10 deletions(-) diff --git a/test/check_ledger.tsv b/test/check_ledger.tsv index 112d1bc7..4f275128 100644 --- a/test/check_ledger.tsv +++ b/test/check_ledger.tsv @@ -1228,13 +1228,13 @@ native_join_vector_agg native_join_vector_agg unique join fold answer equals GUC native_join_vector_agg native_join_vector_agg unique join fold answer equals heap 15;16;17;18;19 never - native_join_vector_agg native_join_vector_agg unique join uses core Agg when GUC off 15;16;17;18;19 never - native_join_vector_agg native_join_vector_agg unique join uses vectorized agg when GUC on 15;16;17;18;19 2026-09-12 drop JOINREL fold -parallel_am_scan parallel_am_scan a parallel table-AM scan returns the same row count as serial 18 never - -parallel_am_scan parallel_am_scan premise: ANALYZE printed a rows= line per launched worker 18 never - -parallel_am_scan parallel_am_scan premise: EXPLAIN ANALYZE launched two workers 18 never - -parallel_am_scan parallel_am_scan premise: the parallel plan has Gather 18 never - -parallel_am_scan parallel_am_scan premise: the parallel plan is still a Seq Scan, not a custom scan 18 never - -parallel_am_scan parallel_am_scan premise: the parallel plan uses two workers 18 never - -parallel_am_scan parallel_am_scan premise: the serial plan is not a columnar custom scan 18 never - -parallel_am_scan parallel_am_scan premise: the table holds every inserted row 18 never - -parallel_am_scan parallel_am_scan premise: with the custom scan off the serial plan is a Seq Scan 18 never - -parallel_am_scan parallel_am_scan workers share the table-AM scan, it is not a single claimer 18 never - +parallel_am_scan parallel_am_scan a parallel table-AM scan returns the same row count as serial 15;16;17;18 never - +parallel_am_scan parallel_am_scan premise: ANALYZE printed a rows= line per launched worker 15;16;17;18 never - +parallel_am_scan parallel_am_scan premise: EXPLAIN ANALYZE launched two workers 15;16;17;18 never - +parallel_am_scan parallel_am_scan premise: the parallel plan has Gather 15;16;17;18 never - +parallel_am_scan parallel_am_scan premise: the parallel plan is still a Seq Scan, not a custom scan 15;16;17;18 never - +parallel_am_scan parallel_am_scan premise: the parallel plan uses two workers 15;16;17;18 never - +parallel_am_scan parallel_am_scan premise: the serial plan is not a columnar custom scan 15;16;17;18 never - +parallel_am_scan parallel_am_scan premise: the table holds every inserted row 15;16;17;18 never - +parallel_am_scan parallel_am_scan premise: with the custom scan off the serial plan is a Seq Scan 15;16;17;18 never - +parallel_am_scan parallel_am_scan workers share the table-AM scan, it is not a single claimer 15;16;17;18 never - From 4ec9a3127e9933f4ea4b8408c70f6b1931c09705 Mon Sep 17 00:00:00 2001 From: OffgridwithJD Date: Thu, 17 Sep 2026 14:55:29 +0000 Subject: [PATCH 3/3] test: a parallel index build must cover the whole table (#1068 review) The scan arms drive pgcolumnar_next_group_index through a parallel seq scan. A parallel index build reaches the same shared claim through table_beginscan_parallel, and it was untested -- which matters because it is the consumer where a claim bug is silent. A scan that double-claims returns duplicate rows and someone notices; an index that SKIPS a group is simply missing entries, and every query using it quietly returns fewer rows. THE WORKER COUNT IS NOT A pgcolumnar GUC, and max_parallel_maintenance_workers alone will not produce one. Measured: max_parallel_maintenance_workers = 0 "building index ... serially" max_parallel_maintenance_workers = 2 "with request for 1 parallel workers" max_parallel_maintenance_workers = 8 "with request for 1 parallel workers" + ALTER TABLE ... SET (parallel_workers = 8) "with request for 8 parallel workers" That GUC is a gate, not a dial. Core sizes the request in plan_create_index_workers() from relpages, and a columnar table reports 69 pages for 2,000,000 rows, so the size heuristic grants ONE worker however large the fixture. The table's parallel_workers reloption is the only thing that produces real parallelism here. Both arms say so in a comment, because a bigger table is what the next person will reach for. PROVED IN BOTH DIRECTIONS by mutating the claim stride, a realistic off-by-one in the shared counter: as proposed 14 passed + 0 failed stride 1 -> 2 FAIL a parallel index build indexes every row of the table: got [25000|612512500] want [50000|1250025000] The SUM is doing real work there: a group read twice cancelling a group skipped leaves the count right and the sum wrong. Both harnesses assert the same four names and observe independently -- the shell suite reads PGC_LOGFILE with awk and compares a concatenated string, the pytest twin reads the cluster's own server.log through the pgc_cluster fixture and compares a tuple. Neither invokes the other. compare_to_bash grades them one-for-one: 128/128. Ledger rows re-derived from runs on all five majors rather than by editing the majors field. Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_012RSw4qMHS7ByE7PY8Ns4cs --- test/check_ledger.tsv | 20 ++--- test/parallel_am_scan.sh | 80 +++++++++++++++++ test/pytest/TESTS.md | 1 + test/pytest/test_parallel_am_scan.py | 123 +++++++++++++++++++++++++++ 4 files changed, 214 insertions(+), 10 deletions(-) diff --git a/test/check_ledger.tsv b/test/check_ledger.tsv index 4f275128..51ff4521 100644 --- a/test/check_ledger.tsv +++ b/test/check_ledger.tsv @@ -1228,13 +1228,13 @@ native_join_vector_agg native_join_vector_agg unique join fold answer equals GUC native_join_vector_agg native_join_vector_agg unique join fold answer equals heap 15;16;17;18;19 never - native_join_vector_agg native_join_vector_agg unique join uses core Agg when GUC off 15;16;17;18;19 never - native_join_vector_agg native_join_vector_agg unique join uses vectorized agg when GUC on 15;16;17;18;19 2026-09-12 drop JOINREL fold -parallel_am_scan parallel_am_scan a parallel table-AM scan returns the same row count as serial 15;16;17;18 never - -parallel_am_scan parallel_am_scan premise: ANALYZE printed a rows= line per launched worker 15;16;17;18 never - -parallel_am_scan parallel_am_scan premise: EXPLAIN ANALYZE launched two workers 15;16;17;18 never - -parallel_am_scan parallel_am_scan premise: the parallel plan has Gather 15;16;17;18 never - -parallel_am_scan parallel_am_scan premise: the parallel plan is still a Seq Scan, not a custom scan 15;16;17;18 never - -parallel_am_scan parallel_am_scan premise: the parallel plan uses two workers 15;16;17;18 never - -parallel_am_scan parallel_am_scan premise: the serial plan is not a columnar custom scan 15;16;17;18 never - -parallel_am_scan parallel_am_scan premise: the table holds every inserted row 15;16;17;18 never - -parallel_am_scan parallel_am_scan premise: with the custom scan off the serial plan is a Seq Scan 15;16;17;18 never - -parallel_am_scan parallel_am_scan workers share the table-AM scan, it is not a single claimer 15;16;17;18 never - +parallel_am_scan parallel_am_scan a parallel table-AM scan returns the same row count as serial 15;16;17;18;19 never - +parallel_am_scan parallel_am_scan premise: ANALYZE printed a rows= line per launched worker 15;16;17;18;19 never - +parallel_am_scan parallel_am_scan premise: EXPLAIN ANALYZE launched two workers 15;16;17;18;19 never - +parallel_am_scan parallel_am_scan premise: the parallel plan has Gather 15;16;17;18;19 never - +parallel_am_scan parallel_am_scan premise: the parallel plan is still a Seq Scan, not a custom scan 15;16;17;18;19 never - +parallel_am_scan parallel_am_scan premise: the parallel plan uses two workers 15;16;17;18;19 never - +parallel_am_scan parallel_am_scan premise: the serial plan is not a columnar custom scan 15;16;17;18;19 never - +parallel_am_scan parallel_am_scan premise: the table holds every inserted row 15;16;17;18;19 never - +parallel_am_scan parallel_am_scan premise: with the custom scan off the serial plan is a Seq Scan 15;16;17;18;19 never - +parallel_am_scan parallel_am_scan workers share the table-AM scan, it is not a single claimer 15;16;17;18;19 never - diff --git a/test/parallel_am_scan.sh b/test/parallel_am_scan.sh index 7e38cf87..b79fdda8 100755 --- a/test/parallel_am_scan.sh +++ b/test/parallel_am_scan.sh @@ -105,4 +105,84 @@ check "a parallel table-AM scan returns the same row count as serial" \ check "workers share the table-AM scan, it is not a single claimer" \ "$n_busy" "2" +# ---- a parallel INDEX BUILD is the other consumer of the shared claim -------- +# +# Everything above drives the shared group claim through a parallel SEQ SCAN. +# A parallel index build reaches the same pgcolumnar_next_group_index through +# table_beginscan_parallel, and it is the consumer where a claim bug is silent: +# a scan that double-claims returns duplicate rows and someone notices, while an +# index that SKIPS a group is simply missing entries and every query using it +# quietly returns fewer rows. +# +# THE WORKER COUNT IS NOT A pgcolumnar GUC, AND max_parallel_maintenance_workers +# ALONE WILL NOT PRODUCE ONE. That GUC is a gate -- 0 builds serially -- but core +# sizes the request in plan_create_index_workers() from relpages, and a columnar +# table reports very few pages for many rows (measured: 69 pages for 2,000,000), +# so the size heuristic grants ONE worker however large the table is. The table's +# `parallel_workers` reloption is the only dial that produces real parallelism +# here, which is why it is set below and why a bigger fixture would not help. +psql_run "ALTER TABLE pam SET (parallel_workers = 4);" + +# Own reader: `q` runs psql without -q, so a multi-statement call prints a SET +# line per SET and the value under test would be whatever came last. This takes +# the final line after dropping those. grep -v and tail both read to EOF, so +# neither can SIGPIPE the writer (#486). +_pam_read() { + env PATH="$PGC_BINDIR:$PATH" psql -h 127.0.0.1 -p "$PGC_PORT" -U postgres \ + -d "$PGC_DB" -Atq -c "$1" 2>/dev/null | grep -v '^SET$' | tail -1 +} + +_pam_mark="pam_build_$$" +q "DO \$\$ BEGIN RAISE LOG '$_pam_mark'; END \$\$;" >/dev/null +q "DROP INDEX IF EXISTS pam_idx;" >/dev/null +q "SET log_min_messages = debug1; + SET max_parallel_maintenance_workers = 4; + SET min_parallel_table_scan_size = 0; + CREATE INDEX pam_idx ON pam (id);" >/dev/null + +# Scoped to a marker this run wrote, so a build from an earlier run in the same +# cluster cannot answer for this one. +_pam_req="$(awk -v m="$_pam_mark" ' + p && /with request for/ { print; exit } + $0 ~ m { p = 1 } + ' "${PGC_LOGFILE:-/dev/null}")" +_pam_nreq="$(printf '%s' "$_pam_req" | tr -dc '0-9 ' | awk '{print $1+0}')" + +check "premise: the index build requested parallel workers" \ + "$([ -n "$_pam_req" ] && [ "${_pam_nreq:-0}" -ge 2 ] && echo yes || echo "no (${_pam_nreq:-none})")" "yes" + +# The plan is classified from a captured string with `case`, not a pipe into an +# early-exit reader. +_pam_planout="$(_pam_read "SET enable_seqscan=off; + SET pgcolumnar.enable_custom_scan=off; + EXPLAIN (COSTS OFF) SELECT count(*) FROM pam WHERE id > 0;")" +_pam_planall="$(env PATH="$PGC_BINDIR:$PATH" psql -h 127.0.0.1 -p "$PGC_PORT" -U postgres \ + -d "$PGC_DB" -Atq -c "SET enable_seqscan=off; + SET pgcolumnar.enable_custom_scan=off; + EXPLAIN (COSTS OFF) SELECT count(*) FROM pam WHERE id > 0;" 2>/dev/null)" +case "$_pam_planall" in + *"Index Only Scan"*) _pam_node="Index Only Scan" ;; + *"Index Scan"*) _pam_node="Index Scan" ;; + *"Seq Scan"*) _pam_node="Seq Scan" ;; + *) _pam_node="" ;; +esac + +check "premise: the comparison reads the table through the index" \ + "$_pam_node" "Index Only Scan" + +# THE PROPERTY: the index built in parallel describes the whole table. Compared +# as an aggregate through the index against the same aggregate through a +# sequential scan -- a count alone would miss a group read twice and a group +# skipped cancelling out, which the sum does not. +_pam_idx="$(_pam_read "SET enable_seqscan=off; SET pgcolumnar.enable_custom_scan=off; + SELECT count(*) || '|' || coalesce(sum(id),0) FROM pam WHERE id > 0;")" +_pam_seq="$(_pam_read "SET enable_indexscan=off; SET enable_bitmapscan=off; + SELECT count(*) || '|' || coalesce(sum(id),0) FROM pam WHERE id > 0;")" + +check "premise: both sides of the comparison returned a value" \ + "$([ -n "$_pam_idx" ] && [ -n "$_pam_seq" ] && echo yes || echo no)" "yes" + +check "a parallel index build indexes every row of the table" \ + "$_pam_idx" "$_pam_seq" + pgc_summary diff --git a/test/pytest/TESTS.md b/test/pytest/TESTS.md index c7d75dc6..b99d31cb 100644 --- a/test/pytest/TESTS.md +++ b/test/pytest/TESTS.md @@ -4358,6 +4358,7 @@ this file uses `ampar`, 80000 rows, groups of 200. Assertion names match. | test | what it holds | | --- | --- | | `test_parallel_am_scan` | the serial plan is a Seq Scan, not a custom scan; the parallel plan is a Seq Scan under Gather with two workers launched; a parallel AM scan returns the same count as serial; both launched workers produced rows | +| `test_a_parallel_index_build_covers_the_whole_table` | a parallel index build requests workers and indexes every row -- compared as count and SUM through the index against a sequential scan, because a group read twice cancelling a group skipped leaves the count right | The load-bearing assertion is `workers share the table-AM scan, it is not a single claimer`. It is unreachable while `phs_nallocated` is first-wins, and diff --git a/test/pytest/test_parallel_am_scan.py b/test/pytest/test_parallel_am_scan.py index c7dc8688..4715f319 100644 --- a/test/pytest/test_parallel_am_scan.py +++ b/test/pytest/test_parallel_am_scan.py @@ -13,6 +13,11 @@ before either finishes the table. """ +import pathlib +import re +import uuid + + def _nodes(plan): stack = [plan[0]["Plan"]] @@ -137,3 +142,121 @@ def test_parallel_am_scan(pgc_conn, expect): 2, "workers share the table-AM scan, it is not a single claimer", ) + + +# ---- a parallel INDEX BUILD is the other consumer of the shared claim -------- +# +# The scan arms above drive pgcolumnar_next_group_index through a parallel seq +# scan. A parallel index build reaches the same claim through +# table_beginscan_parallel, and it is the consumer where a claim bug is silent: +# a scan that double-claims returns duplicate rows and someone notices, while an +# index that SKIPS a group is simply missing entries. +# +# THE WORKER COUNT IS NOT A pgcolumnar GUC, and max_parallel_maintenance_workers +# alone will not produce one. That GUC is a gate -- 0 builds serially -- but core +# sizes the request in plan_create_index_workers() from relpages, and a columnar +# table reports very few pages for many rows (measured: 69 pages for 2,000,000), +# so the heuristic grants ONE worker however large the fixture. The table's +# `parallel_workers` reloption is the only dial. +# +# Independent of test/parallel_am_scan.sh: that suite reads PGC_LOGFILE with awk +# and compares a concatenated string; this one reads the cluster's own +# server.log through the pgc_cluster fixture and compares a tuple. Same seam, +# separate observers, neither invoking the other. +def _requested_workers(cluster, marker): + """Workers the build asked for, from the first request line after MARKER. + + Scoped to a marker this test wrote, so a build from an earlier test in the + same session cluster cannot answer for this one. + """ + log = pathlib.Path(cluster.datadir) / "server.log" + if not log.is_file(): + return None + seen = False + for line in log.read_text(encoding="utf-8", errors="replace").splitlines(): + if not seen: + if marker in line: + seen = True + continue + m = re.search(r"with request for (\d+) parallel worker", line) + if m: + return int(m.group(1)) + return None + + +def _scan_node_of(plan): + stack = [plan[0]["Plan"]] + while stack: + node = stack.pop(0) + t = node.get("Node Type", "") + if t in ("Index Only Scan", "Index Scan", "Seq Scan", "Custom Scan"): + return t + stack.extend(node.get("Plans") or ()) + return "" + + +def test_a_parallel_index_build_covers_the_whole_table(pgc_cluster, pgc_conn, expect): + """The index built in parallel must describe every row.""" + marker = "pamidx_" + uuid.uuid4().hex[:12] + with pgc_conn.cursor() as cur: + cur.execute("CREATE TABLE pamidx (id int, v int) USING pgcolumnar") + cur.execute( + "INSERT INTO pamidx SELECT g, g % 1000 FROM generate_series(1, 300000) g" + # single %, not %%: psycopg only un-doubles when it is + # interpolating, and this call passes no parameters. + ) + cur.execute("ALTER TABLE pamidx SET (parallel_workers = 4)") + cur.execute("ANALYZE pamidx") + # The marker is interpolated, not bound: a placeholder inside a DO + # body cannot be typed ("could not determine data type of parameter + # $1"), because the body is a string literal to the server. Safe to + # interpolate and asserted so -- it is uuid4 hex generated here. + assert marker.replace("_", "").isalnum(), marker + cur.execute("DO $pgcmark$ BEGIN RAISE LOG '%s'; END $pgcmark$" % marker) + cur.execute("SET log_min_messages = debug1") + cur.execute("SET max_parallel_maintenance_workers = 4") + cur.execute("SET min_parallel_table_scan_size = 0") + cur.execute("CREATE INDEX pamidx_id ON pamidx (id)") + cur.execute("RESET log_min_messages") + + req = _requested_workers(pgc_cluster, marker) + expect.num( + 1 if (req is not None and req >= 2) else 0, + 1, + "premise: the index build requested parallel workers", + ) + + with pgc_conn.cursor() as cur: + cur.execute("SET enable_seqscan = off") + cur.execute("SET pgcolumnar.enable_custom_scan = off") + cur.execute( + "EXPLAIN (FORMAT JSON, COSTS OFF) " + "SELECT count(*) FROM pamidx WHERE id > 0" + ) + plan = cur.fetchone()[0] + cur.execute("SELECT count(*), coalesce(sum(id), 0) FROM pamidx WHERE id > 0") + via_index = cur.fetchone() + expect.text( + _scan_node_of(plan), + "Index Only Scan", + "premise: the comparison reads the table through the index", + ) + + with pgc_conn.cursor() as cur: + cur.execute("SET enable_indexscan = off") + cur.execute("SET enable_bitmapscan = off") + cur.execute("SELECT count(*), coalesce(sum(id), 0) FROM pamidx WHERE id > 0") + via_seq = cur.fetchone() + + expect.num( + 1 if (via_index is not None and via_seq is not None) else 0, + 1, + "premise: both sides of the comparison returned a value", + ) + # The SUM, not just the count: a group read twice cancelling a group skipped + # leaves the count right and the sum wrong. + expect.rows( + [via_index], + [via_seq], + "a parallel index build indexes every row of the table", + )