diff --git a/datafusion/common/src/config.rs b/datafusion/common/src/config.rs index 913845a4219e4..14d99025fcd94 100644 --- a/datafusion/common/src/config.rs +++ b/datafusion/common/src/config.rs @@ -1184,6 +1184,16 @@ config_namespace! { /// /// Disabled by default, set to a number greater than 0 for enabling it. pub hash_join_buffering_capacity: usize, default = 0 + + /// (experimental) When true, `FilterExec` measures the selectivity + /// and the evaluation time of each conjunct of an `AND` predicate on + /// the first batches of each partition. Then it evaluates first the + /// conjuncts that remove the most rows per unit of time. The query + /// result does not change, but a fallible conjunct can see different + /// rows: for example, a new order of `b <> 0 AND 1 / b > 2` can cause + /// or prevent a division by zero error. Predicates with volatile + /// expressions are never reordered. + pub adaptive_filter_reordering: bool, default = false } } diff --git a/datafusion/physical-expr/src/expressions/binary.rs b/datafusion/physical-expr/src/expressions/binary.rs index a0a6548518737..484e11e857e7c 100644 --- a/datafusion/physical-expr/src/expressions/binary.rs +++ b/datafusion/physical-expr/src/expressions/binary.rs @@ -1236,7 +1236,11 @@ enum ShortCircuitStrategy { /// the side that cannot short-circuit the operator is rare: /// - for `AND`, when the proportion of `true` is less than or equal to 0.2 /// - for `OR`, when the proportion of `false` is less than or equal to 0.2 -const PRE_SELECTION_THRESHOLD: f32 = 0.2; +/// +/// Public only so that the adaptive conjunct ordering of `FilterExec` can +/// model the same rule. Not part of the stable API. +#[doc(hidden)] +pub const PRE_SELECTION_THRESHOLD: f32 = 0.2; /// Checks if a logical operator (`AND`/`OR`) can short-circuit evaluation based on the left-hand side (lhs) result. /// diff --git a/datafusion/physical-expr/src/expressions/mod.rs b/datafusion/physical-expr/src/expressions/mod.rs index f2f9285de560a..1ddd6e9cea9bd 100644 --- a/datafusion/physical-expr/src/expressions/mod.rs +++ b/datafusion/physical-expr/src/expressions/mod.rs @@ -41,7 +41,7 @@ pub use crate::PhysicalSortExpr; /// Module with some convenient methods used in expression building pub use crate::aggregate::stats::StatsType; -pub use binary::{BinaryExpr, binary, similar_to}; +pub use binary::{BinaryExpr, PRE_SELECTION_THRESHOLD, binary, similar_to}; pub use case::{CaseExpr, case}; pub use cast::{CastExpr, cast}; pub use column::{Column, col, with_new_schema}; diff --git a/datafusion/physical-expr/src/filter_stats.rs b/datafusion/physical-expr/src/filter_stats.rs new file mode 100644 index 0000000000000..f5f5c4f03f8b1 --- /dev/null +++ b/datafusion/physical-expr/src/filter_stats.rs @@ -0,0 +1,199 @@ +// 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. + +//! Measurements of filter evaluations at runtime. +//! +//! Operators that adapt to the data measure the filters that they evaluate: +//! the rows in, the rows that pass and the evaluation time. This module has +//! the shared parts: +//! +//! * [`Clock`]: a monotonic clock that tests can replace, so that decisions +//! that use time are deterministic in tests. [`SystemClock`] is the real +//! clock and [`ManualClock`] is a clock that only moves when a test moves +//! it. +//! * [`FilterCost`]: the counts and the time of one filter, and the values +//! derived from them (cost for each row, rows removed for each +//! nanosecond). +//! +//! For example, an operator can use them to pause a filter that costs more +//! than it saves, or to change the order of the conjuncts of a predicate. + +use std::fmt::Debug; +use std::sync::Arc; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::time::Duration; + +use datafusion_common::instant::Instant; + +/// A monotonic clock in nanoseconds. +/// +/// Production code uses [`SystemClock`]. Tests use [`ManualClock`] (or their +/// own implementation), thus decisions that use time are deterministic in +/// tests. +pub trait Clock: Debug + Send + Sync { + /// Nanoseconds since an arbitrary fixed point. The value never + /// decreases. + fn now_nanos(&self) -> u64; +} + +/// The real monotonic [`Clock`]. +#[derive(Debug, Clone, Copy)] +pub struct SystemClock { + start: Instant, +} + +impl SystemClock { + /// Creates a clock whose zero is now. + pub fn new() -> Self { + Self { + start: Instant::now(), + } + } + + /// A shared [`SystemClock`], as a trait object. + pub fn shared() -> Arc { + Arc::new(Self::new()) + } +} + +impl Default for SystemClock { + fn default() -> Self { + Self::new() + } +} + +impl Clock for SystemClock { + fn now_nanos(&self) -> u64 { + u64::try_from(self.start.elapsed().as_nanos()).unwrap_or(u64::MAX) + } +} + +/// A [`Clock`] that moves only when [`Self::advance`] is called. For tests. +#[derive(Debug, Default)] +pub struct ManualClock { + nanos: AtomicU64, +} + +impl ManualClock { + /// Creates a clock at zero. + pub fn new() -> Self { + Self::default() + } + + /// Moves the clock forward by `nanos` nanoseconds. + pub fn advance(&self, nanos: u64) { + self.nanos.fetch_add(nanos, Ordering::Relaxed); + } +} + +impl Clock for ManualClock { + fn now_nanos(&self) -> u64 { + self.nanos.load(Ordering::Relaxed) + } +} + +/// Returns the nanoseconds in `elapsed`, saturated to `u64::MAX`. +pub fn duration_nanos(elapsed: Duration) -> u64 { + u64::try_from(elapsed.as_nanos()).unwrap_or(u64::MAX) +} + +/// The measurements of one filter: the rows that it was evaluated on, the +/// rows that passed it, and the evaluation time. +#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)] +pub struct FilterCost { + /// Rows that the filter was evaluated on. + pub rows_in: u64, + /// Rows that passed the filter (`true`; `null` does not pass). + pub rows_out: u64, + /// Evaluation time in nanoseconds. + pub nanos: u64, +} + +impl FilterCost { + /// Adds the result of one evaluation. `rows_out` larger than `rows_in` + /// is used as `rows_in`. + pub fn add(&mut self, rows_in: u64, rows_out: u64, nanos: u64) { + self.rows_in = self.rows_in.saturating_add(rows_in); + self.rows_out = self.rows_out.saturating_add(rows_out.min(rows_in)); + self.nanos = self.nanos.saturating_add(nanos); + } + + /// Rows that the filter removed. + pub fn rows_removed(&self) -> u64 { + self.rows_in.saturating_sub(self.rows_out) + } + + /// Nanoseconds for each evaluated row, or `None` if the filter was not + /// evaluated on any row. + pub fn nanos_per_row(&self) -> Option { + (self.rows_in > 0).then(|| self.nanos as f64 / self.rows_in as f64) + } + + /// Rows removed for each nanosecond, `(1 + rows_in - rows_out) / nanos`, + /// or `None` if the filter was not evaluated on any row. A larger value + /// is a better filter to evaluate first. This is the ranking key of + /// Velox (Pedreira et al., VLDB 2022). The `1 +` ranks a filter that + /// removes no rows by its cost, and a zero time is used as 1 ns. + pub fn rows_removed_per_nano(&self) -> Option { + (self.rows_in > 0) + .then(|| (1 + self.rows_removed()) as f64 / self.nanos.max(1) as f64) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn filter_cost_derived_values() { + let empty = FilterCost::default(); + assert_eq!(empty.nanos_per_row(), None); + assert_eq!(empty.rows_removed_per_nano(), None); + + let mut cost = FilterCost::default(); + cost.add(100, 25, 1_000); + cost.add(100, 200, 1_000); + assert_eq!(cost.rows_in, 200); + // `rows_out` is at most `rows_in` for each evaluation. + assert_eq!(cost.rows_out, 125); + assert_eq!(cost.rows_removed(), 75); + assert_eq!(cost.nanos_per_row(), Some(10.0)); + assert_eq!(cost.rows_removed_per_nano(), Some(76.0 / 2_000.0)); + + // A zero time is used as 1 ns. + let mut free = FilterCost::default(); + free.add(10, 0, 0); + assert_eq!(free.rows_removed_per_nano(), Some(11.0)); + } + + #[test] + fn manual_clock_moves_only_when_advanced() { + let clock = ManualClock::new(); + assert_eq!(clock.now_nanos(), 0); + clock.advance(5); + clock.advance(7); + assert_eq!(clock.now_nanos(), 12); + } + + #[test] + fn system_clock_is_monotonic() { + let clock = SystemClock::new(); + let first = clock.now_nanos(); + assert!(clock.now_nanos() >= first); + assert_eq!(duration_nanos(Duration::from_micros(3)), 3_000); + } +} diff --git a/datafusion/physical-expr/src/lib.rs b/datafusion/physical-expr/src/lib.rs index 80e9f88b510ed..75d9d7ea70c2b 100644 --- a/datafusion/physical-expr/src/lib.rs +++ b/datafusion/physical-expr/src/lib.rs @@ -34,6 +34,7 @@ pub mod binary_map { pub mod async_scalar_function; pub mod equivalence; pub mod expressions; +pub mod filter_stats; pub mod higher_order_function; pub mod intervals; mod partitioning; diff --git a/datafusion/physical-plan/src/filter.rs b/datafusion/physical-plan/src/filter.rs index f71d67d04e39c..adc08e46fd0c5 100644 --- a/datafusion/physical-plan/src/filter.rs +++ b/datafusion/physical-plan/src/filter.rs @@ -46,9 +46,12 @@ use crate::stream::EmptyRecordBatchStream; use crate::{ChildrenPropertiesMode, ReplaceChildrenOptions, validate_child_count}; use crate::{ DisplayFormatType, ExecutionPlan, - metrics::{BaselineMetrics, ExecutionPlanMetricsSet, MetricsSet, RatioMetrics}, + metrics::{ + BaselineMetrics, Count, ExecutionPlanMetricsSet, MetricsSet, RatioMetrics, + }, }; +use arrow::array::ArrayRef; use arrow::compute::filter_record_batch; use arrow::datatypes::{DataType, SchemaRef}; use arrow::record_batch::RecordBatch; @@ -77,6 +80,10 @@ use datafusion_physical_expr_common::physical_expr::fmt_sql; use futures::stream::{Stream, StreamExt}; use log::trace; +mod conjunct_order; +use conjunct_order::ConjunctOrder; +use datafusion_physical_expr::filter_stats::SystemClock; + const FILTER_EXEC_DEFAULT_SELECTIVITY: u8 = 20; const FILTER_EXEC_DEFAULT_BATCH_SIZE: usize = 8192; @@ -645,10 +652,26 @@ impl ExecutionPlan for FilterExec { context.session_id(), context.task_id() ); - let metrics = FilterExecMetrics::new(&self.metrics, partition); + let mut metrics = FilterExecMetrics::new(&self.metrics, partition); + let conjunct_order = if context + .session_config() + .options() + .execution + .adaptive_filter_reordering + { + ConjunctOrder::try_new(&self.predicate, &SystemClock::shared()) + } else { + None + }; + if conjunct_order.is_some() { + metrics.adaptive_reorders = Some( + MetricBuilder::new(&self.metrics).counter("adaptive_reorders", partition), + ); + } Ok(Box::pin(FilterExecStream { schema: self.schema(), predicate: Arc::clone(&self.predicate), + conjunct_order, input: self.input.execute(partition, context)?, metrics, projection: self.projection.clone(), @@ -1356,6 +1379,9 @@ struct FilterExecStream { schema: SchemaRef, /// The expression to filter on. This expression must evaluate to a boolean value. predicate: Arc, + /// Evaluates `predicate` with adaptive conjunct ordering, if + /// `datafusion.execution.adaptive_filter_reordering` is true. + conjunct_order: Option, /// The input partition to filter. input: SendableRecordBatchStream, /// Runtime metrics recording @@ -1372,6 +1398,9 @@ struct FilterExecMetrics { baseline_metrics: BaselineMetrics, /// Selectivity of the filter, calculated as output_rows / input_rows selectivity: RatioMetrics, + /// Number of streams that changed the order of the conjuncts. Present + /// only when `datafusion.execution.adaptive_filter_reordering` is true. + adaptive_reorders: Option, // Remember to update `docs/source/user-guide/metrics.md` when adding new metrics, // or modifying metrics comments } @@ -1383,10 +1412,33 @@ impl FilterExecMetrics { selectivity: MetricBuilder::new(metrics) .with_type(MetricType::Summary) .ratio_metrics("selectivity", partition), + adaptive_reorders: None, } } } +impl FilterExecStream { + /// Evaluates `predicate` on `batch`, through `conjunct_order` if it is + /// set, and returns the selection mask. + fn evaluate_predicate(&mut self, batch: &RecordBatch) -> Result { + let mask = match &mut self.conjunct_order { + Some(order) => { + let was_measuring = order.is_measuring(); + let mask = order.evaluate(batch)?; + if was_measuring + && order.new_order().is_some() + && let Some(count) = &self.metrics.adaptive_reorders + { + count.add(1); + } + mask + } + None => self.predicate.evaluate(batch)?, + }; + mask.into_array(batch.num_rows()) + } +} + pub fn batch_filter( batch: &RecordBatch, predicate: &Arc, @@ -1451,9 +1503,7 @@ 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 status = self.evaluate_predicate(&batch) .and_then(|array| { Ok(match self.projection.as_ref() { Some(projection) => { @@ -4215,4 +4265,52 @@ mod tests { ); Ok(()) } + + /// With `adaptive_filter_reordering`, `FilterExec` returns the same rows + /// and reports the `adaptive_reorders` metric. Whether it reorders + /// depends on the wall clock, thus the scenario tests of the decision are + /// in `conjunct_order.rs`, with a mock clock. + #[tokio::test] + async fn adaptive_filter_reordering_returns_same_rows() -> Result<()> { + use datafusion_execution::config::SessionConfig; + + let schema = Arc::new(Schema::new(vec![Field::new("a", DataType::Int32, false)])); + let batch = RecordBatch::try_new( + Arc::clone(&schema), + vec![Arc::new(Int32Array::from_iter_values(0..100))], + )?; + let input: Arc = test::TestMemoryExec::try_new_exec( + &[vec![batch; 20]], + Arc::clone(&schema), + None, + )?; + let a = col("a", &schema)?; + // The selective conjunct is last. + let predicate = binary( + binary( + binary(Arc::clone(&a), Operator::GtEq, lit(0i32), &schema)?, + Operator::And, + binary(Arc::clone(&a), Operator::NotEq, lit(50i32), &schema)?, + &schema, + )?, + Operator::And, + binary(Arc::clone(&a), Operator::Lt, lit(5i32), &schema)?, + &schema, + )?; + let single = binary(a, Operator::Lt, lit(5i32), &schema)?; + + for (predicate, has_metric) in [(predicate, true), (single, false)] { + let mut config = SessionConfig::new(); + config.options_mut().execution.adaptive_filter_reordering = true; + let context = Arc::new(TaskContext::default().with_session_config(config)); + let filter = Arc::new(FilterExec::try_new(predicate, Arc::clone(&input))?); + let batches = collect(filter.execute(0, context)?).await?; + let rows: usize = batches.iter().map(|b| b.num_rows()).sum(); + assert_eq!(rows, 20 * 5); + let reorders = filter.metrics().unwrap().sum_by_name("adaptive_reorders"); + assert_eq!(reorders.is_some(), has_metric); + assert!(reorders.is_none_or(|r| r.as_usize() <= 1)); + } + Ok(()) + } } diff --git a/datafusion/physical-plan/src/filter/conjunct_order.rs b/datafusion/physical-plan/src/filter/conjunct_order.rs new file mode 100644 index 0000000000000..de4d6573f726b --- /dev/null +++ b/datafusion/physical-plan/src/filter/conjunct_order.rs @@ -0,0 +1,733 @@ +// 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. + +//! Adaptive conjunct ordering for [`FilterExec`], controlled by +//! `datafusion.execution.adaptive_filter_reordering`. +//! +//! The order of the conjuncts of an `AND` predicate is important. The `AND` +//! of [`BinaryExpr`] evaluates its right side only on the rows that its left +//! side keeps, if the left side keeps few rows (see +//! [`PRE_SELECTION_THRESHOLD`]). For example, in +//! `regexp_like(s, 'a') AND regexp_like(s, 'rare')`, the second conjunct +//! removes most rows. If it is first, the first conjunct is evaluated only on +//! the few rows that remain. +//! +//! [`ConjunctOrder`] finds a better order for one stream: +//! +//! 1. For the first [`WARMUP_BATCHES`] batches, it evaluates the predicate as +//! written. Each conjunct is wrapped in a [`MeasuredConjunct`] that counts +//! the rows in, the rows out and the evaluation time (a [`FilterCost`], +//! measured with a [`Clock`]). The shape of the `AND` tree does not +//! change, thus each +//! conjunct sees the same rows as without the measurement. +//! 2. Then it ranks the conjuncts by rows removed per nanosecond, and +//! estimates the cost of the new order and of the written order. It uses +//! the new order only if its cost is at least [`MIN_GAIN`] smaller. +//! 3. It evaluates the new order as an ordinary right-nested `AND` chain of +//! [`BinaryExpr`], or the written predicate unchanged. +//! +//! The decision is made one time. The state is local to the stream, thus +//! executions and partitions do not share it. +//! +//! [`FilterExec`]: crate::filter::FilterExec + +use std::fmt::{self, Debug, Display, Formatter}; +use std::hash::{Hash, Hasher}; +use std::sync::Arc; +use std::sync::atomic::{AtomicU64, Ordering::Relaxed}; + +use arrow::array::{Array, RecordBatch}; +use arrow::datatypes::{FieldRef, Schema}; +use datafusion_common::Result; +use datafusion_common::cast::as_boolean_array; +use datafusion_expr::{ColumnarValue, Operator}; +use datafusion_physical_expr::PhysicalExpr; +use datafusion_physical_expr::expressions::{BinaryExpr, PRE_SELECTION_THRESHOLD}; +use datafusion_physical_expr::filter_stats::{Clock, FilterCost}; +use datafusion_physical_expr::split_conjunction; +use datafusion_physical_expr_common::physical_expr::is_volatile; + +/// Number of non-empty batches that a stream measures before it decides. +const WARMUP_BATCHES: usize = 8; + +/// A new order must be at least this fraction cheaper than the written order. +const MIN_GAIN: f64 = 0.05; + +/// The measurements of one conjunct. +#[derive(Debug, Default, Clone, Copy, PartialEq)] +struct ConjunctStats { + /// Rows in, rows out (`true`; `null` does not pass) and evaluation time. + cost: FilterCost, + /// Rows that `AND` gives to the conjuncts after this one, if this one + /// is the left side. See [`rows_passed_on`]. + rows_passed_on: u64, +} + +impl ConjunctStats { + /// Nanoseconds for each row, or `None` if the conjunct was not + /// evaluated. + fn cost_per_row(&self) -> Option { + self.cost.nanos_per_row() + } + + /// The rank key: removed rows per nanosecond, see + /// [`FilterCost::rows_removed_per_nano`]. `None` if the conjunct was not + /// evaluated. + fn rank_key(&self) -> Option { + self.cost.rows_removed_per_nano() + } + + /// Fraction of the rows that the conjuncts after this one see. + fn passed_on_fraction(&self) -> f64 { + if self.cost.rows_in == 0 { + 1.0 + } else { + self.rows_passed_on as f64 / self.cost.rows_in as f64 + } + } +} + +/// Rows that the `AND` of [`BinaryExpr`] gives to its right side, when its +/// left side returns `rows_out` passing rows and `null_count` nulls out of +/// `rows_in` rows. This is the same rule as `check_short_circuit` in +/// `binary.rs`: a left side with nulls never pre-selects, a left side with +/// no passing row stops the evaluation, and a left side that keeps at most +/// [`PRE_SELECTION_THRESHOLD`] of the rows pre-selects them. +fn rows_passed_on(rows_in: usize, rows_out: usize, null_count: usize) -> usize { + if rows_in == 0 || null_count > 0 { + rows_in + } else if rows_out == 0 { + 0 + } else if rows_out as f32 / rows_in as f32 <= PRE_SELECTION_THRESHOLD { + rows_out + } else { + rows_in + } +} + +/// Wraps a conjunct and records [`ConjunctStats`] for it. The result of the +/// conjunct does not change. +#[derive(Debug)] +struct MeasuredConjunct { + inner: Arc, + clock: Arc, + rows_in: AtomicU64, + rows_out: AtomicU64, + rows_passed_on: AtomicU64, + nanos: AtomicU64, +} + +impl MeasuredConjunct { + fn new(inner: Arc, clock: Arc) -> Self { + Self { + inner, + clock, + rows_in: AtomicU64::new(0), + rows_out: AtomicU64::new(0), + rows_passed_on: AtomicU64::new(0), + nanos: AtomicU64::new(0), + } + } + + fn stats(&self) -> ConjunctStats { + ConjunctStats { + cost: FilterCost { + rows_in: self.rows_in.load(Relaxed), + rows_out: self.rows_out.load(Relaxed), + nanos: self.nanos.load(Relaxed), + }, + rows_passed_on: self.rows_passed_on.load(Relaxed), + } + } +} + +impl PartialEq for MeasuredConjunct { + fn eq(&self, other: &Self) -> bool { + self.inner.eq(&other.inner) + } +} + +impl Eq for MeasuredConjunct {} + +impl Hash for MeasuredConjunct { + fn hash(&self, state: &mut H) { + self.inner.hash(state); + } +} + +impl Display for MeasuredConjunct { + fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result { + write!(f, "{}", self.inner) + } +} + +impl PhysicalExpr for MeasuredConjunct { + fn return_field(&self, input_schema: &Schema) -> Result { + self.inner.return_field(input_schema) + } + + fn evaluate(&self, batch: &RecordBatch) -> Result { + let rows_in = batch.num_rows(); + let start = self.clock.now_nanos(); + let result = self.inner.evaluate(batch)?.into_array(rows_in)?; + let nanos = self.clock.now_nanos().saturating_sub(start); + let mask = as_boolean_array(&result)?; + let rows_out = mask.true_count(); + let passed_on = rows_passed_on(rows_in, rows_out, mask.null_count()); + self.rows_in.fetch_add(rows_in as u64, Relaxed); + self.rows_out.fetch_add(rows_out as u64, Relaxed); + self.rows_passed_on.fetch_add(passed_on as u64, Relaxed); + self.nanos.fetch_add(nanos, Relaxed); + Ok(ColumnarValue::Array(result)) + } + + fn children(&self) -> Vec<&Arc> { + vec![&self.inner] + } + + fn with_new_children( + self: Arc, + mut children: Vec>, + ) -> Result> { + datafusion_common::assert_eq_or_internal_err!( + children.len(), + 1, + "MeasuredConjunct: expected 1 child" + ); + Ok(Arc::new(Self::new( + children.remove(0), + Arc::clone(&self.clock), + ))) + } + + fn fmt_sql(&self, f: &mut Formatter<'_>) -> fmt::Result { + self.inner.fmt_sql(f) + } +} + +/// The state of a [`ConjunctOrder`]. +#[derive(Debug)] +enum State { + /// Evaluates the written predicate with measured conjuncts. + Measuring { + /// The written predicate, with each conjunct wrapped. + predicate: Arc, + /// The wrapped conjuncts, in the order of [`split_conjunction`]. + measured: Vec>, + /// Non-empty batches measured so far. + batches: usize, + }, + /// The decision is made. + Decided { + predicate: Arc, + /// The new order (indexes into the written conjuncts), or `None` if + /// the written predicate is kept. + order: Option>, + }, +} + +/// Evaluates an `AND` predicate for one stream, and changes the order of its +/// conjuncts when that is clearly faster. See the [module +/// documentation](self). +#[derive(Debug)] +pub(crate) struct ConjunctOrder { + written: Arc, + state: State, +} + +impl ConjunctOrder { + /// Returns `None` if the order of the conjuncts of `predicate` cannot + /// change: `predicate` has fewer than two conjuncts, or a conjunct is + /// volatile. + pub(crate) fn try_new( + predicate: &Arc, + clock: &Arc, + ) -> Option { + let conjuncts = split_conjunction(predicate); + if conjuncts.len() < 2 || conjuncts.iter().any(|c| is_volatile(c)) { + return None; + } + let mut measured = Vec::with_capacity(conjuncts.len()); + let wrapped = wrap_conjuncts(predicate, clock, &mut measured); + Some(Self { + written: Arc::clone(predicate), + state: State::Measuring { + predicate: wrapped, + measured, + batches: 0, + }, + }) + } + + /// Evaluates the predicate on `batch`. + pub(crate) fn evaluate(&mut self, batch: &RecordBatch) -> Result { + match &mut self.state { + State::Decided { predicate, .. } => predicate.evaluate(batch), + State::Measuring { + predicate, + measured, + batches, + } => { + let result = predicate.evaluate(batch)?; + if batch.num_rows() > 0 { + *batches += 1; + } + if *batches >= WARMUP_BATCHES { + let stats: Vec<_> = measured.iter().map(|m| m.stats()).collect(); + self.state = decide(&self.written, &stats); + } + Ok(result) + } + } + } + + /// The new order, if the decision is made and it changed the order. + pub(crate) fn new_order(&self) -> Option<&[usize]> { + match &self.state { + State::Decided { order, .. } => order.as_deref(), + State::Measuring { .. } => None, + } + } + + /// True while the stream measures the conjuncts. + pub(crate) fn is_measuring(&self) -> bool { + matches!(self.state, State::Measuring { .. }) + } +} + +/// Returns `predicate` with each conjunct of its root `AND` chain wrapped in +/// a [`MeasuredConjunct`]. The shape of the `AND` tree does not change. The +/// wrappers are added to `measured` in the order of [`split_conjunction`]. +fn wrap_conjuncts( + predicate: &Arc, + clock: &Arc, + measured: &mut Vec>, +) -> Arc { + if let Some(binary) = predicate.downcast_ref::() + && *binary.op() == Operator::And + { + let left = wrap_conjuncts(binary.left(), clock, measured); + let right = wrap_conjuncts(binary.right(), clock, measured); + return Arc::new(BinaryExpr::new(left, Operator::And, right)); + } + let wrapper = Arc::new(MeasuredConjunct::new( + Arc::clone(predicate), + Arc::clone(clock), + )); + measured.push(Arc::clone(&wrapper)); + wrapper +} + +/// Decides the order from the measurements. +fn decide(written: &Arc, stats: &[ConjunctStats]) -> State { + let written_order: Vec = (0..stats.len()).collect(); + let order = rank(stats); + let is_better = order != written_order + && expected_cost(stats, &order) + < (1.0 - MIN_GAIN) * expected_cost(stats, &written_order); + let conjuncts = split_conjunction(written); + match right_nested_and(&conjuncts, &order) { + Some(predicate) if is_better => State::Decided { + predicate, + order: Some(order), + }, + // Keep the written tree. A new tree with the same order can change + // where `AND` pre-selects, thus which rows a conjunct sees. + _ => State::Decided { + predicate: Arc::clone(written), + order: None, + }, + } +} + +/// Indexes of the conjuncts, by rank key from high to low. Conjuncts that were +/// not evaluated are last. The sort is stable. +fn rank(stats: &[ConjunctStats]) -> Vec { + let mut order: Vec = (0..stats.len()).collect(); + order.sort_by(|&a, &b| match (stats[a].rank_key(), stats[b].rank_key()) { + (Some(a), Some(b)) => b.total_cmp(&a), + (Some(_), None) => std::cmp::Ordering::Less, + (None, Some(_)) => std::cmp::Ordering::Greater, + (None, None) => std::cmp::Ordering::Equal, + }); + order +} + +/// Estimated nanoseconds for each input row when the conjuncts are +/// evaluated in `order` as a right-nested `AND` chain. Each conjunct costs +/// its cost for each row, times the fraction of the rows that the conjuncts +/// before it pass on (see [`rows_passed_on`]). The conjuncts are assumed to +/// be independent. A conjunct that was not evaluated costs nothing. +fn expected_cost(stats: &[ConjunctStats], order: &[usize]) -> f64 { + let mut fraction = 1.0; + let mut cost = 0.0; + for stats in order.iter().map(|&i| &stats[i]) { + if let Some(cost_per_row) = stats.cost_per_row() { + cost += fraction * cost_per_row; + fraction *= stats.passed_on_fraction(); + } + } + cost +} + +/// `c[o0] AND (c[o1] AND (... AND c[on]))`, or `None` if `order` is empty. +/// A right-nested chain keeps the rows that the `AND` pre-selects for all +/// the conjuncts after it. +fn right_nested_and( + conjuncts: &[&Arc], + order: &[usize], +) -> Option> { + order + .iter() + .rev() + .map(|&i| Arc::clone(conjuncts[i])) + .reduce(|right, left| Arc::new(BinaryExpr::new(left, Operator::And, right))) +} + +#[cfg(test)] +mod tests { + use super::*; + use arrow::array::{BooleanArray, Int32Array}; + use arrow::datatypes::{DataType, Field}; + use datafusion_physical_expr::expressions::{binary, col, lit}; + use datafusion_physical_expr::filter_stats::ManualClock; + + /// A conjunct that moves the manual clock by `nanos_per_row` for each row + /// that it evaluates. The clock does not move otherwise. + #[derive(Debug)] + struct Costly { + inner: Arc, + clock: Arc, + nanos_per_row: u64, + volatile: bool, + } + + impl PartialEq for Costly { + fn eq(&self, other: &Self) -> bool { + self.inner.eq(&other.inner) && self.nanos_per_row == other.nanos_per_row + } + } + + impl Eq for Costly {} + + impl Hash for Costly { + fn hash(&self, state: &mut H) { + self.inner.hash(state); + } + } + + impl Display for Costly { + fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result { + write!(f, "{}", self.inner) + } + } + + impl PhysicalExpr for Costly { + fn return_field(&self, input_schema: &Schema) -> Result { + self.inner.return_field(input_schema) + } + + fn evaluate(&self, batch: &RecordBatch) -> Result { + self.clock + .advance(self.nanos_per_row * batch.num_rows() as u64); + self.inner.evaluate(batch) + } + + fn children(&self) -> Vec<&Arc> { + vec![&self.inner] + } + + fn with_new_children( + self: Arc, + _children: Vec>, + ) -> Result> { + Ok(self) + } + + fn fmt_sql(&self, f: &mut Formatter<'_>) -> fmt::Result { + self.inner.fmt_sql(f) + } + + fn is_volatile_node(&self) -> bool { + self.volatile + } + } + + const ROWS: i32 = 100; + + /// Column `a` is `0..100`. Column `b` is `a` in the even rows and null + /// in the odd rows. + fn schema() -> Schema { + Schema::new(vec![ + Field::new("a", DataType::Int32, false), + Field::new("b", DataType::Int32, true), + ]) + } + + fn batch(rows: i32) -> RecordBatch { + RecordBatch::try_new( + Arc::new(schema()), + vec![ + Arc::new(Int32Array::from_iter_values(0..rows)), + Arc::new(Int32Array::from_iter( + (0..rows).map(|v| (v % 2 == 0).then_some(v)), + )), + ], + ) + .unwrap() + } + + struct Scenario { + clock: Arc, + } + + impl Scenario { + fn new() -> Self { + Self { + clock: Arc::new(ManualClock::new()), + } + } + + /// `column op value`, which costs `nanos_per_row` for each row. + fn conjunct( + &self, + column: &str, + op: Operator, + value: i32, + nanos_per_row: u64, + ) -> Arc { + let schema = schema(); + let inner = + binary(col(column, &schema).unwrap(), op, lit(value), &schema).unwrap(); + Arc::new(Costly { + inner, + clock: Arc::clone(&self.clock), + nanos_per_row, + volatile: false, + }) + } + + fn order(&self, predicate: &Arc) -> ConjunctOrder { + let clock = Arc::clone(&self.clock) as Arc; + ConjunctOrder::try_new(predicate, &clock).expect("can reorder") + } + } + + fn and( + left: Arc, + right: Arc, + ) -> Arc { + Arc::new(BinaryExpr::new(left, Operator::And, right)) + } + + fn mask(value: ColumnarValue) -> BooleanArray { + as_boolean_array(&value.into_array(ROWS as usize).unwrap()) + .unwrap() + .clone() + } + + /// Evaluates `batches` batches. Returns, for each batch, the strategy + /// that was used: `measure`, `written` or `order [..]`. Also checks that + /// each result is the same as the result of `predicate`. + fn run( + order: &mut ConjunctOrder, + predicate: &Arc, + batches: usize, + ) -> Vec { + (0..batches) + .map(|_| { + let strategy = if order.is_measuring() { + "measure".to_string() + } else if let Some(new_order) = order.new_order() { + format!("order {new_order:?}") + } else { + "written".to_string() + }; + let batch = batch(ROWS); + let got = mask(order.evaluate(&batch).unwrap()); + let want = mask(predicate.evaluate(&batch).unwrap()); + assert_eq!(got, want); + strategy + }) + .collect() + } + + /// `measure` for the warm-up batches, then `then` for 2 batches. + fn trace(then: &str) -> Vec { + let mut trace = vec!["measure".to_string(); WARMUP_BATCHES]; + trace.extend([then.to_string(), then.to_string()]); + trace + } + + fn decided_predicate(order: &ConjunctOrder) -> &Arc { + match &order.state { + State::Decided { predicate, .. } => predicate, + State::Measuring { .. } => panic!("still measuring"), + } + } + + /// `a >= 0` keeps all rows and `a < 5` keeps 5% of the rows, at the + /// same cost. `a < 5` goes first. + #[test] + fn selective_conjunct_moves_first() { + let s = Scenario::new(); + let predicate = and( + s.conjunct("a", Operator::GtEq, 0, 10), + s.conjunct("a", Operator::Lt, 5, 10), + ); + let mut order = s.order(&predicate); + assert_eq!(run(&mut order, &predicate, 10), trace("order [1, 0]")); + assert_eq!( + decided_predicate(&order).to_string(), + "a@0 < 5 AND a@0 >= 0" + ); + } + + /// The written order is already the best: the written predicate is kept + /// unchanged (the same `Arc`). + #[test] + fn best_written_order_is_kept() { + let s = Scenario::new(); + let predicate = and( + s.conjunct("a", Operator::Lt, 5, 10), + s.conjunct("a", Operator::GtEq, 0, 10), + ); + let mut order = s.order(&predicate); + assert_eq!(run(&mut order, &predicate, 10), trace("written")); + assert!(Arc::ptr_eq(decided_predicate(&order), &predicate)); + } + + /// `a >= 10` is cheap (1 ns for each row) and removes 10% of the rows. + /// `a < 20` is expensive (100 ns for each row) and removes 80%. The + /// cheap conjunct removes more rows per nanosecond, thus it stays first. + #[test] + fn cheap_conjunct_stays_before_expensive_selective_conjunct() { + let s = Scenario::new(); + let predicate = and( + s.conjunct("a", Operator::GtEq, 10, 1), + s.conjunct("a", Operator::Lt, 20, 100), + ); + let mut order = s.order(&predicate); + assert_eq!(run(&mut order, &predicate, 10), trace("written")); + } + + /// `a < 30` keeps 30% of the rows. This is more than + /// `PRE_SELECTION_THRESHOLD`, thus `AND` evaluates the next conjunct on + /// all rows. The new order is not cheaper, thus it is not used. + #[test] + fn no_reorder_when_and_cannot_pre_select() { + let s = Scenario::new(); + let predicate = and( + s.conjunct("a", Operator::GtEq, 10, 10), + s.conjunct("a", Operator::Lt, 30, 10), + ); + let mut order = s.order(&predicate); + assert_eq!(run(&mut order, &predicate, 10), trace("written")); + + // With `a < 10` (10% of the rows), `AND` pre-selects. + let predicate = and( + s.conjunct("a", Operator::GtEq, 10, 10), + s.conjunct("a", Operator::Lt, 10, 10), + ); + let mut order = s.order(&predicate); + assert_eq!(run(&mut order, &predicate, 10), trace("order [1, 0]")); + } + + /// `b < 10` keeps 5% of the rows, but it returns nulls. `AND` never + /// pre-selects with nulls, thus the new order is not cheaper. + #[test] + fn nulls_prevent_reorder() { + let s = Scenario::new(); + let predicate = and( + s.conjunct("a", Operator::GtEq, 0, 10), + s.conjunct("b", Operator::Lt, 10, 10), + ); + let mut order = s.order(&predicate); + assert_eq!(run(&mut order, &predicate, 10), trace("written")); + } + + /// The new order is a right-nested `AND` chain. + #[test] + fn new_order_is_right_nested() { + let s = Scenario::new(); + let predicate = and( + and( + s.conjunct("a", Operator::GtEq, 0, 10), + s.conjunct("a", Operator::GtEq, 10, 10), + ), + s.conjunct("a", Operator::Lt, 5, 10), + ); + let mut order = s.order(&predicate); + assert_eq!(run(&mut order, &predicate, 10), trace("order [2, 1, 0]")); + let root = decided_predicate(&order) + .downcast_ref::() + .unwrap(); + assert_eq!(root.left().to_string(), "a@0 < 5"); + let right = root.right().downcast_ref::().unwrap(); + assert_eq!(right.left().to_string(), "a@0 >= 10"); + assert_eq!(right.right().to_string(), "a@0 >= 0"); + } + + /// Empty batches do not count as measured batches. + #[test] + fn empty_batches_are_not_measured() { + let s = Scenario::new(); + let predicate = and( + s.conjunct("a", Operator::GtEq, 0, 10), + s.conjunct("a", Operator::Lt, 5, 10), + ); + let mut order = s.order(&predicate); + for _ in 0..2 * WARMUP_BATCHES { + order.evaluate(&batch(0)).unwrap(); + } + assert!(order.is_measuring()); + } + + /// A predicate with fewer than two conjuncts, or with a volatile + /// conjunct, is not reordered. + #[test] + fn some_predicates_are_not_reordered() { + let s = Scenario::new(); + let clock = || Arc::clone(&s.clock) as Arc; + let single = s.conjunct("a", Operator::Lt, 5, 10); + assert!(ConjunctOrder::try_new(&single, &clock()).is_none()); + + let volatile: Arc = Arc::new(Costly { + inner: s.conjunct("a", Operator::Lt, 5, 10), + clock: Arc::clone(&s.clock), + nanos_per_row: 10, + volatile: true, + }); + let predicate = and(s.conjunct("a", Operator::GtEq, 0, 10), volatile); + assert!(ConjunctOrder::try_new(&predicate, &clock()).is_none()); + } + + #[test] + fn rows_passed_on_matches_binary_expr() { + // No passing row: the evaluation stops. + assert_eq!(rows_passed_on(100, 0, 0), 0); + // At or below the threshold: pre-selection. + assert_eq!(rows_passed_on(100, 20, 0), 20); + // Above the threshold: all rows. + assert_eq!(rows_passed_on(100, 21, 0), 100); + // All rows pass. + assert_eq!(rows_passed_on(100, 100, 0), 100); + // Nulls: never pre-selects. + assert_eq!(rows_passed_on(100, 5, 1), 100); + assert_eq!(rows_passed_on(100, 0, 1), 100); + } +} diff --git a/datafusion/sqllogictest/test_files/adaptive_filter_reordering.slt b/datafusion/sqllogictest/test_files/adaptive_filter_reordering.slt new file mode 100644 index 0000000000000..51e5f3665917e --- /dev/null +++ b/datafusion/sqllogictest/test_files/adaptive_filter_reordering.slt @@ -0,0 +1,81 @@ +# 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 `datafusion.execution.adaptive_filter_reordering`. The decision to +# reorder depends on measured time, thus these tests only check that the +# results do not change. The decision is tested with a mock clock in +# `datafusion/physical-plan/src/filter/conjunct_order.rs`. + +statement ok +set datafusion.execution.batch_size = 64; + +statement ok +CREATE TABLE t AS +SELECT value AS a, value % 7 AS b, CAST(value AS VARCHAR) AS s +FROM generate_series(1, 10000); + +# The most selective conjunct is last. +query IIT rowsort +SELECT a, b, s FROM t WHERE b <> 3 AND s LIKE '%1%' AND a < 20; +---- +1 1 1 +11 4 11 +12 5 12 +13 6 13 +14 0 14 +15 1 15 +16 2 16 +18 4 18 +19 5 19 + +statement ok +set datafusion.execution.adaptive_filter_reordering = true; + +query IIT rowsort +SELECT a, b, s FROM t WHERE b <> 3 AND s LIKE '%1%' AND a < 20; +---- +1 1 1 +11 4 11 +12 5 12 +13 6 13 +14 0 14 +15 1 15 +16 2 16 +18 4 18 +19 5 19 + +# Nullable conjuncts +query I +SELECT count(*) FROM t +WHERE CASE WHEN b = 0 THEN NULL ELSE a END > 5000 AND s LIKE '%9%' AND b < 3; +---- +596 + +statement ok +set datafusion.execution.adaptive_filter_reordering = false; + +query I +SELECT count(*) FROM t +WHERE CASE WHEN b = 0 THEN NULL ELSE a END > 5000 AND s LIKE '%9%' AND b < 3; +---- +596 + +statement ok +DROP TABLE t; + +statement ok +set datafusion.execution.batch_size = 8192; diff --git a/datafusion/sqllogictest/test_files/information_schema.slt b/datafusion/sqllogictest/test_files/information_schema.slt index 113f63c71257c..68084bd7188a9 100644 --- a/datafusion/sqllogictest/test_files/information_schema.slt +++ b/datafusion/sqllogictest/test_files/information_schema.slt @@ -213,6 +213,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 @@ -374,6 +375,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 true, `FilterExec` measures the selectivity and the evaluation time of each conjunct of an `AND` predicate on the first batches of each partition. Then it evaluates first the conjuncts that remove the most rows per unit of time. The query result does not change, but a fallible conjunct can see different rows: for example, a new order of `b <> 0 AND 1 / b > 2` can cause or prevent a division by zero error. Predicates with 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 ef4ea00129f29..f3cdfa56dca16 100644 --- a/docs/source/user-guide/configs.md +++ b/docs/source/user-guide/configs.md @@ -145,6 +145,7 @@ The following configuration settings are available: | 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. | +| datafusion.execution.adaptive_filter_reordering | false | (experimental) When true, `FilterExec` measures the selectivity and the evaluation time of each conjunct of an `AND` predicate on the first batches of each partition. Then it evaluates first the conjuncts that remove the most rows per unit of time. The query result does not change, but a fallible conjunct can see different rows: for example, a new order of `b <> 0 AND 1 / b > 2` can cause or prevent a division by zero error. Predicates with volatile expressions are never reordered. | | datafusion.optimizer.enable_distinct_aggregation_soft_limit | true | When set to true, the optimizer will push a limit operation into grouped aggregations which have no aggregate expressions, as a soft limit, emitting groups once the limit is reached, before all rows in the group are read. | | datafusion.optimizer.enable_round_robin_repartition | true | When set to true, the physical plan optimizer will try to add round robin repartitioning to increase parallelism to leverage more CPU cores | | datafusion.optimizer.enable_topk_aggregation | true | When set to true, the optimizer will attempt to perform limit operations during aggregations, if possible | diff --git a/docs/source/user-guide/metrics.md b/docs/source/user-guide/metrics.md index 111df66ccc08a..e0bd6d0820d45 100644 --- a/docs/source/user-guide/metrics.md +++ b/docs/source/user-guide/metrics.md @@ -38,9 +38,10 @@ DataFusion operators expose runtime metrics so you can understand where time is ### FilterExec -| Metric | Description | -| ----------- | ----------------------------------------------------------------- | -| selectivity | Selectivity of the filter, calculated as output_rows / input_rows | +| Metric | Description | +| ----------------- | ----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | +| selectivity | Selectivity of the filter, calculated as output_rows / input_rows | +| adaptive_reorders | Number of partition streams that changed the order of the conjuncts of the predicate. Only present when `datafusion.execution.adaptive_filter_reordering` is true and the predicate has two or more conjuncts and no volatile expression. | ### HashJoinExec