Repository navigation
Conversation
|
run benchmarks |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing parquet-post-scan-filter (7dd85fc) to c8b784a (merge-base) diff using: clickbench_partitioned File an issue against this benchmark runner |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing parquet-post-scan-filter (7dd85fc) to c8b784a (merge-base) diff using: tpcds File an issue against this benchmark runner |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing parquet-post-scan-filter (7dd85fc) to c8b784a (merge-base) diff using: tpch File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: CPU Details (lscpu)Details
Resource Usagetpch — base (merge-base)
tpch — branch
File an issue against this benchmark runner |
7dd85fc to
2fffad2
Compare
|
🤖 Benchmark completed (GKE) | trigger Instance: CPU Details (lscpu)Details
Resource Usagetpcds — base (merge-base)
tpcds — branch
File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: CPU Details (lscpu)Details
Resource Usageclickbench_partitioned — base (merge-base)
clickbench_partitioned — branch
File an issue against this benchmark runner |
1404755 to
cca69df
Compare
|
run benchmarks |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing parquet-post-scan-filter (cca69df) to ad7d6ea (merge-base) diff using: tpcds File an issue against this benchmark runner |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing parquet-post-scan-filter (cca69df) to ad7d6ea (merge-base) diff using: tpch File an issue against this benchmark runner |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing parquet-post-scan-filter (cca69df) to ad7d6ea (merge-base) diff using: clickbench_partitioned File an issue against this benchmark runner |
|
Thank you for opening this pull request! Reviewer note: cargo-semver-checks reported the current version number is not SemVer-compatible with the changes in this pull request (compared against the base branch). Details |
|
🤖 Benchmark completed (GKE) | trigger Instance: CPU Details (lscpu)Details
Resource Usagetpch — base (merge-base)
tpch — branch
File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: CPU Details (lscpu)Details
Resource Usagetpcds — base (merge-base)
tpcds — branch
File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: CPU Details (lscpu)Details
Resource Usageclickbench_partitioned — base (merge-base)
clickbench_partitioned — branch
File an issue against this benchmark runner |
…filter The pushdown=false path in the parquet opener split the whole predicate into 'post_scan_conjuncts' — a per-batch FilterExec-equivalent — which included any dynamic filter conjuncts (HashJoin bounds, TopK threshold, aggregate dynamic filter). For join-heavy TPC-H / TPC-DS this dominates cost: HashJoin's Partitioned- mode dynamic filter is a 'CASE hash(col) % N WHEN pid THEN bounds ELSE lit(false) END' — per-row hash + modulo + CASE branch — and it prunes almost nothing on high-match-rate joins where the downstream hash lookup would eliminate the same rows anyway. Local TPC-H SF1 Q9 profile showed 1.1% self-time in 'expressions::case::PartialResultIndex::merge_n' and 1.3% in 'arrow_select::filter::filter_native' on the PR, both at 0% on main — driving Q9 from 40ms → 80ms (2.09x on CI, 1.79x locally). This commit filters DynamicFilterPhysicalExpr-containing conjuncts out of 'post_scan_conjuncts'. Effects: - RowGroupPruner (added by apache#22450) still sees the full predicate via prepared.predicate, so RG-level dynamic pruning continues to fire on bounds/threshold updates. - pushdown_filters=true path unchanged — dynamic filters still go through the arrow-rs RowFilter. - Downstream operator does the exact equivalent: HashJoin's hash lookup filters rows the bounds would have filtered; TopK's sort heap filters rows the threshold would have filtered. No wrong results. Local TPC-H SF1 (release-nonlto, 3 iters): - baseline (HEAD~2, pre-apache#22384) avg: 27.65 ms - PR + this fix avg: 26.24 ms (net 5% ahead of baseline) - Q9 individually: 34.54 → 34.33 ms (matches baseline, was 80.74 before) Also regenerates push_down_filter_parquet.slt for the membership-off default (from the earlier 'split membership from bounds' commit).
…r (root fix)
Reverts the tactical fix from the previous commit and cures the same
regression at its source. The prior commit skipped
DynamicFilterPhysicalExpr-containing conjuncts from PostScanFilter for
pushdown_filters=false; that recovered TPC-H but killed TPC-DS Q72
(4.91x faster -> no change) by removing row-level pruning for
CollectLeft's cheap bounds too.
Root cause: on PartitionMode::Partitioned, SharedBuildAccumulator emitted
a per-partition 'CASE hash(col) % N WHEN pid THEN bounds ELSE
lit(false) END' as the dynamic filter. On the probe scan this evaluates
hash + modulo + CASE branch per row -- the profile hotspot
(expressions::case::PartialResultIndex::merge_n at 1.11% self-time on
TPC-H Q9 vs 0% on main).
The routing existed to keep the per-partition bounds exact -- a probe
row X with hash(X) % N == P would only be checked against partition P's
bounds. That's exact but redundant with the downstream hash lookup
(which is also per-partition and exact). The lookup filters exactly
what CASE was filtering, at a lower per-row cost, so the CASE routing
buys nothing on the probe scan.
This commit, when the membership gate is off (the production default
after the split-membership commit), emits the union of per-partition
bounds instead:
col >= min(min_0, ..., min_{N-1}) AND col <= max(max_0, ..., max_{N-1})
Same shape as PartitionMode::CollectLeft. A probe row can pass the
union and still miss its build partition, but the exact hash lookup
downstream drops it -- no wrong results. Empty partitions contribute
nothing to the union; if every partition is empty the filter is
lit(false); a canceled partition falls back to lit(true) (permissive,
we lack the info to safely narrow). Membership-opt-in retains the
historical CASE hash-routed form so InListExpr / HashTableLookupExpr
can be applied to the correct partition's build values.
With the expensive per-row form gone, the pushdown_filters=false
PostScanFilter is cheap again, so opener/mod.rs no longer needs to
filter dynamic conjuncts out of post_scan_conjuncts -- reverted.
Local TPC-H SF1 (release-nonlto, 5 iters, warm):
- baseline (HEAD~2, pre-apache#22384) Q9: 40 ms (avg 46)
- PR before any fix Q9: 80 ms (avg 82) -- 2.09x regression
- PR + this fix Q9: 35 ms (avg 52) -- back at baseline, no CASE hotspot
- Full TPC-H avg: baseline 27.65, this fix 26.99 ms (net -2%)
Snapshot in filter_pushdown.rs regenerated to reflect the union form
(the old snapshot's 'CASE hash_repartition % 12 WHEN 5 ...' is gone;
new form is 'a@0 >= aa AND a@0 <= ab AND b@1 >= ba AND b@1 <= bb').
- filter_pushdown.rs: enable enable_hash_join_dynamic_membership_filter in
test_hashjoin_hash_table_pushdown_{collect_left,partitioned} (they
specifically exercise HashTableLookupExpr; membership default is now false).
- filter_pushdown.rs: refresh test_hashjoin_dynamic_filter_pushdown_collect_left
and the force_hash_collisions branch snapshots (bounds only, no IN (SET)
since membership is off by default).
- clickbench.slt, preserve_file_partitioning.slt, projection_pushdown.slt,
repartition_subset_satisfaction.slt: regenerated for post-apache#22384 plan
display + membership-off default.
- configs.md: prettier reformat (trailing whitespace).
Tactical fix on top of apache#22384 to address the benchmark regressions that paper reported (TPCH SF1: +27% total, Q17 2.09x slower, Q3/5/7/8/9/12/ 13/14/18/20 all 1.24-1.67x slower). Approach was suggested by @adriangb in the apache#23420 discussion: "splitting out the min/max range dynamic filters that HashJoinExec pushes down from the hash table ones and then we could turn off the hash table ones by default". Why: HashJoin's build-side dynamic filter today publishes a combined `bounds AND membership` expression to the probe scan. - Bounds (`col >= min AND col <= max`) is 2 comparisons per row (~2ns) and drives the RG-level statistics pruning that is by far the largest contribution. - Membership (`InListExpr` over the build keys, or a hash-table lookup for large builds) is a per-row hash-set / hash-table probe (~50-100ns). apache#22384's contract change (`try_pushdown_filters` always accepts pushable filters, PostScanFilter picks up whatever the RowFilter cannot place) means the combined expression now runs on every scanned batch even with `pushdown_filters = false`, where previously it was silently propagated through source.predicate but never row-evaluated. On multi-join queries with high match rate, the membership check pays the hash cost twice (once in the scan, once inside HashJoin) with no selectivity win — that's exactly the "not earning their keep" case @adriangb described. What: a new config knob `datafusion.optimizer.enable_hash_join_dynamic_membership_filter` (default `false`) gates the membership creation. When off, `SharedBuildAccumulator` skips `create_membership_predicate` in both the CollectLeft and Partitioned finalize paths and publishes only the bounds portion. RG pruning is unaffected. Highly-selective joins with big build sides that used to see 2-3x wins from membership pruning can restore the historical behavior by flipping the knob to `true`. Tests: two new unit tests in `shared_bounds.rs`: - `collect_left_with_gate_off_publishes_bounds_only` drives an accumulator with the gate off, asserts the published expression contains no `InListExpr` and its top op is `AND` (bounds). - `collect_left_with_gate_on_publishes_bounds_and_membership` the inverse, guards against accidentally regressing the wiring. All 398 pre-existing `joins::hash_join` tests still pass. All 1559 `datafusion-physical-plan` lib tests pass. Full `information_schema.slt` passes with the new option listed. Draft while we run benchmarks to quantify how much of the apache#22384 regression this closes. Companion to apache#22384 (adriangb's foundation), follow-up to apache#23532 (DynamicFilter cache — Layer 1 of the regression fix).
…filter The pushdown=false path in the parquet opener split the whole predicate into 'post_scan_conjuncts' — a per-batch FilterExec-equivalent — which included any dynamic filter conjuncts (HashJoin bounds, TopK threshold, aggregate dynamic filter). For join-heavy TPC-H / TPC-DS this dominates cost: HashJoin's Partitioned- mode dynamic filter is a 'CASE hash(col) % N WHEN pid THEN bounds ELSE lit(false) END' — per-row hash + modulo + CASE branch — and it prunes almost nothing on high-match-rate joins where the downstream hash lookup would eliminate the same rows anyway. Local TPC-H SF1 Q9 profile showed 1.1% self-time in 'expressions::case::PartialResultIndex::merge_n' and 1.3% in 'arrow_select::filter::filter_native' on the PR, both at 0% on main — driving Q9 from 40ms → 80ms (2.09x on CI, 1.79x locally). This commit filters DynamicFilterPhysicalExpr-containing conjuncts out of 'post_scan_conjuncts'. Effects: - RowGroupPruner (added by apache#22450) still sees the full predicate via prepared.predicate, so RG-level dynamic pruning continues to fire on bounds/threshold updates. - pushdown_filters=true path unchanged — dynamic filters still go through the arrow-rs RowFilter. - Downstream operator does the exact equivalent: HashJoin's hash lookup filters rows the bounds would have filtered; TopK's sort heap filters rows the threshold would have filtered. No wrong results. Local TPC-H SF1 (release-nonlto, 3 iters): - baseline (HEAD~2, pre-apache#22384) avg: 27.65 ms - PR + this fix avg: 26.24 ms (net 5% ahead of baseline) - Q9 individually: 34.54 → 34.33 ms (matches baseline, was 80.74 before) Also regenerates push_down_filter_parquet.slt for the membership-off default (from the earlier 'split membership from bounds' commit).
…shold of AND The post-scan filter copied the working batch (all its columns) to the surviving rows when a conjunct kept at most 80% of them. The caller then copies the surviving rows again when it applies the final mask. For a cheap range predicate that keeps about half of the rows, the first copy costs much more than the evaluation that it saves on the next conjunct. TPC-DS Q82, `inventory` scan with `inv_quantity_on_hand BETWEEN 100 AND 500`, SF1, all 11.7M rows (EXPLAIN ANALYZE, 3 runs): | | post-scan filter eval | scan compute | |---|---|---| | threshold 0.8 | 28.4 to 29.2 ms | 104 to 106 ms | | threshold 0.2 | 4.5 to 4.6 ms | 74 to 78 ms | | main, `FilterExec` above the scan | 11.7 to 14.1 ms (`FilterExec`) | 64 to 66 ms | The threshold is now `PRE_SELECTION_THRESHOLD` (0.2), the threshold of `AND` in `BinaryExpr` that a `FilterExec` uses. Thus the post-scan filter makes the same copy decision as the `FilterExec` that it replaces. The old comment gave a `CASE` dynamic filter as the reason for 0.8; the partitioned hash join dynamic filters are no longer `CASE` expressions, and optional filters have gates and a measured order. Without a compaction, the working batch also has the rows that the earlier conjuncts removed. The measurements of a later conjunct (its placement statistics and its gate) now count these rows as passing, as they do after a compaction. Before, a later conjunct got the credit for rows that it did not remove. Wall time, min of 16 runs, pushdown on: TPC-DS Q82 1.04x -> 0.96x of main. TPC-H Q14 1.08x -> 0.94x, Q15 1.09x -> 1.00x, Q12 0.94x -> 0.84x. No other TPC-H, TPC-DS or ClickBench query changed outside the A/A noise. A threshold of 0 (never compact) is slower on TPC-H Q20 (1.14x) and TPC-DS Q42, Q52, Q55 (1.2x), thus the compaction stays. PR: apache#22384 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…path The Parquet scan evaluates rejected and non-pushed-down required conjuncts after the decode (apache#22384). Optional conjuncts never go to this post-scan filter. These tests show this behaviour through `ParquetSource::try_pushdown_filters` and all `optional_filter_mode` values: - `optional_filter_is_not_evaluated_post_scan`: with `pushdown_filters = false`, only the required conjuncts run post-scan. The optional conjuncts are used only for statistics pruning. - `rejected_optional_filter_is_not_evaluated_post_scan`: with `pushdown_filters = true`, an optional conjunct that the row filter cannot evaluate (a whole-struct `IS NOT NULL`) is not used. The same conjunct as a required conjunct runs post-scan. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Add `datafusion.execution.adaptive_filter_placement` (default `false`). With `pushdown_filters = true`, the Parquet scan decides for each conjunct where to evaluate it, at file open and again at each row group boundary: - Required conjunct: `RowFilter` (late materialization) or the post-scan filter from apache#22384. - Optional conjunct in the `adaptive` optional filter mode: `RowFilter`, or `Skip` while its gate is paused. A skipped conjunct is not in the `RowFilter`, thus its columns are not decoded. The scan counts down the pause of the gate with the batches of the skipped row groups. The decision for a required conjunct compares the decode time that a row filter saves with the extra fetch latency of a row filter stage: benefit = skippable fraction * unread output bytes per row * decode ns per byte cost = mean fetch latency / rows of the next row group The skippable fraction counts rows in 64-row windows where no row passes (the decoder only skips long runs of removed rows). The measurements are pooled over all files and partitions of the scan. Before enough rows are measured, a conjunct that reads all output columns starts in the post-scan filter. When the placement changes, the stream rebuilds the decoder with `ParquetPushDecoder::into_builder` (new `RowFilter` and projection mask) and builds a new `DecoderProjection`. A file with adaptive placement always uses the batch coalescer and the stream-level `LIMIT`. The placement does not change for a file with a live row selection, or when the new post-scan conjuncts change the narrowed batch schema. The logic is in the new `filter_placement` module: `model` (pure decision), `stats` (pooled measurements) and `FilePlacement` (per-file state). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Tests and sqllogictest plans that landed on main after apache#22384 was opened still expect the old "scan only uses the predicate for pruning" behaviour. With the scan now applying every accepted filter: - Two opener tests (`test_prune_all_null_column_equality_from_file_statistics`, `test_no_prune_when_missing_column_collapses_mixed_predicate`) now expect only the matching rows. The missing-column test also checks `post_scan_rows_pruned` so it still proves the file was read, not pruned. - `string_in_list_pruning.rs` measured unpruned rows with the scan's `output_rows`. It now uses the post-scan matched + pruned counters, which count the rows that the scan decoded. - Regenerated plans in `dynamic_filter_pushdown_config.slt`, `filter_without_sort_exec.slt`, `push_down_filter_parquet.slt`, `range_partitioning.slt` and `range_sorted_time_bin_agg.slt`: the `FilterExec` above parquet scans is gone and the new `post_scan_rows_*` metrics appear. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…stribution Move the check "a round-robin repartition of this input is useful for its row count" from `EnforceDistribution` into `repartition::round_robin_beneficial_for_rows`. The behavior does not change. The next commit uses the same check in the file scan, so that the scan and the optimizer make the same decision. PR: apache#22384 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…et partitions The Parquet scan accepts all pushable filters, thus `FilterPushdown` removes the `FilterExec`. For a scan of one small file (one partition, too small to split into byte ranges), main puts a round-robin `RepartitionExec` between the scan and the `FilterExec`, and a `CoalescePartitionsExec` above them. Without the `FilterExec`, the optimizer adds neither. The filter then runs in one partition, and the scans of the build sides of the hash joins run one after the other in the task of the probe side, not in parallel tasks. On TPC-DS SF1 this made short queries 5% to 30% slower than main (for example Q37 1.26x). `FileScanConfig::try_pushdown_filters` now makes the same decision as `EnforceDistribution`: if the scan has fewer than `target_partitions` partitions, `repartitioned` cannot give more, and a round-robin repartition is useful for the rows that the scan reads, the filters stay above the scan (`PushedDown::No`). The scan still gets them, through the new `FileSource::try_pushdown_pruning_filters`, and uses them only to prune files, row groups and pages. This is what main does with all filters when `pushdown_filters` is false. The plan is then the plan of main for these scans. - Only the filters of a `FilterExec` stay above the scan. A dynamic filter (of a join, a TopK or an aggregate) has no `FilterExec` above the scan, thus the scan applies it as before. - The default of `try_pushdown_pruning_filters` returns `None`: other file sources get their filters as before. - The Parquet scan applies all conjuncts of its predicate or none of them. A scan with a pruning-only predicate uses later filters only to prune too. - `ParquetScanExecNode` gets `pruning_only_predicate`, thus a decoded scan does not apply its predicate again. - An exact row count of at most one batch keeps the filter in the scan: a round-robin repartition cannot split one batch. Tests: - unit tests for the decision in `file_scan_config` and for the pruning-only predicate of `ParquetSource`; - a proto round trip of the pruning-only predicate; - a sqllogictest plan pin in `parquet_filter_pushdown.slt`; - `parquet_statistics.slt` (no statistics, thus unknown rows): the plan is the plan of main again; - two Parquet integration tests that check the filter in the scan use one target partition. PR: apache#22384 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…path The Parquet scan evaluates rejected and non-pushed-down required conjuncts after the decode (apache#22384). Optional conjuncts never go to this post-scan filter. These tests show this behaviour through `ParquetSource::try_pushdown_filters` and all `optional_filter_mode` values: - `optional_filter_is_not_evaluated_post_scan`: with `pushdown_filters = false`, only the required conjuncts run post-scan. The optional conjuncts are used only for statistics pruning. - `rejected_optional_filter_is_not_evaluated_post_scan`: with `pushdown_filters = true`, an optional conjunct that the row filter cannot evaluate (a whole-struct `IS NOT NULL`) is not used. The same conjunct as a required conjunct runs post-scan. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Add `datafusion.execution.adaptive_filter_placement` (default `false`). With `pushdown_filters = true`, the Parquet scan decides for each conjunct where to evaluate it, at file open and again at each row group boundary: - Required conjunct: `RowFilter` (late materialization) or the post-scan filter from apache#22384. - Optional conjunct in the `adaptive` optional filter mode: `RowFilter`, or `Skip` while its gate is paused. A skipped conjunct is not in the `RowFilter`, thus its columns are not decoded. The scan counts down the pause of the gate with the batches of the skipped row groups. The decision for a required conjunct compares the decode time that a row filter saves with the extra fetch latency of a row filter stage: benefit = skippable fraction * unread output bytes per row * decode ns per byte cost = mean fetch latency / rows of the next row group The skippable fraction counts rows in 64-row windows where no row passes (the decoder only skips long runs of removed rows). The measurements are pooled over all files and partitions of the scan. Before enough rows are measured, a conjunct that reads all output columns starts in the post-scan filter. When the placement changes, the stream rebuilds the decoder with `ParquetPushDecoder::into_builder` (new `RowFilter` and projection mask) and builds a new `DecoderProjection`. A file with adaptive placement always uses the batch coalescer and the stream-level `LIMIT`. The placement does not change for a file with a live row selection, or when the new post-scan conjuncts change the narrowed batch schema. The logic is in the new `filter_placement` module: `model` (pure decision), `stats` (pooled measurements) and `FilePlacement` (per-file state). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Conflict: one expected plan in cte.slt. apache#25780 changes the plan of main (the `FilterExec` above the scan keeps its order); this PR removes the `FilterExec` (the scan accepts the filter). This PR's plan is kept. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
apache#25780 on main added `FileSource::exact_filter`: the part of the filter that every output row satisfies, the only part that the scan derives equivalences from. Its `ParquetSource` version returns the pushable conjuncts when `pushdown_filters` is on, and nothing otherwise. This PR changes both cases: - A pruning-only predicate (a filter that stays in a `FilterExec` above a scan that cannot give the target partitions) is used only to prune, also with `pushdown_filters = true`. `exact_filter` returned it, thus the scan claimed that `a` is constant for `a = 5`, the order-preserving repartition merged on `b` only, and `ORDER BY b LIMIT 1` returned 2 instead of 1. It now returns `None`. - With `pushdown_filters = false` the scan applies the accepted conjuncts in the post-scan filter. They are exact, thus `exact_filter` returns them. This keeps the plans of this PR (for example no `SortExec` for `ORDER BY b` with `b = 2`). Tests: a new case in `push_down_filter_parquet.slt` (a plan pin and two results that were wrong: `2` for `LIMIT 1`, and `5 2 / 5 1` for the order) and `exact_filter` checks in the pruning-only unit test. The plan of the apache#25780 case with `pushdown_filters = false` changes: the scan applies `a = 5` and there is no `FilterExec`. PR: apache#22384 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
apache#22384 now has apache/main (with apache#25780, `FileSource::exact_filter`) and its fix for a pruning-only predicate. No conflict. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
apache#25722 now has the current apache#22384, apache/main (apache#25780, `FileSource::exact_filter`) and the fixes for pruning-only predicates and optional conjuncts in `ParquetSource::exact_filter`. No conflict. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Tests and sqllogictest plans that landed on main after apache#22384 was opened still expect the old "scan only uses the predicate for pruning" behaviour. With the scan now applying every accepted filter: - Two opener tests (`test_prune_all_null_column_equality_from_file_statistics`, `test_no_prune_when_missing_column_collapses_mixed_predicate`) now expect only the matching rows. The missing-column test also checks `post_scan_rows_pruned` so it still proves the file was read, not pruned. - `string_in_list_pruning.rs` measured unpruned rows with the scan's `output_rows`. It now uses the post-scan matched + pruned counters, which count the rows that the scan decoded. - Regenerated plans in `dynamic_filter_pushdown_config.slt`, `filter_without_sort_exec.slt`, `push_down_filter_parquet.slt`, `range_partitioning.slt` and `range_sorted_time_bin_agg.slt`: the `FilterExec` above parquet scans is gone and the new `post_scan_rows_*` metrics appear. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…stribution Move the check "a round-robin repartition of this input is useful for its row count" from `EnforceDistribution` into `repartition::round_robin_beneficial_for_rows`. The behavior does not change. The next commit uses the same check in the file scan, so that the scan and the optimizer make the same decision. PR: apache#22384 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…et partitions The Parquet scan accepts all pushable filters, thus `FilterPushdown` removes the `FilterExec`. For a scan of one small file (one partition, too small to split into byte ranges), main puts a round-robin `RepartitionExec` between the scan and the `FilterExec`, and a `CoalescePartitionsExec` above them. Without the `FilterExec`, the optimizer adds neither. The filter then runs in one partition, and the scans of the build sides of the hash joins run one after the other in the task of the probe side, not in parallel tasks. On TPC-DS SF1 this made short queries 5% to 30% slower than main (for example Q37 1.26x). `FileScanConfig::try_pushdown_filters` now makes the same decision as `EnforceDistribution`: if the scan has fewer than `target_partitions` partitions, `repartitioned` cannot give more, and a round-robin repartition is useful for the rows that the scan reads, the filters stay above the scan (`PushedDown::No`). The scan still gets them, through the new `FileSource::try_pushdown_pruning_filters`, and uses them only to prune files, row groups and pages. This is what main does with all filters when `pushdown_filters` is false. The plan is then the plan of main for these scans. - Only the filters of a `FilterExec` stay above the scan. A dynamic filter (of a join, a TopK or an aggregate) has no `FilterExec` above the scan, thus the scan applies it as before. - The default of `try_pushdown_pruning_filters` returns `None`: other file sources get their filters as before. - The Parquet scan applies all conjuncts of its predicate or none of them. A scan with a pruning-only predicate uses later filters only to prune too. - `ParquetScanExecNode` gets `pruning_only_predicate`, thus a decoded scan does not apply its predicate again. - An exact row count of at most one batch keeps the filter in the scan: a round-robin repartition cannot split one batch. Tests: - unit tests for the decision in `file_scan_config` and for the pruning-only predicate of `ParquetSource`; - a proto round trip of the pruning-only predicate; - a sqllogictest plan pin in `parquet_filter_pushdown.slt`; - `parquet_statistics.slt` (no statistics, thus unknown rows): the plan is the plan of main again; - two Parquet integration tests that check the filter in the scan use one target partition. PR: apache#22384 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…path The Parquet scan evaluates rejected and non-pushed-down required conjuncts after the decode (apache#22384). Optional conjuncts never go to this post-scan filter. These tests show this behaviour through `ParquetSource::try_pushdown_filters` and all `optional_filter_mode` values: - `optional_filter_is_not_evaluated_post_scan`: with `pushdown_filters = false`, only the required conjuncts run post-scan. The optional conjuncts are used only for statistics pruning. - `rejected_optional_filter_is_not_evaluated_post_scan`: with `pushdown_filters = true`, an optional conjunct that the row filter cannot evaluate (a whole-struct `IS NOT NULL`) is not used. The same conjunct as a required conjunct runs post-scan. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Add `datafusion.execution.adaptive_filter_placement` (default `false`). With `pushdown_filters = true`, the Parquet scan decides for each conjunct where to evaluate it, at file open and again at each row group boundary: - Required conjunct: `RowFilter` (late materialization) or the post-scan filter from apache#22384. - Optional conjunct in the `adaptive` optional filter mode: `RowFilter`, or `Skip` while its gate is paused. A skipped conjunct is not in the `RowFilter`, thus its columns are not decoded. The scan counts down the pause of the gate with the batches of the skipped row groups. The decision for a required conjunct compares the decode time that a row filter saves with the extra fetch latency of a row filter stage: benefit = skippable fraction * unread output bytes per row * decode ns per byte cost = mean fetch latency / rows of the next row group The skippable fraction counts rows in 64-row windows where no row passes (the decoder only skips long runs of removed rows). The measurements are pooled over all files and partitions of the scan. Before enough rows are measured, a conjunct that reads all output columns starts in the post-scan filter. When the placement changes, the stream rebuilds the decoder with `ParquetPushDecoder::into_builder` (new `RowFilter` and projection mask) and builds a new `DecoderProjection`. A file with adaptive placement always uses the batch coalescer and the stream-level `LIMIT`. The placement does not change for a file with a live row selection, or when the new post-scan conjuncts change the narrowed batch schema. The logic is in the new `filter_placement` module: `model` (pure decision), `stats` (pooled measurements) and `FilePlacement` (per-file state). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
apache#25780 on main added `FileSource::exact_filter`: the part of the filter that every output row satisfies, the only part that the scan derives equivalences from. Its `ParquetSource` version returns the pushable conjuncts when `pushdown_filters` is on, and nothing otherwise. This PR changes both cases: - A pruning-only predicate (a filter that stays in a `FilterExec` above a scan that cannot give the target partitions) is used only to prune, also with `pushdown_filters = true`. `exact_filter` returned it, thus the scan claimed that `a` is constant for `a = 5`, the order-preserving repartition merged on `b` only, and `ORDER BY b LIMIT 1` returned 2 instead of 1. It now returns `None`. - With `pushdown_filters = false` the scan applies the accepted conjuncts in the post-scan filter. They are exact, thus `exact_filter` returns them. This keeps the plans of this PR (for example no `SortExec` for `ORDER BY b` with `b = 2`). Tests: a new case in `push_down_filter_parquet.slt` (a plan pin and two results that were wrong: `2` for `LIMIT 1`, and `5 2 / 5 1` for the order) and `exact_filter` checks in the pruning-only unit test. The plan of the apache#25780 case with `pushdown_filters = false` changes: the scan applies `a = 5` and there is no `FilterExec`. PR: apache#22384 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The new `parquet_statistics.slt` case of apache#25795 on main pins a `FilterExec` above the scan. With apache#22384 the scan accepts the filter (one file of two rows keeps the filter in the scan), thus the plan is the scan alone. Its statistics are still `Rows=Inexact(2)`, not empty, and the query result does not change. Integration of apache#22384 and main in the final state branch. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Conflict in parquet_statistics.slt: apache#25863 on main changes the statistics of the `FilterExec` in six expected plans (distinct count `Inexact(1)`). With this PR the scan applies the filter and there is no `FilterExec` in these plans, thus this PR's plans are kept. The new `mul_wrap` case of apache#25795 on main pins a `FilterExec` above the scan. With this PR the scan applies the filter (one file of two rows keeps it in the scan), thus the plan is the scan alone. Its statistics are `Rows=Inexact(2)`, not empty, and the query result is unchanged. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Tests and sqllogictest plans that landed on main after apache#22384 was opened still expect the old "scan only uses the predicate for pruning" behaviour. With the scan now applying every accepted filter: - Two opener tests (`test_prune_all_null_column_equality_from_file_statistics`, `test_no_prune_when_missing_column_collapses_mixed_predicate`) now expect only the matching rows. The missing-column test also checks `post_scan_rows_pruned` so it still proves the file was read, not pruned. - `string_in_list_pruning.rs` measured unpruned rows with the scan's `output_rows`. It now uses the post-scan matched + pruned counters, which count the rows that the scan decoded. - Regenerated plans in `dynamic_filter_pushdown_config.slt`, `filter_without_sort_exec.slt`, `push_down_filter_parquet.slt`, `range_partitioning.slt` and `range_sorted_time_bin_agg.slt`: the `FilterExec` above parquet scans is gone and the new `post_scan_rows_*` metrics appear. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…stribution Move the check "a round-robin repartition of this input is useful for its row count" from `EnforceDistribution` into `repartition::round_robin_beneficial_for_rows`. The behavior does not change. The next commit uses the same check in the file scan, so that the scan and the optimizer make the same decision. PR: apache#22384 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…et partitions The Parquet scan accepts all pushable filters, thus `FilterPushdown` removes the `FilterExec`. For a scan of one small file (one partition, too small to split into byte ranges), main puts a round-robin `RepartitionExec` between the scan and the `FilterExec`, and a `CoalescePartitionsExec` above them. Without the `FilterExec`, the optimizer adds neither. The filter then runs in one partition, and the scans of the build sides of the hash joins run one after the other in the task of the probe side, not in parallel tasks. On TPC-DS SF1 this made short queries 5% to 30% slower than main (for example Q37 1.26x). `FileScanConfig::try_pushdown_filters` now makes the same decision as `EnforceDistribution`: if the scan has fewer than `target_partitions` partitions, `repartitioned` cannot give more, and a round-robin repartition is useful for the rows that the scan reads, the filters stay above the scan (`PushedDown::No`). The scan still gets them, through the new `FileSource::try_pushdown_pruning_filters`, and uses them only to prune files, row groups and pages. This is what main does with all filters when `pushdown_filters` is false. The plan is then the plan of main for these scans. - Only the filters of a `FilterExec` stay above the scan. A dynamic filter (of a join, a TopK or an aggregate) has no `FilterExec` above the scan, thus the scan applies it as before. - The default of `try_pushdown_pruning_filters` returns `None`: other file sources get their filters as before. - The Parquet scan applies all conjuncts of its predicate or none of them. A scan with a pruning-only predicate uses later filters only to prune too. - `ParquetScanExecNode` gets `pruning_only_predicate`, thus a decoded scan does not apply its predicate again. - An exact row count of at most one batch keeps the filter in the scan: a round-robin repartition cannot split one batch. Tests: - unit tests for the decision in `file_scan_config` and for the pruning-only predicate of `ParquetSource`; - a proto round trip of the pruning-only predicate; - a sqllogictest plan pin in `parquet_filter_pushdown.slt`; - `parquet_statistics.slt` (no statistics, thus unknown rows): the plan is the plan of main again; - two Parquet integration tests that check the filter in the scan use one target partition. PR: apache#22384 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…path The Parquet scan evaluates rejected and non-pushed-down required conjuncts after the decode (apache#22384). Optional conjuncts never go to this post-scan filter. These tests show this behaviour through `ParquetSource::try_pushdown_filters` and all `optional_filter_mode` values: - `optional_filter_is_not_evaluated_post_scan`: with `pushdown_filters = false`, only the required conjuncts run post-scan. The optional conjuncts are used only for statistics pruning. - `rejected_optional_filter_is_not_evaluated_post_scan`: with `pushdown_filters = true`, an optional conjunct that the row filter cannot evaluate (a whole-struct `IS NOT NULL`) is not used. The same conjunct as a required conjunct runs post-scan. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Add `datafusion.execution.adaptive_filter_placement` (default `false`). With `pushdown_filters = true`, the Parquet scan decides for each conjunct where to evaluate it, at file open and again at each row group boundary: - Required conjunct: `RowFilter` (late materialization) or the post-scan filter from apache#22384. - Optional conjunct in the `adaptive` optional filter mode: `RowFilter`, or `Skip` while its gate is paused. A skipped conjunct is not in the `RowFilter`, thus its columns are not decoded. The scan counts down the pause of the gate with the batches of the skipped row groups. The decision for a required conjunct compares the decode time that a row filter saves with the extra fetch latency of a row filter stage: benefit = skippable fraction * unread output bytes per row * decode ns per byte cost = mean fetch latency / rows of the next row group The skippable fraction counts rows in 64-row windows where no row passes (the decoder only skips long runs of removed rows). The measurements are pooled over all files and partitions of the scan. Before enough rows are measured, a conjunct that reads all output columns starts in the post-scan filter. When the placement changes, the stream rebuilds the decoder with `ParquetPushDecoder::into_builder` (new `RowFilter` and projection mask) and builds a new `DecoderProjection`. A file with adaptive placement always uses the batch coalescer and the stream-level `LIMIT`. The placement does not change for a file with a live row selection, or when the new post-scan conjuncts change the narrowed batch schema. The logic is in the new `filter_placement` module: `model` (pure decision), `stats` (pooled measurements) and `FilePlacement` (per-file state). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
apache#25780 on main added `FileSource::exact_filter`: the part of the filter that every output row satisfies, the only part that the scan derives equivalences from. Its `ParquetSource` version returns the pushable conjuncts when `pushdown_filters` is on, and nothing otherwise. This PR changes both cases: - A pruning-only predicate (a filter that stays in a `FilterExec` above a scan that cannot give the target partitions) is used only to prune, also with `pushdown_filters = true`. `exact_filter` returned it, thus the scan claimed that `a` is constant for `a = 5`, the order-preserving repartition merged on `b` only, and `ORDER BY b LIMIT 1` returned 2 instead of 1. It now returns `None`. - With `pushdown_filters = false` the scan applies the accepted conjuncts in the post-scan filter. They are exact, thus `exact_filter` returns them. This keeps the plans of this PR (for example no `SortExec` for `ORDER BY b` with `b = 2`). Tests: a new case in `push_down_filter_parquet.slt` (a plan pin and two results that were wrong: `2` for `LIMIT 1`, and `5 2 / 5 1` for the order) and `exact_filter` checks in the pruning-only unit test. The plan of the apache#25780 case with `pushdown_filters = false` changes: the scan applies `a = 5` and there is no `FilterExec`. PR: apache#22384 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The new `parquet_statistics.slt` case of apache#25795 on main pins a `FilterExec` above the scan. With apache#22384 the scan accepts the filter (one file of two rows keeps the filter in the scan), thus the plan is the scan alone. Its statistics are still `Rows=Inexact(2)`, not empty, and the query result does not change. Integration of apache#22384 and main in the final state branch. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Which issue does this PR close?
Rationale for this change
The Parquet scan today gives the predicate to the source for row-group / page / bloom pruning, but only applies it row-level via
RowFilterwhenpushdown_filters=true. With pushdown off, aFilterExecis left above the scan to do row-level filtering. This is the substrate the adaptive-filter work (#22237 / #22144) builds on, but it also has a real correctness bug onmainthat's worth fixing on its own.build_row_filter(row_filter.rs:994-1083, see its own doc comment at1009-1014) silently drops conjuncts thatFilterCandidateBuilder::buildreturnsOk(None)for, andRowFilterGenerator::buildswallows whole-build errors. By the timebuild_row_filterruns,ParquetSource::try_pushdown_filtershas already accepted the filter and the parentFilterExechas been removed — so those dropped conjuncts are never applied anywhere and the query returns wrong results. The most reproducible trigger is the per-file expr adapter rewriting a predicate that was pushable at table schema time into somethingPushdownCheckerrejects at physical file schema time (schema evolution / coercion, whole-struct refs introduced by the rewrite, etc.).This PR makes the Parquet scan always own its pushable filters and guarantees every accepted conjunct is applied — either by the parquet
RowFilteror by a new in-scan post-scan filter evaluated on decoded batches (the in-scan equivalent of aFilterExec). Nothing is silently dropped.What changes are included in this PR?
row_filter.rs— never drop conjuncts.build_row_filternow returnsResult<(Option<RowFilter>, Vec<Arc<dyn PhysicalExpr>>)>— the second element is the conjuncts it could not place.RowFilterGeneratorexposes them viarejected_conjuncts(); on whole-file build errors it routes every conjunct through that list (no silent error swallowing).post_scan_filter.rs(new module) — encapsulates the projection widening + rebasing + filter evaluation behind a small API:PostScanFilter— evaluates a predicate on decoded batches; SQLWHEREsemantics (NULLdrops the row); records rows-pruned / matched / time.DecoderProjection::build(projection, post_scan_conjuncts, schemas, …)— widens the decoder projection over (user projection ∪ post-scan conjunct columns), rebases the projection and conjuncts onto the decoder's stream schema, and returns theProjectionMask,Projector,replace_schemaflag, and the rebasedPostScanFilter. Empty conjuncts list = the prior projection-only behaviour, so the opener routes every file through this one call.ParquetSource::try_pushdown_filters— always returns the per-filterYes/Nodiscriminant based oncan_expr_be_pushed_down_with_schemas, regardless of thepushdown_filtersconfig. The flag still records whether theRowFilter(vs. post-scan) path is used downstream.opener/mod.rs::build_stream— orchestrates: builds theRowFilterGeneratoronly whenpushdown_filters=true; computespost_scan_conjuncts(rejected conjuncts when pushdown is on, full split-conjunction of the predicate when off); callsDecoderProjection::build; routes theLIMITtoremaining_limitinstead of a decoder limit whenever the post-scan filter is present (decoder-local limit + post-scan filter is unsafe — the decoder would stop before the post-scan rejected enough rows). The prior inlinebuild_projection_read_plan/reassign_expr_columns/make_projectorblock is replaced by the singleDecoderProjection::buildcall — net simplification.push_decoder.rs—PushDecoderStreamStatecarries anOption<PostScanFilter>; in theDecodeResult::Dataarm it applies the filter, skips empty batches, then enforcesremaining_limitand projects.DecoderBuilderConfigis fedprojection_mask: &ProjectionMaskdirectly (no longer the fullParquetReadPlan).metrics.rs— newpost_scan_rows_pruned/post_scan_rows_matchedcounters andpost_scan_filter_eval_timeTime, mirroring the existingpushdown_rows_*/row_pushdown_eval_timesoEXPLAIN ANALYZEkeeps surfacing filter cost once theFilterExecis gone.Filters stay above a scan that cannot give the target partitions
A filter that the scan applies runs in the partitions of the scan. On main, a scan of one small file (one partition, too small to split into byte ranges) has a round-robin
RepartitionExecand aFilterExecabove it. TheFilterExecruns intarget_partitionspartitions, and theCoalescePartitionsExecabove it starts the scan in its own task. If the scan applies the filter, the optimizer adds none of these operators. The filter then runs in one partition, and the build sides of the hash joins are scanned one after the other, not in parallel.FileScanConfig::try_pushdown_filtersnow makes the same decision asEnforceDistribution(the check is shared, seerepartition::round_robin_beneficial_for_rows):target_partitions, cannot split, can have more rows than one batchFilterExec→RepartitionExec: RoundRobinBatch→ scan (pruning only)target_partitionspartitions, or can splitFilterExec→ scan (pruning only)FilterExec→RepartitionExec: RoundRobinBatch→ scanFileSource::try_pushdown_pruning_filters, and uses them only to prune files, row groups and pages. This is what main does with all filters whenpushdown_filtersis false. The default implementation returnsNone, thus other file sources get their filters as before.FilterExecstay above the scan. A dynamic filter (of a join, a TopK or an aggregate) has noFilterExecabove the scan, thus the scan applies it.ParquetScanExecNodehas a new fieldpruning_only_predicate, thus a decoded scan does not apply its predicate again.FileSource::exact_filter, fix: derive scan equivalences only from filters the source applies exactly #25780 on main): the scan reports its accepted conjuncts as exact withpushdown_filters = falsetoo (it applies them in the post-scan filter), and nothing for a pruning-only predicate. Without the second rule,a = 5of a pruning-only predicate made the repartition below theFilterExecmerge onbonly, andORDER BY b LIMIT 1returned a wrong row (new case inpush_down_filter_parquet.slt, and a unit test).Effect on TPC-DS SF1 with the default configuration (
pushdown_filters = false), minimum of 12 runs, ratio to main (lower is better):TPC-H SF1 does not change outside the noise.
The adaptive-filter machinery from #22237 / #22144 (
SelectivityTracker,FilterIdtagging, per-conjunct pruning stats,StrategySwapmid-stream swaps,OptionalFilterPhysicalExpr, customarrow-rsbranch, the threefilter_pushdown_*config knobs) is intentionally not included — this PR is the standalone substrate they would build on.Are these changes tested?
Yes.
build_row_filter_surfaces_rejected_struct_conjunct(row_filter.rs) asserts the new API contract directly —build_row_filterno longer drops the rejected conjunct.rejected_struct_conjunct_runs_post_scan_not_dropped(opener/mod.rs) is an end-to-end test: withpushdown_filters=trueand as IS NOT NULLpredicate over a struct column where row 1 is NULL,mainreturns 3 rows (conjunct silently dropped, predicate relaxed) and this PR returns the correct 2.file_scan_configand for the pruning-only predicate ofParquetSource, a protobuf round trip ofpruning_only_predicate, and a plan test inparquet_filter_pushdown.slt.parquet_statistics.slt(a table without statistics) has the plan of main again. Two Parquet integration tests that check the filter in the scan use one target partition..sltfiles are regenerated (clickbench, push_down_filter_parquet, projection_pushdown, parquet*, etc.) — theFilterExecabove parquet scans is gone from those plans. Spurious whitespace-only churn from--completewas reverted.Are there any user-facing changes?
pushdown_filters=true. Before this PR they were silently relaxed.FilterExecno longer appears above aDataSourceExecfor pushable filters on a parquet source. The predicate appears aspredicate=…on theDataSourceExec. Query results are unchanged. Exception: for a scan that cannot givetarget_partitionspartitions, theFilterExecand the round-robinRepartitionExecstay above the scan, as on main (see the table above).ParquetFileMetrics—post_scan_rows_pruned,post_scan_rows_matched,post_scan_filter_eval_time— appear inEXPLAIN ANALYZEoutput for parquet scans.build_row_filterreturn type changes fromResult<Option<RowFilter>>toResult<(Option<RowFilter>, Vec<Arc<dyn PhysicalExpr>>)>; callers must apply the rejected conjuncts (or they'll have the same drop-on-floor bug this PR fixes). New provided methodFileSource::try_pushdown_pruning_filters(defaultNone). New fieldpruning_only_predicatein theParquetScanExecNodeprotobuf message.Draft. Happy to split into a stack (refactor → row_filter fix → opener orchestration + tests) if reviewers prefer.
🤖 Generated with Claude Code