Repository navigation
Conversation
…s nothing Page-index pruning installed a row selection for each row group that it examined, also when the selection skipped no rows. The opener treats any row selection as live, so a select-all selection disabled runtime (dynamic filter) row-group pruning and statistics-based row-group reordering for the whole file. Keep the row group as a full scan when the page index skips no rows. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…IN list In partitioned mode the hash join pushes a routed filter: CASE hash_repartition % N WHEN i THEN bounds_i AND key IN (list_i) ... END When every non-empty build partition pushes an InList and no partition is canceled, the routing is redundant. Routing is a deterministic function of the join keys, so a build row with key K is in the partition that a probe row with key K routes to. A test of K against the union of all lists therefore accepts the same rows as the routed CASE. The per-partition bounds reject no additional rows either, because every key in a list is inside the bounds of its partition. The filter is now `key IN (union)`. This removes the per-row routing hash from the probe side, and the pruning code can use an InList (up to `max_in_list_size`), which it cannot do with a CASE. The union is capped at 1 MiB. Each partition's list is limited independently by `hash_join_inlist_pushdown_max_size`, so the union grows with the partition count; past the cap the routed CASE, where a probe row checks only one list, stays in use. Partitions that push a hash table, and builds with a canceled partition, also keep the CASE. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The per-partition InList arrays hold one entry per build row, and the collapsed union concatenated them. On TPC-DS SF1 with 12 partitions, Q65 pushed 54,000 entries for 6 distinct keys, and Q18 pushed 10,848 entries for 1,178 distinct keys. The pruning code uses an IN list only up to `max_in_list_size` (20) entries, so these filters gave no pruning term, and the long lists made the statistics evaluation for each file range slow. The union is now deduplicated and sorted with the arrow row format before the IN list is built. This works for single and struct (multi-column) keys and dictionaries, keeps one NULL, and makes the list order deterministic. The 1 MiB cap now applies to the deduplicated union. The collapsed filter also gets one range per key column, `col >= min AND col <= max`, from the combined bounds of all non-empty partitions, before the IN list. Every key in the union is inside this range, so the filter stays exact, and the pruning code can use the range when the list has more than `max_in_list_size` entries. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
… as separate filters
A partitioned hash join pushed one dynamic filter to the probe side:
DynamicFilter [ CASE hash_repartition % N
WHEN i THEN bounds_i AND membership_i ... END ]
The pruning code cannot use a `CASE`, thus the build-side bounds did not
prune files, row groups or pages, even when the probe side is clustered by
the join key.
A partitioned join now pushes two dynamic filters:
DynamicFilter [ bounds ] AND DynamicFilter [ membership ]
- Bounds: the union of the bounds of all partitions (new `bounds_union`
module, adapted from the prototype in apache#24235). It does not need routing,
so the pruning code can use it. With hash partitioning each column gets
one range. With range partitioning the disjoint ranges of the partitions
are kept (up to 8 for each column, OR'd), so the filter also rejects keys
in the gaps between them.
- Membership: the routed `CASE` without the per-partition bounds (they
reject no row that the membership check of the partition accepts), or
the collapsed IN list when every partition pushes an IN list.
The bounds stay in the `CASE` (as before) when the union cannot describe
the build side: a canceled partition, or no usable bounds. An empty build
sets both filters to `false`. The NULL escape of null-equal and null-aware
joins wraps each filter.
A collect-left join does not change: it pushes one filter that holds
`bounds AND membership`. It has no routing `CASE`, so the pruning code can
already use its bounds.
Each filter has its own expression id. The join keeps each filter that
reached a consumer, and both get `update()` and `mark_complete()`. The
proto gets a new `dynamic_filter_bounds` field; a plan without it restores
one filter that holds both the bounds and the membership check, as before.
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Add the `hj_ordered_subset` SQL benchmark suite. The probe table (`events`, 100 days x 200k rows, 100k-row row groups) is sorted by the join key, and the build side matches 1 day, 10 days, or a scattered 1% of the key range (control). The bounds of the hash join dynamic filter can prune the probe row groups outside the matched range, and the membership check passes almost every remaining row. Subgroups `partitioned` (Q01-Q03) and `collect_left` (Q04-Q06) force the HashJoinExec mode, which `expect_plan` checks. An assert checks that every build row matches exactly one event. The load SQL writes the data inline; HJOS_DAYS, HJOS_ROWS_PER_DAY and HJOS_RG_SIZE size it. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…d conjuncts post-scan
Rebased onto main. The original commit 1 of this PR ("extract
DecoderProjection from build_stream") landed independently on main as
current `decoder_projection` / single-decoder + `rg_plan` model.
Two changes, both applied inside the parquet scan so the parent
`FilterExec` can be removed unconditionally for pushable filters:
1. Never drop conjuncts the `RowFilter` cannot place.
`build_row_filter` previously `.flatten()`-ed away conjuncts that
`FilterCandidateBuilder::build` rejected (whole-struct references,
per-file physical-schema mismatches) and swallowed whole-build
errors. By the time it runs, `try_pushdown_filters` has already
removed the `FilterExec`, so those conjuncts were applied nowhere —
wrong results. `build_row_filter` now returns
`(Option<RowFilter>, Vec<rejected>)`, `RowFilterGenerator` exposes
`rejected_conjuncts()`, and a whole-file build error routes every
conjunct to the rejected list rather than relaxing the predicate.
2. Always accept pushable filters and run the remainder post-scan.
`try_pushdown_filters` reports each pushable filter as accepted so the
`FilterExec` is always removed; the scan owns the predicate. The
opener routes conjuncts to two places, applying every one:
- pushdown_filters=true -> row-filterable conjuncts via the parquet
`RowFilter`; rejected conjuncts via the in-scan post-scan filter.
- pushdown_filters=false -> the whole predicate runs as a post-scan
filter on decoded batches (behaviorally identical to `FilterExec`).
Implementation:
- `DecoderProjection` (main's `decoder_projection` module) grows a
`post_scan_conjuncts` parameter: it widens the decoder mask over
(user projection ∪ post-scan filter columns), rebases the conjuncts
onto the stream schema, and returns a `PostScanFilter` applied to
every decoded batch with SQL `WHERE` semantics. Virtual-column
conjuncts are stripped from the read-plan mask (they aren't file
columns) but kept in the post-scan predicate, which sees the
reader-appended virtual columns.
- `PushDecoderStreamState` applies the post-scan filter in the decoded-
batch arm, skips empty batches, and re-introduces a stream-level
`remaining_limit` (main enforces LIMIT decoder-locally, which is
unsafe once a post-scan filter can reject rows). The opener routes the
limit to `remaining_limit` iff a post-scan filter is present.
- New `post_scan_rows_pruned` / `post_scan_rows_matched` counters and
`post_scan_filter_eval_time` on `ParquetFileMetrics`.
Tests:
- `build_row_filter_surfaces_rejected_struct_conjunct` (row_filter.rs)
asserts the rejected struct conjunct is returned, not dropped.
- `rejected_struct_conjunct_runs_post_scan_not_dropped` (opener) is
end-to-end: `s IS NOT NULL` over a struct column with pushdown on
returns 2 (was 3 before the fix).
- Parquet `.slt` files regenerated: `FilterExec` above `DataSourceExec`
gone, predicate on the scan, `post_scan_rows_*` metrics on EXPLAIN
ANALYZE. Opener / core insta / page_pruning assertions updated for the
now-applied predicate.
Squashed follow-up commits (see their original messages on the
adriangb/parquet-post-scan-filter branch):
- perf(parquet-datasource): narrow to the projector's columns before
filtering
- perf(parquet-datasource): coalesce post-scan filter output back to
batch_size
- perf(parquet-datasource): compact between conjuncts in the post-scan
filter
Co-Authored-By: Claude Opus 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>
When the TopK dynamic filter prunes every remaining row group at a row group boundary, `rebuild_decoder_at_boundary` returns `Ok(true)` and the stream calls `finish()`. With a batch coalescer (a post-scan filter is present, for example with the default `pushdown_filters = false`), `finish()` flushes the coalescer and returns a batch. The next poll then went back to the decoder, which still pointed at a row group that the plan had dropped, and `sync_rg_plan_to_decoder_frontier` failed with "push decoder frontier RG N is not in rg_plan; decoder and plan have diverged". ClickBench Q23, Q24 and Q26 fail with this error. After the flush, the stream now only drains the coalescer. The new sqllogictest in `dynamic_row_group_pruning.slt` fails without the fix. 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>
…orrectness Add `OptionalFilterPhysicalExpr`, a transparent wrapper that marks a filter as optional: a consumer can skip it without changing the query result. A consumer can skip it only when the wrapper is a direct conjunct of the root AND chain of its predicate. In all other positions the wrapper is transparent, because `evaluate()` always evaluates the inner expression. `snapshot()` returns the inner expression, so pruning sees through it. Also add: - `split_optional` and `is_optional_filter` helpers in `physical_expr::utils` for consumers - `PhysicalOptionalFilterNode` proto message (field 29 in `PhysicalExprNode`) with self-encoding `try_to_proto`/`try_from_proto` No producer uses the wrapper yet, so there is no behavior change. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
With `pushdown_filters = false`, the scan evaluates all accepted filters for each row after the decode. This includes the optional filters (hash join, TopK and aggregate dynamic filters, which the producers wrap in `Optional(...)`). The join dynamic filter evaluated for each row in the scan caused a 1.16x TPC-H regression. Optional filters are not needed for correctness. Thus the post-scan filter now never gets an optional conjunct (a root `AND` conjunct found with `split_optional`): - `pushdown_filters = false`: only the required conjuncts run post-scan. Optional conjuncts are used only for statistics, page index, bloom filter and file pruning. - `pushdown_filters = true`: an optional conjunct that the row filter rejects for a file is not used for that file. A required conjunct that is rejected still runs post-scan. - A whole-file row filter build error sends only the required conjuncts to the post-scan filter. Accepted optional conjuncts with `pushdown_filters = true` stay row filter predicates, as before. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Add the `datafusion_physical_expr::filter_stats` module with the shared primitives that adaptive filter code uses to measure filters at runtime: - `Clock`: a monotonic clock in nanoseconds that tests can replace. `SystemClock` is the real clock. `ManualClock` moves only when a test moves it, thus decisions that use time are deterministic in tests. - `FilterCost`: the rows in, the rows out and the evaluation time of one filter, and the derived cost for each row and rows removed for each nanosecond. - `duration_nanos`: a `Duration` in nanoseconds, saturated to `u64::MAX`. No code uses the module yet, thus behavior does not change. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Add adaptive conjunct reordering behind `datafusion.execution.adaptive_filter_reordering` (default `false`): each stream measures its conjuncts over a warm-up, ranks them by rows dropped per nanosecond, and adopts a new order (built as a plain `BinaryExpr` AND chain) only if the estimated cost, using `BinaryExpr`'s pre-selection rule, is at least 5% lower. Predicates with volatile expressions are never reordered. The decision is made one time for each stream, and streams do not share state. The measurements use `Clock` and `FilterCost` of `datafusion_physical_expr::filter_stats`, so tests use a `ManualClock`. `PRE_SELECTION_THRESHOLD` is exported (doc-hidden) so that the cost model uses the same rule as `BinaryExpr`. New metric: `adaptive_reorders`. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
… than they save Add a runtime gate that pauses optional filters (filters that are not needed for correctness, such as hash join and TopK dynamic filters) when they cost more than they save. The gate is a per-stream state machine (Evaluate / Paused with exponential backoff) that restarts evaluation when the filter changes. The gate finds the dynamic filters one time with `DynamicFilterTracking::classify` and then polls their subscriptions, so a check does not walk the filter tree. Gates do not share state. At the end of each window of evaluated batches the gate pauses the filter if the window removed no rows, or if its evaluation time is larger than the work that the removed rows save: `(rows_in - rows_out) * saving_ns_per_row`. The saving for each row is the configured minimum plus an optional value that the consumer measures and updates (`MeasuredRowSaving`). The cost rule has a margin (pause above 1.1x the saving, resume below 0.9x) so that a filter does not switch on and off when cost and saving are almost equal. Add the `datafusion.execution.optional_filter_min_saving_ns_per_row` option (default 20). No operator uses the gate yet, so behavior does not change. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
A consumer can now give the gate a fixed cost for each evaluated row in addition to the evaluation time, with `MeasuredRowSaving::set_overhead_ns_per_row`. The gate adds it to the cost of each window. The Parquet scan uses it for the fixed cost of a row filter stage, which is larger than the evaluation time of a cheap predicate. PR: apache#25674 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Each gate paid for its own first window and its own probes. A scan opens its files at the same time, thus a filter that removes nothing cost one window in each file (TPC-H Q9: five join filters that remove no rows and a CASE routing filter that costs 83 ns for each row, in each of 12 files). `SharedGateVerdict` holds the last pause (or end of a pause) of the gates of one plan site in one atomic word. A gate without evidence of its own (before its first decision, after a filter change, after a pause) uses a pause that another gate published after the last verdict that it saw: a new gate starts paused, and a gate in its first window or in a probe window stops and pauses. Thus usually only one gate probes after a pause. A gate that keeps the filter does not use the pauses of other gates (skewed data). A filter change clears the shared pause, because it was measured on the old filter. PR: apache#25674 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
A gate decided after `sample_batches` batches, whatever their size. After a selective row filter a batch can have 2 to 7 rows, and the fixed cost of each call then looks like 600 to 8000 ns for each row: ClickBench Q23 paused the TopK filter on `EventTime` (0.4 ns for each row on full batches) on such windows, and the shared verdict spread these pauses to the other files. All decisions (pause, keep, probe) now need a window of at least `sample_batches` batches and `MIN_OBSERVED_ROWS` rows. The constant moves to `filter_stats`, so that the gate and the Parquet filter placement use the same sample size. All published shared pauses come from such windows. PR: apache#25674 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The gate assumed that each removed row saves `min_saving_ns_per_row` (20 ns) after the filter. For hash join dynamic filters this is the probe work of the join, and it is much smaller in star joins with small dimension tables: 3.5 to 8 ns for each probe row on the TPC-DS SF1 `date_dim` joins (Q65, Q67), 2 ns on Q90, 17 ns on TPC-H Q9. Filters that cost 3 to 7 ns for each row and remove 80% of the rows thus stayed on, and cost more than the join work they saved (8-18% slower than `pruning_only` on the bot). `RemovedRowWork` (in `filter_stats`) is the work that the producer of a filter does for each row that the filter removes, as the producer measures it. Each `DynamicFilterPhysicalExpr` has one, shared by all its derived filters. The gate uses the smallest measured work of the dynamic filters in its filter as the saving of a removed row, and `min_saving_ns_per_row` only until the producer has measured `MIN_OBSERVED_ROWS` rows (a prior). PR: apache#25674 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Add `OptionalFilterMode` and the `datafusion.execution.optional_filter_mode` option (default `always`): - `always`: evaluate optional filters like any other pushed-down filter (today's behavior). - `adaptive`: evaluate each optional filter behind an `OptionalFilterGate`, which pauses it while it costs more than it saves or removes no rows. - `pruning_only`: use optional filters only for statistics pruning. The Parquet scan and `FilterExec` read this option. This commit alone does not change behavior. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…heir cost Make the Parquet scan the first consumer of optional filters (conjuncts wrapped in `OptionalFilterPhysicalExpr`) when row level filter pushdown (`datafusion.execution.parquet.pushdown_filters`) is enabled. The `datafusion.execution.optional_filter_mode` option controls the row filter: - `always` (default): optional conjuncts are normal RowFilter predicates. Behavior does not change. - `adaptive`: each optional conjunct is a separate RowFilter predicate, after all required predicates, with an `OptionalFilterGate`. When the gate skips a batch, the predicate lets all rows pass without evaluation. Each file has its own gates; gates do not share state. - `pruning_only`: optional conjuncts are not in the RowFilter. In `adaptive` mode, the gate times each evaluation of the predicate and pauses a filter that removes no rows or that costs more than it saves. The saving of a removed row is the configured minimum (`datafusion.execution.optional_filter_min_saving_ns_per_row`) plus the decode time of the output columns that the filter does not read. The scan estimates this decode time from the compressed size of those column chunks (file metadata) and a decode speed in ns per compressed byte that it measures over the whole scan (time to produce each output batch from the decoder, which does not include the row filter). `ParquetSource::try_pushdown_filters` reads the mode and the minimum saving from the session configuration. In all modes, required conjuncts do not change, and statistics pruning (files, row groups, pages) uses optional conjuncts as before. An optional conjunct that cannot be pushed down for a file is dropped for that file. Add the lazily registered metrics `optional_filter_rows_skipped`, `optional_filter_pauses` and `optional_filter_eval_time`. 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>
`FilterExec` now handles optional conjuncts (see `OptionalFilterPhysicalExpr`) according to `datafusion.execution.optional_filter_mode`: - `always` (default): the predicate is not split; today's behavior. - `pruning_only`: optional conjuncts are not evaluated. - `adaptive`: required conjuncts run as one ordinary `BinaryExpr` AND chain, and each optional conjunct runs behind its own `OptionalFilterGate`, so optional conjuncts that remove no rows, or that cost more than they save, are paused. The gates measure the evaluation time; the saving for each removed row is `datafusion.execution.optional_filter_min_saving_ns_per_row`. Each stream has its own gates; streams and executions do not share state. Optional conjuncts do not feed equivalence classes or constants. New metrics: `optional_filter_rows_skipped`, `optional_filter_pauses`, `optional_filter_eval_time`. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
A gate decides only on windows of at least `MIN_OBSERVED_ROWS` rows. Each test batch now repeats `0..100` so that half a batch has at least `MIN_OBSERVED_ROWS` rows (the smallest window of a gate). PR: apache#25683 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
`HashJoinExec`, `SortExec` (TopK) and `AggregateExec` now push their dynamic filters down as `Optional(DynamicFilter)`. The operator itself still removes the rows that the filter would remove, so the filter is only a performance hint. The producer keeps its own unwrapped `DynamicFilterPhysicalExpr` for `update()` and `mark_complete()`; only the pushed copy is wrapped. A partitioned hash join wraps each of its two pushed filters (bounds and membership) separately, so both stay direct conjuncts of the scan predicate. There is no behavior change. `OptionalFilterPhysicalExpr` evaluates its inner expression and `snapshot()` removes the wrapper, so scans and pruning use the filter as before. Only the EXPLAIN text changes, from `DynamicFilter [...]` to `Optional(DynamicFilter [...])`. Hash join key transfer rewrites a parent filter below the wrapper, so a transferred optional dynamic filter stays optional. A transferred required filter stays required: for inner and semi joins it is an exact replacement of the parent filter. Direct downcasts that must see through the wrapper: - `HashJoinExec::consumed_dynamic_filter` unwraps its self filters before it looks for a consumer, so the join still produces its filters. - `NestedLoopJoinExec::gather_filters_for_pushdown` uses the new `as_dynamic_filter` helper to route parent dynamic filters. Add `as_dynamic_filter` and `debug_assert_optional_on_root_chain` to `physical_expr::utils`. `ParquetSource::try_pushdown_filters` now checks in debug builds that optional filters are direct conjuncts of the root AND chain. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Update the expected plans in sqllogictest files. Pushed-down dynamic filters now show as `Optional(DynamicFilter [...])`. No query result changes. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…h removed row The hash join records, in the `RemovedRowWork` of each dynamic filter that it produces, the rows of each probe batch and the time of the work that it does for every probe row, match or no match: the evaluation and the hashes of the join keys and the hash table lookup. A row that the filter removes before the join does not get this work. The work for a matched row (the output) is not in it: the filter does not remove matched rows. PR: apache#25681 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
A row that a dynamic filter removes is a probe row without a match. Its saving is the work that such a row gets: the evaluation and the hashes of the join keys and the hash table lookup. The check of the candidates (`equal_rows_arr`) and the output indices are work for the matches only. While the filter is on, most probe rows that reach the join are matches, thus this work made the measured saving too large. PR: apache#25681 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>
…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#25727 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The tests of apache#25729 (page index select-all) and the test of apache#25727 for the coalesced rows at row group boundaries show the TopK dynamic filter in the scan predicate. With apache#25681 the pushed filter is `Optional(DynamicFilter [...])`. Integration of apache#25729, apache#25727 and apache#25681 in the final state branch. 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>
`ParquetSource::exact_filter` (apache#25780) returns the conjuncts that every output row satisfies. An optional conjunct is not one of them: the scan does not evaluate it after the decode (this PR), and it drops it when the `RowFilter` cannot evaluate it. `exact_filter` now skips optional conjuncts, so that the scan claims no equivalence from them. Test: `exact_filter_excludes_optional_conjuncts` (failed before this change: `a@0 = 1 AND Optional(b@1 = 2)` with `pushdown_filters = false`). PR: apache#25722 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
`ParquetSource::exact_filter` (apache#25780) returns the conjuncts that every output row satisfies. An optional conjunct is not one of them: in the `adaptive` mode its gate can skip it, in the `pruning_only` mode the scan does not evaluate it, and the scan drops it when the `RowFilter` cannot evaluate it. `exact_filter` now skips optional conjuncts in all modes, so that the scan claims no equivalence from them. Test: `exact_filter_excludes_optional_conjuncts` (failed before this change: `a@0 = 1 AND Optional(b@1 = 2)`). PR: apache#25682 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
With adaptive filter placement, an optional conjunct can be in the post-scan filter (its gate asks for each batch) or skipped. Say so in the documentation of `ParquetSource::exact_filter`, which already excludes optional conjuncts. PR: apache#25727 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>
…e default Change two defaults: - `datafusion.execution.optional_filter_mode` = `adaptive` (was `always`). The Parquet scan and `FilterExec` pause optional filters (hash join, TopK and aggregate dynamic filters) when they remove no rows or cost more than they save. - `datafusion.execution.adaptive_filter_placement` = `true` (was `false`). With `pushdown_filters = true`, each filter conjunct starts after the decode, and the Parquet scan makes it a row filter only when the measurements show that this saves more than it costs. The Parquet scan uses the two options only when `datafusion.execution.parquet.pushdown_filters` is true. `ParquetSource::new` keeps adaptive filter placement off, because the source has no setter for it: only `try_pushdown_filters` applies the session value. Tests that check exact row filter metrics (predicate cache, runtime row group pruning, optional filters in the row filter, EXPLAIN ANALYZE categories) set the previous values, because the adaptive decisions use time measurements, and because adaptive filter placement coalesces small batches, which delays TopK dynamic filters on tiny row groups. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The previous commit made the scan hand the coalesced rows to the TopK at each row group boundary. Thus the runtime row group pruning tests do not need the previous defaults any more: - `dynamic_row_group_pruning.slt` runs with the defaults. A new case checks that the scan prunes the same row groups without adaptive filter placement. - The Parquet test harness keeps `optional_filter_mode = 'always'`, because its tests check exact metrics, but uses the default adaptive filter placement. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
With the adaptive filter placement, an optional filter starts in the post-scan filter like a required conjunct, thus the `EXPLAIN ANALYZE` metrics move from `pushdown_rows_*` and `predicate_cache_*` to `post_scan_rows_*`. The results do not change. PR: apache#25727 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
EXPERIMENT: pin every arrow crate to pydantic/arrow-rs claude/push-decoder-batch-granular-scan-plan-60.0.0: the 60.0.0 release plus batch-granular decoding in ParquetPushDecoder (apache/arrow-rs#6946) and ParquetPushDecoder::scan_plan (apache/arrow-rs#10555). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
… read-ahead Opt-in with datafusion.execution.parquet.read_ahead_bytes (default unset). If set, the scan builds the same ParquetPushDecoder as the default path with FetchGranularity::Batch and drives it with try_decode. ReadAhead fetches ranges from scan_plan() in the background within that many bytes. datafusion.execution.parquet.read_ahead_conditional (default false) also reads ahead ranges that a pushed-down filter can make unnecessary. Filtered scans, runtime row-group pruning and the per-row-group filter toggle use the same path as main. Squashed from the history in branch claude/datafusion-rowgroup-buffering-history (pydantic/datafusion). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…ranges Skipping ranges that a pushed-down filter can make unnecessary made filtered scans on object storage 2-16x slower. Read-ahead now fetches them, bounded by the read-ahead window. Proto field 40 is reserved. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Each read-ahead stream registers a `ParquetReadAhead[partition]` consumer. The reservation follows the bytes the decoder holds plus the bytes in flight. Bytes the decoder asks for are always reserved (grow). Speculative read-ahead takes only what the pool can grant (try_grow); ranges that do not fit stay pending. `FileSource::create_morselizer_with_context` gives the source the scan's `TaskContext`. The default delegates to `create_morselizer`. `ParquetSource` uses it to get the memory pool. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Member
Author
|
run benchmark tpch tpcds clickbench_partitioned env:
DATAFUSION_EXECUTION_PARQUET_PUSHDOWN_FILTERS: "true"
SIMULATE_LATENCY: "true"
changed:
env:
DATAFUSION_EXECUTION_PARQUET_READ_AHEAD_BYTES: "104857600" |
|
Macroscope skipped reviewing this pull request. Per-review cost limit exceeded (workspace setting). This review would cost an estimated $56.68, which exceeds your per-review limit of $10.00. The top 3 files driving up this estimate:
Tip To get this pull request reviewed, you can:
|
Member
Author
|
Moved to adriangb#17 (the benchmark bot does not watch this repo). |
|
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 |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Combines apache#24086 (batch-granular Parquet scans with
read_ahead_bytes) and apache#25752 (optional filters stack) for benchmarking. Do not merge.The merge resolves conflicts in
datafusion/datasource-parquet/src/push_decoder.rs: the read-ahead path (transition_streaming) and the default path share the post-scan filter, coalescer, limit and row group boundary logic (pending_output,push_decoded_batch,handle_row_group_boundary).Local checks:
datafusion-datasource-parquetunit tests pass, parquet sqllogictests pass, and the full sqllogictest suite withread_ahead_bytesforced to 64 KiB returns the same query results (only predicate cache metrics differ).🤖 Generated with Claude Code