Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
24 commits
Select commit Hold shift + click to select a range
3e6453e
perf: skip functional dependency work when inputs carry none
dwsmith1983 Sep 23, 2026
d1e18cd
Merge branch 'main' into perf/projection-func-deps-early-return
dwsmith1983 Sep 24, 2026
b0aef89
Merge branch 'main' into perf/projection-func-deps-early-return
dwsmith1983 Sep 24, 2026
104e181
Merge branch 'main' into perf/projection-func-deps-early-return
dwsmith1983 Sep 25, 2026
45385e2
Merge branch 'main' into perf/projection-func-deps-early-return
dwsmith1983 Sep 25, 2026
dbf9904
Merge branch 'main' into perf/projection-func-deps-early-return
dwsmith1983 Sep 26, 2026
e7c5130
Merge branch 'main' into perf/projection-func-deps-early-return
dwsmith1983 Sep 27, 2026
30dc6e0
Merge branch 'main' into perf/projection-func-deps-early-return
dwsmith1983 Sep 27, 2026
83c8800
Merge branch 'main' into perf/projection-func-deps-early-return
dwsmith1983 Sep 27, 2026
7dd5380
Merge branch 'main' into perf/projection-func-deps-early-return
dwsmith1983 Sep 28, 2026
dbb31b2
Merge branch 'main' into perf/projection-func-deps-early-return
dwsmith1983 Sep 28, 2026
763710b
Merge branch 'main' into perf/projection-func-deps-early-return
dwsmith1983 Sep 29, 2026
4fffa49
Merge remote-tracking branch 'origin/main' into perf/projection-func-…
dwsmith1983 Sep 29, 2026
e581903
Merge branch 'main' into perf/projection-func-deps-early-return
dwsmith1983 Sep 29, 2026
85d79df
Merge branch 'main' into perf/projection-func-deps-early-return
dwsmith1983 Sep 30, 2026
9ca7906
Merge remote-tracking branch 'origin/main' into perf/projection-func-…
dwsmith1983 Sep 30, 2026
d91f49b
Merge branch 'main' into perf/projection-func-deps-early-return
dwsmith1983 Oct 1, 2026
4d2016f
Merge branch 'main' into perf/projection-func-deps-early-return
dwsmith1983 Oct 1, 2026
5705883
Merge remote-tracking branch 'origin/main' into perf/projection-func-…
dwsmith1983 Oct 1, 2026
e773326
Merge branch 'main' into perf/projection-func-deps-early-return
dwsmith1983 Oct 2, 2026
7cfb210
perf: skip get_target_functional_dependencies work without dependencies
dwsmith1983 Oct 2, 2026
d53a472
Merge branch 'main' into perf/projection-func-deps-early-return
dwsmith1983 Oct 3, 2026
42f4bec
Merge branch 'main' into perf/projection-func-deps-early-return
dwsmith1983 Oct 3, 2026
e964336
Merge branch 'main' into perf/projection-func-deps-early-return
dwsmith1983 Oct 3, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
122 changes: 65 additions & 57 deletions datafusion/common/src/functional_dependencies.rs
Original file line number Diff line number Diff line change
Expand Up @@ -458,71 +458,76 @@ 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::<Vec<_>>();
// Get functional dependencies of the schema:
let func_dependencies = aggr_input_schema.functional_dependencies();
for FunctionalDependence {
source_indices,
nullable,
null_equality,
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::<Vec<_>>();

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());
}
}
// 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() {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The same "no dependencies, skip" check is now at 3 sites (here, add_group_by_exprs_from_dependencies, and calc_func_dependencies_for_project). Could get_target_functional_dependencies also return None early, before it calls schema.field_names()? Then future callers get the fast path too:

let dependencies = schema.functional_dependencies();
if dependencies.is_empty() {
    return None;
}

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done, it now returns None up front when the schema has no dependencies. I kept the other two checks: the one in add_group_by_exprs_from_dependencies also skips building the GROUP BY field names before the call, and calc_func_dependencies_for_project doesn't go through this function at all.

let aggr_input_fields = aggr_input_schema.field_names();
// Compute once: this does not change in the loop.
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,
null_equality,
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() {
// GROUP BY treats NULLs as equal: a determinant covering the
// complete grouping key gets at most one output row per NULL too.
let output_null_equality =
if new_source_indices.len() == group_by_expr_names.len() {
NullEquality::NullEqualsNull
} else {
*null_equality
};
aggregate_func_dependencies.push(
FunctionalDependence::new(
new_source_indices,
target_indices.clone(),
*nullable,
)
.with_mode(mode)
.with_null_equality(output_null_equality),
// 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
.iter()
.map(|&idx| &aggr_input_fields[idx])
.collect::<Vec<_>>();

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() {
// GROUP BY treats NULLs as equal: a determinant covering the
// complete grouping key gets at most one output row per NULL too.
let output_null_equality =
if new_source_indices.len() == group_by_expr_names.len() {
NullEquality::NullEqualsNull
} else {
*null_equality
};
aggregate_func_dependencies.push(
FunctionalDependence::new(
new_source_indices,
target_indices.clone(),
*nullable,
)
.with_mode(mode)
.with_null_equality(output_null_equality),
);
}
}
}

Expand Down Expand Up @@ -565,8 +570,11 @@ pub fn get_target_functional_dependencies(
schema: &DFSchema,
group_by_expr_names: &[String],
) -> Option<Vec<usize>> {
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,
Expand Down
32 changes: 32 additions & 0 deletions datafusion/expr/src/logical_plan/builder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1979,6 +1979,14 @@ pub fn add_group_by_exprs_from_dependencies(
mut group_expr: Vec<Expr>,
schema: &DFSchemaRef,
) -> Result<Vec<Expr>> {
// 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
Expand Down Expand Up @@ -3106,6 +3114,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
Expand Down
170 changes: 147 additions & 23 deletions datafusion/expr/src/logical_plan/plan.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4381,7 +4381,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() {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Question: with this early return, the Expr::Wildcard branch no longer calls exprlist_to_fields(...)? when the input has no dependencies. So an error from that call does not come from here anymore. I think projection schema construction gives the same error before this point. Can you confirm?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes. projection_schema is the only caller, and it runs exprlist_to_fields(exprs, input)? over the full expression list (wildcard included) before it calls calc_func_dependencies_for_project. to_field is deterministic, so any error the per-wildcard call could raise has already surfaced there. The early return only skips repeating that work.

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
Expand All @@ -4403,39 +4425,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::<Vec<_>>(),
)
}
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::<Result<Vec<_>>>()?
.into_iter()
.flatten()
.collect::<Vec<_>>();

Ok(input
.schema()
.functional_dependencies()
Ok(input_func_dependencies
.project_functional_dependencies(&proj_indices, exprs.len()))
}

Expand Down Expand Up @@ -5477,6 +5480,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))])
}
Expand Down
Loading