From 3e6453e755bb40d9229cef0f4541834984876c11 Mon Sep 17 00:00:00 2001 From: Dustin Smith Date: Wed, 23 Sep 2026 22:15:01 +0700 Subject: [PATCH 1/2] perf: skip functional dependency work when inputs carry none Most projections and aggregates sit on inputs with no primary key or unique constraint, so their functional dependency set is empty. The projection, aggregate, and GROUP BY paths now check that up front and return early instead of resolving every expression against the input fields. Projection lookups also move from a linear scan over field names to a hash map built once, so a real constraint still resolves correctly without rescanning per expression. The aggregate path only drops the per-dependence loop, so the GROUP BY dependency is still produced when one exists. Measured on an interleaved run against the base commit, noise floor about two percent: logical_select_all_from_1000 -25%, physical_select_all_from_1000 -12%, physical_plan_clickbench_all -7%, logical_wide_aggregate_100_exprs -3%. TPC-H planning unchanged. --- .../common/src/functional_dependencies.rs | 102 ++++++----- datafusion/expr/src/logical_plan/builder.rs | 32 ++++ datafusion/expr/src/logical_plan/plan.rs | 170 +++++++++++++++--- 3 files changed, 235 insertions(+), 69 deletions(-) diff --git a/datafusion/common/src/functional_dependencies.rs b/datafusion/common/src/functional_dependencies.rs index e8275aac2da4c..54b4cefd5dae2 100644 --- a/datafusion/common/src/functional_dependencies.rs +++ b/datafusion/common/src/functional_dependencies.rs @@ -423,61 +423,71 @@ pub fn aggregate_functional_dependencies( aggr_schema: &DFSchema, ) -> FunctionalDependencies { let mut aggregate_func_dependencies = vec![]; - let aggr_input_fields = aggr_input_schema.field_names(); let aggr_fields = aggr_schema.fields(); // Association covers the whole table: let target_indices = (0..aggr_schema.fields().len()).collect::>(); // Get functional dependencies of the schema: let func_dependencies = aggr_input_schema.functional_dependencies(); - for FunctionalDependence { - source_indices, - nullable, - mode, - .. - } in &func_dependencies.deps - { - // Keep source indices in a `HashSet` to prevent duplicate entries: - let mut new_source_indices = vec![]; - let mut new_source_field_names = vec![]; - let source_field_names = source_indices - .iter() - .map(|&idx| &aggr_input_fields[idx]) - .collect::>(); - - for (idx, group_by_expr_name) in group_by_expr_names.iter().enumerate() { - // When one of the input determinant expressions matches with - // the GROUP BY expression, add the index of the GROUP BY - // expression as a new determinant key: - if source_field_names.contains(&group_by_expr_name) { - new_source_indices.push(idx); - new_source_field_names.push(group_by_expr_name.clone()); - } - } + // If the input carries no functional dependencies, the loop below can + // never turn one into an aggregate dependency (it only re-expresses + // dependencies that already exist on the input), so skip building the + // input field names and resolving target indices for it entirely. The + // GROUP BY-key dependency added after this block does not depend on the + // input's functional dependencies, so it still runs unconditionally. + if !func_dependencies.is_empty() { + let aggr_input_fields = aggr_input_schema.field_names(); + // Loop-invariant: does not depend on the per-dependence loop + // variables, so compute it once instead of on every iteration. let existing_target_indices = get_target_functional_dependencies(aggr_input_schema, group_by_expr_names); - let new_target_indices = get_target_functional_dependencies( - aggr_input_schema, - &new_source_field_names, - ); - let mode = if existing_target_indices == new_target_indices - && new_target_indices.is_some() + for FunctionalDependence { + source_indices, + nullable, + mode, + .. + } in &func_dependencies.deps { - // If dependency covers all GROUP BY expressions, mode will be `Single`: - Dependency::Single - } else { - // Otherwise, existing mode is preserved: - *mode - }; - // All of the composite indices occur in the GROUP BY expression: - if new_source_indices.len() == source_indices.len() { - aggregate_func_dependencies.push( - FunctionalDependence::new( - new_source_indices, - target_indices.clone(), - *nullable, - ) - .with_mode(mode), + // Keep source indices in a `HashSet` to prevent duplicate entries: + let mut new_source_indices = vec![]; + let mut new_source_field_names = vec![]; + let source_field_names = source_indices + .iter() + .map(|&idx| &aggr_input_fields[idx]) + .collect::>(); + + for (idx, group_by_expr_name) in group_by_expr_names.iter().enumerate() { + // When one of the input determinant expressions matches with + // the GROUP BY expression, add the index of the GROUP BY + // expression as a new determinant key: + if source_field_names.contains(&group_by_expr_name) { + new_source_indices.push(idx); + new_source_field_names.push(group_by_expr_name.clone()); + } + } + let new_target_indices = get_target_functional_dependencies( + aggr_input_schema, + &new_source_field_names, ); + let mode = if existing_target_indices == new_target_indices + && new_target_indices.is_some() + { + // If dependency covers all GROUP BY expressions, mode will be `Single`: + Dependency::Single + } else { + // Otherwise, existing mode is preserved: + *mode + }; + // All of the composite indices occur in the GROUP BY expression: + if new_source_indices.len() == source_indices.len() { + aggregate_func_dependencies.push( + FunctionalDependence::new( + new_source_indices, + target_indices.clone(), + *nullable, + ) + .with_mode(mode), + ); + } } } diff --git a/datafusion/expr/src/logical_plan/builder.rs b/datafusion/expr/src/logical_plan/builder.rs index 8600a3cd38677..fdcaabee2938c 100644 --- a/datafusion/expr/src/logical_plan/builder.rs +++ b/datafusion/expr/src/logical_plan/builder.rs @@ -1982,6 +1982,14 @@ pub fn add_group_by_exprs_from_dependencies( mut group_expr: Vec, schema: &DFSchemaRef, ) -> Result> { + // With no functional dependencies on the input schema, + // `get_target_functional_dependencies` below can never resolve any + // target indices, so skip formatting the GROUP BY expression names and + // return the GROUP BY list unchanged. + if schema.functional_dependencies().is_empty() { + return Ok(group_expr); + } + // Names of the fields produced by the GROUP BY exprs for example, `GROUP BY // c1 + 1` produces an output field named `"c1 + 1"` let mut group_by_field_names = group_expr @@ -3089,6 +3097,30 @@ mod tests { Ok(()) } + #[test] + fn plan_builder_aggregate_with_implicit_group_by_exprs_no_constraints() -> Result<()> + { + // With no PRIMARY KEY / UNIQUE constraint on the input, there are no + // functional dependencies to expand the GROUP BY with, so the GROUP + // BY list must be left exactly as given. + let table_source = table_source(&employee_schema()); + + let options = + LogicalPlanBuilderOptions::new().with_add_implicit_group_by_exprs(true); + let plan = + LogicalPlanBuilder::scan("employee_csv", table_source, Some(vec![0, 3, 4]))? + .with_options(options) + .aggregate(vec![col("id")], vec![sum(col("salary"))])? + .build()?; + + assert_snapshot!(plan, @r" + Aggregate: groupBy=[[employee_csv.id]], aggr=[[sum(employee_csv.salary)]] + TableScan: employee_csv projection=[id, state, salary] + "); + + Ok(()) + } + #[test] fn plan_builder_aggregate_rejects_nested_aggregates() -> Result<()> { // https://github.com/apache/datafusion/issues/23812 diff --git a/datafusion/expr/src/logical_plan/plan.rs b/datafusion/expr/src/logical_plan/plan.rs index 96e7feebf6f45..f13480cd95abc 100644 --- a/datafusion/expr/src/logical_plan/plan.rs +++ b/datafusion/expr/src/logical_plan/plan.rs @@ -4357,7 +4357,29 @@ fn calc_func_dependencies_for_project( // Sentinel for projection outputs that do not map back to any input field. const COMPUTED_EXPR_INDEX: usize = usize::MAX; + let input_func_dependencies = input.schema().functional_dependencies(); + // Projecting an empty set of dependencies always yields an empty set, so + // skip resolving projection expressions against the input fields. This is + // the common case because table sources carry no constraints by default. + if input_func_dependencies.is_empty() { + return Ok(FunctionalDependencies::empty()); + } + + // Map each input field name to its first index so that projection + // expressions resolve with a hash lookup instead of a linear scan. let input_fields = input.schema().field_names(); + let mut input_index_by_name: HashMap<&str, usize> = + HashMap::with_capacity(input_fields.len()); + for (index, name) in input_fields.iter().enumerate() { + input_index_by_name.entry(name.as_str()).or_insert(index); + } + let input_index = |name: &str| { + input_index_by_name + .get(name) + .copied() + .unwrap_or(COMPUTED_EXPR_INDEX) + }; + // Map each projection output position to its input column index. // A projection expression can produce multiple output columns, such as `*`. let proj_indices = exprs @@ -4379,39 +4401,20 @@ fn calc_func_dependencies_for_project( let flat_name = qualifier .map(|t| format!("{}.{}", t, f.name())) .unwrap_or_else(|| f.name().clone()); - input_fields - .iter() - .position(|item| *item == flat_name) - .unwrap_or(COMPUTED_EXPR_INDEX) + input_index(&flat_name) }) .collect::>(), ) } - Expr::Alias(alias) => { - let name = format!("{}", alias.expr); - let input_index = input_fields - .iter() - .position(|item| *item == name) - .unwrap_or(COMPUTED_EXPR_INDEX); - Ok(vec![input_index]) - } - _ => { - let name = format!("{expr}"); - let input_index = input_fields - .iter() - .position(|item| *item == name) - .unwrap_or(COMPUTED_EXPR_INDEX); - Ok(vec![input_index]) - } + Expr::Alias(alias) => Ok(vec![input_index(&format!("{}", alias.expr))]), + _ => Ok(vec![input_index(&format!("{expr}"))]), }) .collect::>>()? .into_iter() .flatten() .collect::>(); - Ok(input - .schema() - .functional_dependencies() + Ok(input_func_dependencies .project_functional_dependencies(&proj_indices, exprs.len())) } @@ -5448,6 +5451,127 @@ mod tests { Ok(()) } + #[test] + fn projection_with_alias_preserves_pk() -> Result<()> { + let constraints = + Constraints::new_unverified(vec![Constraint::PrimaryKey(vec![0])]); + let source = Arc::new( + LogicalTableSource::new(Arc::new(employee_schema())) + .with_constraints(constraints), + ); + let plan = LogicalPlanBuilder::scan("employee_csv", source, None)? + .project(vec![col("id").alias("emp_id"), col("first_name")])? + .build()?; + + let deps = plan.schema().functional_dependencies(); + assert_eq!(deps.len(), 1); + assert_eq!(deps[0].source_indices, vec![0]); + assert_eq!(deps[0].target_indices, vec![0, 1]); + + Ok(()) + } + + #[test] + fn projection_over_unconstrained_table_has_no_dependencies() -> Result<()> { + let source = Arc::new(LogicalTableSource::new(Arc::new(employee_schema()))); + let plan = LogicalPlanBuilder::scan("employee_csv", source, None)? + .project(vec![col("id"), col("first_name")])? + .build()?; + + let deps = plan.schema().functional_dependencies(); + assert!(deps.is_empty()); + + Ok(()) + } + + #[test] + fn projection_duplicate_flattened_name_uses_first_input_index() -> Result<()> { + // Build an input schema where a qualified field (`orders`.`id`) and an + // unqualified field that is literally named `"orders.id"` flatten to + // the exact same lookup key that `calc_func_dependencies_for_project` + // uses to resolve projection expressions against input fields. This is + // the only way two entries of `DFSchema::field_names()` can collide + // (`DFSchema::check_names` otherwise forbids duplicate names), and it + // pins that the hash-map based lookup resolves such a collision to the + // *first* matching index, exactly like the linear `position()` scan it + // replaces. + let schema = DFSchema::new_with_metadata( + vec![ + ( + Some(TableReference::bare("orders")), + Arc::new(Field::new("id", DataType::Int32, false)), + ), + ( + None, + Arc::new(Field::new("orders.id", DataType::Int32, false)), + ), + ], + Metadata::default(), + )? + .with_functional_dependencies(FunctionalDependencies::new(vec![ + FunctionalDependence::new(vec![0], vec![0, 1], false) + .with_mode(Dependency::Single), + ]))?; + let input = LogicalPlan::EmptyRelation(EmptyRelation { + produce_one_row: true, + schema: Arc::new(schema), + }); + + // References the *unqualified* second field, whose flattened name + // ("orders.id") collides with the first (qualified) field's. + let exprs = vec![Expr::Column(Column::new_unqualified("orders.id"))]; + let deps = calc_func_dependencies_for_project(&exprs, &input)?; + + assert_eq!(deps.len(), 1); + assert_eq!(deps[0].source_indices, vec![0]); + + Ok(()) + } + + #[test] + fn aggregate_group_by_on_primary_key_reports_single_dependency() -> Result<()> { + let constraints = + Constraints::new_unverified(vec![Constraint::PrimaryKey(vec![0])]); + let source = Arc::new( + LogicalTableSource::new(Arc::new(employee_schema())) + .with_constraints(constraints), + ); + let plan = LogicalPlanBuilder::scan("employee_csv", source, None)? + .aggregate(vec![col("id")], vec![count(lit(true))])? + .build()?; + + let deps = plan.schema().functional_dependencies(); + assert_eq!(deps.len(), 1); + assert_eq!(deps[0].source_indices, vec![0]); + assert_eq!(deps[0].mode, Dependency::Single); + + Ok(()) + } + + #[test] + fn aggregate_group_by_without_constraints_still_reports_single_dependency() + -> Result<()> { + // Grouping guarantees uniqueness of the GROUP BY key regardless of + // whether the input table carries any PRIMARY KEY / UNIQUE + // constraints, so `aggregate_functional_dependencies` must still + // report a `Single` dependency spanning the whole aggregate output. + // This pins that behavior so the early return added for the (far + // more common) case of an input with no functional dependencies at + // all cannot be mistakenly widened to also skip this GROUP BY-only + // dependency. + let source = Arc::new(LogicalTableSource::new(Arc::new(employee_schema()))); + let plan = LogicalPlanBuilder::scan("employee_csv", source, None)? + .aggregate(vec![col("id")], vec![count(lit(true))])? + .build()?; + + let deps = plan.schema().functional_dependencies(); + assert_eq!(deps.len(), 1); + assert_eq!(deps[0].source_indices, vec![0]); + assert_eq!(deps[0].mode, Dependency::Single); + + Ok(()) + } + fn i32_split_point(value: i32) -> SplitPoint { SplitPoint::new(vec![ScalarValue::Int32(Some(value))]) } From 7cfb21029b3028786bff2498b4415f26c49a8720 Mon Sep 17 00:00:00 2001 From: Dustin Smith Date: Sat, 3 Oct 2026 00:20:13 +0700 Subject: [PATCH 2/2] perf: skip get_target_functional_dependencies work without dependencies Return None before collecting field names when the schema has no functional dependencies, so every caller gets the fast path. Also shorten the comments in aggregate_functional_dependencies and fix one that described the source indices as a HashSet. --- .../common/src/functional_dependencies.rs | 18 ++++++++---------- 1 file changed, 8 insertions(+), 10 deletions(-) diff --git a/datafusion/common/src/functional_dependencies.rs b/datafusion/common/src/functional_dependencies.rs index c54de35f3ad20..5ce3ed3b2b94d 100644 --- a/datafusion/common/src/functional_dependencies.rs +++ b/datafusion/common/src/functional_dependencies.rs @@ -463,16 +463,11 @@ pub fn aggregate_functional_dependencies( let target_indices = (0..aggr_schema.fields().len()).collect::>(); // Get functional dependencies of the schema: let func_dependencies = aggr_input_schema.functional_dependencies(); - // If the input carries no functional dependencies, the loop below can - // never turn one into an aggregate dependency (it only re-expresses - // dependencies that already exist on the input), so skip building the - // input field names and resolving target indices for it entirely. The - // GROUP BY-key dependency added after this block does not depend on the - // input's functional dependencies, so it still runs unconditionally. + // The loop below only re-expresses input dependencies. Skip it when the + // input has none. The GROUP BY-key dependency below always runs. if !func_dependencies.is_empty() { let aggr_input_fields = aggr_input_schema.field_names(); - // Loop-invariant: does not depend on the per-dependence loop - // variables, so compute it once instead of on every iteration. + // Compute once: this does not change in the loop. let existing_target_indices = get_target_functional_dependencies(aggr_input_schema, group_by_expr_names); for FunctionalDependence { @@ -483,7 +478,7 @@ pub fn aggregate_functional_dependencies( .. } in &func_dependencies.deps { - // Keep source indices in a `HashSet` to prevent duplicate entries: + // Indices into the GROUP BY list for this determinant: let mut new_source_indices = vec![]; let mut new_source_field_names = vec![]; let source_field_names = source_indices @@ -575,8 +570,11 @@ pub fn get_target_functional_dependencies( schema: &DFSchema, group_by_expr_names: &[String], ) -> Option> { - let mut combined_target_indices = HashSet::new(); let dependencies = schema.functional_dependencies(); + if dependencies.is_empty() { + return None; + } + let mut combined_target_indices = HashSet::new(); let field_names = schema.field_names(); for FunctionalDependence { source_indices,