Skip to content
186 changes: 146 additions & 40 deletions datafusion/physical-plan/src/filter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1125,39 +1125,44 @@ fn unique_match_limit(
predicate: &Arc<dyn PhysicalExpr>,
statistics: &Statistics,
) -> Option<usize> {
let mut limit: Option<usize> = None;
for expr in split_conjunction(predicate) {
let Some((index, values)) = restricted_column(expr) else {
continue;
};
let holds_once = statistics
let holds_once_fn = |index: usize| {
statistics
.column_statistics
.get(index)
.is_some_and(|column| holds_each_value_once(column, &statistics.num_rows));
if !holds_once {
continue;
}
limit = Some(limit.map_or(values, |limit: usize| limit.min(values)));
}
limit
.is_some_and(|column| holds_each_value_once(column, &statistics.num_rows))
};
split_conjunction(predicate)
.into_iter()
.filter_map(|expr| restricted_column(expr, holds_once_fn))
.min()
}

/// The column an expression restricts to a fixed set of values, and how many values
/// that is. NULL is never one of them: it matches nothing.
fn restricted_column(expr: &Arc<dyn PhysicalExpr>) -> Option<(usize, usize)> {
/// How many values an expression restricts a column to, when it restricts a column
/// for which `holds_once_fn` is true to a fixed set of values. NULL is never one of
/// them: it matches nothing.
///
/// The column is checked before the values are counted, so IN lists on columns that
/// may hold a value more than once are never scanned.
fn restricted_column(
expr: &Arc<dyn PhysicalExpr>,
holds_once_fn: impl Fn(usize) -> bool,
) -> Option<usize> {
if let Some(in_list) = expr.downcast_ref::<InListExpr>() {
if in_list.negated() {
return None;
}
let column = in_list.expr().downcast_ref::<Column>()?;
let mut values: Vec<&ScalarValue> = vec![];
if !holds_once_fn(column.index()) {
return None;
}
let mut values: HashSet<&ScalarValue> = HashSet::new();
for expr in in_list.list() {
let value = expr.downcast_ref::<Literal>()?.value();
if !value.is_null() && !values.contains(&value) {
values.push(value);
if !value.is_null() {
values.insert(value);
}
}
return Some((column.index(), values.len()));
return Some(values.len());
}

let binary = expr.downcast_ref::<BinaryExpr>()?;
Expand All @@ -1173,7 +1178,7 @@ fn restricted_column(expr: &Arc<dyn PhysicalExpr>) -> Option<(usize, usize)> {
_ => return None,
};
let value = literal.downcast_ref::<Literal>()?.value();
(!value.is_null()).then_some((column.index(), 1))
(!value.is_null() && holds_once_fn(column.index())).then_some(1)
}

/// Whether the column has as many distinct values as it has non-null rows, so each
Expand Down Expand Up @@ -1921,11 +1926,11 @@ mod tests {
}

/// An equality on a column that holds each value once matches one row at most,
/// including on a type interval analysis cannot read, where the default
/// selectivity would otherwise apply.
/// including on a type neither interval analysis nor the string estimates can
/// read, where the default selectivity would otherwise apply.
#[tokio::test]
async fn test_filter_statistics_equality_on_a_unique_column() -> Result<()> {
let schema = Schema::new(vec![Field::new("id", DataType::Utf8, true)]);
let schema = Schema::new(vec![Field::new("id", DataType::Binary, true)]);
let unique = ColumnStatistics {
null_count: Precision::Exact(0),
distinct_count: Precision::Exact(100),
Expand All @@ -1940,8 +1945,12 @@ mod tests {
},
schema.clone(),
));
let predicate =
binary(col("id", &schema)?, Operator::Eq, lit("seven"), &schema)?;
let predicate = binary(
col("id", &schema)?,
Operator::Eq,
lit(ScalarValue::Binary(Some(b"seven".to_vec()))),
&schema,
)?;
let filter: Arc<dyn ExecutionPlan> =
Arc::new(FilterExec::try_new(predicate, input)?);
Ok(StatisticsContext::new()
Expand All @@ -1965,6 +1974,15 @@ mod tests {
assert_eq!(
rows(ColumnStatistics {
distinct_count: Precision::Absent,
..unique.clone()
})?,
Precision::Inexact(20)
);

// Nor without a null count.
assert_eq!(
rows(ColumnStatistics {
null_count: Precision::Absent,
..unique
})?,
Precision::Inexact(20)
Expand All @@ -1973,44 +1991,132 @@ mod tests {
Ok(())
}

/// Asking a unique column for three values matches three rows at most.
/// Asking a unique column for N distinct values matches N rows at most.
#[tokio::test]
async fn test_filter_statistics_in_list_on_a_unique_column() -> Result<()> {
use datafusion_physical_expr::expressions::in_list;

let null = || lit(ScalarValue::Utf8(None));

let schema = Schema::new(vec![Field::new("id", DataType::Utf8, true)]);
let rows = |list: Vec<&str>, negated: bool| -> Result<Precision<usize>> {
let rows = |list: Vec<Arc<dyn PhysicalExpr>>,
negated: bool|
-> Result<Precision<usize>> {
let input = Arc::new(StatisticsExec::new(
Statistics {
num_rows: Precision::Exact(100),
total_byte_size: Precision::Exact(800),
num_rows: Precision::Exact(100_000),
total_byte_size: Precision::Absent,
column_statistics: vec![ColumnStatistics {
null_count: Precision::Exact(0),
distinct_count: Precision::Exact(100),
distinct_count: Precision::Exact(100_000),
..Default::default()
}],
},
schema.clone(),
));
let predicate = in_list(
col("id", &schema)?,
list.into_iter().map(|value| lit(value) as _).collect(),
&negated,
&schema,
)?;
let predicate = in_list(col("id", &schema)?, list, &negated, &schema)?;
let filter: Arc<dyn ExecutionPlan> =
Arc::new(FilterExec::try_new(predicate, input)?);
Ok(StatisticsContext::new()
.compute(filter.as_ref(), &StatisticsArgs::new())?
.num_rows)
};

assert_eq!(rows(vec!["a", "b", "c"], false)?, Precision::Inexact(3));
assert_eq!(
rows(vec![lit("a"), lit("b"), lit("c")], false)?,
Precision::Inexact(3)
);
// Repeats ask for the same row twice.
assert_eq!(rows(vec!["a", "b", "a"], false)?, Precision::Inexact(2));
assert_eq!(
rows(vec![lit("a"), lit("b"), lit("a")], false)?,
Precision::Inexact(2)
);
// NULL matches nothing.
assert_eq!(rows(vec![null(); 32], false)?, Precision::Inexact(0));
// `NOT IN` selects nearly everything, so the default applies.
assert_eq!(rows(vec!["a", "b", "c"], true)?, Precision::Inexact(20));
assert_eq!(
rows(vec![lit("a"), lit("b"), lit("c")], true)?,
Precision::Inexact(20_000)
);
// A list member that is not a literal leaves the set open.
assert_eq!(
rows(vec![lit("a"), col("id", &schema)?], false)?,
Precision::Inexact(20_000)
);

// A large list, with non-adjacent repeats and NULLs that must not raise the
// cap.
let distinct_values = 4096;
let values: Vec<_> = (0..distinct_values)
.map(|i| lit(format!("value_{i}")))
.collect();
let repeated_values = values
.iter()
.chain(values.iter().rev())
.cloned()
.chain(std::iter::repeat_n(null(), 32))
.collect();
assert_eq!(rows(values, false)?, Precision::Inexact(distinct_values));
assert_eq!(
rows(repeated_values, false)?,
Precision::Inexact(distinct_values)
);

Ok(())
}

/// An IN list caps the estimate only when the column it restricts holds each
/// value once; a unique column elsewhere in the table does not count.
#[test]
fn test_filter_statistics_in_list_requires_unique_target_column() -> Result<()> {
use datafusion_physical_expr::expressions::in_list;

let schema = Schema::new(vec![
Field::new("other", DataType::Utf8, true),
Field::new("id", DataType::Utf8, true),
]);
let unique = ColumnStatistics {
null_count: Precision::Exact(0),
distinct_count: Precision::Exact(100_000),
..Default::default()
};
let non_unique = ColumnStatistics {
distinct_count: Precision::Exact(50_000),
..unique.clone()
};
let cases = [
(
"only an unrelated column is unique",
vec![unique.clone(), non_unique.clone()],
20_000,
),
("the second column is unique", vec![non_unique, unique], 3),
];

for (description, column_statistics, expected_rows) in cases {
let input = Arc::new(StatisticsExec::new(
Statistics {
num_rows: Precision::Exact(100_000),
total_byte_size: Precision::Absent,
column_statistics,
},
schema.clone(),
));
let predicate = in_list(
col("id", &schema)?,
vec![lit("a"), lit("b"), lit("a"), lit("c")],
&false,
&schema,
)?;
let filter = FilterExec::try_new(predicate, input)?;
let statistics =
StatisticsContext::new().compute(&filter, &StatisticsArgs::new())?;
assert_eq!(
statistics.num_rows,
Precision::Inexact(expected_rows),
"{description}"
);
}
Ok(())
}

Expand Down
Loading