From 3083d3e7870bf22ee1e3654f0760211aa14780c3 Mon Sep 17 00:00:00 2001 From: Jake Wallin Date: Mon, 28 Sep 2026 11:46:19 -0400 Subject: [PATCH 1/2] fix: Estimate string equality and LIKE filter selectivity --- datafusion/common/src/rounding.rs | 2 +- .../tests/custom_sources_cases/statistics.rs | 122 +++++++ datafusion/physical-plan/src/filter.rs | 309 +++++++++++++++++- 3 files changed, 425 insertions(+), 8 deletions(-) diff --git a/datafusion/common/src/rounding.rs b/datafusion/common/src/rounding.rs index 1796143d7cf1a..b3b514e8631c8 100644 --- a/datafusion/common/src/rounding.rs +++ b/datafusion/common/src/rounding.rs @@ -254,7 +254,7 @@ where } } _ => {} - }; + } Ok(result) } diff --git a/datafusion/core/tests/custom_sources_cases/statistics.rs b/datafusion/core/tests/custom_sources_cases/statistics.rs index 7701ca0279125..8fb2f0e5d1f7d 100644 --- a/datafusion/core/tests/custom_sources_cases/statistics.rs +++ b/datafusion/core/tests/custom_sources_cases/statistics.rs @@ -285,6 +285,128 @@ async fn sql_filter() -> Result<()> { Ok(()) } +fn string_filter_ctx( + data_type: DataType, + distinct_count: Precision, +) -> Result { + init_ctx( + Statistics { + num_rows: Precision::Exact(1000), + total_byte_size: Precision::Absent, + column_statistics: vec![ColumnStatistics { + null_count: Precision::Exact(200), + distinct_count, + ..ColumnStatistics::new_unknown() + }], + }, + Schema::new(vec![Field::new("c1", data_type, true)]), + ) +} + +async fn string_filter_rows( + ctx: &SessionContext, + predicate: &str, +) -> Result> { + let plan = ctx + .sql(&format!("SELECT * FROM stats_table WHERE {predicate}")) + .await? + .create_physical_plan() + .await?; + Ok(StatisticsContext::new() + .compute(plan.as_ref(), &StatisticsArgs::new())? + .num_rows) +} + +#[tokio::test] +async fn sql_string_filter_selectivity() -> Result<()> { + // There are 800 non-null rows and 20 distinct strings. Equality estimates + // 40 matches, and negated predicates exclude nulls as well as matches. + // LIKE estimates distinguish unanchored patterns from patterns anchored + // at one or both ends; escaped wildcards are literal characters. + let cases = [ + ("c1 = 'foo'", 40), + ("'foo' = c1", 40), + ("c1 != 'foo'", 760), + ("'foo' != c1", 760), + ("c1 LIKE 'foo'", 40), + ("c1 NOT LIKE 'foo'", 760), + ("c1 LIKE '%foo%'", 160), + ("c1 NOT LIKE '%foo%'", 640), + ("c1 LIKE 'foo%'", 80), + ("c1 NOT LIKE 'foo%'", 720), + ("c1 LIKE '%foo'", 80), + ("c1 NOT LIKE '%foo'", 720), + ("c1 LIKE 'f_o'", 40), + ("c1 NOT LIKE 'f_o'", 760), + (r"c1 LIKE 'foo\%'", 40), + (r"c1 NOT LIKE 'foo\%'", 760), + (r"c1 LIKE 'foo\%%'", 80), + (r"c1 NOT LIKE 'foo\%%'", 720), + (r"c1 LIKE 'foo\'", 40), + (r"c1 NOT LIKE 'foo\'", 760), + (r"c1 LIKE '%foo\'", 80), + (r"c1 NOT LIKE '%foo\'", 720), + // Case-sensitive NDV cannot estimate a case-insensitive equality. + ("c1 ILIKE 'foo'", 160), + ("c1 NOT ILIKE 'foo'", 640), + ]; + for data_type in [DataType::Utf8, DataType::LargeUtf8, DataType::Utf8View] { + let ctx = string_filter_ctx(data_type.clone(), Precision::Exact(20))?; + for (predicate, expected_rows) in cases { + assert_eq!( + string_filter_rows(&ctx, predicate).await?, + Precision::Inexact(expected_rows), + "{data_type:?}: {predicate}" + ); + } + } + Ok(()) +} + +#[tokio::test] +async fn sql_string_filter_without_distinct_count() -> Result<()> { + let ctx = string_filter_ctx(DataType::Utf8View, Precision::Absent)?; + // Without NDV, use the configured default (20%) over non-null rows. + for (predicate, expected_rows) in [ + ("c1 = 'foo'", 160), + ("c1 != 'foo'", 640), + ("c1 LIKE 'foo'", 160), + ("c1 NOT LIKE 'foo'", 640), + ] { + assert_eq!( + string_filter_rows(&ctx, predicate).await?, + Precision::Inexact(expected_rows), + "{predicate}" + ); + } + Ok(()) +} + +#[tokio::test] +async fn sql_string_filter_custom_selectivity() -> Result<()> { + let ctx = string_filter_ctx(DataType::Utf8View, Precision::Absent)?; + ctx.sql("SET datafusion.optimizer.default_filter_selectivity = 40") + .await? + .collect() + .await?; + + for (predicate, expected_rows) in [ + ("c1 = 'foo'", 320), + ("c1 != 'foo'", 480), + ("c1 LIKE '%foo%'", 320), + ("c1 NOT LIKE '%foo%'", 480), + ("c1 LIKE 'foo%'", 160), + ("c1 NOT LIKE 'foo%'", 640), + ] { + assert_eq!( + string_filter_rows(&ctx, predicate).await?, + Precision::Inexact(expected_rows), + "{predicate}" + ); + } + Ok(()) +} + #[tokio::test] async fn sql_limit() -> Result<()> { let (stats, schema) = fully_defined(); diff --git a/datafusion/physical-plan/src/filter.rs b/datafusion/physical-plan/src/filter.rs index f71d67d04e39c..48c0cc1257e37 100644 --- a/datafusion/physical-plan/src/filter.rs +++ b/datafusion/physical-plan/src/filter.rs @@ -64,7 +64,7 @@ use datafusion_execution::TaskContext; use datafusion_expr::Operator; use datafusion_physical_expr::equivalence::ProjectionMapping; use datafusion_physical_expr::expressions::{ - BinaryExpr, Column, InListExpr, IsNotNullExpr, Literal, lit, + BinaryExpr, Column, InListExpr, IsNotNullExpr, LikeExpr, Literal, lit, }; use datafusion_physical_expr::intervals::utils::check_support; use datafusion_physical_expr::utils::collect_columns; @@ -346,8 +346,8 @@ impl FilterExec { } /// Calculates `Statistics` for `FilterExec` by applying the filter's - /// selectivity (default, or estimated from interval analysis) to the input - /// statistics. + /// selectivity (estimated from interval analysis or string statistics, with + /// a configurable fallback) to the input statistics. /// /// The estimated output row count is used to keep the per-column statistics /// consistent with it: @@ -414,10 +414,17 @@ impl FilterExec { ); (selectivity, filtered_num_rows, cs) } else { - // Without interval boundaries, use the default selectivity and - // apply the row-count constraints that still follow from the - // filter predicate. - let selectivity = default_selectivity as f64 / 100.0; + // String predicates cannot use interval analysis. Estimate them + // from distinct counts and pattern form when possible, then + // apply the row-count constraints implied by the predicate. + let default_selectivity = default_selectivity as f64 / 100.0; + let selectivity = string_selectivity( + predicate, + schema, + &input_stats, + default_selectivity, + ) + .unwrap_or(default_selectivity); let filtered_num_rows = input_num_rows.with_estimated_selectivity(selectivity); let mut cs = input_stats.to_inexact().column_statistics; @@ -1181,6 +1188,17 @@ fn collect_null_rejecting_columns(predicate: &Arc) -> HashSet< let mut columns = HashSet::new(); for expr in split_conjunction(predicate) { + // LIKE (including ILIKE and their negations) returns NULL if either + // operand is NULL. + if let Some(like) = expr.downcast_ref::() { + for operand in [like.expr(), like.pattern()] { + if let Some(col) = operand.downcast_ref::() { + columns.insert(col.index()); + } + } + continue; + } + // `col IS NOT NULL` keeps only rows where `col` is non-null. if let Some(is_not_null) = expr.downcast_ref::() { if let Some(col) = is_not_null.arg().downcast_ref::() { @@ -1207,6 +1225,135 @@ fn collect_null_rejecting_columns(predicate: &Arc) -> HashSet< columns } +/// Estimates a single string comparison against a literal. Compound predicates +/// retain the fallback: multiplying estimates for predicates on the same column +/// would count null rejection (and possibly the same restriction) repeatedly. +fn string_selectivity( + predicate: &Arc, + schema: &SchemaRef, + input_stats: &Statistics, + default_selectivity: f64, +) -> Option { + let (column, literal, negated, like) = + if let Some(binary) = predicate.downcast_ref::() { + if !matches!(binary.op(), Operator::Eq | Operator::NotEq) { + return None; + } + let operands = [ + (binary.left(), binary.right()), + (binary.right(), binary.left()), + ]; + let (column, literal) = operands.into_iter().find_map(|(left, right)| { + Some(( + left.downcast_ref::()?, + right.downcast_ref::()?, + )) + })?; + (column, literal, *binary.op() == Operator::NotEq, None) + } else { + let like = predicate.downcast_ref::()?; + ( + like.expr().downcast_ref::()?, + like.pattern().downcast_ref::()?, + like.negated(), + Some(like), + ) + }; + + if !matches!( + schema.field(column.index()).data_type(), + DataType::Utf8 | DataType::LargeUtf8 | DataType::Utf8View + ) { + return None; + } + let Some(value) = literal.value().try_as_str()? else { + // Both a comparison and its negation evaluate to NULL here. + return Some(0.0); + }; + let column_stats = &input_stats.column_statistics[column.index()]; + let equality_selectivity = column_stats + .distinct_count + .get_value() + .filter(|&&ndv| ndv > 0) + .map_or(default_selectivity, |&ndv| 1.0 / ndv as f64); + let positive_selectivity = match like { + Some(like) => like_selectivity( + value, + like.case_insensitive(), + equality_selectivity, + default_selectivity, + ), + None => equality_selectivity, + }; + let selectivity = if negated { + 1.0 - positive_selectivity + } else { + positive_selectivity + }; + + // Apply the complement within the non-null rows, not across all input + // rows. Unknown null counts use the usual assumption of no nulls. + match ( + input_stats.num_rows.get_value(), + column_stats.null_count.get_value(), + ) { + (Some(&rows), Some(&nulls)) if rows > 0 => { + Some(rows.saturating_sub(nulls) as f64 * selectivity / rows as f64) + } + (Some(0), _) => Some(0.0), + _ => Some(selectivity), + } +} + +/// Classifies a LIKE pattern without treating escaped wildcards as operators. +/// In the absence of string histograms, wildcard patterns use the configured +/// fallback, halved for each anchored end. Thus prefix/suffix patterns are +/// estimated to match fewer rows than an unanchored contains pattern. These are +/// heuristics, not bounds on the number of matches. +fn like_selectivity( + pattern: &str, + case_insensitive: bool, + equality_selectivity: f64, + default_selectivity: f64, +) -> f64 { + let mut chars = pattern.chars(); + let mut has_wildcard = false; + let mut all_percent = !pattern.is_empty(); + let mut starts_with_percent = false; + let mut ends_with_percent = false; + let mut first = true; + while let Some(ch) = chars.next() { + let percent = ch == '%'; + if ch == '\\' { + // Arrow treats a trailing backslash as a literal backslash too. + chars.next(); + } else if matches!(ch, '%' | '_') { + has_wildcard = true; + } + all_percent &= percent; + if first { + starts_with_percent = percent; + first = false; + } + ends_with_percent = percent; + } + if all_percent { + 1.0 + } else if !has_wildcard { + // Case folding may merge several distinct input values, so the raw + // NDV does not give the selectivity of a literal ILIKE pattern. + if case_insensitive { + default_selectivity + } else { + equality_selectivity + } + } else { + let start_factor = if starts_with_percent { 1.0 } else { 0.5 }; + let end_factor = if ends_with_percent { 1.0 } else { 0.5 }; + default_selectivity * start_factor * end_factor + } +} + /// Converts an interval bound to a [`Precision`] value. NULL bounds (which /// represent "unbounded" in the interval type) map to [`Precision::Absent`]. fn interval_bound_to_precision( @@ -4174,6 +4321,154 @@ mod tests { Ok(()) } + fn like_filter_statistics( + num_rows: Precision, + column: ColumnStatistics, + pattern: Option<&str>, + negated: bool, + ) -> Result> { + let schema = Schema::new(vec![Field::new("name", DataType::Utf8, true)]); + let input = Arc::new(StatisticsExec::new( + Statistics { + num_rows, + total_byte_size: Precision::Absent, + column_statistics: vec![column], + }, + schema.clone(), + )); + let predicate = Arc::new(LikeExpr::new( + negated, + false, + col("name", &schema)?, + lit(ScalarValue::Utf8(pattern.map(str::to_owned))), + )); + let filter = FilterExec::try_new(predicate, input)?; + StatisticsContext::new().compute(&filter, &StatisticsArgs::new()) + } + + #[test] + fn test_filter_statistics_like_literal_patterns() -> Result<()> { + // SQL simplification rewrites these patterns before FilterExec sees + // them. Construct LIKE directly to exercise physical-plan callers too. + for (pattern, positive_rows, negative_rows) in [ + (Some("%"), 800, 0), + (Some("%%%"), 800, 0), + (Some(""), 40, 760), + (Some(r"\%"), 40, 760), + (Some(r"\_"), 40, 760), + (Some(r"\%%"), 80, 720), + (Some(r"%\%"), 80, 720), + (None, 0, 0), + ] { + for (negated, expected_rows) in + [(false, positive_rows), (true, negative_rows)] + { + let statistics = like_filter_statistics( + Precision::Exact(1000), + ColumnStatistics { + null_count: Precision::Exact(200), + distinct_count: Precision::Exact(20), + ..Default::default() + }, + pattern, + negated, + )?; + assert_eq!( + statistics.num_rows, + Precision::Inexact(expected_rows), + "pattern={pattern:?}, negated={negated}" + ); + assert_eq!( + statistics.column_statistics[0].null_count, + Precision::Exact(0) + ); + } + } + Ok(()) + } + + #[test] + fn test_filter_statistics_like_missing_or_zero_counts() -> Result<()> { + use Precision::{Absent, Exact, Inexact}; + + for (rows, nulls, distinct, positive_rows, negative_rows) in [ + (Exact(1000), Absent, Exact(20), Inexact(50), Inexact(950)), + ( + Exact(1000), + Exact(200), + Exact(0), + Inexact(160), + Inexact(640), + ), + (Exact(1000), Exact(1000), Exact(0), Inexact(0), Inexact(0)), + (Exact(0), Exact(0), Exact(0), Exact(0), Exact(0)), + (Inexact(0), Absent, Absent, Inexact(0), Inexact(0)), + (Absent, Absent, Exact(20), Absent, Absent), + ] { + for (negated, expected_rows) in + [(false, positive_rows), (true, negative_rows)] + { + let statistics = like_filter_statistics( + rows, + ColumnStatistics { + null_count: nulls, + distinct_count: distinct, + ..Default::default() + }, + Some(""), + negated, + )?; + assert_eq!( + statistics.num_rows, expected_rows, + "rows={rows:?}, nulls={nulls:?}, distinct={distinct:?}, negated={negated}" + ); + assert_eq!(statistics.column_statistics[0].null_count, Exact(0)); + } + } + Ok(()) + } + + #[test] + fn test_filter_statistics_like_column_pattern_rejects_nulls() -> Result<()> { + let schema = Schema::new(vec![ + Field::new("name", DataType::Utf8, true), + Field::new("pattern", DataType::Utf8, true), + ]); + for negated in [false, true] { + let input = Arc::new(StatisticsExec::new( + Statistics { + num_rows: Precision::Exact(1000), + total_byte_size: Precision::Absent, + column_statistics: [200, 300] + .into_iter() + .map(|nulls| ColumnStatistics { + null_count: Precision::Exact(nulls), + distinct_count: Precision::Exact(20), + ..Default::default() + }) + .collect(), + }, + schema.clone(), + )); + let predicate = Arc::new(LikeExpr::new( + negated, + false, + col("name", &schema)?, + col("pattern", &schema)?, + )); + let filter = FilterExec::try_new(predicate, input)?; + let statistics = + StatisticsContext::new().compute(&filter, &StatisticsArgs::new())?; + // A varying pattern retains the fallback estimate, but neither + // operand can be NULL in any surviving row. + assert_eq!(statistics.num_rows, Precision::Inexact(200)); + for column in &statistics.column_statistics { + assert_eq!(column.null_count, Precision::Exact(0)); + } + } + Ok(()) + } + #[tokio::test] async fn test_filter_statistics_is_not_null_rejects_nulls() -> Result<()> { let schema = Schema::new(vec![Field::new("name", DataType::Utf8, true)]); From a5b133207382036b0037c320322ad086da506b8a Mon Sep 17 00:00:00 2001 From: Jake Wallin Date: Mon, 28 Sep 2026 23:37:33 -0400 Subject: [PATCH 2/2] test: Cover string selectivity fallback and clarify LIKE test setup --- datafusion/physical-plan/src/filter.rs | 53 +++++++++++++++++--------- 1 file changed, 36 insertions(+), 17 deletions(-) diff --git a/datafusion/physical-plan/src/filter.rs b/datafusion/physical-plan/src/filter.rs index 48c0cc1257e37..9080e987a1da1 100644 --- a/datafusion/physical-plan/src/filter.rs +++ b/datafusion/physical-plan/src/filter.rs @@ -4321,12 +4321,30 @@ mod tests { Ok(()) } + #[test] + fn test_string_selectivity_non_string_literal() { + let schema = + Arc::new(Schema::new(vec![Field::new("name", DataType::Utf8, true)])); + let predicate: Arc = Arc::new(BinaryExpr::new( + Arc::new(Column::new("name", 0)), + Operator::Eq, + lit(42_i32), + )); + let statistics = Statistics::new_unknown(&schema); + + // SQL coerces comparison operands, so exercise this defensive fallback directly. + assert_eq!( + string_selectivity(&predicate, &schema, &statistics, 0.2), + None + ); + } + fn like_filter_statistics( num_rows: Precision, column: ColumnStatistics, pattern: Option<&str>, negated: bool, - ) -> Result> { + ) -> Arc { let schema = Schema::new(vec![Field::new("name", DataType::Utf8, true)]); let input = Arc::new(StatisticsExec::new( Statistics { @@ -4339,15 +4357,17 @@ mod tests { let predicate = Arc::new(LikeExpr::new( negated, false, - col("name", &schema)?, + col("name", &schema).expect("name column exists"), lit(ScalarValue::Utf8(pattern.map(str::to_owned))), )); - let filter = FilterExec::try_new(predicate, input)?; - StatisticsContext::new().compute(&filter, &StatisticsArgs::new()) + let filter = FilterExec::try_new(predicate, input).expect("valid LIKE filter"); + StatisticsContext::new() + .compute(&filter, &StatisticsArgs::new()) + .expect("LIKE filter statistics") } #[test] - fn test_filter_statistics_like_literal_patterns() -> Result<()> { + fn test_filter_statistics_like_literal_patterns() { // SQL simplification rewrites these patterns before FilterExec sees // them. Construct LIKE directly to exercise physical-plan callers too. for (pattern, positive_rows, negative_rows) in [ @@ -4372,7 +4392,7 @@ mod tests { }, pattern, negated, - )?; + ); assert_eq!( statistics.num_rows, Precision::Inexact(expected_rows), @@ -4384,11 +4404,10 @@ mod tests { ); } } - Ok(()) } #[test] - fn test_filter_statistics_like_missing_or_zero_counts() -> Result<()> { + fn test_filter_statistics_like_missing_or_zero_counts() { use Precision::{Absent, Exact, Inexact}; for (rows, nulls, distinct, positive_rows, negative_rows) in [ @@ -4417,7 +4436,7 @@ mod tests { }, Some(""), negated, - )?; + ); assert_eq!( statistics.num_rows, expected_rows, "rows={rows:?}, nulls={nulls:?}, distinct={distinct:?}, negated={negated}" @@ -4425,11 +4444,10 @@ mod tests { assert_eq!(statistics.column_statistics[0].null_count, Exact(0)); } } - Ok(()) } #[test] - fn test_filter_statistics_like_column_pattern_rejects_nulls() -> Result<()> { + fn test_filter_statistics_like_column_pattern_rejects_nulls() { let schema = Schema::new(vec![ Field::new("name", DataType::Utf8, true), Field::new("pattern", DataType::Utf8, true), @@ -4453,12 +4471,14 @@ mod tests { let predicate = Arc::new(LikeExpr::new( negated, false, - col("name", &schema)?, - col("pattern", &schema)?, + col("name", &schema).expect("name column exists"), + col("pattern", &schema).expect("pattern column exists"), )); - let filter = FilterExec::try_new(predicate, input)?; - let statistics = - StatisticsContext::new().compute(&filter, &StatisticsArgs::new())?; + let filter = + FilterExec::try_new(predicate, input).expect("valid LIKE filter"); + let statistics = StatisticsContext::new() + .compute(&filter, &StatisticsArgs::new()) + .expect("LIKE filter statistics"); // A varying pattern retains the fallback estimate, but neither // operand can be NULL in any surviving row. assert_eq!(statistics.num_rows, Precision::Inexact(200)); @@ -4466,7 +4486,6 @@ mod tests { assert_eq!(column.null_count, Precision::Exact(0)); } } - Ok(()) } #[tokio::test]