Skip to content

Commit 87974dc

Browse files
adriangbclaude
andcommitted
feat(gate): use the measured work of the filter producer as the saving
The gate assumed that each removed row saves `min_saving_ns_per_row` (20 ns) after the filter. For hash join dynamic filters this is the probe work of the join, and it is much smaller in star joins with small dimension tables: 3.5 to 8 ns for each probe row on the TPC-DS SF1 `date_dim` joins (Q65, Q67), 2 ns on Q90, 17 ns on TPC-H Q9. Filters that cost 3 to 7 ns for each row and remove 80% of the rows thus stayed on, and cost more than the join work they saved (8-18% slower than `pruning_only` on the bot). `RemovedRowWork` (in `filter_stats`) is the work that the producer of a filter does for each row that the filter removes, as the producer measures it. Each `DynamicFilterPhysicalExpr` has one, shared by all its derived filters. The gate uses the smallest measured work of the dynamic filters in its filter as the saving of a removed row, and `min_saving_ns_per_row` only until the producer has measured `MIN_OBSERVED_ROWS` rows (a prior). PR: #25674 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
1 parent 0e3905c commit 87974dc

3 files changed

Lines changed: 146 additions & 11 deletions

File tree

‎datafusion/physical-expr/src/expressions/dynamic_filters/mod.rs‎

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@ use std::{fmt::Display, hash::Hash, sync::Arc};
2121
use tokio::sync::watch;
2222

2323
use crate::PhysicalExpr;
24+
use crate::filter_stats::RemovedRowWork;
2425
use arrow::datatypes::{DataType, Schema};
2526
#[cfg(feature = "proto")]
2627
use datafusion_common::internal_datafusion_err;
@@ -94,6 +95,10 @@ pub struct DynamicFilterPhysicalExpr {
9495
/// But this can have overhead in production, so it's only included in our tests.
9596
data_type: Arc<RwLock<Option<DataType>>>,
9697
nullable: Arc<RwLock<Option<bool>>>,
98+
/// The work that the producer of the filter does for each row that the
99+
/// filter removes, see [`Self::removed_row_work`]. Shared by all
100+
/// derived filters.
101+
removed_row_work: Arc<RemovedRowWork>,
97102
}
98103

99104
impl std::fmt::Debug for DynamicFilterPhysicalExpr {
@@ -223,9 +228,20 @@ impl DynamicFilterPhysicalExpr {
223228
state_watch,
224229
data_type: Arc::new(RwLock::new(None)),
225230
nullable: Arc::new(RwLock::new(None)),
231+
removed_row_work: Arc::default(),
226232
}
227233
}
228234

235+
/// The work that the producer of this filter does for each row that
236+
/// the filter removes before it (for example the probe of a hash join),
237+
/// as the producer measures it. A consumer that decides if the filter
238+
/// is worth its cost uses it as the saving of a removed row (see
239+
/// [`OptionalFilterGate`](crate::optional_filter_gate::OptionalFilterGate)).
240+
/// The producer records its work in it.
241+
pub fn removed_row_work(&self) -> &Arc<RemovedRowWork> {
242+
&self.removed_row_work
243+
}
244+
229245
fn remap_children(
230246
children: &[Arc<dyn PhysicalExpr>],
231247
remapped_children: Option<&Vec<Arc<dyn PhysicalExpr>>>,
@@ -490,6 +506,7 @@ impl DynamicFilterPhysicalExpr {
490506
state_watch,
491507
data_type: Arc::new(RwLock::new(None)),
492508
nullable: Arc::new(RwLock::new(None)),
509+
removed_row_work: Arc::default(),
493510
}
494511
}
495512
}
@@ -518,6 +535,7 @@ impl PhysicalExpr for DynamicFilterPhysicalExpr {
518535
state_watch: self.state_watch.clone(),
519536
data_type: Arc::clone(&self.data_type),
520537
nullable: Arc::clone(&self.nullable),
538+
removed_row_work: Arc::clone(&self.removed_row_work),
521539
}))
522540
}
523541

