adriangb opened a new issue, #25035:
URL: https://github.com/apache/datafusion/issues/25035
### Is your feature request related to a problem or challenge?
`BinaryExpr::evaluate`
(`datafusion/physical-expr/src/expressions/binary.rs:542`)
evaluates its LHS, then `check_short_circuit` (`binary.rs:1188`) returns
`PreSelection` for `AND` when the LHS boolean array has no nulls and
`true_count / len <= PRE_SELECTION_THRESHOLD` (0.2, `binary.rs:1167`).
Pre-selection
does `filter_record_batch(batch, &mask)` on the original batch
(`binary.rs:562`),
evaluates the RHS on the survivors, and `pre_selection_scatter`s back to
full length
(`binary.rs:1307`).
Conjunctions are built **left-deep** — `datafusion_expr::utils::conjunction`
is
`reduce(Expr::and)` (`datafusion/expr/src/utils.rs:1296`),
`datafusion_physical_expr::utils::conjunction_opt` is a left fold
(`datafusion/physical-expr/src/utils/mod.rs:122`), and SQL's `AND` is
left-associative. So at each level the LHS is the **accumulated prefix**,
not a
single conjunct. Two consequences:
1. **The threshold is tested against the wrong quantity.** With several
individually unselective conjuncts, no single conjunct is selective, but
the
prefix multiplies below 20% partway down the chain. Pre-selection then
fires to
save evaluating a couple of cheap comparisons on the remaining rows — a
losing
trade.
2. **Left-nesting repeats the work.** Every level that trips the threshold
re-filters the *original* batch and scatters back to full length, instead
of
compacting once and keeping survivors compacted. And
`filter_record_batch` copies
**every column of the batch**, not just the ones the RHS reads, so the
cost
scales with batch width.
Measured on `BinaryExpr` directly at `262936eef5` (8192-row batches, k
conjuncts
`ci < cutoff` on independent Int32 columns, left-deep vs right-deep, median
of 7
interleaved rounds). `fires` counts prefix levels below the threshold:
| k | per-conjunct pass rate | left-deep | right-deep | right/left | fires |
|---|---|---|---|---|---|
| 4 | 0.50 | 13.3 us | 5.0 us | 0.376 | 1 |
| 8 | 0.50 | 31.8 us | 10.9 us | 0.344 | 5 |
| 8 | 0.70 | 41.0 us | 10.8 us | 0.263 | 3 |
| 8 | 0.90 | 11.0 us | 10.8 us | 0.988 | 0 |
| 16 | 0.70 | 95.8 us | 24.4 us | 0.254 | 11 |
| 16 | 0.90 | 25.5 us | 24.6 us | 0.967 | 0 |
Every `fires = 0` row is within 4% of parity, so the two shapes are
equivalent when
the threshold never trips.
Width amplifies it. Same 8 conjuncts, with unreferenced payload columns
added to the
batch:
| extra columns | pass rate | left-deep | right-deep | right/left |
|---|---|---|---|---|
| 0 | 0.70 | 40.3 us | 11.1 us | 0.277 |
| 8 | 0.70 | 50.7 us | 11.0 us | 0.218 |
| 32 | 0.70 | 81.9 us | 11.3 us | 0.138 |
End to end this costs TPC-H Q06 — the canonical 4-conjunct `lineitem` scan
filter —
about 17%.
(Noticed while working on runtime conjunct reordering in #22698; the problem
above
is in stock DataFusion and independent of that PR.)
### Describe the solution you'd like
**Option A: build conjunctions right-deep.** Recommend against. Besides ~180
expected-output lines across ~30 `.slt` files, the InList merge rules at
`datafusion/optimizer/src/simplify_expressions/expr_simplifier.rs:1888-1975`
require
the two InLists to be direct siblings and would silently stop firing; proto's
linearized `operands` form reconstructs left-deep (`binary.rs:973`), so the
shape
does not survive serialization; and it changes which rows a fallible
conjunct sees
(`b <> 0 AND 1/b > 2`), a user-visible semantic change. (`split_conjunction`
itself
is shape-agnostic in both crates, so that part would have been fine.)
**Option B: evaluate `AND` chains n-ary.** Keep the tree shape, but when
evaluating
an `AND`, flatten the chain, evaluate conjuncts against a working batch, and
compact
only when the accumulated mask crosses the threshold — what a right-deep
tree would
give without changing the tree. EXPLAIN, proto and the simplifier rules are
all
untouched. Three-valued logic is preserved by refusing to compact while the
accumulated mask has nulls, matching today's behaviour. `OR` is left alone.
This is how Velox evaluates conjunctions, and the point is worth making
explicitly
because Velox is where the `time / (1 + n_in - n_out)` conjunct-ordering
metric
comes from: that metric is applied to a *flat list* of conjuncts evaluated
against a
*narrowing* row set, not to a binary tree. Verified against the paper and the
current source:
- **One n-ary node, flattened at compile time.** Pedreira et al., *Velox:
Meta's
Unified Execution Engine*, PVLDB 15(12), 2022, §4.3.1 "Adaptive Conjunct
Reordering" (p. 3377): "AND(AND(AND(a, b), c), AND(d, e)) is flattened to
a single
AND(a, b, c, d, e) node during compilation". In source, `ConjunctExpr`
(`velox/expression/ConjunctExpr.h`) takes a `std::vector<ExprPtr>&&
inputs`;
`ExprCompiler.cpp` treats `and`/`or` as flattenable in `shouldFlatten` and
calls
`expression::utils::flattenInput` (`velox/expression/ExprUtils.cpp`), which
recursively folds nested same-name calls into one input list.
- **Later conjuncts only see surviving rows.**
`ConjunctExpr::evalSpecialForm` (`ConjunctExpr.cpp`) evaluates each input
with
`inputs_[inputOrder_[i]]->eval(*activeRows, ...)`, then `updateResult`
removes
decided rows from `activeRows` and the loop exits when `countSelected()`
reaches 0.
Velox passes the shrinking `SelectivityVector` *down* into the child
rather than
copying the batch, so there is no filter-and-scatter at all; DataFusion's
kernels
take whole arrays, so compacting once and keeping the survivors compacted
is the
closest equivalent.
- **Three-valued logic is handled per row, without stopping evaluation on
null.**
In `updateAnd` the drop set is `testFalse = ~testValue & testPresent` and
only
`active &= ~testFalse` — a definite `false` retires the row; a `null`
input flips a
so-far-`true` result to null but the row stays active, so a later `false`
can
still decide it. That is the same rule option B needs (a null in the
accumulated
mask must not be treated as eliminated).
- **Reordering is a separate concern layered on top.** `maybeReorderInputs`
sorts
`inputOrder_` by `SelectivityInfo::timeToDropValue()`
(`velox/common/base/SelectivityInfo.h`: `timeClocks_ / (numIn_ -
numOut_)`, with a
guard when nothing was dropped — the paper writes it as
`time / (1 + n_in - n_out)`), gated by
`adaptiveFilterReorderingEnabled()`. The
n-ary, narrowing evaluation works with reordering off; it is the shape
that makes
reordering meaningful, not the other way round.
DuckDB does the same thing: `ExpressionExecutor::Select` on a
`BoundConjunctionExpression`
(`src/execution/expression_executor/execute_conjunction.cpp`)
walks the children in `permutation` order and, once rows are filtered out,
switches
`current_sel = true_sel` so each subsequent child evaluates only the passing
tuples;
its `AdaptiveFilter` (`src/execution/adaptive_filter.cpp`) reorders by
randomized
adjacent swaps kept or reverted on measured runtime. I did not check
ClickHouse.
**Option C: make the compaction decision cost-aware.** Option B alone is
necessary
but not sufficient. With B, shape stops mattering (left/right parity within
0.4%),
but at 40 columns / 0.70 pass rate it still lands at 36.9 us against 11.3 us
for a
tree that never compacts — B removes the *repeated* filtering but still
performs one
unprofitable compaction, because the threshold cannot see that the remaining
conjuncts are cheap and the batch is wide. A complete fix needs B plus a
decision
that weighs the filter and scatter against the RHS's actual cost and the
batch
width.
### Reproducing
Micro-benchmark: an integration test in `datafusion/physical-expr/tests/`
that builds
the same conjuncts as left-deep and right-deep `BinaryExpr` trees and times
both
interleaved. Sweep k in {2,4,8,16}, per-conjunct pass rate in
{0.5,0.7,0.9,0.99},
batch sizes {1024,8192}, plus unreferenced payload columns. The existing
`datafusion/physical-expr/benches/binary_op.rs` (`cargo bench --bench
binary_op`,
`benchmark_binary_op_in_short_circuit`) already exercises this code path and
would
be the natural home for a regression benchmark.
End to end:
```
cd benchmarks
../target/release/benchmark_runner tpch -i 5
../target/release/benchmark_runner predicate_eval -i 5
```
### Measured results
Prototype of option B behind an env switch so one binary runs both
configurations;
6-8 interleaved A/B/A' rounds x 5 iterations, AC power, medians, warm-up
iteration
dropped.
| suite | query | baseline | prototype | ratio | rounds faster | A/A drift |
|---|---|---|---|---|---|---|
| TPC-H SF1 | Q06 | 14.80 ms | 12.23 ms | 0.826 | 6/6 | 0.980 |
| TPC-H SF1 | all other 21 | | | 0.95-1.04 | coin flip | median 1.4% |
| predicate_eval 1M | cardinality_q33_k16 | 2.724 ms | 2.157 ms | 0.792 |
7/8 | 1.014 |
| predicate_eval 1M | all other 25 | | | 0.95-1.12 | coin flip | median 2.6%
|
Q06 and cardinality_q33_k16 are the only results outside the noise floor.
Both are
exactly the predicted shape: several conjuncts whose prefix crosses the
threshold
partway down. Only queries whose runtime is dominated by a multi-conjunct
filter move
at all; the micro-benchmark ratios are filter-evaluation-only numbers.
Correctness of the prototype: all 510 sqllogictest files pass, all
`datafusion-physical-expr` unit tests pass, and all 26 `predicate_eval`
queries
validate against result CSVs persisted from unmodified main.
### Describe alternatives you've considered
Options A and C above.
### Additional context
Related but distinct: #15631 optimizes the short-circuit *check* itself
(`count_ones` versus testing for any set bit), not the quantity being tested.
Separately: `predicate_eval.benchmark.template` has no `result` directive, so
`--result-mode validate` currently verifies nothing for that suite.
Sources for the prior-art section:
https://www.vldb.org/pvldb/vol15/p3372-pedreira.pdf
(§4.3.1), https://github.com/facebookincubator/velox
(`velox/expression/ConjunctExpr.{h,cpp}`,
`velox/expression/ExprCompiler.cpp`, `velox/expression/ExprUtils.cpp`,
`velox/common/base/SelectivityInfo.h`), https://github.com/duckdb/duckdb
(`src/execution/expression_executor/execute_conjunction.cpp`,
`src/execution/adaptive_filter.cpp`); all read at `main` on 2026-09-07.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]