use std::collections::{BTreeMap, HashSet};
use std::sync::Arc;
use datafusion::arrow::datatypes::SchemaRef;
use datafusion::physical_expr::expressions::{Column, DynamicFilterPhysicalExpr};
use datafusion::physical_plan::PhysicalExpr;
#[derive(Clone, Debug)]
pub(crate) struct DeltaRetainedDynamicFilter {
pub(crate) physical_expr: Arc<dyn PhysicalExpr>,
pub(crate) partition_columns: Vec<DeltaDynamicFilterColumn>,
pub(crate) provider_schema: SchemaRef,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub(crate) struct DeltaDynamicFilterColumn {
pub(crate) name: String,
pub(crate) index: usize,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) enum DeltaDynamicFilterOutcome {
Accepted,
Rejected,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) enum DeltaDynamicFilterRejectionReason {
NotDynamicFilter,
NoReferencedColumns,
InternalColumn,
UnknownColumn,
DataColumn,
MixedPartitionAndData,
}
#[derive(Clone, Debug)]
pub(crate) struct DeltaDynamicFilterDecision {
pub(crate) outcome: DeltaDynamicFilterOutcome,
pub(crate) retained_filter: Option<DeltaRetainedDynamicFilter>,
#[allow(dead_code)]
pub(crate) rejection_reason: Option<DeltaDynamicFilterRejectionReason>,
}
#[derive(Clone, Debug, Default)]
pub(crate) struct DeltaDynamicFilterPlan {
pub(crate) decisions: Vec<DeltaDynamicFilterDecision>,
pub(crate) accepted_filters: Vec<DeltaRetainedDynamicFilter>,
}
impl DeltaDynamicFilterPlan {
#[must_use]
pub(crate) fn from_filters(
filters: &[Arc<dyn PhysicalExpr>],
provider_schema: &SchemaRef,
partition_columns: &[String],
) -> Self {
let decisions = filters
.iter()
.map(|filter| classify_dynamic_filter(filter, provider_schema, partition_columns))
.collect::<Vec<_>>();
let accepted_filters = decisions
.iter()
.filter_map(|decision| decision.retained_filter.clone())
.collect();
Self {
decisions,
accepted_filters,
}
}
#[must_use]
pub(crate) fn has_accepted_filters(&self) -> bool {
!self.accepted_filters.is_empty()
}
}
fn classify_dynamic_filter(
filter: &Arc<dyn PhysicalExpr>,
provider_schema: &SchemaRef,
partition_columns: &[String],
) -> DeltaDynamicFilterDecision {
if !contains_dynamic_filter(filter.as_ref()) {
return rejected(DeltaDynamicFilterRejectionReason::NotDynamicFilter);
}
let references = collect_column_references(filter.as_ref(), provider_schema);
if references.has_internal_column {
return rejected(DeltaDynamicFilterRejectionReason::InternalColumn);
}
if references.has_unknown_column {
return rejected(DeltaDynamicFilterRejectionReason::UnknownColumn);
}
if references.columns.is_empty() {
return rejected(DeltaDynamicFilterRejectionReason::NoReferencedColumns);
}
let partition_column_set = partition_columns
.iter()
.map(String::as_str)
.collect::<HashSet<_>>();
let (partition_columns, data_column_count): (Vec<_>, usize) =
references.columns.into_values().fold(
(Vec::new(), 0),
|(mut partition_columns, data_count), column| {
if partition_column_set.contains(column.name.as_str()) {
partition_columns.push(column);
(partition_columns, data_count)
} else {
(partition_columns, data_count + 1)
}
},
);
match (partition_columns.is_empty(), data_column_count == 0) {
(false, true) => DeltaDynamicFilterDecision {
outcome: DeltaDynamicFilterOutcome::Accepted,
retained_filter: Some(DeltaRetainedDynamicFilter {
physical_expr: Arc::clone(filter),
partition_columns,
provider_schema: Arc::clone(provider_schema),
}),
rejection_reason: None,
},
(false, false) => rejected(DeltaDynamicFilterRejectionReason::MixedPartitionAndData),
(true, false) => rejected(DeltaDynamicFilterRejectionReason::DataColumn),
(true, true) => rejected(DeltaDynamicFilterRejectionReason::NoReferencedColumns),
}
}
fn rejected(reason: DeltaDynamicFilterRejectionReason) -> DeltaDynamicFilterDecision {
DeltaDynamicFilterDecision {
outcome: DeltaDynamicFilterOutcome::Rejected,
retained_filter: None,
rejection_reason: Some(reason),
}
}
fn contains_dynamic_filter(expr: &dyn PhysicalExpr) -> bool {
expr.is::<DynamicFilterPhysicalExpr>()
|| expr
.children()
.into_iter()
.any(|child| contains_dynamic_filter(child.as_ref()))
}
#[derive(Default)]
struct ColumnReferences {
columns: BTreeMap<(usize, String), DeltaDynamicFilterColumn>,
has_internal_column: bool,
has_unknown_column: bool,
}
fn collect_column_references(
expr: &dyn PhysicalExpr,
provider_schema: &SchemaRef,
) -> ColumnReferences {
let mut references = ColumnReferences::default();
collect_column_references_into(expr, provider_schema, &mut references);
references
}
fn collect_column_references_into(
expr: &dyn PhysicalExpr,
provider_schema: &SchemaRef,
references: &mut ColumnReferences,
) {
if let Some(column) = expr.downcast_ref::<Column>() {
collect_column_reference(column, provider_schema, references);
}
for child in expr.children() {
collect_column_references_into(child.as_ref(), provider_schema, references);
}
}
fn collect_column_reference(
column: &Column,
provider_schema: &SchemaRef,
references: &mut ColumnReferences,
) {
if column.name().starts_with("__delta_arrow_reader_")
|| column.name().starts_with("__delta_funnel_")
{
references.has_internal_column = true;
return;
}
let Some(field) = provider_schema.fields().get(column.index()) else {
references.has_unknown_column = true;
return;
};
if field.name() != column.name() {
references.has_unknown_column = true;
return;
}
references.columns.insert(
(column.index(), column.name().to_owned()),
DeltaDynamicFilterColumn {
name: column.name().to_owned(),
index: column.index(),
},
);
}
#[cfg(test)]
mod tests {
use datafusion::arrow::datatypes::{DataType, Field, Schema};
use datafusion::logical_expr::Operator;
use datafusion::physical_expr::expressions::{BinaryExpr, lit};
use super::*;
fn test_schema() -> SchemaRef {
Arc::new(Schema::new(vec![
Field::new("id", DataType::Int32, false),
Field::new("customer_name", DataType::Utf8, true),
Field::new("region", DataType::Utf8, true),
Field::new("event_date", DataType::Date32, true),
]))
}
fn column(name: &str, index: usize) -> Arc<dyn PhysicalExpr> {
Arc::new(Column::new(name, index))
}
fn dynamic_filter(children: Vec<Arc<dyn PhysicalExpr>>) -> Arc<dyn PhysicalExpr> {
Arc::new(DynamicFilterPhysicalExpr::new(children, lit(true)))
}
fn plan_for(filter: Arc<dyn PhysicalExpr>) -> DeltaDynamicFilterPlan {
DeltaDynamicFilterPlan::from_filters(
&[filter],
&test_schema(),
&["region".to_owned(), "event_date".to_owned()],
)
}
#[test]
fn partition_dynamic_filter_is_accepted() {
let plan = plan_for(dynamic_filter(vec![column("region", 2)]));
assert!(plan.has_accepted_filters());
assert_eq!(
plan.decisions[0].outcome,
DeltaDynamicFilterOutcome::Accepted
);
assert_eq!(
plan.accepted_filters[0].partition_columns,
vec![DeltaDynamicFilterColumn {
name: "region".to_owned(),
index: 2,
}]
);
}
#[test]
fn no_column_dynamic_filter_is_rejected() {
let plan = plan_for(dynamic_filter(Vec::new()));
assert!(!plan.has_accepted_filters());
assert_eq!(
plan.decisions[0].rejection_reason,
Some(DeltaDynamicFilterRejectionReason::NoReferencedColumns)
);
}
#[test]
fn multi_partition_dynamic_filter_retains_sorted_column_mappings() {
let plan = plan_for(dynamic_filter(vec![
column("event_date", 3),
column("region", 2),
]));
assert_eq!(
plan.accepted_filters[0].partition_columns,
vec![
DeltaDynamicFilterColumn {
name: "region".to_owned(),
index: 2,
},
DeltaDynamicFilterColumn {
name: "event_date".to_owned(),
index: 3,
},
]
);
}
#[test]
fn data_column_dynamic_filter_is_rejected() {
let plan = plan_for(dynamic_filter(vec![column("id", 0)]));
assert!(!plan.has_accepted_filters());
assert_eq!(
plan.decisions[0].rejection_reason,
Some(DeltaDynamicFilterRejectionReason::DataColumn)
);
}
#[test]
fn unknown_column_dynamic_filter_is_rejected() {
let plan = plan_for(dynamic_filter(vec![column("ghost", 99)]));
assert!(!plan.has_accepted_filters());
assert_eq!(
plan.decisions[0].rejection_reason,
Some(DeltaDynamicFilterRejectionReason::UnknownColumn)
);
}
#[test]
fn mixed_partition_and_data_dynamic_filter_is_rejected() {
let plan = plan_for(dynamic_filter(vec![column("region", 2), column("id", 0)]));
assert!(!plan.has_accepted_filters());
assert_eq!(
plan.decisions[0].rejection_reason,
Some(DeltaDynamicFilterRejectionReason::MixedPartitionAndData)
);
}
#[test]
fn dynamic_filter_wrapped_with_data_column_is_rejected_as_mixed() {
let dynamic = dynamic_filter(vec![column("region", 2)]);
let wrapped = Arc::new(BinaryExpr::new(dynamic, Operator::And, column("id", 0)));
let plan = plan_for(wrapped);
assert!(!plan.has_accepted_filters());
assert_eq!(
plan.decisions[0].rejection_reason,
Some(DeltaDynamicFilterRejectionReason::MixedPartitionAndData)
);
}
#[test]
fn non_dynamic_filter_is_rejected() {
let filter = Arc::new(BinaryExpr::new(
column("region", 2),
Operator::Eq,
lit("us-west"),
));
let plan = plan_for(filter);
assert!(!plan.has_accepted_filters());
assert_eq!(
plan.decisions[0].rejection_reason,
Some(DeltaDynamicFilterRejectionReason::NotDynamicFilter)
);
}
#[test]
fn internal_column_dynamic_filter_is_rejected() {
let plan = plan_for(dynamic_filter(vec![column("__delta_funnel_row_index", 0)]));
assert!(!plan.has_accepted_filters());
assert_eq!(
plan.decisions[0].rejection_reason,
Some(DeltaDynamicFilterRejectionReason::InternalColumn)
);
}
}