@@ -614,6 +632,7 @@ impl PhysicalExpr for DynamicFilterPhysicalExpr {
614632
state_watch: _, // Runtime channel, recreated from inner state by from_parts().
615633
data_type: _, // Cached test invariant, recomputed from the expression.
616634
nullable: _, // Cached test invariant, recomputed from the expression.
635+
removed_row_work: _, // Runtime measurement of the producer.
617636
} = self;
618637

619638
let children = children
@@ -837,6 +856,7 @@ impl DynamicFilterPhysicalExpr {
837856
state_watch: self.state_watch.clone(),
838857
data_type: Arc::clone(&self.data_type),
839858
nullable: Arc::clone(&self.nullable),
859+
removed_row_work: Arc::clone(&self.removed_row_work),
840860
}
841861
}
842862
}

‎datafusion/physical-expr/src/filter_stats.rs‎

Lines changed: 53 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -163,10 +163,63 @@ impl FilterCost {
163163
}
164164
}
165165

166+
/// The work, in nanoseconds for each row, that an operator does on the rows
167+
/// that a filter removes before them, as measured by that operator.
168+
///
169+
/// For example, a hash join computes the hashes of the join keys of each
170+
/// probe row and looks them up in its hash table. A row that the dynamic
171+
/// filter of the join removes in the scan does not get this work, thus this
172+
/// is the saving of a removed row. The work that depends on a match (the
173+
/// output of a matched row) is not in it: the filter does not remove
174+
/// matched rows. The producer of a dynamic filter measures it and the
175+
/// consumers of the filter read it (see
176+
/// [`DynamicFilterPhysicalExpr::removed_row_work`]).
177+
///
178+
/// Lock-free.
179+
///
180+
/// [`DynamicFilterPhysicalExpr::removed_row_work`]: crate::expressions::DynamicFilterPhysicalExpr::removed_row_work
181+
#[derive(Debug, Default)]
182+
pub struct RemovedRowWork {
183+
rows: AtomicU64,
184+
nanos: AtomicU64,
185+
}
186+
187+
impl RemovedRowWork {
188+
/// Creates an empty measurement.
189+
pub fn new() -> Self {
190+
Self::default()
191+
}
192+
193+
/// Adds `rows` rows and `nanos` nanoseconds of work. The two can be
194+
/// recorded separately (for example the rows when a batch arrives and
195+
/// the time of each step).
196+
pub fn record(&self, rows: u64, nanos: u64) {
197+
self.rows.fetch_add(rows, Ordering::Relaxed);
198+
self.nanos.fetch_add(nanos, Ordering::Relaxed);
199+
}
200+
201+
/// The work for each row, or `None` before [`MIN_OBSERVED_ROWS`] rows.
202+
pub fn ns_per_row(&self) -> Option<f64> {
203+
let rows = self.rows.load(Ordering::Relaxed);
204+
(rows >= MIN_OBSERVED_ROWS)
205+
.then(|| self.nanos.load(Ordering::Relaxed) as f64 / rows as f64)
206+
}
207+
}
208+
166209
#[cfg(test)]
167210
mod tests {
168211
use super::*;
169212

213+
#[test]
214+
fn removed_row_work_needs_min_observed_rows() {
215+
let work = RemovedRowWork::new();
216+
work.record(MIN_OBSERVED_ROWS - 1, 0);
217+
work.record(0, 3 * (MIN_OBSERVED_ROWS - 1));
218+
assert_eq!(work.ns_per_row(), None);
219+
work.record(1, 3);
220+
assert_eq!(work.ns_per_row(), Some(3.0));
221+
}
222+
170223
#[test]
171224
fn filter_cost_derived_values() {
172225
let empty = FilterCost::default();

‎datafusion/physical-expr/src/optional_filter_gate.rs‎

Lines changed: 73 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -73,12 +73,20 @@
7373
//! ```text
7474
//! cost_ns = evaluation time + rows_in * measured overhead
7575
//! saving_ns = (rows_in - rows_out) * saving_ns_per_row
76-
//! saving_ns_per_row = min_saving_ns_per_row + measured saving
76+
//! saving_ns_per_row = producer work + measured saving
7777
//! ```
7878
//!
79-
//! `min_saving_ns_per_row` comes from the configuration. It is the work
80-
//! that a removed row saves after the filter, for example a hash table
81-
//! probe in a join. The *measured saving* and the *measured overhead* are
79+
//! The *producer work* is the work that a removed row saves after the
80+
//! filter, in the operator that produced the filter, for example the
81+
//! hash and the hash table lookup of a probe row in a hash join. The
82+
//! producer measures it ([`RemovedRowWork`], see
83+
//! [`DynamicFilterPhysicalExpr::removed_row_work`]): a hash join with a
84+
//! small build side does 2 to 8 ns of work for each probe row (TPC-DS
85+
//! SF1 star joins), a join with a large build side much more. Until the
86+
//! producer has measured [`MIN_OBSERVED_ROWS`] rows, the gate uses
87+
//! `min_saving_ns_per_row` from the configuration (a prior). With more
88+
//! than one dynamic filter in the filter, the smallest measured work is
89+
//! used. The *measured saving* and the *measured overhead* are
8290
//! optional: a consumer that can measure more work that a removed row
8391
//! saves (the Parquet scan measures the decode time of the columns that
8492
//! the filter does not read), or a fixed cost for each evaluated row in
@@ -146,9 +154,11 @@ use std::time::Duration;
146154
use datafusion_common::config::ExecutionOptions;
147155
use datafusion_physical_expr_common::physical_expr::PhysicalExpr;
148156

149-
use crate::expressions::DynamicFilterTracking;
157+
use datafusion_common::tree_node::{TreeNode, TreeNodeRecursion};
158+
159+
use crate::expressions::{DynamicFilterPhysicalExpr, DynamicFilterTracking};
150160
use crate::filter_stats::{
151-
Clock, FilterCost, MIN_OBSERVED_ROWS, SystemClock, duration_nanos,
161+
Clock, FilterCost, MIN_OBSERVED_ROWS, RemovedRowWork, SystemClock, duration_nanos,
152162
};
153163

154164
/// A running filter is paused by the cost rule only if its cost is larger
@@ -172,7 +182,8 @@ pub struct OptionalFilterGateConfig {
172182
/// `initial_pause_batches` are used as `initial_pause_batches`.
173183
pub max_pause_batches: usize,
174184
/// Work, in nanoseconds, that each row removed by the filter saves after
175-
/// the filter, at the least. The gate adds the saving that the consumer
185+
/// the filter, until the producer of the filter has measured it (see
186+
/// [`RemovedRowWork`]). The gate adds the saving that the consumer
176187
/// measures (see [`MeasuredRowSaving`]). The gate pauses a filter whose
177188
/// evaluation time is larger than the saving of the rows that it
178189
/// removes. Negative values are used as 0.
@@ -421,9 +432,11 @@ pub struct OptionalFilterGate {
421432
config: OptionalFilterGateConfig,
422433
/// The clock that consumers use to measure the evaluation time.
423434
clock: Arc<dyn Clock>,
424-
/// The saving that the consumer measures, added to
425-
/// `config.min_saving_ns_per_row`.
435+
/// The saving that the consumer measures, added to the producer work.
426436
measured_saving: Option<Arc<MeasuredRowSaving>>,
437+
/// The work that the producers of the dynamic filters in `filter` do
438+
/// for each removed row, see [`RemovedRowWork`].
439+
producer_work: Vec<Arc<RemovedRowWork>>,
427440
state: GateState,
428441
/// True while the current window is a probe after a pause. The cost
429442
/// rule then uses [`RESUME_COST_MARGIN`].
@@ -457,12 +470,22 @@ impl OptionalFilterGate {
457470
pub fn new(filter: Arc<dyn PhysicalExpr>, config: OptionalFilterGateConfig) -> Self {
458471
let config = config.normalized();
459472
let tracking = DynamicFilterTracking::classify(&filter);
473+
let mut producer_work = vec![];
474+
filter
475+
.apply(|expr| {
476+
if let Some(dynamic) = expr.downcast_ref::<DynamicFilterPhysicalExpr>() {
477+
producer_work.push(Arc::clone(dynamic.removed_row_work()));
478+
}
479+
Ok(TreeNodeRecursion::Continue)
480+
})
481+
.expect("the closure is infallible");
460482
Self {
461483
filter,
462484
tracking,
463485
config,
464486
clock: SystemClock::shared(),
465487
measured_saving: None,
488+
producer_work,
466489
state: GateState::new_window(),
467490
probing: false,
468491
backoff: config.initial_pause_batches,
@@ -513,13 +536,26 @@ impl OptionalFilterGate {
513536
}
514537

515538
/// The work, in nanoseconds, that the gate assumes each removed row
516-
/// saves now: the configured minimum plus the measured saving.
539+
/// saves now: the work of the producer (measured, or the configured
540+
/// `min_saving_ns_per_row` before the measurement) plus the measured
541+
/// saving of the consumer.
517542
pub fn saving_ns_per_row(&self) -> f64 {
518543
let measured = self
519544
.measured_saving
520545
.as_ref()
521546
.map_or(0.0, |saving| saving.ns_per_row());
522-
self.config.min_saving_ns_per_row + measured
547+
self.producer_work_ns_per_row() + measured
548+
}
549+
550+
/// The work of the producer for each removed row: the smallest measured
551+
/// [`RemovedRowWork`] of the dynamic filters in the filter, or
552+
/// `min_saving_ns_per_row` if none is measured yet.
553+
fn producer_work_ns_per_row(&self) -> f64 {
554+
self.producer_work
555+
.iter()
556+
.filter_map(|work| work.ns_per_row())
557+
.reduce(f64::min)
558+
.unwrap_or(self.config.min_saving_ns_per_row)
523559
}
524560

525561
/// The work, in nanoseconds for each evaluated row, that the gate adds
@@ -1374,6 +1410,32 @@ mod tests {
13741410
assert!(!shared_gate(filter, &shared).is_paused());
13751411
}
13761412

1413+
/// The producer of a dynamic filter measures its work for each row that
1414+
/// the filter removes: once measured, it replaces the configured
1415+
/// `min_saving_ns_per_row`.
1416+
#[test]
1417+
fn producer_work_replaces_configured_saving() {
1418+
let (dynamic, filter) = dynamic_filter();
1419+
let mut gate = gate_with(filter);
1420+
assert_eq!(gate.saving_ns_per_row(), 20.0);
1421+
// Removes 80% at 5 ns for each row: 5 < 0.8 * 20 * 1.1, it stays on.
1422+
for _ in 0..4 {
1423+
assert_eq!(feed_timed(&mut gate, 0.2, 5.0), GateDecision::Evaluate);
1424+
}
1425+
assert!(!gate.is_paused());
1426+
1427+
// The producer measures 4 ns for each removed row: 0.8 * 4 < 5.
1428+
let work = dynamic.removed_row_work();
1429+
work.record(MIN_OBSERVED_ROWS, 4 * MIN_OBSERVED_ROWS);
1430+
assert_eq!(gate.saving_ns_per_row(), 4.0);
1431+
assert_eq!(feed_timed(&mut gate, 0.2, 5.0), GateDecision::Evaluate);
1432+
assert_eq!(feed_timed(&mut gate, 0.2, 5.0), GateDecision::Evaluate);
1433+
assert!(gate.is_paused());
1434+
1435+
// A filter without dynamic filters uses the configuration.
1436+
assert_eq!(new_gate().saving_ns_per_row(), 20.0);
1437+
}
1438+
13771439
/// A clock that moves by a fixed step each time it is read.
13781440
#[derive(Debug)]
13791441
struct SteppingClock {

0 commit comments

Comments
 (0)