diff --git a/CHANGELOG.md b/CHANGELOG.md index 1b7ca1fa..3a3c1a75 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -615,87 +615,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. -- An index fetch pinned once per projected column, while a sequential scan - already coalesced adjacent chunk ranges into one read. - - `pgcolumnar_fetch_row` issued two `PgColumnarReadLogicalData` calls per - column (validity bitmap, then the value stream). The scan path - (`pgcolumnar_native_read_projected`) sorts those ranges and merges the ones - that touch. Adjacent columns in a row group are laid out back to back, so a - wide btree fetch of a small group pinned the same pages once per column. - - Measured on PostgreSQL 18 with `EXPLAIN (ANALYZE, BUFFERS)` executor pins - (planning excluded): 16 int columns, one row via the index, 64 pins for one - column and 94 for sixteen -- exactly two extra pins per extra column. After - the fetch path coalesces the same way the scan does, both counts are 61. - New twins `native_fetch_coalesce` and `test_native_fetch_coalesce.py`. - - THE VALIDITY COPY IS BOUNDED BY THE CHUNK BEFORE IT RUNS. The coalesced path - copies `validityBytes` out of a span buffer that is only guaranteed to hold - `page_length` bytes for the chunk being served, and the test reconciling the - two ran three lines AFTER the copy. A chunk whose catalog `page_length` was - smaller than its validity bitmap therefore read past the allocation. - - Reproduced against a build with `-fsanitize=address`, by poisoning - `pgcolumnar.column_chunk.page_length` on the last chunk by `page_offset` and - issuing a plain index-scan `SELECT`: - - AddressSanitizer: heap-buffer-overflow - READ of size 625, 0 bytes after a 2640-byte region - pgcolumnar_fetch_coalesce_read (the memcpy) - pgcolumnar_fetch_row - printtup - - The backend died and the cluster entered recovery. Main cannot have this - shape: its non-coalesced fill reads straight from storage into an - exactly-sized destination, so no in-memory extent exists to exceed. The span - buffer and the copy out of it are both new here. - - Hoisting the `page_length >= validityBytes` test above the copy closes it. An - inconsistent chunk is left for the non-coalesced path, which refuses it. - - The regression arm is an ORDERING pin, not a behavioural one, and that is - deliberate: reading ~117 bytes past a palloc'd span reads adjacent heap and - returns quietly without a sanitizer, so a behavioural arm would report PASS - on the broken code. Both harnesses assert the order, each reading the source - its own way -- awk over line numbers in the shell suite, a regex over - character offsets in the pytest twin. Proved by MOVING the guard below the - copy rather than deleting it, which leaves both statements present and - reddens only the ordering arm. - - AND A CHUNK THE CHECKED DECODE PATH WOULD REFUSE IS LEFT FOR IT, so the refusal - keeps its SQLSTATE. The range-building loop now defers any chunk whose - page_length is under the validity bitmap or whose value stream would not fit a - uint32. - - Without that, this change SHADOWS #1063's typed refusal. `pgcolumnar_fetch_row` - calls the coalescing helper before the per-column loop reaches - `pgcolumnar_chunk_value_bytes`, and the helper builds its ranges straight from - `page_length`, so a poisoned length spans ~4GB and palloc raises first. - Measured on the two composed: - - without the defer native_chunk_length_bound 5 passed + 1 failed - ERROR: invalid memory alloc request size 4294971754 - with the defer native_chunk_length_bound 6 passed + 0 failed (XX001) - with the defer native_fetch_coalesce 7 passed + 0 failed - - The last line matters: the wide-fetch pin still passes, so deferring the - inconsistent chunk is not disabling coalescing to make a test green. - - Reported by @jdatcmd, who composed the two branches rather than reading them. - - THE ORDERING ARM IS ANCHORED ON THE CONTAINMENT TEST, because the function now - holds two guards with the same text -- the deferral above and the bound on the - copy. An unanchored search finds the first, which is in the wrong loop, and the - arm would then pass with the bound deleted. The containment test belongs only - to the distribution loop. Proved by deleting ONLY that guard and leaving the - deferral: both arms redden. - - The two source patterns use bracket expressions rather than backslash-escaped - parens. `awk -v` processes escapes in the value and `\(` is undefined, so mawk - keeps the backslash and matches while gawk strips it -- silently not matching - for one pattern, and exiting fatally on `Unmatched (` for the other. CI runners - carry gawk. Verified identical under both. +- 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, @@ -2420,6 +2351,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 e623bcf6..8a15c286 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) { /* @@ -3192,20 +3185,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 802e1339..65151499 100644 --- a/test/check_ledger.tsv +++ b/test/check_ledger.tsv @@ -1241,3 +1241,17 @@ 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 index build indexes every row of the table 15;16;17;18;19 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: both sides of the comparison returned a value 15;16;17;18;19 never - +parallel_am_scan parallel_am_scan premise: the comparison reads the table through the index 15;16;17;18;19 never - +parallel_am_scan parallel_am_scan premise: the index build requested parallel 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/check_ledger_budget.txt b/test/check_ledger_budget.txt index 20b5d09c..979fb74f 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 1235 +checks_never_observed_red 1249 diff --git a/test/parallel_am_scan.sh b/test/parallel_am_scan.sh new file mode 100755 index 00000000..d2176a0c --- /dev/null +++ b/test/parallel_am_scan.sh @@ -0,0 +1,195 @@ +#!/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" + +# ---- 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" + +# FAIL CLOSED. If the build errored, BOTH reads come back empty and comparing +# "" against "" reports PASS -- measured: mutating the shared claim so every +# participant walks its own index made the build fail, both reads returned +# nothing, and this check passed on a tree where the property was broken. The +# premise above caught it, but a headline check that says PASS when it measured +# nothing is worse than no check. Distinct sentinels cannot collide. +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 3eb4fabe..369f00a0 100644 --- a/test/pytest/TESTS.md +++ b/test/pytest/TESTS.md @@ -91,6 +91,7 @@ behaviour, the source of that number is named. - [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_native_chunk_length_bound.py: a truncated chunk length cannot fetch](#44-test_native_chunk_length_boundpy-a-truncated-chunk-length-cannot-fetch) - [45. test_native_fetch_coalesce.py: index fetch I/O is not per-column](#45-test_native_fetch_coalescepy-index-fetch-io-is-not-per-column) +- [46. test_parallel_am_scan.py: a table-AM parallel scan must share work](#46-test_parallel_am_scanpy-a-table-am-parallel-scan-must-share-work) ## 1. How to read a test in here @@ -4376,3 +4377,26 @@ own observations. Assertion names match the shell suite. | --- | --- | | `test_native_fetch_coalesce` | a point lookup uses the index and returns the projected values; executor pins for one column and for every column are both measurable, and the wide fetch does not pin once per column | | `test_the_validity_copy_is_bounded_before_the_chunk_is_read` | the bound on the validity copy precedes the copy, read as positions in the coalescing helper rather than as the presence of both statements -- the overread it guards had both | +## 46. 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 | +| `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 +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 7d4728ba..099600a5 100644 --- a/test/pytest/expected_tests.txt +++ b/test/pytest/expected_tests.txt @@ -237,10 +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. -# 410 -> 411 when test_native_fetch_coalesce.py landed on the rebased tree. Re-derived by collection, never by adding one to 410. -# 411 -> 412 when the coalescing suite gained an ordering arm for the validity -# copy. Re-derived by collecting the complement of NO_CLUSTER, the way the job -# builds FILES, never by adding one to 411: `412 tests collected`. guard_tests -# was re-derived in the same run and did NOT move -- 346 -- which is the expected -# answer for a file the cluster leg owns, and checking it was the point. -cluster_tests 413 +# cluster_tests re-derived by collection on the rebased tree, never by adding a delta measured on another tree. +cluster_tests 415 diff --git a/test/pytest/test_compare_to_bash.py b/test/pytest/test_compare_to_bash.py index 0e39dc16..5d7b8930 100644 --- a/test/pytest/test_compare_to_bash.py +++ b/test/pytest/test_compare_to_bash.py @@ -86,7 +86,7 @@ # 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_chunk_length_bound", "native_fetch_coalesce", "native_ownership", "native_projection", + "native_chunk_length_bound", "native_fetch_coalesce", "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..132ad5bb --- /dev/null +++ b/test/pytest/test_parallel_am_scan.py @@ -0,0 +1,265 @@ +"""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. +""" + +import pathlib +import re +import uuid + + + +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", + ) + + +# ---- 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. + # FAIL CLOSED, same reason as the shell twin: if the build errored both + # fetches return None and [None] == [None] would report PASS on a broken + # tree. Distinct sentinels cannot collide. + expect.rows( + [via_index if via_index is not None else ("the index read returned nothing",)], + [via_seq if via_seq is not None else ("the sequential read returned nothing",)], + "a parallel index build indexes every row of the table", + ) diff --git a/test/run_all_versions.sh b/test/run_all_versions.sh index 05741bfd..d755f8ef 100755 --- a/test/run_all_versions.sh +++ b/test/run_all_versions.sh @@ -225,6 +225,7 @@ SUITES=( objstore_tls_read objstore_userinfo parallel + parallel_am_scan parallel_copy parallel_copy_dedup parallel_degree