Skip to content

Commit 1eb4f7c

Browse files
adriangbclaude
andcommitted
feat: add OptionalFilterGate and optional_filter_mode config
Add a runtime gate that pauses optional filters (filters that are not needed for correctness, such as hash join and TopK dynamic filters) when they remove too few rows. The gate is a per-stream state machine (Evaluate / Paused with exponential backoff) that restarts evaluation when the filter changes. The gate finds the dynamic filters one time with `DynamicFilterTracking::classify` and then polls their subscriptions, so a check does not walk the filter tree. Gates of one plan site share lock-free pooled statistics, which seed the state of new gates. Add a crate-private generation tag to `DynamicFilterTracker` (the sum of the latest observed generations of the dynamic filters, including complete ones), so the pooled statistics can tag verdicts by filter version without calling `snapshot_generation`. Add the `datafusion.execution.optional_filter_mode` (always | adaptive | pruning_only, default always) and `datafusion.execution.optional_filter_max_pass_ratio` (default 0.8) options. No operator uses the gate yet, so behavior does not change. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
1 parent 0413381 commit 1eb4f7c

7 files changed

Lines changed: 1194 additions & 7 deletions

File tree

‎datafusion/common/src/config.rs‎

Lines changed: 75 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -881,6 +881,61 @@ impl Display for MapKeyDedupPolicy {
881881
}
882882
}
883883

