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