Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
65 commits
Select commit Hold shift + click to select a range
446765a
fix: do not install a Parquet row selection when the page index prune…
adriangb Sep 25, 2026
7607cb8
perf: collapse partitioned hash join InList dynamic filters into one …
adriangb Sep 24, 2026
d2a3eb3
perf: dedupe collapsed hash join IN lists and add combined bounds
adriangb Sep 24, 2026
205c2e7
perf: push partitioned hash join dynamic filter bounds and membership…
adriangb Sep 24, 2026
bbbd3a3
bench: add ordered-subset hash join dynamic filter benchmark
adriangb Sep 24, 2026
9219137
fix: rename test to pass the typos check
adriangb Sep 25, 2026
e74402b
feat(parquet-datasource): always accept pushable filters, run rejecte…
adriangb Sep 24, 2026
588e894
Update tests added on main for always-accepted parquet filters
adriangb Sep 24, 2026
7af48ce
fix: stop the Parquet push decoder stream after the final flush
adriangb Sep 25, 2026
c5ce4b2
refactor(physical-plan): share the round-robin row check of EnforceDi…
adriangb Sep 26, 2026
bbb6193
fix(datasource): keep a filter above a scan that cannot give the targ…
adriangb Sep 26, 2026
b003a46
feat: add OptionalFilterPhysicalExpr to mark filters not needed for c…
adriangb Sep 24, 2026
8f894aa
feat: keep optional filters pruning-only in the Parquet post-scan path
adriangb Sep 24, 2026
d0bd6e9
feat: add filter_stats (FilterCost, Clock) for filter measurements
adriangb Sep 24, 2026
5c81d84
feat: adaptive conjunct reordering in FilterExec
adriangb Sep 24, 2026
0bbc09e
feat: add OptionalFilterGate to pause optional filters that cost more…
adriangb Sep 24, 2026
81cafbd
feat(gate): add a measured overhead for each evaluated row to the cost
adriangb Sep 25, 2026
3eecaed
feat(gate): share pauses between the gates of one plan site
adriangb Sep 25, 2026
68d5038
feat(gate): decide only on windows of at least MIN_OBSERVED_ROWS rows
adriangb Sep 25, 2026
7b412f4
feat(gate): use the measured work of the filter producer as the saving
adriangb Sep 25, 2026
b9c2579
feat: add optional_filter_mode config
adriangb Sep 24, 2026
d2b303e
feat: let the Parquet scan skip optional filters that are not worth t…
adriangb Sep 24, 2026
9bd9ec1
test: end-to-end tests for optional filters in the Parquet post-scan …
adriangb Sep 25, 2026
3802295
feat: let FilterExec skip optional filters that are not worth their cost
adriangb Sep 24, 2026
06fa789
test(filter): size the FilterExec optional filter tests for gate windows
adriangb Sep 25, 2026
da18a81
feat: mark hash join, TopK and aggregate dynamic filters as optional
adriangb Sep 24, 2026
586f5e9
test: update plan snapshots for optional dynamic filters
adriangb Sep 24, 2026
9c3922b
feat(hash-join): measure the work that a dynamic filter saves for eac…
adriangb Sep 25, 2026
76e4b9d
fix(hash-join): measure only the work of a probe row without a match
adriangb Sep 26, 2026
c5ed30c
feat(parquet): adaptive placement of scan filter conjuncts
adriangb Sep 24, 2026
19dc188
test: sqllogictests for adaptive filter placement
adriangb Sep 24, 2026
8a166f4
feat(parquet): carry row selections across placement rebuilds
adriangb Sep 25, 2026
6d81b2f
feat(parquet): measure skippable rows on the stage input
adriangb Sep 25, 2026
ad398d8
feat(parquet): start filters post-scan and add a row filter stage cost
adriangb Sep 25, 2026
647af41
fix(parquet): hand coalesced rows downstream at row group boundaries
adriangb Sep 25, 2026
e34984d
fix(parquet): count the actual rows in optional_filter_rows_skipped
adriangb Sep 25, 2026
3d5fbb4
fix(parquet): count only skippable rows in the saving of an optional …
adriangb Sep 25, 2026
a455e18
feat(parquet): share the pauses of optional filter gates in a scan
adriangb Sep 25, 2026
56a99da
feat(parquet): place and order required and optional conjuncts with o…
adriangb Sep 25, 2026
e7e9928
feat(parquet): add hysteresis to the placement changes of a file
adriangb Sep 25, 2026
5791de6
feat(parquet): measure unmeasured conjuncts first and forget old filters
adriangb Sep 25, 2026
72f6691
test(parquet): gate windows of MIN_OBSERVED_ROWS rows in row filter t…
adriangb Sep 25, 2026
6276115
refactor(parquet): use the shared MIN_OBSERVED_ROWS in the placement …
adriangb Sep 25, 2026
0b80d81
feat(parquet): let measured evidence decide the placement, drop the s…
adriangb Sep 26, 2026
79f9869
feat(parquet): order row filter predicates by measured rows removed p…
adriangb Sep 25, 2026
9616f5c
refactor(parquet): share the ranking and the hysteresis with the row …
adriangb Sep 25, 2026
cd05bb2
feat(parquet): add the measured post-scan copy to the benefit of a ro…
adriangb Sep 25, 2026
1060d95
test(parquet): size the optional filter integration tests for gate wi…
adriangb Sep 25, 2026
6a3a7da
feat(parquet): order the post-scan conjuncts by their measurements at…
adriangb Sep 25, 2026
22b411f
fix(parquet): fix rustdoc links of the adaptive filter placement
adriangb Sep 26, 2026
4fac55a
fix(parquet): compact the post-scan batch with the pre-selection thre…
adriangb Sep 26, 2026
3318ef1
test: optional dynamic filters in runtime row group pruning tests
adriangb Sep 27, 2026
c32df68
fix(parquet): no scan equivalences from a pruning-only predicate
adriangb Sep 28, 2026
45dcc5e
fix(parquet): no scan equivalences from optional conjuncts
adriangb Sep 28, 2026
a7b6bc5
fix(parquet): no scan equivalences from optional conjuncts
adriangb Sep 28, 2026
2cd2a65
docs(parquet): optional conjuncts of the placement are not exact
adriangb Sep 28, 2026
bd93482
test: the scan applies the filter of the wrapping multiplication test
adriangb Sep 28, 2026
b88a5b5
feat: make adaptive optional filters and adaptive filter placement th…
adriangb Sep 25, 2026
3e5f53c
test: run runtime row group pruning tests with the new defaults
adriangb Sep 25, 2026
f1036e8
test: optional filters start in the post-scan filter in sqllogictests
adriangb Sep 25, 2026
f679435
chore: pin arrow-rs with batch-granular push decoding and scan_plan
adriangb Sep 25, 2026
754d844
POC: batch-granular Parquet scans with the push decoder and scan_plan…
adriangb Sep 25, 2026
667c4dc
remove read_ahead_conditional; read-ahead always fetches conditional …
adriangb Sep 27, 2026
02429ca
account Parquet read-ahead in the MemoryPool
adriangb Sep 27, 2026
07caaa3
Merge PR 24086 (read-ahead) into PR 25752 (optional filters)
adriangb Sep 29, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
58 changes: 20 additions & 38 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

