use std::cmp::Ordering;
use std::collections::HashSet;
use std::sync::{Arc, LazyLock};
use std::time::Instant;
use tracing::{debug, error};
use crate::actions::visitors::SelectionVectorVisitor;
use crate::actions::{MAX_VALUES, MIN_VALUES, NULL_COUNT, NUM_RECORDS};
use crate::error::DeltaResult;
use crate::expressions::{
col, column_name, column_pred, lit, BinaryPredicateOp, ColumnName, Expression as Expr,
ExpressionRef, JunctionPredicateOp, OpaquePredicateOpRef, Predicate as Pred, PredicateRef,
Scalar,
};
use crate::kernel_predicates::{
DataSkippingPredicateEvaluator, KernelPredicateEvaluator, KernelPredicateEvaluatorDefaults,
};
use crate::scan::data_skipping::stats_schema::is_skipping_eligible_datatype;
use crate::scan::log_replay::PARTITION_VALUES_PARSED_NAME;
use crate::scan::metrics::ScanMetrics;
use crate::schema::{lazy_schema_ref, schema_ref, DataType, PrimitiveType, SchemaRef};
use crate::table_configuration::TableConfiguration;
use crate::utils::require;
use crate::{Engine, EngineData, Error, ExpressionEvaluator, PredicateEvaluator, RowVisitor as _};
pub(crate) mod stats_schema;
#[cfg(test)]
mod tests;
use delta_kernel_derive::internal_api;
#[cfg(test)]
pub(crate) fn as_data_skipping_predicate(pred: &Pred) -> Option<Pred> {
let stats_columns = all_referenced_columns(pred);
DataSkippingPredicateCreator::new(&Default::default(), &stats_columns).eval(pred)
}
#[cfg(test)]
pub(crate) fn all_referenced_columns(pred: &Pred) -> HashSet<ColumnName> {
pred.references().into_iter().cloned().collect()
}
#[cfg(test)]
pub(crate) fn as_data_skipping_predicate_with_partitions(
pred: &Pred,
partition_columns: &HashSet<ColumnName>,
) -> Option<Pred> {
let stats_columns = all_referenced_columns(pred);
DataSkippingPredicateCreator::new(partition_columns, &stats_columns).eval(pred)
}
#[cfg(test)]
fn as_sql_data_skipping_predicate(
pred: &Pred,
partition_columns: &HashSet<ColumnName>,
) -> Option<Pred> {
let stats_columns = all_referenced_columns(pred);
as_sql_data_skipping_predicate_with_stats_columns(pred, partition_columns, &stats_columns)
}
pub(crate) fn as_sql_data_skipping_predicate_with_stats_columns(
pred: &Pred,
partition_columns: &HashSet<ColumnName>,
stats_columns: &HashSet<ColumnName>,
) -> Option<Pred> {
DataSkippingPredicateCreator::new(partition_columns, stats_columns).eval_sql_where(pred)
}
#[internal_api]
pub(crate) struct DataSkippingFilter {
stats_evaluator: Arc<dyn ExpressionEvaluator>,
skipping_evaluator: Arc<dyn PredicateEvaluator>,
filter_evaluator: Arc<dyn PredicateEvaluator>,
metrics: Option<Arc<ScanMetrics>>,
}
impl DataSkippingFilter {
#[allow(clippy::too_many_arguments)]
pub(crate) fn new(
engine: &dyn Engine,
predicate: Option<PredicateRef>,
stats_schema: Option<&SchemaRef>,
stats_expr: ExpressionRef,
partition_schema: Option<&SchemaRef>,
partition_expr: ExpressionRef,
is_add_expr: ExpressionRef,
input_schema: SchemaRef,
stats_columns: &HashSet<ColumnName>,
metrics: Option<Arc<ScanMetrics>>,
) -> Option<Self> {
static FILTER_PRED: LazyLock<PredicateRef> =
LazyLock::new(|| Arc::new(col!("output").distinct(lit(false))));
static FILTER_SCHEMA: LazyLock<SchemaRef> = lazy_schema_ref! {
nullable "output": BOOLEAN,
};
let predicate = predicate?;
debug!("Creating a data skipping filter for {:#?}", predicate);
let (unified_schema, unified_expr, partition_columns) =
Self::build_unified_schema_and_expr(
stats_schema,
stats_expr,
partition_schema,
partition_expr,
is_add_expr,
)?;
let stats_evaluator = engine
.evaluation_handler()
.new_expression_evaluator(
input_schema,
unified_expr,
unified_schema.as_ref().clone().into(),
)
.inspect_err(|e| error!("Failed to create stats evaluator: {e}"))
.ok()?;
let skipping_evaluator = engine
.evaluation_handler()
.new_predicate_evaluator(
unified_schema.clone(),
Arc::new(as_sql_data_skipping_predicate_with_stats_columns(
&predicate,
&partition_columns,
stats_columns,
)?),
)
.inspect_err(|e| error!("Failed to create skipping evaluator: {e}"))
.ok()?;
let filter_evaluator = engine
.evaluation_handler()
.new_predicate_evaluator(FILTER_SCHEMA.clone(), FILTER_PRED.clone())
.inspect_err(|e| error!("Failed to create filter evaluator: {e}"))
.ok()?;
Some(Self {
stats_evaluator,
skipping_evaluator,
filter_evaluator,
metrics,
})
}
pub(crate) fn for_raw_action_batch(
engine: &dyn Engine,
physical_predicate: PredicateRef,
table_configuration: &TableConfiguration,
input_schema: SchemaRef,
) -> Option<Self> {
let predicate_refs: Vec<ColumnName> = physical_predicate
.references()
.into_iter()
.cloned()
.collect();
let physical_stats_columns = table_configuration.physical_stats_columns_set(None);
let physical_stats_schema = table_configuration
.build_expected_stats_schemas(None, Some(&predicate_refs))
.ok()?
.physical;
let partition_schema = table_configuration.predicate_partition_schema(&predicate_refs);
let stats_expr = Arc::new(Expr::parse_json(
col!("add.stats"),
physical_stats_schema.clone(),
));
let partition_expr = Arc::new(Expr::map_to_struct(col!("add.partitionValues")));
let is_add_expr = Arc::new(Pred::is_not_null(col!("add.path")).into());
Self::new(
engine,
Some(physical_predicate),
Some(&physical_stats_schema),
stats_expr,
partition_schema.as_ref(),
partition_expr,
is_add_expr,
input_schema,
&physical_stats_columns,
None,
)
}
fn build_unified_schema_and_expr(
physical_stats_schema: Option<&SchemaRef>,
stats_expr: ExpressionRef,
physical_partition_schema: Option<&SchemaRef>,
partition_expr: ExpressionRef,
is_add_expr: ExpressionRef,
) -> Option<(SchemaRef, ExpressionRef, HashSet<ColumnName>)> {
let partition_columns: HashSet<ColumnName> = physical_partition_schema
.map(|s| {
s.fields()
.filter(|f| {
matches!(
f.data_type(),
DataType::Primitive(primitive)
if is_skipping_eligible_datatype(primitive)
|| matches!(
primitive,
PrimitiveType::Boolean | PrimitiveType::Binary
)
)
})
.map(|f| ColumnName::new([f.name()]))
.collect()
})
.unwrap_or_default();
let unified_schema = match (physical_stats_schema, physical_partition_schema) {
(Some(stats), Some(ps)) => schema_ref! {
nullable "stats_parsed": (stats.as_ref().clone()),
nullable "partitionValues_parsed": (ps.as_ref().clone()),
not_null "is_add": BOOLEAN,
},
(Some(stats), None) => schema_ref! {
nullable "stats_parsed": (stats.as_ref().clone()),
not_null "is_add": BOOLEAN,
},
(None, Some(ps)) => schema_ref! {
nullable "partitionValues_parsed": (ps.as_ref().clone()),
not_null "is_add": BOOLEAN,
},
(None, None) => return None,
};
let unified_expr = match (
physical_stats_schema.is_some(),
physical_partition_schema.is_some(),
) {
(true, true) => Arc::new(Expr::struct_from([stats_expr, partition_expr, is_add_expr])),
(true, false) => Arc::new(Expr::struct_from([stats_expr, is_add_expr])),
(false, true) => Arc::new(Expr::struct_from([partition_expr, is_add_expr])),
(false, false) => return None,
};
Some((unified_schema, unified_expr, partition_columns))
}
pub(crate) fn apply(&self, batch: &dyn EngineData) -> DeltaResult<Vec<bool>> {
let start_time = Instant::now();
let batch_len = batch.len();
let file_stats = self.stats_evaluator.evaluate(batch)?;
require!(
file_stats.len() == batch_len,
Error::internal_error(format!(
"stats evaluator output length {} != batch length {}",
file_stats.len(),
batch_len
))
);
let skipping_predicate = self.skipping_evaluator.evaluate(&*file_stats)?;
require!(
skipping_predicate.len() == batch_len,
Error::internal_error(format!(
"skipping evaluator output length {} != batch length {}",
skipping_predicate.len(),
batch_len
))
);
let selection_vector = self
.filter_evaluator
.evaluate(skipping_predicate.as_ref())?;
debug_assert_eq!(selection_vector.len(), batch_len);
require!(
selection_vector.len() == batch_len,
Error::internal_error(format!(
"filter evaluator output length {} != batch length {}",
selection_vector.len(),
batch_len
))
);
let mut visitor = SelectionVectorVisitor::default();
visitor.visit_rows_of(selection_vector.as_ref())?;
if visitor.num_filtered > 0 {
debug!(
"data skipping filtered {}/{batch_len} rows from batch",
visitor.num_filtered
);
}
if let Some(metrics) = self.metrics.as_ref() {
metrics.add_predicate_filtered(visitor.num_filtered);
metrics.add_predicate_eval_time_ns(start_time.elapsed().as_nanos() as u64)
}
Ok(visitor.selection_vector)
}
}
pub(crate) fn as_checkpoint_skipping_predicate(
pred: &Pred,
physical_partition_columns: &HashSet<ColumnName>,
physical_floating_partition_columns: &HashSet<ColumnName>,
physical_stats_columns: &HashSet<ColumnName>,
) -> Option<Pred> {
CheckpointDataSkippingPredicateCreator {
data_skipping_columns: DataSkippingColumns {
physical_partition_columns,
physical_stats_columns,
},
physical_floating_partition_columns,
}
.eval(pred)
}
fn comparison_predicate(ord: Ordering, col: Expr, val: &Scalar, inverted: bool) -> Pred {
let pred_fn = match (ord, inverted) {
(Ordering::Less, false) => Pred::lt,
(Ordering::Less, true) => Pred::ge,
(Ordering::Equal, false) => Pred::eq,
(Ordering::Equal, true) => Pred::ne,
(Ordering::Greater, false) => Pred::gt,
(Ordering::Greater, true) => Pred::le,
};
pred_fn(col, val.clone())
}
fn collect_junction_preds(
mut op: JunctionPredicateOp,
preds: &mut dyn Iterator<Item = Option<Pred>>,
inverted: bool,
) -> Pred {
if inverted {
op = op.invert();
}
let mut keep_null = true;
let preds: Vec<_> = preds
.flat_map(|p| match p {
Some(pred) => Some(pred),
None => keep_null.then(|| {
keep_null = false;
Pred::NULL
}),
})
.collect();
Pred::junction(op, preds)
}
fn adjust_scalar_for_max_stat_truncation(val: &Scalar) -> Scalar {
match val {
Scalar::Timestamp(micros) => Scalar::Timestamp(micros.saturating_sub(999)),
Scalar::TimestampNtz(micros) => Scalar::TimestampNtz(micros.saturating_sub(999)),
other => other.clone(),
}
}
fn partition_value_expr(col: &ColumnName) -> Expr {
Expr::from(column_name!(PARTITION_VALUES_PARSED_NAME).join(col))
}
fn is_partition_value_reference(expr: &Expr) -> bool {
matches!(expr, Expr::Column(name)
if name.path().first().is_some_and(|f| f == PARTITION_VALUES_PARSED_NAME))
}
fn has_min_max_stats(data_type: &DataType) -> bool {
matches!(data_type, DataType::Primitive(ptype) if is_skipping_eligible_datatype(ptype))
}
struct DataSkippingColumns<'a> {
physical_partition_columns: &'a HashSet<ColumnName>,
physical_stats_columns: &'a HashSet<ColumnName>,
}
impl DataSkippingColumns<'_> {
fn is_partition_column(&self, col: &ColumnName) -> bool {
self.physical_partition_columns.contains(col)
}
fn is_stats_column(&self, col: &ColumnName) -> bool {
self.physical_stats_columns.contains(col)
}
fn min_stat(&self, col: &ColumnName, data_type: &DataType) -> Option<Expr> {
if self.is_partition_column(col) {
Some(partition_value_expr(col))
} else {
(self.is_stats_column(col) && has_min_max_stats(data_type))
.then(|| Expr::from(column_name!("stats_parsed", MIN_VALUES).join(col)))
}
}
fn max_stat(&self, col: &ColumnName, data_type: &DataType) -> Option<Expr> {
if self.is_partition_column(col) {
Some(partition_value_expr(col))
} else {
(self.is_stats_column(col) && has_min_max_stats(data_type))
.then(|| Expr::from(column_name!("stats_parsed", MAX_VALUES).join(col)))
}
}
fn max_stat_comparison_value(&self, col: &ColumnName, val: &Scalar) -> Scalar {
if self.is_partition_column(col) {
val.clone()
} else {
adjust_scalar_for_max_stat_truncation(val)
}
}
fn nullcount_stat(&self, col: &ColumnName) -> Option<Expr> {
if self.is_partition_column(col) {
None
} else {
self.is_stats_column(col)
.then(|| Expr::from(column_name!("stats_parsed", NULL_COUNT).join(col)))
}
}
fn rowcount_stat(&self) -> Expr {
col!("stats_parsed", NUM_RECORDS)
}
}
struct DataSkippingPredicateCreator<'a> {
data_skipping_columns: DataSkippingColumns<'a>,
}
impl<'a> DataSkippingPredicateCreator<'a> {
fn new(
physical_partition_columns: &'a HashSet<ColumnName>,
physical_stats_columns: &'a HashSet<ColumnName>,
) -> Self {
Self {
data_skipping_columns: DataSkippingColumns {
physical_partition_columns,
physical_stats_columns,
},
}
}
fn guard_for_removes(&self, pred: Pred) -> Pred {
Pred::or(Pred::not(column_pred!("is_add")), pred)
}
}
impl DataSkippingPredicateEvaluator for DataSkippingPredicateCreator<'_> {
type Output = Pred;
type ColumnStat = Expr;
fn get_min_stat(&self, col: &ColumnName, data_type: &DataType) -> Option<Expr> {
self.data_skipping_columns.min_stat(col, data_type)
}
fn get_max_stat(&self, col: &ColumnName, data_type: &DataType) -> Option<Expr> {
self.data_skipping_columns.max_stat(col, data_type)
}
fn partial_cmp_max_stat(
&self,
col: &ColumnName,
val: &Scalar,
ord: Ordering,
inverted: bool,
) -> Option<Pred> {
let max = self.get_max_stat(col, &val.data_type())?;
let adjusted = self
.data_skipping_columns
.max_stat_comparison_value(col, val);
self.eval_partial_cmp(ord, max, &adjusted, inverted)
}
fn get_nullcount_stat(&self, col: &ColumnName) -> Option<Expr> {
self.data_skipping_columns.nullcount_stat(col)
}
fn get_rowcount_stat(&self) -> Option<Expr> {
Some(self.data_skipping_columns.rowcount_stat())
}
fn eval_partial_cmp(
&self,
ord: Ordering,
col: Expr,
val: &Scalar,
inverted: bool,
) -> Option<Pred> {
let is_partition = is_partition_value_reference(&col);
let cmp = comparison_predicate(ord, col, val, inverted);
Some(if is_partition {
self.guard_for_removes(cmp)
} else {
cmp
})
}
fn eval_pred_scalar(&self, val: &Scalar, inverted: bool) -> Option<Pred> {
KernelPredicateEvaluatorDefaults::eval_pred_scalar(val, inverted).map(Pred::literal)
}
fn eval_pred_scalar_is_null(&self, val: &Scalar, inverted: bool) -> Option<Pred> {
KernelPredicateEvaluatorDefaults::eval_pred_scalar_is_null(val, inverted).map(Pred::literal)
}
fn eval_pred_is_null(&self, col: &ColumnName, inverted: bool) -> Option<Pred> {
if self.data_skipping_columns.is_partition_column(col) {
let pv_expr = partition_value_expr(col);
let pred = if inverted {
Pred::is_not_null(pv_expr)
} else {
Pred::is_null(pv_expr)
};
Some(self.guard_for_removes(pred))
} else {
let safe_to_skip = match inverted {
true => self.get_rowcount_stat()?, false => lit(0i64), };
Some(Pred::ne(self.get_nullcount_stat(col)?, safe_to_skip))
}
}
fn eval_pred_binary_scalars(
&self,
op: BinaryPredicateOp,
left: &Scalar,
right: &Scalar,
inverted: bool,
) -> Option<Pred> {
KernelPredicateEvaluatorDefaults::eval_pred_binary_scalars(op, left, right, inverted)
.map(Pred::literal)
}
fn eval_pred_cast(
&self,
_op: BinaryPredicateOp,
_col: &ColumnName,
_target: &DataType,
_val: &Scalar,
_inverted: bool,
) -> Option<Pred> {
None
}
fn eval_pred_opaque(
&self,
op: &OpaquePredicateOpRef,
exprs: &[Expr],
inverted: bool,
) -> Option<Pred> {
let pred = op.as_data_skipping_predicate(self, exprs, inverted)?;
Some(self.guard_for_removes(pred))
}
fn finish_eval_pred_junction(
&self,
op: JunctionPredicateOp,
preds: &mut dyn Iterator<Item = Option<Pred>>,
inverted: bool,
) -> Option<Pred> {
Some(collect_junction_preds(op, preds, inverted))
}
}
struct CheckpointDataSkippingPredicateCreator<'a> {
data_skipping_columns: DataSkippingColumns<'a>,
physical_floating_partition_columns: &'a HashSet<ColumnName>,
}
impl CheckpointDataSkippingPredicateCreator<'_> {
fn partition_min_max_may_omit_values(&self, col: &ColumnName) -> bool {
self.physical_floating_partition_columns.contains(col)
}
}
impl DataSkippingPredicateEvaluator for CheckpointDataSkippingPredicateCreator<'_> {
type Output = Pred;
type ColumnStat = Expr;
fn get_min_stat(&self, col: &ColumnName, data_type: &DataType) -> Option<Expr> {
if self.partition_min_max_may_omit_values(col) {
return None;
}
self.data_skipping_columns.min_stat(col, data_type)
}
fn get_max_stat(&self, col: &ColumnName, data_type: &DataType) -> Option<Expr> {
if self.partition_min_max_may_omit_values(col) {
return None;
}
self.data_skipping_columns.max_stat(col, data_type)
}
fn get_nullcount_stat(&self, col: &ColumnName) -> Option<Expr> {
self.data_skipping_columns.nullcount_stat(col)
}
fn get_rowcount_stat(&self) -> Option<Expr> {
Some(self.data_skipping_columns.rowcount_stat())
}
fn partial_cmp_max_stat(
&self,
col: &ColumnName,
val: &Scalar,
ord: Ordering,
inverted: bool,
) -> Option<Pred> {
let max = self.get_max_stat(col, &val.data_type())?;
let adjusted = self
.data_skipping_columns
.max_stat_comparison_value(col, val);
self.eval_partial_cmp(ord, max, &adjusted, inverted)
}
fn eval_partial_cmp(
&self,
ord: Ordering,
col: Expr,
val: &Scalar,
inverted: bool,
) -> Option<Pred> {
let comparison = comparison_predicate(ord, col.clone(), val, inverted);
Some(if is_partition_value_reference(&col) {
comparison
} else {
Pred::or(Pred::is_null(col), comparison)
})
}
fn eval_pred_scalar(&self, val: &Scalar, inverted: bool) -> Option<Pred> {
KernelPredicateEvaluatorDefaults::eval_pred_scalar(val, inverted).map(Pred::literal)
}
fn eval_pred_scalar_is_null(&self, val: &Scalar, inverted: bool) -> Option<Pred> {
KernelPredicateEvaluatorDefaults::eval_pred_scalar_is_null(val, inverted).map(Pred::literal)
}
fn eval_pred_cast(
&self,
op: BinaryPredicateOp,
col: &ColumnName,
target: &DataType,
val: &Scalar,
inverted: bool,
) -> Option<Pred> {
if !self.data_skipping_columns.is_partition_column(col) {
return None;
}
let ord = match op {
BinaryPredicateOp::LessThan => Ordering::Less,
BinaryPredicateOp::Equal => Ordering::Equal,
BinaryPredicateOp::GreaterThan => Ordering::Greater,
BinaryPredicateOp::Distinct | BinaryPredicateOp::In => return None,
};
let cast = Expr::cast(partition_value_expr(col), target.clone());
Some(comparison_predicate(ord, cast, val, inverted))
}
fn eval_pred_is_null(&self, col: &ColumnName, inverted: bool) -> Option<Pred> {
if self.data_skipping_columns.is_partition_column(col) {
let partition_value = partition_value_expr(col);
return Some(if inverted {
Pred::is_not_null(partition_value)
} else {
Pred::is_null(partition_value)
});
}
if inverted {
return None; }
let nullcount = self.get_nullcount_stat(col)?;
let comparison = Pred::ne(nullcount.clone(), lit(0i64));
Some(Pred::or(Pred::is_null(nullcount), comparison))
}
fn eval_pred_binary_scalars(
&self,
op: BinaryPredicateOp,
left: &Scalar,
right: &Scalar,
inverted: bool,
) -> Option<Pred> {
KernelPredicateEvaluatorDefaults::eval_pred_binary_scalars(op, left, right, inverted)
.map(Pred::literal)
}
fn eval_pred_opaque(
&self,
_op: &OpaquePredicateOpRef,
_exprs: &[Expr],
_inverted: bool,
) -> Option<Pred> {
None
}
fn finish_eval_pred_junction(
&self,
op: JunctionPredicateOp,
preds: &mut dyn Iterator<Item = Option<Pred>>,
inverted: bool,
) -> Option<Pred> {
Some(collect_junction_preds(op, preds, inverted))
}
}