884+
/// How DataFusion evaluates optional filters.
885+
///
886+
/// Optional filters are filters that are not needed for correctness, such as
887+
/// dynamic filters pushed down by hash joins and TopK. See
888+
/// [`ExecutionOptions::optional_filter_mode`].
889+
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
890+
pub enum OptionalFilterMode {
891+
/// Evaluate optional filters like any other pushed-down filter.
892+
#[default]
893+
Always,
894+
/// Pause optional filters that remove too few rows. Try them again at
895+
/// intervals to find out if they became selective.
896+
Adaptive,
897+
/// Use optional filters only for statistics pruning (for example of
898+
/// files, row groups and pages). Never evaluate them row by row.
899+
PruningOnly,
900+
}
901+
902+
impl FromStr for OptionalFilterMode {
903+
type Err = DataFusionError;
904+
905+
fn from_str(s: &str) -> Result<Self, Self::Err> {
906+
match s.to_ascii_lowercase().as_str() {
907+
"always" => Ok(Self::Always),
908+
"adaptive" => Ok(Self::Adaptive),
909+
"pruning_only" => Ok(Self::PruningOnly),
910+
other => Err(DataFusionError::Configuration(format!(
911+
"Invalid optional filter mode: {other}. Expected one of: always, adaptive, pruning_only"
912+
))),
913+
}
914+
}
915+
}
916+
917+
impl ConfigField for OptionalFilterMode {
918+
fn visit<V: Visit>(&self, v: &mut V, key: &str, description: &'static str) {
919+
v.some(key, self, description)
920+
}
921+
922+
fn set(&mut self, _: &str, value: &str) -> Result<()> {
923+
*self = OptionalFilterMode::from_str(value)?;
924+
Ok(())
925+
}
926+
}
927+
928+
impl Display for OptionalFilterMode {
929+
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
930+
let str = match self {
931+
Self::Always => "always",
932+
Self::Adaptive => "adaptive",
933+
Self::PruningOnly => "pruning_only",
934+
};
935+
write!(f, "{str}")
936+
}
937+
}
938+
884939
impl From<SpillCompression> for Option<CompressionType> {
885940
fn from(c: SpillCompression) -> Self {
886941
match c {
@@ -1182,6 +1237,26 @@ config_namespace! {
11821237
///
11831238
/// Disabled by default, set to a number greater than 0 for enabling it.
11841239
pub hash_join_buffering_capacity: usize, default = 0
1240+
1241+
/// Controls how DataFusion evaluates filters that are not needed for
1242+
/// correctness, such as the dynamic filters that hash joins and TopK
1243+
/// push down into scans. `always` evaluates these filters like any
1244+
/// other pushed-down filter. `adaptive` pauses these filters when they
1245+
/// remove too few rows (see
1246+
/// `datafusion.execution.optional_filter_max_pass_ratio`), and tries
1247+
/// them again at intervals. `pruning_only` uses these filters only to
1248+
/// prune files, row groups and pages with statistics, and never
1249+
/// evaluates them row by row.
1250+
///
1251+
/// This option is most important when
1252+
/// `datafusion.execution.parquet.pushdown_filters` is true, because then
1253+
/// the Parquet reader evaluates pushed-down filters row by row.
1254+
pub optional_filter_mode: OptionalFilterMode, default = OptionalFilterMode::Always
1255+
1256+
/// When `datafusion.execution.optional_filter_mode` is `adaptive`, pause
1257+
/// an optional filter when more than this fraction of the rows pass it.
1258+
/// Use a value between 0.0 and 1.0.
1259+
pub optional_filter_max_pass_ratio: f64, default = 0.8
11851260
}
11861261
}
11871262

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

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -755,6 +755,13 @@ pub(crate) struct DynamicFilterSubscription {
755755
}
756756

757757
impl DynamicFilterSubscription {
758+
/// The latest generation of the filter that this subscription observed:
759+
/// the generation at [`DynamicFilterPhysicalExpr::subscribe`] time, or
760+
/// the generation of the last change that [`Self::observe`] reported.
761+
pub(crate) fn last_generation(&self) -> u64 {
762+
self.last_generation
763+
}
764+
758765
/// Observe the latest state of the filter.
759766
///
760767
/// Reports whether the filter's expression advanced since the previous call

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

Lines changed: 135 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -69,29 +69,56 @@ impl DynamicFilterTracking {
6969
/// Walk `predicate` once and classify its dynamic-filter content,
7070
/// subscribing to every filter that is not yet complete.
7171
pub fn classify(predicate: &Arc<dyn PhysicalExpr>) -> Self {
72+
Self::classify_with_generation_tag(predicate).0
73+
}
74+
75+
/// Same as [`Self::classify`], but also returns the generation tag of
76+
/// `predicate` at classification time.
77+
///
78+
/// The tag is the wrapping sum of the generations of all dynamic filters
79+
/// in `predicate` (0 for a [`Self::Static`] predicate). Two callers that
80+
/// classify the same predicate get the same tag if no filter changed
81+
/// between the two calls. For [`Self::Watching`], use
82+
/// [`DynamicFilterTracker::generation_tag`] after
83+
/// [`DynamicFilterTracker::changed`] returns `true` to get the new tag.
84+
/// For the other variants the tag never changes.
85+
pub(crate) fn classify_with_generation_tag(
86+
predicate: &Arc<dyn PhysicalExpr>,
87+
) -> (Self, u64) {
7288
let mut subscriptions = Vec::new();
89+
let mut completed_generations = 0u64;
7390
let mut found_any = false;
7491
predicate
7592
.apply(|expr| {
7693
if let Some(filter) = expr.downcast_ref::<DynamicFilterPhysicalExpr>() {
7794
found_any = true;
7895
// Already-complete filters can never change again, so there
79-
// is no point subscribing to them.
80-
if !filter.is_complete() {
96+
// is no point subscribing to them. Keep their (final)
97+
// generation for the generation tag.
98+
if filter.is_complete() {
99+
completed_generations = completed_generations
100+
.wrapping_add(filter.current_generation());
101+
} else {
81102
subscriptions.push(filter.subscribe());
82103
}
83104
}
84105
Ok(TreeNodeRecursion::Continue)
85106
})
86107
.expect("traversal closure is infallible");
87108

88-
if !found_any {
109+
let tracker = DynamicFilterTracker {
110+
subscriptions,
111+
completed_generations,
112+
};
113+
let generation_tag = tracker.generation_tag();
114+
let tracking = if !found_any {
89115
DynamicFilterTracking::Static
90-
} else if subscriptions.is_empty() {
116+
} else if tracker.subscriptions.is_empty() {
91117
DynamicFilterTracking::AllComplete
92118
} else {
93-
DynamicFilterTracking::Watching(DynamicFilterTracker { subscriptions })
94-
}
119+
DynamicFilterTracking::Watching(tracker)
120+
};
121+
(tracking, generation_tag)
95122
}
96123

97124
/// `true` if the predicate contains any dynamic filter (complete or not),
@@ -123,6 +150,10 @@ pub struct DynamicFilterTracker {
123150
/// Subscriptions to the not-yet-complete dynamic filters. Entries are
124151
/// dropped as their filters complete, so the set only shrinks.
125152
subscriptions: Vec<DynamicFilterSubscription>,
153+
/// Wrapping sum of the final generations of the complete filters: the
154+
/// filters that were complete at classification time and the filters
155+
/// whose subscriptions were dropped. See [`Self::generation_tag`].
156+
completed_generations: u64,
126157
}
127158

128159
impl DynamicFilterTracker {
@@ -134,14 +165,41 @@ impl DynamicFilterTracker {
134165
/// returns `false`.
135166
pub fn changed(&mut self) -> bool {
136167
let mut changed = false;
168+
let completed_generations = &mut self.completed_generations;
137169
self.subscriptions.retain_mut(|subscription| {
138170
let change = subscription.observe();
139171
changed |= change.changed;
172+
if change.complete {
173+
// Keep the final generation so the generation tag does not
174+
// change when the subscription is dropped.
175+
*completed_generations =
176+
completed_generations.wrapping_add(subscription.last_generation());
177+
}
140178
// Keep the subscription only while the filter can still change.
141179
!change.complete
142180
});
143181
changed
144182
}
183+
184+
/// A tag that identifies the generations of the watched filters, as
185+
/// observed by the last call to [`Self::changed`] (or at classification
186+
/// time, before the first call).
187+
///
188+
/// The tag is the wrapping sum of the latest observed generation of every
189+
/// dynamic filter in the predicate, complete or not. It is equal to the
190+
/// tag that [`DynamicFilterTracking::classify_with_generation_tag`]
191+
/// returns for the same predicate, if no filter changed since the last
192+
/// observation. It does not change when a filter completes without an
193+
/// update. A generation only increases, thus in practice two different
194+
/// tags mean that a filter changed. Computing it does not walk the
195+
/// predicate: it only reads the watched subscriptions.
196+
pub(crate) fn generation_tag(&self) -> u64 {
197+
self.subscriptions
198+
.iter()
199+
.fold(self.completed_generations, |tag, subscription| {
200+
tag.wrapping_add(subscription.last_generation())
201+
})
202+
}
145203
}
146204

147205
#[cfg(test)]
@@ -157,7 +215,7 @@ impl DynamicFilterTracker {
157215
}
158216

159217
/// `true` once every watched filter has completed and been dropped.
160-
fn is_exhausted(&self) -> bool {
218+
pub(crate) fn is_exhausted(&self) -> bool {
161219
self.subscriptions.is_empty()
162220
}
163221
}
@@ -328,4 +386,74 @@ mod tests {
328386
assert!(!tracker.changed());
329387
assert!(tracker.is_exhausted());
330388
}
389+
390+
#[test]
391+
fn generation_tag_of_static_predicate_is_zero() {
392+
let (tracking, tag) =
393+
DynamicFilterTracking::classify_with_generation_tag(&lit(true));
394+
assert!(matches!(tracking, DynamicFilterTracking::Static));
395+
assert_eq!(tag, 0);
396+
}
397+
398+
#[test]
399+
fn generation_tag_follows_observed_updates() {
400+
let (predicate, filter) = dynamic_predicate();
401+
let (mut tracking, initial_tag) =
402+
DynamicFilterTracking::classify_with_generation_tag(&predicate);
403+
let tracker = tracking.watcher().unwrap();
404+
assert_eq!(tracker.generation_tag(), initial_tag);
405+
406+
// The tag changes only when the tracker observes the update.
407+
filter.update(lit(false)).unwrap();
408+
assert_eq!(tracker.generation_tag(), initial_tag);
409+
assert!(tracker.changed());
410+
let updated_tag = tracker.generation_tag();
411+
assert_ne!(updated_tag, initial_tag);
412+
413+
// A new classification of the same predicate gets the same tag.
414+
let (_, tag) = DynamicFilterTracking::classify_with_generation_tag(&predicate);
415+
assert_eq!(tag, updated_tag);
416+
}
417+
418+
#[test]
419+
fn generation_tag_is_stable_when_filter_completes() {
420+
let (predicate, filter) = dynamic_predicate();
421+
let (mut tracking, _) =
422+
DynamicFilterTracking::classify_with_generation_tag(&predicate);
423+
let tracker = tracking.watcher().unwrap();
424+
425+
filter.update(lit(false)).unwrap();
426+
assert!(tracker.changed());
427+
let tag = tracker.generation_tag();
428+
429+
// Completion drops the subscription but keeps its final generation.
430+
filter.mark_complete();
431+
assert!(!tracker.changed());
432+
assert!(tracker.is_exhausted());
433+
assert_eq!(tracker.generation_tag(), tag);
434+
435+
// A predicate classified after completion gets the same tag.
436+
let (tracking, complete_tag) =
437+
DynamicFilterTracking::classify_with_generation_tag(&predicate);
438+
assert!(matches!(tracking, DynamicFilterTracking::AllComplete));
439+
assert_eq!(complete_tag, tag);
440+
}
441+
442+
#[test]
443+
fn generation_tag_of_coalesced_update_and_complete() {
444+
let (predicate, filter) = dynamic_predicate();
445+
let (mut tracking, initial_tag) =
446+
DynamicFilterTracking::classify_with_generation_tag(&predicate);
447+
let tracker = tracking.watcher().unwrap();
448+
449+
filter.update(lit(false)).unwrap();
450+
filter.mark_complete();
451+
assert!(tracker.changed());
452+
assert!(tracker.is_exhausted());
453+
assert_ne!(tracker.generation_tag(), initial_tag);
454+
455+
let (_, complete_tag) =
456+
DynamicFilterTracking::classify_with_generation_tag(&predicate);
457+
assert_eq!(tracker.generation_tag(), complete_tag);
458+
}
331459
}

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -36,6 +36,7 @@ pub mod equivalence;
3636
pub mod expressions;
3737
pub mod higher_order_function;
3838
pub mod intervals;
39+
pub mod optional_filter_gate;
3940
mod partitioning;
4041
mod physical_expr;
4142
pub mod planner;

0 commit comments

Comments
 (0)