use std::sync::Arc;
use arrow::datatypes::{FieldRef, SchemaRef};
use datafusion_common::{Result, internal_datafusion_err, pruning::PrunableStatistics};
use datafusion_datasource::PartitionedFile;
use datafusion_physical_expr::DynamicFilterTracking;
use datafusion_physical_expr_common::physical_expr::PhysicalExpr;
use datafusion_physical_plan::metrics::Count;
use log::debug;
use crate::build_pruning_predicate;
pub struct FilePruner {
predicate: Arc<dyn PhysicalExpr>,
tracking: DynamicFilterTracking,
checked_once: bool,
file_schema: SchemaRef,
file_stats_pruning: PrunableStatistics,
predicate_creation_errors: Count,
}
impl FilePruner {
#[deprecated(
since = "52.0.0",
note = "Use `try_new` instead which returns None if no statistics are available"
)]
#[expect(clippy::needless_pass_by_value)]
pub fn new(
predicate: Arc<dyn PhysicalExpr>,
logical_file_schema: &SchemaRef,
_partition_fields: Vec<FieldRef>,
partitioned_file: PartitionedFile,
predicate_creation_errors: Count,
) -> Result<Self> {
Self::try_new(
predicate,
logical_file_schema,
&partitioned_file,
predicate_creation_errors,
)
.ok_or_else(|| {
internal_datafusion_err!(
"FilePruner::new called on a file without statistics: {:?}",
partitioned_file
)
})
}
pub fn try_new(
predicate: Arc<dyn PhysicalExpr>,
file_schema: &SchemaRef,
partitioned_file: &PartitionedFile,
predicate_creation_errors: Count,
) -> Option<Self> {
let file_stats = partitioned_file.statistics.as_ref()?;
let tracking = DynamicFilterTracking::classify(&predicate);
if !partitioned_file.has_statistics() && !tracking.contains_dynamic_filter() {
return None;
}
let file_stats_pruning =
PrunableStatistics::new(vec![file_stats.clone()], Arc::clone(file_schema));
Some(Self {
predicate,
tracking,
checked_once: false,
file_schema: Arc::clone(file_schema),
file_stats_pruning,
predicate_creation_errors,
})
}
pub fn is_watching(&self) -> bool {
matches!(self.tracking, DynamicFilterTracking::Watching(_))
}
pub fn should_prune(&mut self) -> Result<bool> {
let should_build = if self.checked_once {
self.tracking.watcher().is_some_and(|w| w.changed())
} else {
self.checked_once = true;
true
};
if !should_build {
return Ok(false);
}
let pruning_predicate = build_pruning_predicate(
Arc::clone(&self.predicate),
&self.file_schema,
&self.predicate_creation_errors,
);
let Some(pruning_predicate) = pruning_predicate else {
return Ok(false);
};
match pruning_predicate.prune(&self.file_stats_pruning) {
Ok(values) => {
assert!(values.len() == 1);
if values.into_iter().all(|v| !v) {
return Ok(true);
}
}
Err(e) => {
debug!("Ignoring error building pruning predicate for file: {e}");
self.predicate_creation_errors.add(1);
}
}
Ok(false)
}
}