Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
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
119 changes: 104 additions & 15 deletions datafusion/optimizer/src/extract_leaf_expressions.rs
Original file line number Diff line number Diff line change
Expand Up @@ -714,13 +714,33 @@ fn build_extraction_projection_impl(
})
.collect();

// The names the merged projection already produces. The merge keeps
// every expression of `existing`, so the parent can still read each of
// these names, and a pass-through column that carries one of them makes
// the output schema ambiguous.
let output_names: std::collections::HashSet<&str> = existing
.schema
.fields()
.iter()
.map(|f| f.name().as_str())
.collect();

let input_schema = existing.input.schema();
for col in columns_needed {
// The projection produces this name, so the parent reads it from the
// projection's own output. Do not resolve the name through the rename
// map: a rename such as `t.a AS b` is a computed output, and its
// input column `t.a` beside the output field `a` of a second rename
// gives an ambiguous schema (issue #25446).
if output_names.contains(col.name.as_str()) {
continue;
}
let col_expr = Expr::Column(col.clone());
let resolved = replace_cols_by_name(col_expr, &replace_map)?;
if let Expr::Column(resolved_col) = &resolved
&& !existing_cols.contains(resolved_col)
&& input_schema.has_column(resolved_col)
&& !output_names.contains(resolved_col.name.as_str())
{
proj_exprs.push(Expr::Column(resolved_col.clone()));
}
Expand Down Expand Up @@ -2792,14 +2812,11 @@ mod tests {
## After Pushdown
Projection: __datafusion_extracted_1 AS leaf_udf(x,Utf8("a"))
Filter: x IS NOT NULL
Projection: test.user AS x, leaf_udf(test.user, Utf8("a")) AS __datafusion_extracted_1, test.user
Projection: test.user AS x, leaf_udf(test.user, Utf8("a")) AS __datafusion_extracted_1
TableScan: test projection=[user]

## Optimized
Projection: __datafusion_extracted_1 AS leaf_udf(x,Utf8("a"))
Filter: x IS NOT NULL
Projection: test.user AS x, leaf_udf(test.user, Utf8("a")) AS __datafusion_extracted_1
TableScan: test projection=[user]
(same as after pushdown)
"#)
}

Expand All @@ -2826,14 +2843,11 @@ mod tests {
## After Pushdown
Projection: __datafusion_extracted_1 IS NOT NULL AS leaf_udf(x,Utf8("a")) IS NOT NULL
Filter: x IS NOT NULL
Projection: test.user AS x, leaf_udf(test.user, Utf8("a")) AS __datafusion_extracted_1, test.user
Projection: test.user AS x, leaf_udf(test.user, Utf8("a")) AS __datafusion_extracted_1
TableScan: test projection=[user]

## Optimized
Projection: __datafusion_extracted_1 IS NOT NULL AS leaf_udf(x,Utf8("a")) IS NOT NULL
Filter: x IS NOT NULL
Projection: test.user AS x, leaf_udf(test.user, Utf8("a")) AS __datafusion_extracted_1
TableScan: test projection=[user]
(same as after pushdown)
"#)
}

Expand All @@ -2855,17 +2869,92 @@ mod tests {
## After Extraction
Projection: x
Filter: __datafusion_extracted_1 = Utf8("active")
Projection: test.user AS x, leaf_udf(test.user, Utf8("a")) AS __datafusion_extracted_1, test.user
Projection: test.user AS x, leaf_udf(test.user, Utf8("a")) AS __datafusion_extracted_1
TableScan: test projection=[user]

## After Pushdown
(same as after extraction)

## Optimized
Projection: x
Filter: __datafusion_extracted_1 = Utf8("active")
Projection: test.user AS x, leaf_udf(test.user, Utf8("a")) AS __datafusion_extracted_1
TableScan: test projection=[user]
(same as after pushdown)
"#)
}

/// A projection that swaps two column names, below a filter and a struct
/// field read. The merge must not resolve the parent's column references
/// through the rename map: the input column `test.user` beside the output
/// field `user` of the other rename gives an ambiguous schema (#25446).
#[test]
fn test_extract_above_projection_that_swaps_column_names() -> Result<()> {
let table_scan = test_table_scan_with_struct()?;
let plan = LogicalPlanBuilder::from(table_scan)
.project(vec![col("user").alias("id"), col("id").alias("user")])?
.filter(col("user").gt(lit(0u32)))?
.project(vec![col("id"), col("user"), leaf_udf(col("id"), "name")])?
.build()?;

assert_stages!(plan, @r#"
## Original Plan
Projection: id, user, leaf_udf(id, Utf8("name"))
Filter: user > UInt32(0)
Projection: test.user AS id, test.id AS user
TableScan: test projection=[id, user]

## After Extraction
(same as original)

## After Pushdown
Projection: id, user, __datafusion_extracted_1 AS leaf_udf(id,Utf8("name"))
Filter: user > UInt32(0)
Projection: test.user AS id, test.id AS user, leaf_udf(test.user, Utf8("name")) AS __datafusion_extracted_1
TableScan: test projection=[id, user]

## Optimized
(same as after pushdown)
"#)
}

/// A projection that gives a computed column the name of its own input
/// column. The extraction projection goes below the filter, which puts a
/// column named `id` under the rename, so the two plans have the same set
/// of field names. The recovery projection must stay, or the rename is lost
/// and the query gives wrong results.
#[test]
fn test_extract_above_projection_that_redefines_column_name() -> Result<()> {
let table_scan = test_table_scan_with_struct()?;
let plan = LogicalPlanBuilder::from(table_scan)
.filter(col("id").gt(lit(0u32)))?
.project(vec![(col("id") * lit(10u32)).alias("id"), col("user")])?
.alias("sub")?
.project(vec![col("sub.id"), leaf_udf(col("sub.user"), "name")])?
.build()?;

assert_stages!(plan, @r#"
## Original Plan
Projection: sub.id, leaf_udf(sub.user, Utf8("name"))
SubqueryAlias: sub
Projection: test.id * UInt32(10) AS id, test.user
Filter: test.id > UInt32(0)
TableScan: test projection=[id, user]

## After Extraction
(same as original)

## After Pushdown
Projection: sub.id, __datafusion_extracted_1 AS leaf_udf(sub.user,Utf8("name"))
SubqueryAlias: sub
Projection: test.id * UInt32(10) AS id, test.user, __datafusion_extracted_1
Filter: test.id > UInt32(0)
Projection: leaf_udf(test.user, Utf8("name")) AS __datafusion_extracted_1, test.id, test.user
TableScan: test projection=[id, user]

## Optimized
Projection: sub.id, __datafusion_extracted_1 AS leaf_udf(sub.user,Utf8("name"))
SubqueryAlias: sub
Projection: test.id * UInt32(10) AS id, __datafusion_extracted_1
Filter: test.id > UInt32(0)
Projection: leaf_udf(test.user, Utf8("name")) AS __datafusion_extracted_1, test.id
TableScan: test projection=[id, user]
"#)
}

Expand Down
75 changes: 75 additions & 0 deletions datafusion/sqllogictest/test_files/struct.slt
Original file line number Diff line number Diff line change
Expand Up @@ -1890,3 +1890,78 @@ from (values (1),(2),(1)) group by 1;
----
{a: 1, z: NULL} 2
{a: 2, z: NULL} 1

# A sub-query projection that renames a column to the name of a different
# column of the same input, with a struct field read above it. Leaf projection
# pushdown resolved the parent's column references through the rename map and
# added the renamed input column a second time, which made the output schema
# ambiguous (https://github.com/apache/datafusion/issues/25446).
statement ok
create table rename_swap_struct(a int, b int, s struct<x varchar>) as values (1, 10, {x: 'p'}), (2, 20, {x: 'q'});

# a rename beside a same-name alias, under a filter
query IIT
select b, a, s['x'] from (select rename_swap_struct.a as b, rename_swap_struct.b as a, s from rename_swap_struct) where a > 0 order by b;

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.

Optional: could we change this to where a > 15 and expect only 2 20 q? The current a > 0 case passes with either possible binding because both source columns are positive, while a > 15 would also verify the alias binding at the result level.

----
1 10 p
2 20 q

# the same shape under a limit
query IIT rowsort
select b, a, s['x'] from (select rename_swap_struct.a as b, rename_swap_struct.b as a, s from rename_swap_struct) limit 10;
----
1 10 p
2 20 q

# a swap of two column names
query IIT
select a, c, s['x'] from (select rename_swap_struct.b as a, rename_swap_struct.a as c, s from rename_swap_struct) where a > 0 order by a;
----
10 1 p
20 2 q

# the same shape under an order by
query IIT
select b, a, s['x'] from (select rename_swap_struct.a as b, rename_swap_struct.b as a, s from rename_swap_struct) order by b;
----
1 10 p
2 20 q

# the same shape under a group by
query III
select b, a, count(s['x']) from (select rename_swap_struct.a as b, rename_swap_struct.b as a, s from rename_swap_struct) group by b, a order by b;
----
1 10 1
2 20 1

# the rename stays above the extraction projection, and the struct field is
# still read at the scan
query TT
explain select b, a, s['x'] from (select rename_swap_struct.a as b, rename_swap_struct.b as a, s from rename_swap_struct) where a > 0;
----
logical_plan
01)Projection: rename_swap_struct.a AS b, rename_swap_struct.b AS a, __datafusion_extracted_1 AS rename_swap_struct.s[x]
02)--Filter: rename_swap_struct.b > Int32(0)
03)----Projection: get_field(rename_swap_struct.s, Utf8("x")) AS __datafusion_extracted_1, rename_swap_struct.a, rename_swap_struct.b
04)------TableScan: rename_swap_struct projection=[a, b, s]
physical_plan
01)ProjectionExec: expr=[a@1 as b, b@2 as a, __datafusion_extracted_1@0 as rename_swap_struct.s[x]]
02)--FilterExec: b@2 > 0
03)----ProjectionExec: expr=[get_field(s@2, x) as __datafusion_extracted_1, a@0 as a, b@1 as b]
04)------DataSourceExec: partitions=1, partition_sizes=[1]

statement ok
drop table rename_swap_struct;

# The same swap where the struct field name is also a column name of the table.
statement ok
create table rename_swap_struct_field(a int, b int, s struct<b varchar>) as values (1, 10, {b: 'p'}), (2, 20, {b: 'q'});

query IIT
select a, b, s['b'] from (select b as a, a as b, s from rename_swap_struct_field limit 100) order by a;
----
10 1 p
20 2 q

statement ok
drop table rename_swap_struct_field;
Loading