22 changes: 22 additions & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -406,3 +406,25 @@ debug = false
debug-assertions = false
strip = "debuginfo"
incremental = false

# Pin arrow-rs to pydantic/arrow-rs claude/push-decoder-batch-granular-scan-plan-60.0.0:
# 60.0.0 plus apache/arrow-rs#6946 (batch-granular decoding) and
# apache/arrow-rs#10555 (scan_plan). Patch every arrow crate together.
[patch.crates-io]
arrow-arith = { git = "https://github.com/pydantic/arrow-rs.git", branch = "claude/push-decoder-batch-granular-scan-plan-60.0.0" }
arrow-array = { git = "https://github.com/pydantic/arrow-rs.git", branch = "claude/push-decoder-batch-granular-scan-plan-60.0.0" }
arrow-avro = { git = "https://github.com/pydantic/arrow-rs.git", branch = "claude/push-decoder-batch-granular-scan-plan-60.0.0" }
arrow-buffer = { git = "https://github.com/pydantic/arrow-rs.git", branch = "claude/push-decoder-batch-granular-scan-plan-60.0.0" }
arrow-cast = { git = "https://github.com/pydantic/arrow-rs.git", branch = "claude/push-decoder-batch-granular-scan-plan-60.0.0" }
arrow-csv = { git = "https://github.com/pydantic/arrow-rs.git", branch = "claude/push-decoder-batch-granular-scan-plan-60.0.0" }
arrow-data = { git = "https://github.com/pydantic/arrow-rs.git", branch = "claude/push-decoder-batch-granular-scan-plan-60.0.0" }
arrow-flight = { git = "https://github.com/pydantic/arrow-rs.git", branch = "claude/push-decoder-batch-granular-scan-plan-60.0.0" }
arrow-ipc = { git = "https://github.com/pydantic/arrow-rs.git", branch = "claude/push-decoder-batch-granular-scan-plan-60.0.0" }
arrow-json = { git = "https://github.com/pydantic/arrow-rs.git", branch = "claude/push-decoder-batch-granular-scan-plan-60.0.0" }
arrow-ord = { git = "https://github.com/pydantic/arrow-rs.git", branch = "claude/push-decoder-batch-granular-scan-plan-60.0.0" }
arrow-row = { git = "https://github.com/pydantic/arrow-rs.git", branch = "claude/push-decoder-batch-granular-scan-plan-60.0.0" }
arrow-schema = { git = "https://github.com/pydantic/arrow-rs.git", branch = "claude/push-decoder-batch-granular-scan-plan-60.0.0" }
arrow-select = { git = "https://github.com/pydantic/arrow-rs.git", branch = "claude/push-decoder-batch-granular-scan-plan-60.0.0" }
arrow-string = { git = "https://github.com/pydantic/arrow-rs.git", branch = "claude/push-decoder-batch-granular-scan-plan-60.0.0" }
arrow = { git = "https://github.com/pydantic/arrow-rs.git", branch = "claude/push-decoder-batch-granular-scan-plan-60.0.0" }
parquet = { git = "https://github.com/pydantic/arrow-rs.git", branch = "claude/push-decoder-batch-granular-scan-plan-60.0.0" }
22 changes: 22 additions & 0 deletions benchmarks/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -966,6 +966,28 @@ Several queries are included to test hash joins under various workloads.
./bench.sh run hj
```

## Hash Join Dynamic Filter on an Ordered Subset

This benchmark (`hj_ordered_subset`) measures the dynamic filter that a hash join pushes to its probe-side Parquet scan when the probe table is sorted by the join key and the build side matches an ordered subset of the key range. This is common for time-ordered fact tables: the events of one day, or of a range of days.

The load SQL writes an `events` table of 100 days x 200,000 rows (20M rows, 100,000-row row groups) sorted by `event_id`, and two build tables with one row for each event. The queries join `events` to 1 day (Q01, Q04), to 10 days (Q02, Q05), or to 1% of the events scattered over the whole range (Q03, Q06, the control). The bounds of the dynamic filter (`event_id >= min AND event_id <= max`) can prune the `events` row groups outside the matched range, and the membership check passes almost every row that remains. In the control, the bounds cannot prune anything.

The `partitioned` subgroup (Q01-Q03) forces `HashJoinExec` mode `Partitioned`, and the `collect_left` subgroup (Q04-Q06) forces mode `CollectLeft`.

### Example Run

```bash
# No need to generate data: the suite's load SQL writes the Parquet files

