From 62ce9cd922dbefe74a0d45c4dfe66bf267205f21 Mon Sep 17 00:00:00 2001 From: Victor Date: Wed, 30 Sep 2026 15:54:23 +0200 Subject: [PATCH 1/5] fix: make restricted_column run linearly --- datafusion/physical-plan/src/filter.rs | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/datafusion/physical-plan/src/filter.rs b/datafusion/physical-plan/src/filter.rs index bc1c1ce5d1c96..92671082185b7 100644 --- a/datafusion/physical-plan/src/filter.rs +++ b/datafusion/physical-plan/src/filter.rs @@ -1143,11 +1143,11 @@ fn restricted_column(expr: &Arc) -> Option<(usize, usize)> { return None; } let column = in_list.expr().downcast_ref::()?; - let mut values: Vec<&ScalarValue> = vec![]; + let mut values: HashSet<&ScalarValue> = HashSet::new(); for expr in in_list.list() { let value = expr.downcast_ref::()?.value(); - if !value.is_null() && !values.contains(&value) { - values.push(value); + if !value.is_null() { + values.insert(value); } } return Some((column.index(), values.len())); From a057886bf0661d2d9e5945a8cff291e038d2ff2e Mon Sep 17 00:00:00 2001 From: Victor Date: Wed, 30 Sep 2026 15:56:00 +0200 Subject: [PATCH 2/5] feat: add early return in unique_match_limit --- datafusion/physical-plan/src/filter.rs | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/datafusion/physical-plan/src/filter.rs b/datafusion/physical-plan/src/filter.rs index 92671082185b7..81f59cf818cbb 100644 --- a/datafusion/physical-plan/src/filter.rs +++ b/datafusion/physical-plan/src/filter.rs @@ -1118,6 +1118,13 @@ fn unique_match_limit( predicate: &Arc, statistics: &Statistics, ) -> Option { + if !statistics + .column_statistics + .iter() + .any(|column| holds_each_value_once(column, &statistics.num_rows)) + { + return None; + } let mut limit: Option = None; for expr in split_conjunction(predicate) { let Some((index, values)) = restricted_column(expr) else { From cb25016461a5333abf3f3666275e49719d5f6df8 Mon Sep 17 00:00:00 2001 From: Victor Date: Wed, 30 Sep 2026 16:20:50 +0200 Subject: [PATCH 3/5] chore: add tests --- datafusion/physical-plan/src/filter.rs | 150 +++++++++++++++++++++++++ 1 file changed, 150 insertions(+) diff --git a/datafusion/physical-plan/src/filter.rs b/datafusion/physical-plan/src/filter.rs index 81f59cf818cbb..f216b8dca4443 100644 --- a/datafusion/physical-plan/src/filter.rs +++ b/datafusion/physical-plan/src/filter.rs @@ -1874,6 +1874,156 @@ mod tests { Ok(()) } + #[test] + fn test_filter_statistics_large_in_list() -> Result<()> { + use datafusion_physical_expr::expressions::in_list; + + let schema = Schema::new(vec![Field::new("id", DataType::Utf8, true)]); + let distinct_values = 4096; + let values: Vec> = (0..distinct_values) + .map(|i| lit(format!("value_{i}")) as _) + .collect(); + // Non-adjacent duplicates and NULLs must not increase the row cap. + let repeated_values = values + .iter() + .chain(values.iter().rev()) + .cloned() + .chain(std::iter::repeat_n(lit(ScalarValue::Utf8(None)) as _, 32)) + .collect(); + let cases = [ + ("distinct values", values, false, distinct_values), + ( + "duplicates and nulls", + repeated_values, + false, + distinct_values, + ), + ( + "only nulls", + vec![lit(ScalarValue::Utf8(None)); 32], + false, + 0, + ), + ( + "non-literal list member", + vec![lit("value_0"), col("id", &schema)?], + false, + 20_000, + ), + ("negated list", vec![lit("value_0")], true, 20_000), + ]; + + for (description, list, negated, expected_rows) in cases { + let input = Arc::new(StatisticsExec::new( + Statistics { + 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_000), + ..Default::default() + }], + }, + schema.clone(), + )); + let predicate = in_list(col("id", &schema)?, list, &negated, &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(()) + } + + #[test] + fn test_filter_statistics_in_list_uniqueness_precheck() -> 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 = [ + ( + "no unique columns", + vec![non_unique.clone(), non_unique.clone()], + 20_000, + ), + ( + "unknown statistics", + vec![ColumnStatistics::new_unknown(); 2], + 20_000, + ), + ( + "missing distinct count", + vec![ + non_unique.clone(), + ColumnStatistics { + distinct_count: Precision::Absent, + ..unique.clone() + }, + ], + 20_000, + ), + ( + "missing null count", + vec![ + non_unique.clone(), + ColumnStatistics { + null_count: Precision::Absent, + ..unique.clone() + }, + ], + 20_000, + ), + ( + "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(()) + } + #[tokio::test] async fn test_filter_statistics_basic_expr() -> Result<()> { // Table: From 247c95617e4cd9b27f67c753ca4ae821e1d66de6 Mon Sep 17 00:00:00 2001 From: Victor Date: Fri, 9 Oct 2026 12:29:29 +0200 Subject: [PATCH 4/5] fix: more precise check --- datafusion/physical-plan/src/filter.rs | 48 ++++++++++++-------------- 1 file changed, 23 insertions(+), 25 deletions(-) diff --git a/datafusion/physical-plan/src/filter.rs b/datafusion/physical-plan/src/filter.rs index 8fc663c099084..97aa2ba3408c2 100644 --- a/datafusion/physical-plan/src/filter.rs +++ b/datafusion/physical-plan/src/filter.rs @@ -1125,38 +1125,36 @@ fn unique_match_limit( predicate: &Arc, statistics: &Statistics, ) -> Option { - if !statistics - .column_statistics - .iter() - .any(|column| holds_each_value_once(column, &statistics.num_rows)) - { - return None; - } - let mut limit: Option = 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) -> 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, + holds_once_fn: impl Fn(usize) -> bool, +) -> Option { if let Some(in_list) = expr.downcast_ref::() { if in_list.negated() { return None; } let column = in_list.expr().downcast_ref::()?; + 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::()?.value(); @@ -1164,7 +1162,7 @@ fn restricted_column(expr: &Arc) -> Option<(usize, usize)> { values.insert(value); } } - return Some((column.index(), values.len())); + return Some(values.len()); } let binary = expr.downcast_ref::()?; @@ -1180,7 +1178,7 @@ fn restricted_column(expr: &Arc) -> Option<(usize, usize)> { _ => return None, }; let value = literal.downcast_ref::()?.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 From dd7ac259698f11f66dad6d9d0c679e06d7cb7f09 Mon Sep 17 00:00:00 2001 From: Victor Date: Fri, 9 Oct 2026 12:31:54 +0200 Subject: [PATCH 5/5] chore: update tests --- datafusion/physical-plan/src/filter.rs | 169 +++++++++---------------- 1 file changed, 60 insertions(+), 109 deletions(-) diff --git a/datafusion/physical-plan/src/filter.rs b/datafusion/physical-plan/src/filter.rs index 97aa2ba3408c2..6443784197d0f 100644 --- a/datafusion/physical-plan/src/filter.rs +++ b/datafusion/physical-plan/src/filter.rs @@ -1926,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), @@ -1945,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 = Arc::new(FilterExec::try_new(predicate, input)?); Ok(StatisticsContext::new() @@ -1970,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) @@ -1978,31 +1991,30 @@ 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> { + let rows = |list: Vec>, + negated: bool| + -> Result> { 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 = Arc::new(FilterExec::try_new(predicate, input)?); Ok(StatisticsContext::new() @@ -2010,82 +2022,53 @@ mod tests { .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)); - - Ok(()) - } - - #[test] - fn test_filter_statistics_large_in_list() -> Result<()> { - use datafusion_physical_expr::expressions::in_list; + 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) + ); - let schema = Schema::new(vec![Field::new("id", DataType::Utf8, true)]); + // 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}")) as _) + let values: Vec<_> = (0..distinct_values) + .map(|i| lit(format!("value_{i}"))) .collect(); - // Non-adjacent duplicates and NULLs must not increase the row cap. let repeated_values = values .iter() .chain(values.iter().rev()) .cloned() - .chain(std::iter::repeat_n(lit(ScalarValue::Utf8(None)) as _, 32)) + .chain(std::iter::repeat_n(null(), 32)) .collect(); - let cases = [ - ("distinct values", values, false, distinct_values), - ( - "duplicates and nulls", - repeated_values, - false, - distinct_values, - ), - ( - "only nulls", - vec![lit(ScalarValue::Utf8(None)); 32], - false, - 0, - ), - ( - "non-literal list member", - vec![lit("value_0"), col("id", &schema)?], - false, - 20_000, - ), - ("negated list", vec![lit("value_0")], true, 20_000), - ]; + assert_eq!(rows(values, false)?, Precision::Inexact(distinct_values)); + assert_eq!( + rows(repeated_values, false)?, + Precision::Inexact(distinct_values) + ); - for (description, list, negated, expected_rows) in cases { - let input = Arc::new(StatisticsExec::new( - Statistics { - 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_000), - ..Default::default() - }], - }, - schema.clone(), - )); - let predicate = in_list(col("id", &schema)?, list, &negated, &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(()) } + /// 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_uniqueness_precheck() -> Result<()> { + fn test_filter_statistics_in_list_requires_unique_target_column() -> Result<()> { use datafusion_physical_expr::expressions::in_list; let schema = Schema::new(vec![ @@ -2102,38 +2085,6 @@ mod tests { ..unique.clone() }; let cases = [ - ( - "no unique columns", - vec![non_unique.clone(), non_unique.clone()], - 20_000, - ), - ( - "unknown statistics", - vec![ColumnStatistics::new_unknown(); 2], - 20_000, - ), - ( - "missing distinct count", - vec![ - non_unique.clone(), - ColumnStatistics { - distinct_count: Precision::Absent, - ..unique.clone() - }, - ], - 20_000, - ), - ( - "missing null count", - vec![ - non_unique.clone(), - ColumnStatistics { - null_count: Precision::Absent, - ..unique.clone() - }, - ], - 20_000, - ), ( "only an unrelated column is unique", vec![unique.clone(), non_unique.clone()],