From 65d18ee4ad0913f596c62780c3edcd14921dae29 Mon Sep 17 00:00:00 2001 From: Adrian Garcia Badaracco <1755071+adriangb@users.noreply.github.com> Date: Sun, 28 Jun 2026 16:55:42 +0200 Subject: [PATCH 1/3] feat(physical-plan): adaptive conjunct reordering in FilterExec (compact-once core) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Add runtime, statistics-based conjunct reordering for `FilterExec`, off by default behind `datafusion.execution.adaptive_filter_reordering`. A conjunctive predicate is evaluated through a compact-once loop: conjunct masks are AND-combined and the working batch is physically compacted to the surviving rows once the accumulated mask is selective enough, so a selective conjunct shrinks the batch the conjuncts after it must decode. This compaction — not reordering a fused `BinaryExpr` AND, which does not compact between conjuncts — is the source of the win, and reordering compounds it. Each conjunct is timed and counted on the rows that reach it during a short warm-up; the conjuncts are then ranked by rows discarded per nanosecond (`(1 - pass_rate) / cost_per_row`) and, if the ranked order is materially cheaper than the written one, it is adopted and frozen. Results, plan, and EXPLAIN are unchanged; volatile predicates are never reordered. This is the minimal core. Benchmarks (predicate_eval) confirm it captures the "buried selective conjunct" wins (costsel_q01 ~-14%, width ~-12%) but also that compact-once regresses cheap-predicate conjunctions (cardinality k8 ~+37%) where the compaction overhead is not repaid — a guard that keeps the plain fused evaluation for those is added in the next commit. Cross-stream sharing, drift re-measurement, and confidence-interval statistics are later layers. Tested by unit tests (compact-once result-equivalence in any order, ranking, expected-cost weighting, adopt/keep decisions) and an end-to-end `adaptive_filter.slt` asserting identical results and EXPLAIN with the flag on and off. Co-Authored-By: Claude Opus 4.8 (1M context) Claude-Session: https://claude.ai/code/session_01Lh7i9DyeFWuTFWjogrVNkb --- datafusion/common/src/config.rs | 11 + .../physical-plan/src/adaptive_filter.rs | 538 ++++++++++++++++++ datafusion/physical-plan/src/filter.rs | 26 +- datafusion/physical-plan/src/lib.rs | 1 + .../test_files/adaptive_filter.slt | 80 +++ .../test_files/information_schema.slt | 2 + docs/source/user-guide/configs.md | 1 + 7 files changed, 656 insertions(+), 3 deletions(-) create mode 100644 datafusion/physical-plan/src/adaptive_filter.rs create mode 100644 datafusion/sqllogictest/test_files/adaptive_filter.slt diff --git a/datafusion/common/src/config.rs b/datafusion/common/src/config.rs index cc263dfe3e619..fb0b46d3218df 100644 --- a/datafusion/common/src/config.rs +++ b/datafusion/common/src/config.rs @@ -926,6 +926,17 @@ config_namespace! { /// tables with a highly-selective join filter, but is also slightly slower. pub enforce_batch_size_in_joins: bool, default = false + /// (experimental) When enabled, `FilterExec` adaptively reorders the + /// conjuncts of a conjunctive predicate at runtime. It measures each + /// conjunct's selectivity and evaluation cost on the rows that reach it + /// and runs the conjuncts that discard the most rows per unit of CPU + /// time first, so cheap-and-selective predicates gate expensive ones. + /// Reordering never changes query results (only the evaluation order of + /// a conjunction) but can change observable side effects of fallible + /// predicates, so it is off by default. Predicates containing volatile + /// expressions are never reordered. + pub adaptive_filter_reordering: bool, default = false + /// Size (bytes) of data buffer DataFusion uses when writing output files. /// This affects the size of the data chunks that are uploaded to remote /// object stores (e.g. AWS S3). If very large (>= 100 GiB) output files are being diff --git a/datafusion/physical-plan/src/adaptive_filter.rs b/datafusion/physical-plan/src/adaptive_filter.rs new file mode 100644 index 0000000000000..46e76d6ddc591 --- /dev/null +++ b/datafusion/physical-plan/src/adaptive_filter.rs @@ -0,0 +1,538 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! Runtime-adaptive evaluation of a conjunctive (`AND`) predicate in +//! [`FilterExec`](crate::filter::FilterExec). +//! +//! Predicate evaluation order matters: a selective predicate run first gates +//! the work of the predicates after it. DataFusion's `BinaryExpr` `AND` +//! short-circuit only gates on the *leftmost* conjunct, so a conjunction whose +//! selective member is written last (e.g. +//! `regexp_like(s,'a') AND … AND regexp_like(s,'rare')`) evaluates every +//! predicate against ~every row. +//! +//! ## How it evaluates: the compact-once loop +//! +//! The conjuncts are evaluated sequentially, combining their boolean results +//! with `AND`. The working batch is physically compacted to the surviving rows +//! once the accumulated mask becomes selective enough — so a run of +//! non-selective conjuncts costs only cheap bitwise `AND`s, while a selective +//! conjunct shrinks the batch the conjuncts after it must decode. This +//! compaction is what makes ordering pay off (and is itself a win even without +//! reordering): a left-deep fused `BinaryExpr` `AND` does *not* compact between +//! conjuncts, so it evaluates ~every conjunct on ~every row regardless of order. +//! +//! ## How it orders +//! +//! Each conjunct is timed and counted on exactly the rows it evaluated, giving +//! its marginal selectivity and per-row cost. After a short warm-up the +//! conjuncts are ranked by rows discarded per nanosecond +//! (`(1 - pass_rate) / cost_per_row`, the classic optimal ordering key for +//! independent conjuncts), and if the ranked order is *materially* cheaper than +//! the written one it is adopted. The order then stays fixed. +//! +//! It is **off by default** +//! (`datafusion.execution.adaptive_filter_reordering`) and never changes query +//! results: a conjunction's value is independent of evaluation order. Predicates +//! containing volatile expressions are never reordered (their observable side +//! effects depend on order). +//! +//! This is the minimal core of the adaptive evaluator. Richer policies +//! (A/B-validated adoption, cross-stream sharing, drift re-measurement) build on +//! top of it. + +use std::sync::Arc; + +use arrow::array::{Array, ArrayRef, BooleanArray, BooleanBufferBuilder, UInt32Array}; +use arrow::buffer::BooleanBuffer; +use arrow::compute::kernels::boolean::and; +use arrow::compute::{filter, filter_record_batch, prep_null_mask_filter}; +use arrow::record_batch::RecordBatch; +use datafusion_common::Result; +use datafusion_common::cast::as_boolean_array; +use datafusion_common::instant::Instant; +use datafusion_physical_expr::PhysicalExpr; +use datafusion_physical_expr::utils::split_conjunction; +use datafusion_physical_expr_common::physical_expr::is_volatile; + +/// Batches measured before the order is settled. +const WARMUP_BATCHES: u64 = 8; + +/// Fraction of the conjunction's expected per-row cost below which a reorder is +/// immaterial. A candidate order is adopted only if it is expected to cost less +/// than `(1 - TIE_COST_FRACTION)` of the written order, so interchangeable +/// conjuncts never trigger a reorder. +const TIE_COST_FRACTION: f64 = 0.05; + +/// Physically compact the working batch to the surviving rows only when the +/// accumulated mask keeps at most this fraction of them. Above this, the cost +/// of materializing a barely-smaller batch is not repaid, so we keep evaluating +/// against the full working batch and just `AND` the boolean masks. +const COMPACTION_SELECTIVITY_THRESHOLD: f64 = 0.2; + +/// Per-conjunct measurement: marginal pass rate and per-row evaluation cost, +/// accumulated over the warm-up window on exactly the rows that reached the +/// conjunct. +#[derive(Debug, Default, Clone)] +struct ConjunctStats { + /// Total rows the conjunct was evaluated on. + rows: u64, + /// Rows that passed (non-null `true`, matching SQL filter semantics). + matched: u64, + /// Total evaluation time, nanoseconds. + nanos: u64, +} + +impl ConjunctStats { + fn record(&mut self, matched: u64, rows: u64, nanos: u64) { + self.rows += rows; + self.matched += matched; + self.nanos += nanos; + } + + /// Fraction of rows that pass, or `None` if never evaluated on any row. + fn pass_rate(&self) -> Option { + (self.rows > 0).then(|| self.matched as f64 / self.rows as f64) + } + + /// Per-row evaluation cost in nanoseconds, or `None` if unmeasured. + fn cost_per_row(&self) -> Option { + (self.rows > 0 && self.nanos > 0).then(|| self.nanos as f64 / self.rows as f64) + } + + /// Ranking key: rows discarded per nanosecond of evaluation + /// (`(1 - pass_rate) / cost_per_row`). Maximising this is exactly + /// minimising `cost_per_row / (1 - pass_rate)`, the classic optimal + /// ordering key for independent conjuncts — so a selective-but-expensive + /// predicate correctly sorts ahead of a cheap-but-unselective one. + /// `None` when unmeasured, so such conjuncts sort last. + fn effectiveness(&self) -> Option { + let cost = self.cost_per_row()?; + let pass = self.pass_rate()?; + Some((1.0 - pass) / cost) + } +} + +/// Adaptive evaluator for a single conjunctive predicate, owned per partition +/// stream (single-threaded, no locking). +#[derive(Debug)] +pub(crate) struct AdaptiveConjunction { + /// The split conjuncts. `stats`/`order` indices refer to positions here. + conjuncts: Vec>, + /// Per-conjunct measurements, indexed by conjunct position. + stats: Vec, + /// Evaluation order: indices into `conjuncts`. Starts as the written order + /// and may be permuted once after the warm-up. + order: Vec, + /// Warm-up batches still to measure; `0` means the order has settled. + warmup_left: u64, +} + +impl AdaptiveConjunction { + /// Build an adaptive evaluator for `predicate`, or `None` if adaptive + /// reordering does not apply: + /// + /// - `enabled` is false (the config flag is off); + /// - the predicate has fewer than two `AND` conjuncts (nothing to reorder); + /// - any conjunct is volatile (reordering could change side effects). + pub(crate) fn try_new( + predicate: &Arc, + enabled: bool, + ) -> Option { + if !enabled { + return None; + } + let conjuncts: Vec> = split_conjunction(predicate) + .into_iter() + .map(Arc::clone) + .collect(); + if conjuncts.len() < 2 || conjuncts.iter().any(is_volatile) { + return None; + } + let stats = vec![ConjunctStats::default(); conjuncts.len()]; + let order = (0..conjuncts.len()).collect(); + Some(Self { + conjuncts, + stats, + order, + warmup_left: WARMUP_BATCHES, + }) + } + + /// Evaluate the conjunction against `batch`, returning the boolean mask + /// (over the batch's rows) of rows that passed every conjunct. + /// + /// During the warm-up the conjuncts are measured (on the rows that reach + /// each one); once the window ends the order is settled and subsequent + /// batches evaluate with no instrumentation. + pub(crate) fn evaluate(&mut self, batch: &RecordBatch) -> Result { + if self.warmup_left == 0 { + return eval_conjuncts(&self.conjuncts, &self.order, batch, None); + } + let result = + eval_conjuncts(&self.conjuncts, &self.order, batch, Some(&mut self.stats))?; + self.warmup_left -= 1; + if self.warmup_left == 0 { + self.settle(); + } + Ok(result) + } + + /// Rank the conjuncts by measured effectiveness and, if the resulting order + /// is materially cheaper than the written order, adopt it. + fn settle(&mut self) { + let candidate = rank_by_effectiveness(&self.stats); + if candidate == self.order { + return; + } + let candidate_cost = expected_cost_per_row(&self.stats, &candidate); + let current_cost = expected_cost_per_row(&self.stats, &self.order); + if candidate_cost < (1.0 - TIE_COST_FRACTION) * current_cost { + self.order = candidate; + } + } +} + +/// Evaluate `conjuncts` in `order` against `batch` via the compact-once loop, +/// returning the boolean mask (over the batch's original rows) of rows that +/// passed every conjunct. With `stats`, each conjunct is additionally timed and +/// counted on exactly the rows it evaluated (its marginal selectivity and cost +/// on the current working population). +/// +/// The working batch is physically compacted to the surviving rows only once +/// the accumulated mask becomes selective enough (see +/// [`COMPACTION_SELECTIVITY_THRESHOLD`]); until then masks are combined with a +/// cheap bitwise `AND`, so a run of non-selective conjuncts pays no +/// materialization cost. Unlike a fused `BinaryExpr` chain, survivors stay +/// compacted across the remaining conjuncts instead of being re-evaluated on +/// every row. +fn eval_conjuncts( + conjuncts: &[Arc], + order: &[usize], + batch: &RecordBatch, + mut stats: Option<&mut [ConjunctStats]>, +) -> Result { + let num_rows = batch.num_rows(); + if num_rows == 0 { + return Ok(Arc::new(BooleanArray::from(Vec::::new()))); + } + + // `working` is the batch conjuncts are evaluated against. `acc` is the + // accumulated (`AND`-combined, null-free) result over `working`'s rows since + // the last compaction; `None` means all of them are still live. `live` maps + // `working`'s rows back to original row indices; `None` until a compaction + // first drops rows. + let mut working = batch.clone(); + let mut acc: Option = None; + let mut live: Option = None; + + for &id in order { + let rows_in = working.num_rows(); + + let timer = stats.is_some().then(Instant::now); + let array = conjuncts[id].evaluate(&working)?.into_array(rows_in)?; + let mask = as_boolean_array(&array)?; + // `matched` counts non-null trues (SQL filter semantics). + let matched = mask.true_count() as u64; + + if let (Some(stats), Some(timer)) = (stats.as_deref_mut(), timer) { + let eval_nanos = timer.elapsed().as_nanos() as u64; + stats[id].record(matched, rows_in as u64, eval_nanos); + } + + // An all-true mask leaves the accumulated result untouched. + if matched == rows_in as u64 && mask.null_count() == 0 { + continue; + } + + // Fold this conjunct into the accumulated mask (null -> false). + let mask = if mask.null_count() > 0 { + prep_null_mask_filter(mask) + } else { + mask.clone() + }; + let folded = match &acc { + None => mask, + Some(prev) => and(prev, &mask)?, + }; + + let alive = folded.true_count(); + if alive == 0 { + // Nothing survives; the result is all-false over the original rows. + return Ok(Arc::new(BooleanArray::new( + BooleanBuffer::new_unset(num_rows), + None, + ))); + } + // Compact only when the survivors are a small fraction of the working + // batch — otherwise the copy is not worth it. + if (alive as f64) <= COMPACTION_SELECTIVITY_THRESHOLD * rows_in as f64 { + working = filter_record_batch(&working, &folded)?; + let indices = live.take().unwrap_or_else(|| { + Arc::new(UInt32Array::from_iter_values(0..num_rows as u32)) + }); + live = Some(filter(&indices, &folded)?); + acc = None; + } else { + acc = Some(folded); + } + } + + match live { + // Never compacted: `acc` (or all-true) already covers the original rows. + None => Ok(match acc { + Some(acc) => Arc::new(acc), + None => Arc::new(BooleanArray::new(BooleanBuffer::new_set(num_rows), None)), + }), + // Compacted at least once: scatter the surviving original indices + // (`live`, narrowed by any residual `acc`) into a full-length mask. + Some(indices) => { + let indices = match acc { + Some(acc) => filter(&indices, &acc)?, + None => indices, + }; + let indices = indices + .as_any() + .downcast_ref::() + .expect("u32 live"); + let mut builder = BooleanBufferBuilder::new(num_rows); + builder.append_n(num_rows, false); + for &idx in indices.values() { + builder.set_bit(idx as usize, true); + } + Ok(Arc::new(BooleanArray::new(builder.finish(), None))) + } + } +} + +/// Rank conjunct ids by effectiveness (discards per nanosecond) descending; +/// ids without measurements sort last. Stable, so equal ids keep their order. +fn rank_by_effectiveness(stats: &[ConjunctStats]) -> Vec { + let mut ids: Vec = (0..stats.len()).collect(); + ids.sort_by( + |&a, &b| match (stats[a].effectiveness(), stats[b].effectiveness()) { + (Some(x), Some(y)) => y.partial_cmp(&x).unwrap_or(std::cmp::Ordering::Equal), + (Some(_), None) => std::cmp::Ordering::Less, + (None, Some(_)) => std::cmp::Ordering::Greater, + (None, None) => std::cmp::Ordering::Equal, + }, + ); + ids +} + +/// Expected cost of evaluating the conjuncts in `order`, in nanoseconds per +/// input row: each conjunct's measured per-row cost weighted by the fraction of +/// rows expected to reach it (the product of the pass rates of the conjuncts +/// before it, treated as independent). Unmeasured conjuncts contribute nothing. +fn expected_cost_per_row(stats: &[ConjunctStats], order: &[usize]) -> f64 { + let mut weight = 1.0_f64; + let mut total = 0.0_f64; + for &id in order { + let (Some(cost), Some(pass)) = (stats[id].cost_per_row(), stats[id].pass_rate()) + else { + continue; + }; + total += weight * cost; + weight *= pass; + } + total +} + +#[cfg(test)] +mod tests { + use super::*; + + use arrow::array::Int32Array; + use arrow::datatypes::{DataType, Field, Schema}; + use datafusion_expr::Operator; + use datafusion_physical_expr::expressions::{binary, col, lit}; + + fn schema() -> Arc { + Arc::new(Schema::new(vec![ + Field::new("a", DataType::Int32, false), + Field::new("b", DataType::Int32, false), + ])) + } + + fn batch(schema: &Arc, a: Vec, b: Vec) -> RecordBatch { + RecordBatch::try_new( + Arc::clone(schema), + vec![Arc::new(Int32Array::from(a)), Arc::new(Int32Array::from(b))], + ) + .unwrap() + } + + /// `a > 2 AND b < 5` + fn predicate(schema: &Arc) -> Arc { + let left = + binary(col("a", schema).unwrap(), Operator::Gt, lit(2i32), schema).unwrap(); + let right = + binary(col("b", schema).unwrap(), Operator::Lt, lit(5i32), schema).unwrap(); + binary(left, Operator::And, right, schema).unwrap() + } + + fn conjuncts(schema: &Arc) -> Vec> { + split_conjunction(&predicate(schema)) + .into_iter() + .map(Arc::clone) + .collect() + } + + fn passing_rows(mask: &ArrayRef) -> Vec { + let mask = as_boolean_array(mask).unwrap(); + (0..mask.len()) + .filter(|&i| !mask.is_null(i) && mask.value(i)) + .collect() + } + + fn stats(rows: u64, matched: u64, nanos: u64) -> ConjunctStats { + ConjunctStats { + rows, + matched, + nanos, + } + } + + #[test] + fn disabled_is_not_adaptive() { + let schema = schema(); + assert!(AdaptiveConjunction::try_new(&predicate(&schema), false).is_none()); + } + + #[test] + fn single_conjunct_is_not_adaptive() { + let schema = schema(); + let p = + binary(col("a", &schema).unwrap(), Operator::Gt, lit(2i32), &schema).unwrap(); + assert!(AdaptiveConjunction::try_new(&p, true).is_none()); + } + + #[test] + fn two_conjuncts_are_adaptive() { + let schema = schema(); + let adaptive = AdaptiveConjunction::try_new(&predicate(&schema), true).unwrap(); + assert_eq!(adaptive.conjuncts.len(), 2); + assert_eq!(adaptive.order, vec![0, 1]); + assert_eq!(adaptive.warmup_left, WARMUP_BATCHES); + } + + #[test] + fn ranks_by_discards_per_nanosecond() { + // id 0: cheap (1ns/row), unselective (pass 0.9): eff = 0.1 / 1 = 0.1 + // id 1: expensive (5ns/row), selective (pass 0.01): eff = 0.99 / 5 = 0.198 (first) + // id 2: cheap (1ns/row), very unselective (pass 0.95): eff = 0.05 / 1 = 0.05 (last) + let s = vec![ + stats(1000, 900, 1000), + stats(1000, 10, 5000), + stats(1000, 950, 1000), + ]; + assert_eq!(rank_by_effectiveness(&s), vec![1, 0, 2]); + } + + #[test] + fn unmeasured_conjuncts_sort_last() { + let s = vec![ + stats(0, 0, 0), // unmeasured -> last + stats(1000, 10, 1000), // selective + stats(1000, 900, 1000), // unselective + ]; + assert_eq!(rank_by_effectiveness(&s), vec![1, 2, 0]); + } + + #[test] + fn expected_cost_weights_by_upstream_pass_rate() { + // a: cost 1, pass 0.5 ; b: cost 10, pass 0.5 + let s = vec![stats(1000, 500, 1000), stats(1000, 500, 10_000)]; + // order [0,1]: 1 + 0.5*10 = 6 + assert!((expected_cost_per_row(&s, &[0, 1]) - 6.0).abs() < 1e-9); + // order [1,0]: 10 + 0.5*1 = 10.5 + assert!((expected_cost_per_row(&s, &[1, 0]) - 10.5).abs() < 1e-9); + } + + /// The compact-once loop returns exactly the rows the plain predicate + /// keeps, in any order and whether or not compaction triggers. + #[test] + fn eval_conjuncts_matches_predicate_in_any_order() { + let schema = schema(); + let cs = conjuncts(&schema); + let p = predicate(&schema); + + // A batch where `b < 5` is rare (forces a compaction) and one where it + // is common (no compaction). + for b in [ + (0..100).map(|x| x % 50).collect::>(), // b<5 rare + (0..100).map(|x| x % 3).collect::>(), // b<5 common + ] { + let a: Vec = (0..100).collect(); + let rb = batch(&schema, a, b); + let want = p.evaluate(&rb).unwrap().into_array(rb.num_rows()).unwrap(); + for order in [vec![0, 1], vec![1, 0]] { + let got = eval_conjuncts(&cs, &order, &rb, None).unwrap(); + assert_eq!(passing_rows(&got), passing_rows(&want), "order {order:?}"); + } + } + } + + /// Across the warm-up boundary the mask must always equal the plain + /// predicate's, before and after the order settles. + #[test] + fn evaluate_matches_predicate_across_warmup() { + let schema = schema(); + let p = predicate(&schema); + let mut adaptive = AdaptiveConjunction::try_new(&p, true).unwrap(); + + for round in 0..(WARMUP_BATCHES as i32 + 4) { + let base = round * 10; + let a: Vec = (base..base + 10).collect(); + let b: Vec = (base..base + 10).map(|x| x.rem_euclid(9)).collect(); + let rb = batch(&schema, a, b); + + let got = adaptive.evaluate(&rb).unwrap(); + let want = p.evaluate(&rb).unwrap().into_array(rb.num_rows()).unwrap(); + assert_eq!( + passing_rows(&got), + passing_rows(&want), + "mismatch on round {round}" + ); + } + assert_eq!(adaptive.warmup_left, 0); + } + + /// A reorder is adopted only when materially cheaper; an already-good order + /// is left untouched. + #[test] + fn settle_keeps_order_when_not_materially_better() { + let schema = schema(); + let mut adaptive = + AdaptiveConjunction::try_new(&predicate(&schema), true).unwrap(); + // Two equally cheap, equally selective conjuncts: swapping cannot help. + adaptive.stats = vec![stats(1000, 500, 1000), stats(1000, 500, 1000)]; + adaptive.settle(); + assert_eq!(adaptive.order, vec![0, 1]); + } + + #[test] + fn settle_adopts_materially_cheaper_order() { + let schema = schema(); + let mut adaptive = + AdaptiveConjunction::try_new(&predicate(&schema), true).unwrap(); + // id 1 is far more selective and equally cheap: it should move first. + adaptive.stats = vec![stats(1000, 900, 1000), stats(1000, 10, 1000)]; + adaptive.settle(); + assert_eq!(adaptive.order, vec![1, 0]); + } +} diff --git a/datafusion/physical-plan/src/filter.rs b/datafusion/physical-plan/src/filter.rs index d23dd380423d1..9a254142d3388 100644 --- a/datafusion/physical-plan/src/filter.rs +++ b/datafusion/physical-plan/src/filter.rs @@ -28,6 +28,7 @@ use super::{ ColumnStatistics, DisplayAs, ExecutionPlanProperties, PlanProperties, RecordBatchStream, SendableRecordBatchStream, Statistics, }; +use crate::adaptive_filter::AdaptiveConjunction; use crate::check_if_same_properties; use crate::coalesce::{LimitedBatchCoalescer, PushBatchStatus}; use crate::common::can_project; @@ -571,9 +572,18 @@ impl ExecutionPlan for FilterExec { context.task_id() ); let metrics = FilterExecMetrics::new(&self.metrics, partition); + let adaptive = AdaptiveConjunction::try_new( + &self.predicate, + context + .session_config() + .options() + .execution + .adaptive_filter_reordering, + ); Ok(Box::pin(FilterExecStream { schema: self.schema(), predicate: Arc::clone(&self.predicate), + adaptive, input: self.input.execute(partition, context)?, metrics, projection: self.projection.clone(), @@ -1051,6 +1061,10 @@ struct FilterExecStream { schema: SchemaRef, /// The expression to filter on. This expression must evaluate to a boolean value. predicate: Arc, + /// When set, the predicate is a reorderable conjunction evaluated + /// adaptively (conjuncts measured, then reordered) instead of via + /// `predicate`. + adaptive: Option, /// The input partition to filter. input: SendableRecordBatchStream, /// Runtime metrics recording @@ -1146,9 +1160,15 @@ impl Stream for FilterExecStream { } Some(Ok(batch)) => { let timer = elapsed_compute.timer(); - let status = self.predicate.as_ref() - .evaluate(&batch) - .and_then(|v| v.into_array(batch.num_rows())) + let array = match self.adaptive.as_mut() { + Some(adaptive) => adaptive.evaluate(&batch), + None => self + .predicate + .as_ref() + .evaluate(&batch) + .and_then(|v| v.into_array(batch.num_rows())), + }; + let status = array .and_then(|array| { Ok(match self.projection.as_ref() { Some(projection) => { diff --git a/datafusion/physical-plan/src/lib.rs b/datafusion/physical-plan/src/lib.rs index 6cc6e44c32cc3..66ee0454ff2ff 100644 --- a/datafusion/physical-plan/src/lib.rs +++ b/datafusion/physical-plan/src/lib.rs @@ -61,6 +61,7 @@ mod render_tree; mod topk; mod visitor; +mod adaptive_filter; pub mod aggregates; pub mod analyze; pub mod async_func; diff --git a/datafusion/sqllogictest/test_files/adaptive_filter.slt b/datafusion/sqllogictest/test_files/adaptive_filter.slt new file mode 100644 index 0000000000000..b0e2d958a3be9 --- /dev/null +++ b/datafusion/sqllogictest/test_files/adaptive_filter.slt @@ -0,0 +1,80 @@ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, +# software distributed under the License is distributed on an +# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +# KIND, either express or implied. See the License for the +# specific language governing permissions and limitations +# under the License. + +# Tests for execution.adaptive_filter_reordering. Runtime reordering of a +# conjunction must never change query results, only evaluation order. + +statement ok +CREATE TABLE t AS +SELECT + i AS a, + i % 7 AS b, + arrow_cast(i, 'Utf8') AS s +FROM generate_series(1, 1000) AS tbl(i); + +# Baseline (flag off): multi-conjunct filter mixing a cheap comparison with an +# expensive LIKE. +query I +SELECT count(*) FROM t WHERE b = 3 AND s LIKE '1%'; +---- +17 + +statement ok +SET datafusion.execution.adaptive_filter_reordering = true; + +# Same query with adaptive reordering on must return the same result. +query I +SELECT count(*) FROM t WHERE b = 3 AND s LIKE '1%'; +---- +17 + +# A three-conjunct predicate, including an expensive-but-selective LIKE. +query I +SELECT count(*) FROM t WHERE a > 100 AND s LIKE '5%' AND b <> 0; +---- +86 + +# Full row materialization (not just count) is unchanged by reordering. +query IIT +SELECT a, b, s FROM t WHERE b = 1 AND a > 990 AND s LIKE '99%' ORDER BY a; +---- +995 1 995 + +# EXPLAIN: runtime reordering is invisible to the plan — the FilterExec +# predicate is unchanged whether the flag is on or off. +query TT +EXPLAIN SELECT count(*) FROM t WHERE b = 3 AND s LIKE '1%'; +---- +logical_plan +01)Projection: count(Int64(1)) AS count(*) +02)--Aggregate: groupBy=[[]], aggr=[[count(Int64(1))]] +03)----Projection: +04)------Filter: t.b = Int64(3) AND t.s LIKE Utf8("1%") +05)--------TableScan: t projection=[b, s] +physical_plan +01)ProjectionExec: expr=[count(Int64(1))@0 as count(*)] +02)--AggregateExec: mode=Final, gby=[], aggr=[count(Int64(1))] +03)----CoalescePartitionsExec +04)------AggregateExec: mode=Partial, gby=[], aggr=[count(Int64(1))] +05)--------FilterExec: b@0 = 3 AND s@1 LIKE 1%, projection=[] +06)----------DataSourceExec: partitions=4, partition_sizes=[1, 0, 0, 0] + +statement ok +SET datafusion.execution.adaptive_filter_reordering = false; + +statement ok +DROP TABLE t; diff --git a/datafusion/sqllogictest/test_files/information_schema.slt b/datafusion/sqllogictest/test_files/information_schema.slt index 0932f58a7c03f..2e6b567c77905 100644 --- a/datafusion/sqllogictest/test_files/information_schema.slt +++ b/datafusion/sqllogictest/test_files/information_schema.slt @@ -214,6 +214,7 @@ datafusion.catalog.has_header true datafusion.catalog.information_schema true datafusion.catalog.location NULL datafusion.catalog.newlines_in_values false +datafusion.execution.adaptive_filter_reordering false datafusion.execution.batch_size 8192 datafusion.execution.coalesce_batches true datafusion.execution.collect_statistics true @@ -371,6 +372,7 @@ datafusion.catalog.has_header true Default value for `format.has_header` for `CR datafusion.catalog.information_schema true Should DataFusion provide access to `information_schema` virtual tables for displaying schema information datafusion.catalog.location NULL Location scanned to load tables for `default` schema datafusion.catalog.newlines_in_values false Specifies whether newlines in (quoted) CSV values are supported. This is the default value for `format.newlines_in_values` for `CREATE EXTERNAL TABLE` if not specified explicitly in the statement. Parsing newlines in quoted values may be affected by execution behaviour such as parallel file scanning. Setting this to `true` ensures that newlines in values are parsed successfully, which may reduce performance. +datafusion.execution.adaptive_filter_reordering false (experimental) When enabled, `FilterExec` adaptively reorders the conjuncts of a conjunctive predicate at runtime. It measures each conjunct's selectivity and evaluation cost on the rows that reach it and runs the conjuncts that discard the most rows per unit of CPU time first, so cheap-and-selective predicates gate expensive ones. Reordering never changes query results (only the evaluation order of a conjunction) but can change observable side effects of fallible predicates, so it is off by default. Predicates containing volatile expressions are never reordered. datafusion.execution.batch_size 8192 Default batch size while creating new batches, it's especially useful for buffer-in-memory batches since creating tiny batches would result in too much metadata memory consumption datafusion.execution.coalesce_batches true When set to true, record batches will be examined between each operator and small batches will be coalesced into larger batches. This is helpful when there are highly selective filters or joins that could produce tiny output batches. The target batch size is determined by the configuration setting datafusion.execution.collect_statistics true Should DataFusion collect statistics when first creating a table. Has no effect after the table is created. Defaults to true. diff --git a/docs/source/user-guide/configs.md b/docs/source/user-guide/configs.md index f70daef317216..f8f8d31964370 100644 --- a/docs/source/user-guide/configs.md +++ b/docs/source/user-guide/configs.md @@ -139,6 +139,7 @@ The following configuration settings are available: | datafusion.execution.skip_partial_aggregation_probe_rows_threshold | 100000 | Number of input rows partial aggregation partition should process, before aggregation ratio check and trying to switch to skipping aggregation mode | | datafusion.execution.use_row_number_estimates_to_optimize_partitioning | false | Should DataFusion use row number estimates at the input to decide whether increasing parallelism is beneficial or not. By default, only exact row numbers (not estimates) are used for this decision. Setting this flag to `true` will likely produce better plans. if the source of statistics is accurate. We plan to make this the default in the future. | | datafusion.execution.enforce_batch_size_in_joins | false | Should DataFusion enforce batch size in joins or not. By default, DataFusion will not enforce batch size in joins. Enforcing batch size in joins can reduce memory usage when joining large tables with a highly-selective join filter, but is also slightly slower. | +| datafusion.execution.adaptive_filter_reordering | false | (experimental) When enabled, `FilterExec` adaptively reorders the conjuncts of a conjunctive predicate at runtime. It measures each conjunct's selectivity and evaluation cost on the rows that reach it and runs the conjuncts that discard the most rows per unit of CPU time first, so cheap-and-selective predicates gate expensive ones. Reordering never changes query results (only the evaluation order of a conjunction) but can change observable side effects of fallible predicates, so it is off by default. Predicates containing volatile expressions are never reordered. | | datafusion.execution.objectstore_writer_buffer_size | 10485760 | Size (bytes) of data buffer DataFusion uses when writing output files. This affects the size of the data chunks that are uploaded to remote object stores (e.g. AWS S3). If very large (>= 100 GiB) output files are being written, it may be necessary to increase this size to avoid errors from the remote end point. | | datafusion.execution.enable_ansi_mode | false | Whether to enable ANSI SQL mode. The flag is experimental and relevant only for DataFusion Spark built-in functions When `enable_ansi_mode` is set to `true`, the query engine follows ANSI SQL semantics for expressions, casting, and error handling. This means: - **Strict type coercion rules:** implicit casts between incompatible types are disallowed. - **Standard SQL arithmetic behavior:** operations such as division by zero, numeric overflow, or invalid casts raise runtime errors rather than returning `NULL` or adjusted values. - **Consistent ANSI behavior** for string concatenation, comparisons, and `NULL` handling. When `enable_ansi_mode` is `false` (the default), the engine uses a more permissive, non-ANSI mode designed for user convenience and backward compatibility. In this mode: - Implicit casts between types are allowed (e.g., string to integer when possible). - Arithmetic operations are more lenient — for example, `abs()` on the minimum representable integer value returns the input value instead of raising overflow. - Division by zero or invalid casts may return `NULL` instead of failing. # Default `false` — ANSI SQL mode is disabled by default. | | datafusion.execution.hash_join_buffering_capacity | 0 | How many bytes to buffer in the probe side of hash joins while the build side is concurrently being built. Without this, hash joins will wait until the full materialization of the build side before polling the probe side. This is useful in scenarios where the query is not completely CPU bounded, allowing to do some early work concurrently and reducing the latency of the query. Note that when hash join buffering is enabled, the probe side will start eagerly polling data, not giving time for the producer side of dynamic filters to produce any meaningful predicate. Queries with dynamic filters might see performance degradation. Disabled by default, set to a number greater than 0 for enabling it. | From 9d962a0ffb52a5c080e10a348360b5b8126e942e Mon Sep 17 00:00:00 2001 From: Adrian Garcia Badaracco <1755071+adriangb@users.noreply.github.com> Date: Mon, 29 Jun 2026 00:32:00 +0200 Subject: [PATCH 2/3] feat(physical-plan): pool adaptive measurements across partition streams MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A `FilterExec` is split across many partition streams, each seeing a slice of the data. With per-stream warm-up, every stream pays its own measurement cost, and when each stream is only a handful of batches long that warm-up is most of its work — so the reordering win never materialises (and the warm-up overhead shows up as a regression). Benchmarked: at 12 partitions the costsel_q01 win collapsed from -67% (single stream) to -14%. Share the measurements. `AdaptiveFilterShared` holds a per-conjunct stats pool plus a settled-order epoch, common to every stream of one `FilterExec`. Each stream measures a batch into a local accumulator and folds it into the pool; the first stream to reach `WARMUP_BATCHES` pooled batches decides the order and publishes it by bumping the epoch. Other streams poll the epoch with one relaxed atomic load per batch and adopt the published order without paying warm-up. The warm-up is thus paid ~once per query, not once per stream. Restores the recovered win at default partitioning: width -56..-66%, costsel_q01 -58%, cardinality k16 -26% (was -12%, -14%, -6% without sharing), matching or beating the full design. Steady-state regressions on conjunctions where compaction does not pay (neutral_q61 ~+11%, cardinality k4 ~+5%, costsel_q02/q03 ~+3-5%) remain — a Fused-vs-CompactOnce guard addresses those in the next commit. Co-Authored-By: Claude Opus 4.8 (1M context) Claude-Session: https://claude.ai/code/session_01Lh7i9DyeFWuTFWjogrVNkb --- .../physical-plan/src/adaptive_filter.rs | 233 ++++++++++++++---- datafusion/physical-plan/src/filter.rs | 13 +- 2 files changed, 197 insertions(+), 49 deletions(-) diff --git a/datafusion/physical-plan/src/adaptive_filter.rs b/datafusion/physical-plan/src/adaptive_filter.rs index 46e76d6ddc591..b444b1a7c1704 100644 --- a/datafusion/physical-plan/src/adaptive_filter.rs +++ b/datafusion/physical-plan/src/adaptive_filter.rs @@ -56,6 +56,8 @@ //! top of it. use std::sync::Arc; +use std::sync::Mutex; +use std::sync::atomic::{AtomicU64, Ordering}; use arrow::array::{Array, ArrayRef, BooleanArray, BooleanBufferBuilder, UInt32Array}; use arrow::buffer::BooleanBuffer; @@ -104,6 +106,14 @@ impl ConjunctStats { self.nanos += nanos; } + /// Fold another accumulator's counts into this one (they are plain sums, so + /// merging is addition). Used to pool measurements across partition streams. + fn merge(&mut self, other: &Self) { + self.rows += other.rows; + self.matched += other.matched; + self.nanos += other.nanos; + } + /// Fraction of rows that pass, or `None` if never evaluated on any row. fn pass_rate(&self) -> Option { (self.rows > 0).then(|| self.matched as f64 / self.rows as f64) @@ -127,19 +137,61 @@ impl ConjunctStats { } } +/// State shared by every partition stream of one `FilterExec`, so the streams +/// learn as one: per-conjunct measurements are pooled across streams and the +/// first stream to accumulate enough samples settles the order for all of them. +/// +/// This matters because a `FilterExec` is split across many partition streams, +/// each seeing only a slice of the data. Without sharing, every stream pays its +/// own warm-up — and when each stream is only a handful of batches long, that +/// warm-up is most of its work, so the reordering win never materialises. With +/// sharing the warm-up is paid roughly once per query, not once per stream. +#[derive(Debug, Default)] +pub(crate) struct AdaptiveFilterShared { + /// `0` until an order is published; bumped once when the first stream + /// settles. Streams poll it with one relaxed atomic load per batch. + epoch: AtomicU64, + inner: Mutex, +} + +#[derive(Debug, Default)] +struct SharedInner { + /// Per-conjunct counts pooled across all streams (indexed by conjunct + /// position). Empty until the first measured batch sizes it. + stats: Vec, + /// Measured batches contributed by all streams so far. + measured_batches: u64, + /// The settled evaluation order, once decided; `None` while learning. + settled: Option>, +} + +impl AdaptiveFilterShared { + pub(crate) fn new() -> Self { + Self::default() + } + + /// The published settled order, or `None` if the streams are still learning. + fn settled_order(&self) -> Option> { + self.inner.lock().expect("poisoned").settled.clone() + } +} + /// Adaptive evaluator for a single conjunctive predicate, owned per partition -/// stream (single-threaded, no locking). +/// stream. Measurements are pooled into the shared [`AdaptiveFilterShared`]; +/// the per-stream state is just the current order and how far it has caught up. #[derive(Debug)] pub(crate) struct AdaptiveConjunction { - /// The split conjuncts. `stats`/`order` indices refer to positions here. + /// The split conjuncts. `order` indices refer to positions here. conjuncts: Vec>, - /// Per-conjunct measurements, indexed by conjunct position. - stats: Vec, - /// Evaluation order: indices into `conjuncts`. Starts as the written order - /// and may be permuted once after the warm-up. + /// Measurements and the settled order, shared by every partition stream. + shared: Arc, + /// Evaluation order: indices into `conjuncts`. The written order until a + /// settled order is adopted. order: Vec, - /// Warm-up batches still to measure; `0` means the order has settled. - warmup_left: u64, + /// Shared epoch this stream has caught up to. + epoch_seen: u64, + /// Whether the order is settled (frozen): this stream no longer measures. + settled: bool, } impl AdaptiveConjunction { @@ -149,9 +201,13 @@ impl AdaptiveConjunction { /// - `enabled` is false (the config flag is off); /// - the predicate has fewer than two `AND` conjuncts (nothing to reorder); /// - any conjunct is volatile (reordering could change side effects). + /// + /// `shared` is the state common to all partition streams of the owning + /// `FilterExec`. pub(crate) fn try_new( predicate: &Arc, enabled: bool, + shared: Arc, ) -> Option { if !enabled { return None; @@ -163,47 +219,83 @@ impl AdaptiveConjunction { if conjuncts.len() < 2 || conjuncts.iter().any(is_volatile) { return None; } - let stats = vec![ConjunctStats::default(); conjuncts.len()]; let order = (0..conjuncts.len()).collect(); Some(Self { conjuncts, - stats, + shared, order, - warmup_left: WARMUP_BATCHES, + epoch_seen: 0, + settled: false, }) } /// Evaluate the conjunction against `batch`, returning the boolean mask /// (over the batch's rows) of rows that passed every conjunct. /// - /// During the warm-up the conjuncts are measured (on the rows that reach - /// each one); once the window ends the order is settled and subsequent - /// batches evaluate with no instrumentation. + /// Until the order settles, each batch is measured and its counts pooled + /// into the shared registry; a stream adopts the settled order another + /// stream published as soon as it sees the epoch advance. pub(crate) fn evaluate(&mut self, batch: &RecordBatch) -> Result { - if self.warmup_left == 0 { + // Adopt a settled order another stream published since we last looked: + // one relaxed atomic load per batch, a lock only on the transition. + if !self.settled { + let epoch = self.shared.epoch.load(Ordering::Acquire); + if epoch != self.epoch_seen { + self.epoch_seen = epoch; + if let Some(order) = self.shared.settled_order() { + self.order = order; + self.settled = true; + } + } + } + if self.settled { return eval_conjuncts(&self.conjuncts, &self.order, batch, None); } + + // Measure this batch into a local accumulator, then pool it. + let mut local = vec![ConjunctStats::default(); self.conjuncts.len()]; let result = - eval_conjuncts(&self.conjuncts, &self.order, batch, Some(&mut self.stats))?; - self.warmup_left -= 1; - if self.warmup_left == 0 { - self.settle(); - } + eval_conjuncts(&self.conjuncts, &self.order, batch, Some(&mut local))?; + self.pool_and_maybe_settle(&local); Ok(result) } - /// Rank the conjuncts by measured effectiveness and, if the resulting order - /// is materially cheaper than the written order, adopt it. - fn settle(&mut self) { - let candidate = rank_by_effectiveness(&self.stats); - if candidate == self.order { - return; + /// Merge this batch's measurements into the shared pool and, once enough + /// batches have accrued across all streams, decide and publish the order. + fn pool_and_maybe_settle(&mut self, local: &[ConjunctStats]) { + let mut inner = self.shared.inner.lock().expect("poisoned"); + if inner.stats.len() != local.len() { + inner.stats = vec![ConjunctStats::default(); local.len()]; } - let candidate_cost = expected_cost_per_row(&self.stats, &candidate); - let current_cost = expected_cost_per_row(&self.stats, &self.order); - if candidate_cost < (1.0 - TIE_COST_FRACTION) * current_cost { - self.order = candidate; + for (s, l) in inner.stats.iter_mut().zip(local) { + s.merge(l); + } + inner.measured_batches += 1; + if inner.settled.is_some() || inner.measured_batches < WARMUP_BATCHES { + return; } + let order = settle_order(&inner.stats); + inner.settled = Some(order.clone()); + drop(inner); + self.order = order; + self.settled = true; + self.shared.epoch.fetch_add(1, Ordering::Release); + } +} + +/// Rank the conjuncts by effectiveness and adopt the ranking only if it is +/// materially cheaper than the written order; otherwise keep the written order +/// (so interchangeable conjuncts are never reshuffled). +fn settle_order(stats: &[ConjunctStats]) -> Vec { + let identity: Vec = (0..stats.len()).collect(); + let candidate = rank_by_effectiveness(stats); + if candidate != identity + && expected_cost_per_row(stats, &candidate) + < (1.0 - TIE_COST_FRACTION) * expected_cost_per_row(stats, &identity) + { + candidate + } else { + identity } } @@ -407,10 +499,22 @@ mod tests { } } + /// `try_new` with a fresh, unshared registry. + fn try_new( + predicate: &Arc, + enabled: bool, + ) -> Option { + AdaptiveConjunction::try_new( + predicate, + enabled, + Arc::new(AdaptiveFilterShared::new()), + ) + } + #[test] fn disabled_is_not_adaptive() { let schema = schema(); - assert!(AdaptiveConjunction::try_new(&predicate(&schema), false).is_none()); + assert!(try_new(&predicate(&schema), false).is_none()); } #[test] @@ -418,16 +522,16 @@ mod tests { let schema = schema(); let p = binary(col("a", &schema).unwrap(), Operator::Gt, lit(2i32), &schema).unwrap(); - assert!(AdaptiveConjunction::try_new(&p, true).is_none()); + assert!(try_new(&p, true).is_none()); } #[test] fn two_conjuncts_are_adaptive() { let schema = schema(); - let adaptive = AdaptiveConjunction::try_new(&predicate(&schema), true).unwrap(); + let adaptive = try_new(&predicate(&schema), true).unwrap(); assert_eq!(adaptive.conjuncts.len(), 2); assert_eq!(adaptive.order, vec![0, 1]); - assert_eq!(adaptive.warmup_left, WARMUP_BATCHES); + assert!(!adaptive.settled); } #[test] @@ -493,7 +597,7 @@ mod tests { fn evaluate_matches_predicate_across_warmup() { let schema = schema(); let p = predicate(&schema); - let mut adaptive = AdaptiveConjunction::try_new(&p, true).unwrap(); + let mut adaptive = try_new(&p, true).unwrap(); for round in 0..(WARMUP_BATCHES as i32 + 4) { let base = round * 10; @@ -509,30 +613,63 @@ mod tests { "mismatch on round {round}" ); } - assert_eq!(adaptive.warmup_left, 0); + assert!(adaptive.settled); } /// A reorder is adopted only when materially cheaper; an already-good order /// is left untouched. #[test] fn settle_keeps_order_when_not_materially_better() { - let schema = schema(); - let mut adaptive = - AdaptiveConjunction::try_new(&predicate(&schema), true).unwrap(); // Two equally cheap, equally selective conjuncts: swapping cannot help. - adaptive.stats = vec![stats(1000, 500, 1000), stats(1000, 500, 1000)]; - adaptive.settle(); - assert_eq!(adaptive.order, vec![0, 1]); + let s = vec![stats(1000, 500, 1000), stats(1000, 500, 1000)]; + assert_eq!(settle_order(&s), vec![0, 1]); } #[test] fn settle_adopts_materially_cheaper_order() { - let schema = schema(); - let mut adaptive = - AdaptiveConjunction::try_new(&predicate(&schema), true).unwrap(); // id 1 is far more selective and equally cheap: it should move first. - adaptive.stats = vec![stats(1000, 900, 1000), stats(1000, 10, 1000)]; - adaptive.settle(); - assert_eq!(adaptive.order, vec![1, 0]); + let s = vec![stats(1000, 900, 1000), stats(1000, 10, 1000)]; + assert_eq!(settle_order(&s), vec![1, 0]); + } + + /// Two streams sharing one registry settle the order together: the pooled + /// warm-up is `WARMUP_BATCHES` total across both streams, and once one + /// stream publishes the order the other adopts it on its next batch. + #[test] + fn streams_pool_measurements_and_share_settled_order() { + let schema = schema(); + let p = predicate(&schema); + let shared = Arc::new(AdaptiveFilterShared::new()); + let mut s1 = AdaptiveConjunction::try_new(&p, true, Arc::clone(&shared)).unwrap(); + let mut s2 = AdaptiveConjunction::try_new(&p, true, Arc::clone(&shared)).unwrap(); + + // `b < 5` (conjunct 1) is the selective one; drive both streams with + // batches where it keeps ~1 row in 25. + let mk = |round: i32| { + let base = round * 100; + let a: Vec = (base..base + 100).collect(); + let b: Vec = (base..base + 100).map(|x| x.rem_euclid(25)).collect(); + batch(&schema, a, b) + }; + + // Alternate the two streams for `WARMUP_BATCHES` pooled batches; the + // order settles partway through and both streams must end settled. + for round in 0..(WARMUP_BATCHES as i32) { + let rb = mk(round); + for s in [&mut s1, &mut s2] { + let got = s.evaluate(&rb).unwrap(); + let want = p.evaluate(&rb).unwrap().into_array(rb.num_rows()).unwrap(); + assert_eq!(passing_rows(&got), passing_rows(&want)); + } + } + + assert!(shared.settled_order().is_some()); + // One more batch each lets a not-yet-settled stream adopt the epoch. + s1.evaluate(&mk(99)).unwrap(); + s2.evaluate(&mk(99)).unwrap(); + assert!(s1.settled && s2.settled); + // The selective conjunct was promoted to the front for both. + assert_eq!(s1.order, vec![1, 0]); + assert_eq!(s2.order, vec![1, 0]); } } diff --git a/datafusion/physical-plan/src/filter.rs b/datafusion/physical-plan/src/filter.rs index 9a254142d3388..5c2d5a6b04a64 100644 --- a/datafusion/physical-plan/src/filter.rs +++ b/datafusion/physical-plan/src/filter.rs @@ -28,7 +28,7 @@ use super::{ ColumnStatistics, DisplayAs, ExecutionPlanProperties, PlanProperties, RecordBatchStream, SendableRecordBatchStream, Statistics, }; -use crate::adaptive_filter::AdaptiveConjunction; +use crate::adaptive_filter::{AdaptiveConjunction, AdaptiveFilterShared}; use crate::check_if_same_properties; use crate::coalesce::{LimitedBatchCoalescer, PushBatchStatus}; use crate::common::can_project; @@ -99,6 +99,10 @@ pub struct FilterExec { batch_size: usize, /// Number of rows to fetch fetch: Option, + /// Measurements shared by all partition streams, used by adaptive conjunct + /// reordering (see [`AdaptiveConjunction`]) so the streams learn as one. + /// Fresh per plan node; never affects the plan. + adaptive_stats: Arc, } /// Builder for [`FilterExec`] to set optional parameters @@ -218,6 +222,7 @@ impl FilterExecBuilder { projection: self.projection, batch_size: self.batch_size, fetch: self.fetch, + adaptive_stats: Arc::new(AdaptiveFilterShared::new()), }) } } @@ -291,6 +296,7 @@ impl FilterExec { projection: self.projection.clone(), batch_size, fetch: self.fetch, + adaptive_stats: Arc::clone(&self.adaptive_stats), }) } @@ -579,6 +585,7 @@ impl ExecutionPlan for FilterExec { .options() .execution .adaptive_filter_reordering, + Arc::clone(&self.adaptive_stats), ); Ok(Box::pin(FilterExecStream { schema: self.schema(), @@ -774,6 +781,9 @@ impl ExecutionPlan for FilterExec { projection: self.projection.clone(), batch_size: self.batch_size, fetch: self.fetch, + // The predicate changed; pooled per-conjunct stats no longer + // describe it. + adaptive_stats: Arc::new(AdaptiveFilterShared::new()), }; Some(Arc::new(new) as _) }; @@ -798,6 +808,7 @@ impl ExecutionPlan for FilterExec { projection: self.projection.clone(), batch_size: self.batch_size, fetch, + adaptive_stats: Arc::clone(&self.adaptive_stats), })) } From e3182a090fd22953c5ad5b9d3a447599d5ce7946 Mon Sep 17 00:00:00 2001 From: Adrian Garcia Badaracco <1755071+adriangb@users.noreply.github.com> Date: Mon, 29 Jun 2026 01:41:00 +0200 Subject: [PATCH 3/3] feat(physical-plan): use compact-once only in service of a reorder MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The compact-once loop wins by gating expensive conjuncts behind a selective one, but its per-conjunct bookkeeping (mask AND, true_count, the compaction copy) is pure overhead when there is nothing to gate. On a conjunction of interchangeable predicates — several equally expensive, equally unselective regexps, say — the warm-up settles on the written order (nothing to reorder) yet still paid compact-once on every batch, regressing ~11% vs the plain predicate (predicate_eval neutral_q61). Guard it: compact-once is adopted only when the warm-up actually reorders the conjuncts. When the settled order equals the written order, evaluate the predicate as-is — byte-for-byte the flag-off path, zero overhead. Since every real win reorders (a selective conjunct moves toward the front), this keeps the full win while removing the no-reorder regression. predicate_eval (vs flag off): neutral_q61 +11% -> ~0; wins preserved (costsel_q01 -60%, width -58..-66%, cardinality k16 -31%). A small residual remains on low-cardinality cheap conjunctions that do reorder (k4/k8 ~+3-4%), where compaction's cost is not repaid by gating so few/cheap predicates; the full champion/challenger arbiter regresses these more (~+10%), so a heavier guard is not worth it here. Co-Authored-By: Claude Opus 4.8 (1M context) Claude-Session: https://claude.ai/code/session_01Lh7i9DyeFWuTFWjogrVNkb --- .../physical-plan/src/adaptive_filter.rs | 160 ++++++++++++++---- 1 file changed, 131 insertions(+), 29 deletions(-) diff --git a/datafusion/physical-plan/src/adaptive_filter.rs b/datafusion/physical-plan/src/adaptive_filter.rs index b444b1a7c1704..ca4a683d4e9fa 100644 --- a/datafusion/physical-plan/src/adaptive_filter.rs +++ b/datafusion/physical-plan/src/adaptive_filter.rs @@ -45,15 +45,31 @@ //! independent conjuncts), and if the ranked order is *materially* cheaper than //! the written one it is adopted. The order then stays fixed. //! +//! Compact-once is used **only in service of a reorder**: if the warm-up does +//! not reorder the conjuncts (e.g. they are interchangeable), the written +//! predicate is evaluated as-is, so a conjunction that does not benefit from +//! reordering pays no compact-once overhead and behaves exactly as it would +//! with the feature off. +//! +//! ## How it shares +//! +//! A `FilterExec` is split across many partition streams, each seeing only a +//! slice of the data. Measurements are pooled into a shared +//! [`AdaptiveFilterShared`] so the streams learn as one: the first stream to +//! accumulate enough samples settles the order and publishes it, and the others +//! adopt it (one relaxed atomic load per batch) without each re-paying the +//! warm-up — which is what makes the win materialise when each stream is only a +//! handful of batches long. +//! //! It is **off by default** //! (`datafusion.execution.adaptive_filter_reordering`) and never changes query //! results: a conjunction's value is independent of evaluation order. Predicates //! containing volatile expressions are never reordered (their observable side //! effects depend on order). //! -//! This is the minimal core of the adaptive evaluator. Richer policies -//! (A/B-validated adoption, cross-stream sharing, drift re-measurement) build on -//! top of it. +//! This is the core of the adaptive evaluator. Further policies (drift +//! re-measurement, confidence-interval statistics, A/B-validated adoption for +//! the cases a cost model cannot separate) can build on top of it. use std::sync::Arc; use std::sync::Mutex; @@ -161,8 +177,19 @@ struct SharedInner { stats: Vec, /// Measured batches contributed by all streams so far. measured_batches: u64, - /// The settled evaluation order, once decided; `None` while learning. - settled: Option>, + /// The settled decision, once made; `None` while learning. + settled: Option, +} + +/// The settled outcome of the warm-up: the evaluation order, and whether to run +/// it through the compact-once loop or the plain predicate. +#[derive(Debug, Clone)] +struct Settled { + /// Evaluation order: indices into the conjunct list. + order: Vec, + /// `true` to run `order` through the compact-once loop; `false` to evaluate + /// the written predicate as-is (see [`settle`]). + compact: bool, } impl AdaptiveFilterShared { @@ -170,8 +197,8 @@ impl AdaptiveFilterShared { Self::default() } - /// The published settled order, or `None` if the streams are still learning. - fn settled_order(&self) -> Option> { + /// The published settled decision, or `None` if streams are still learning. + fn settled(&self) -> Option { self.inner.lock().expect("poisoned").settled.clone() } } @@ -183,11 +210,18 @@ impl AdaptiveFilterShared { pub(crate) struct AdaptiveConjunction { /// The split conjuncts. `order` indices refer to positions here. conjuncts: Vec>, - /// Measurements and the settled order, shared by every partition stream. + /// The written predicate, evaluated as-is when the settled order does not + /// reorder it (so a settled non-reorder costs exactly what the flag-off path + /// costs — no compact-once overhead). + predicate: Arc, + /// Measurements and the settled decision, shared by every partition stream. shared: Arc, /// Evaluation order: indices into `conjuncts`. The written order until a /// settled order is adopted. order: Vec, + /// Whether the settled order runs through the compact-once loop; `false` + /// means evaluate [`predicate`](Self::predicate) directly. + compact: bool, /// Shared epoch this stream has caught up to. epoch_seen: u64, /// Whether the order is settled (frozen): this stream no longer measures. @@ -222,8 +256,10 @@ impl AdaptiveConjunction { let order = (0..conjuncts.len()).collect(); Some(Self { conjuncts, + predicate: Arc::clone(predicate), shared, order, + compact: false, epoch_seen: 0, settled: false, }) @@ -242,14 +278,13 @@ impl AdaptiveConjunction { let epoch = self.shared.epoch.load(Ordering::Acquire); if epoch != self.epoch_seen { self.epoch_seen = epoch; - if let Some(order) = self.shared.settled_order() { - self.order = order; - self.settled = true; + if let Some(decision) = self.shared.settled() { + self.adopt(decision); } } } if self.settled { - return eval_conjuncts(&self.conjuncts, &self.order, batch, None); + return self.evaluate_settled(batch); } // Measure this batch into a local accumulator, then pool it. @@ -260,6 +295,23 @@ impl AdaptiveConjunction { Ok(result) } + /// Evaluate the settled arrangement with no instrumentation: the + /// compact-once loop when the order was reordered, or the written predicate + /// directly otherwise (identical to the feature being off). + fn evaluate_settled(&self, batch: &RecordBatch) -> Result { + if self.compact { + eval_conjuncts(&self.conjuncts, &self.order, batch, None) + } else { + self.predicate.evaluate(batch)?.into_array(batch.num_rows()) + } + } + + fn adopt(&mut self, decision: Settled) { + self.order = decision.order; + self.compact = decision.compact; + self.settled = true; + } + /// Merge this batch's measurements into the shared pool and, once enough /// batches have accrued across all streams, decide and publish the order. fn pool_and_maybe_settle(&mut self, local: &[ConjunctStats]) { @@ -274,28 +326,40 @@ impl AdaptiveConjunction { if inner.settled.is_some() || inner.measured_batches < WARMUP_BATCHES { return; } - let order = settle_order(&inner.stats); - inner.settled = Some(order.clone()); + let decision = settle(&inner.stats); + inner.settled = Some(decision.clone()); drop(inner); - self.order = order; - self.settled = true; + self.adopt(decision); self.shared.epoch.fetch_add(1, Ordering::Release); } } +/// Decide the settled arrangement from the pooled measurements. +/// /// Rank the conjuncts by effectiveness and adopt the ranking only if it is -/// materially cheaper than the written order; otherwise keep the written order -/// (so interchangeable conjuncts are never reshuffled). -fn settle_order(stats: &[ConjunctStats]) -> Vec { +/// materially cheaper than the written order. A genuine reorder runs through the +/// compact-once loop (the source of the win); otherwise the written predicate is +/// kept and evaluated as-is. This is the guard: compact-once is only ever used +/// in service of a reorder, so a conjunction that does not benefit from +/// reordering (interchangeable conjuncts, e.g. several equally expensive +/// unselective predicates) pays no compact-once overhead and behaves exactly as +/// it would with the feature off. +fn settle(stats: &[ConjunctStats]) -> Settled { let identity: Vec = (0..stats.len()).collect(); let candidate = rank_by_effectiveness(stats); if candidate != identity && expected_cost_per_row(stats, &candidate) < (1.0 - TIE_COST_FRACTION) * expected_cost_per_row(stats, &identity) { - candidate + Settled { + order: candidate, + compact: true, + } } else { - identity + Settled { + order: identity, + compact: false, + } } } @@ -617,19 +681,55 @@ mod tests { } /// A reorder is adopted only when materially cheaper; an already-good order - /// is left untouched. + /// is left untouched and runs the plain predicate (no compact-once). #[test] fn settle_keeps_order_when_not_materially_better() { - // Two equally cheap, equally selective conjuncts: swapping cannot help. + // Two equally cheap, equally selective conjuncts: swapping cannot help, + // so the written order stands and compact-once is not used. let s = vec![stats(1000, 500, 1000), stats(1000, 500, 1000)]; - assert_eq!(settle_order(&s), vec![0, 1]); + let d = settle(&s); + assert_eq!(d.order, vec![0, 1]); + assert!(!d.compact); } #[test] - fn settle_adopts_materially_cheaper_order() { - // id 1 is far more selective and equally cheap: it should move first. + fn settle_adopts_materially_cheaper_order_with_compaction() { + // id 1 is far more selective and equally cheap: it should move first and + // run through the compact-once loop. let s = vec![stats(1000, 900, 1000), stats(1000, 10, 1000)]; - assert_eq!(settle_order(&s), vec![1, 0]); + let d = settle(&s); + assert_eq!(d.order, vec![1, 0]); + assert!(d.compact); + } + + /// When the order does not change, the settled evaluator runs the plain + /// predicate (compact-once is only used in service of a reorder), so an + /// interchangeable conjunction costs exactly what the flag-off path costs. + #[test] + fn no_reorder_evaluates_plain_predicate() { + let schema = schema(); + // Both conjuncts equally cheap and selective: nothing to reorder. + let left = + binary(col("a", &schema).unwrap(), Operator::Gt, lit(2i32), &schema).unwrap(); + let right = + binary(col("b", &schema).unwrap(), Operator::Gt, lit(2i32), &schema).unwrap(); + let p = binary(left, Operator::And, right, &schema).unwrap(); + let mut adaptive = try_new(&p, true).unwrap(); + + for round in 0..(WARMUP_BATCHES as i32 + 2) { + let base = round * 10; + let a: Vec = (base..base + 10).collect(); + let b: Vec = (base..base + 10).collect(); + let rb = batch(&schema, a, b); + let got = adaptive.evaluate(&rb).unwrap(); + let want = p.evaluate(&rb).unwrap().into_array(rb.num_rows()).unwrap(); + assert_eq!(passing_rows(&got), passing_rows(&want)); + } + assert!(adaptive.settled); + assert!( + !adaptive.compact, + "interchangeable conjuncts stay on the plain predicate" + ); } /// Two streams sharing one registry settle the order together: the pooled @@ -663,13 +763,15 @@ mod tests { } } - assert!(shared.settled_order().is_some()); + assert!(shared.settled().is_some()); // One more batch each lets a not-yet-settled stream adopt the epoch. s1.evaluate(&mk(99)).unwrap(); s2.evaluate(&mk(99)).unwrap(); assert!(s1.settled && s2.settled); - // The selective conjunct was promoted to the front for both. + // The selective conjunct was promoted to the front for both, and the + // reorder runs through the compact-once loop. assert_eq!(s1.order, vec![1, 0]); assert_eq!(s2.order, vec![1, 0]); + assert!(s1.compact && s2.compact); } }