./bench.sh run hj_ordered_subset

# With parquet filter pushdown
DATAFUSION_EXECUTION_PARQUET_PUSHDOWN_FILTERS=true ./bench.sh run hj_ordered_subset

# Smaller data: 50 days x 100,000 rows
HJOS_DAYS=50 HJOS_ROWS_PER_DAY=100000 ./bench.sh run hj_ordered_subset
```

## Null-Aware Join

This benchmark focuses on `NOT IN` subqueries, which plan as null-aware joins: an
Expand Down
34 changes: 34 additions & 0 deletions benchmarks/bench.sh
Original file line number Diff line number Diff line change
Expand Up @@ -119,6 +119,10 @@ spill_views: Sort and GROUP BY queries that spill StringView/BinaryVi
null_aware_join: Null-aware (NOT IN) hash join micro-benchmarks: uncorrelated, non-equality-correlated and equality-correlated
NOT IN across NULL fractions, to measure the per-pair join-filter work the correlated cases do
(data generated inline by the suite's load SQL from range(); knobs: NAJ_ROWS, NAJ_LARGE_ROWS)
hj_ordered_subset: Hash join dynamic filter on a probe table sorted by the join key, where the build side matches an ordered subset
(1 day, 10 days) of the key range, so the filter bounds can prune probe row groups; plus a scattered 1% control
(subgroups via BENCH_SUBGROUP: partitioned, collect_left)
(data generated inline by the suite's load SQL; knobs: HJOS_DAYS, HJOS_ROWS_PER_DAY, HJOS_RG_SIZE)
projection_subquery: IN / NOT IN / EXISTS subqueries in the SELECT list (see https://github.com/apache/datafusion/issues/25341); each query projects the
boolean subquery result and aggregates it, so the cost is the decorrelation plan and not the output size
(q07 correlates on '<' instead of '=', so it keeps the nested-loop plan and acts as the control)
Expand Down Expand Up @@ -295,6 +299,10 @@ main() {
# Data is generated inline by the suite's load SQL.
echo "projection_subquery: no external data to generate"
;;
hj_ordered_subset)
# Data is generated inline by the suite's load SQL (COPY).
echo "hj_ordered_subset: no external data to generate"
;;
asof_join)
data_asof_join
;;
Expand Down Expand Up @@ -552,6 +560,9 @@ main() {
projection_subquery)
run_projection_subquery
;;
hj_ordered_subset)
run_hj_ordered_subset
;;
asof_join)
run_asof_join
;;
Expand Down Expand Up @@ -1009,6 +1020,29 @@ run_null_aware_join() {
bash -c "$SQL_CARGO_COMMAND"
}

# Runs the hj_ordered_subset suite: a hash join whose probe table is sorted by
# the join key and whose build side matches an ordered subset of the key range,
# so the bounds of the join's dynamic filter can prune probe row groups. Data is
# generated inline by the load SQL (COPY), so there is no data step. Set
# DATAFUSION_EXECUTION_PARQUET_PUSHDOWN_FILTERS to vary how the scan uses the
# filter.
# Knobs (string-substituted into the load SQL, not engine config):
# BENCH_SUBGROUP run one subgroup (partitioned, collect_left)
# HJOS_DAYS days in the events table (default 100, at least 50)
# HJOS_ROWS_PER_DAY events per day (default 200_000)
# HJOS_RG_SIZE parquet row-group size (default 100_000)
run_hj_ordered_subset() {
echo "Running hj_ordered_subset benchmark (subgroup=${BENCH_SUBGROUP:-all}, days=${HJOS_DAYS:-100}, rows_per_day=${HJOS_ROWS_PER_DAY:-200000}, rg_size=${HJOS_RG_SIZE:-100000})..."
debug_run env BENCH_NAME=hj_ordered_subset \
BENCH_RESULTS_FILE="$(sql_results_file hj_ordered_subset)" \
${BENCH_SUBGROUP:+BENCH_SUBGROUP="${BENCH_SUBGROUP}"} \
HJOS_DAYS="${HJOS_DAYS:-100}" \
HJOS_ROWS_PER_DAY="${HJOS_ROWS_PER_DAY:-200000}" \
HJOS_RG_SIZE="${HJOS_RG_SIZE:-100000}" \
${QUERY:+BENCH_QUERY="${QUERY}"} \
bash -c "$SQL_CARGO_COMMAND"
}

# Runs the projection_subquery suite: IN / NOT IN / EXISTS subqueries that sit
# in the SELECT list instead of a filter (see
# https://github.com/apache/datafusion/issues/25341). The load SQL builds the
Expand Down
1 change: 1 addition & 0 deletions benchmarks/sql_benchmarks/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@ in the community:
| `clickbench_sorted` | ClickBench benchmark using a pre-sorted hits file. |
| `h2o` | The `h2o` benchmark |
| `hj` | Hash join benchmark |
| `hj_ordered_subset` | Hash join dynamic filter on a probe table sorted by the join key, where the build side matches an ordered subset of the key range (1 day, 10 days) or a scattered 1% (control). Subgroups (`--subgroup`): `partitioned`, `collect_left`. Size the data with `HJOS_DAYS`, `HJOS_ROWS_PER_DAY` and `HJOS_RG_SIZE`. |
| `imdb` | IMDb benchmark |
| `nlj` | Nested‑loop join benchmark |
| `null_aware_join` | Null-aware (`NOT IN`) hash join micro-benchmarks. Q01-Q03 are uncorrelated `NOT IN` across NULL fractions and are linear in the table size (`NAJ_LARGE_ROWS`, default `1000000`). Q04-Q08 are correlated, so the correlation predicate stays behind as a join filter that the join applies per candidate (build row × probe row) pair while deciding which outer rows are UNKNOWN; without an equality correlation there are no scope keys to narrow those pairs, so their cost grows with the square of `NAJ_ROWS` (default `10000`). Q08 adds an equality correlation, which turns those pairs into a hash lookup. All tables are built inline from `range()`, so there is no data step. |
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
# Partitioned join, build side = 1 day (1% of the event_id range, ordered).
# The bounds of the dynamic filter can prune ~99% of the events row groups.
subgroup partitioned

template sql_benchmarks/hj_ordered_subset/hj_ordered_subset.benchmark.template
NAME=Q01_partitioned_ordered_1_day
JOIN_MODE=partitioned
PLAN_MODE=Partitioned
BUILD_TABLE=labels_by_day
BUILD_FILTER=l.day = 42
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
# Partitioned join, build side = 10 days (10% of the event_id range, ordered).
# The bounds of the dynamic filter can prune ~90% of the events row groups.
subgroup partitioned

template sql_benchmarks/hj_ordered_subset/hj_ordered_subset.benchmark.template
NAME=Q02_partitioned_ordered_10_days
JOIN_MODE=partitioned
PLAN_MODE=Partitioned
BUILD_TABLE=labels_by_day
BUILD_FILTER=l.day BETWEEN 40 AND 49
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
# Control: partitioned join, build side = 1% of the events, scattered over the
# whole event_id range. The bounds cannot prune anything; only the membership
# check of the dynamic filter rejects rows.
subgroup partitioned

template sql_benchmarks/hj_ordered_subset/hj_ordered_subset.benchmark.template
NAME=Q03_partitioned_scattered_1_pct
JOIN_MODE=partitioned
PLAN_MODE=Partitioned
BUILD_TABLE=labels_by_bucket
BUILD_FILTER=l.bucket = 42
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
# Same as Q01 with a CollectLeft join.
subgroup collect_left

template sql_benchmarks/hj_ordered_subset/hj_ordered_subset.benchmark.template
NAME=Q04_collect_left_ordered_1_day
JOIN_MODE=collect_left
PLAN_MODE=CollectLeft
BUILD_TABLE=labels_by_day
BUILD_FILTER=l.day = 42
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
# Same as Q02 with a CollectLeft join.
subgroup collect_left

template sql_benchmarks/hj_ordered_subset/hj_ordered_subset.benchmark.template
NAME=Q05_collect_left_ordered_10_days
JOIN_MODE=collect_left
PLAN_MODE=CollectLeft
BUILD_TABLE=labels_by_day
BUILD_FILTER=l.day BETWEEN 40 AND 49
Loading
Loading