This is an automated email from the ASF dual-hosted git repository.
github-merge-queue[bot] pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/datafusion.git
The following commit(s) were added to refs/heads/main by this push:
new 2424daeb18 fix: preserve computed names in derived SQL projections
(#25521)
2424daeb18 is described below
commit 2424daeb18d526c30d5351655289af71d6c05844
Author: Stefan Wang <[email protected]>
AuthorDate: Tue Sep 22 03:24:13 2026 +0000
fix: preserve computed names in derived SQL projections (#25521)
## Which issue does this PR close?
Related to https://github.com/apache/datafusion/issues/21066.
## Rationale for this change
An outer query can reference a computed result by its logical output
name while the generated inner query leaves that result unnamed. The
database cannot resolve the outer reference. For example, BigQuery
output refers to `sum_40j1_46j1_id_41` without defining it in the
derived table.
## What changes are included in this PR?
Derived projections now emit matching output names whether or not the
table has an alias. The unparser waits until the final table alias is
known, so a table-column alias list still supplies the names. Explicit
expression aliases and plain column references keep their existing SQL.
## What is the testing strategy for this PR?
The regression exercises PostgreSQL and BigQuery output through the
unparser. Existing column-alias and UNNEST snapshots remain unchanged.
Roundtrip tests still compare the whole logical plan after removing only
unqualified, metadata-free column self-aliases.
### Testing Done
SQLite 3.51.0 executes the emitted BigQuery SQL without rewriting it.
This uses local SQLite, not the Google BigQuery service. The pre-fix
query below was captured from
https://github.com/apache/datafusion/commit/4abdeef1e03e7d37be7a0dd1683985e9e3c8a7d1;
the fixed query comes directly from the regression run.
| # | Scenario | Command | Result |
|---|---|---|---|
| 1 | Recorded pre-fix SQL | `python3 - "$log"` below | `no such column:
sum_40j1_46j1_id_41` |
| 2 | Freshly emitted fixed SQL | `python3 - "$log"` below | `[(3,)]` |
Run from this checkout:
```bash
set -euo pipefail
log=$(mktemp)
CARGO_BUILD_JOBS=2 cargo test --locked --profile ci \
-p datafusion -p datafusion-sql --features
datafusion-sql/recursive_protection \
--test sql_integration unparse_preserves_derived_aggregate_output_name \
-- --nocapture > "$log" 2>&1
cat "$log"
python3 - "$log" <<'PY'
import sqlite3
import sys
from pathlib import Path
before = (
"SELECT `sum_40j1_46j1_id_41` FROM "
"(SELECT sum(`j1`.`j1_id`) AS `visible`, sum(`j1`.`j1_id`) FROM `j1`)"
)
emitted: list[str] = [
line.removeprefix("BIGQUERY_SQL=")
for line in Path(sys.argv[1]).read_text().splitlines()
if line.startswith("BIGQUERY_SQL=")
]
assert len(emitted) == 1, emitted
print(f"SQLite version: {sqlite3.sqlite_version}")
with sqlite3.connect(database=":memory:") as connection:
connection.executescript(
"CREATE TABLE j1 (j1_id INTEGER); INSERT INTO j1 VALUES (1), (2);"
)
print(f"before SQL: {before}")
try:
connection.execute(before).fetchall()
except sqlite3.OperationalError as error:
assert str(error) == "no such column: sum_40j1_46j1_id_41", error
print(f"before error: {error}")
else:
raise AssertionError("The recorded pre-fix query unexpectedly
succeeded.")
print(f"after SQL: {emitted[0]}")
rows = connection.execute(emitted[0]).fetchall()
print(f"after rows: {rows}")
assert rows == [(3,)], rows
PY
```
<details>
<summary>Raw logs</summary>
```text
SQLite version: 3.51.0
before SQL: SELECT `sum_40j1_46j1_id_41` FROM (SELECT sum(`j1`.`j1_id`) AS
`visible`, sum(`j1`.`j1_id`) FROM `j1`)
before error: no such column: sum_40j1_46j1_id_41
after SQL: SELECT `sum_40j1_46j1_id_41` FROM (SELECT sum(`j1`.`j1_id`) AS
`visible`, sum(`j1`.`j1_id`) AS `sum_40j1_46j1_id_41` FROM `j1`)
after rows: [(3,)]
```
</details>
## Are there any user-facing changes?
Generated derived projections expose the computed column names used by
outer queries, including when the dialect does not require a table
alias. SQL input validation and public APIs are unchanged.
---------
Signed-off-by: 1fanwang <[email protected]>
---
datafusion/sql/src/unparser/ast.rs | 45 ++++++++++++++++++++----
datafusion/sql/src/unparser/plan.rs | 30 ++++++++++------
datafusion/sql/tests/cases/plan_to_sql.rs | 57 ++++++++++++++++++++++++++++++-
3 files changed, 115 insertions(+), 17 deletions(-)
diff --git a/datafusion/sql/src/unparser/ast.rs
b/datafusion/sql/src/unparser/ast.rs
index c335d3ee4c..38b51afe3e 100644
--- a/datafusion/sql/src/unparser/ast.rs
+++ b/datafusion/sql/src/unparser/ast.rs
@@ -677,6 +677,7 @@ pub struct DerivedRelationBuilder {
lateral: Option<bool>,
subquery: Option<Box<ast::Query>>,
alias: Option<ast::TableAlias>,
+ projection_names: Vec<ast::Ident>,
}
impl DerivedRelationBuilder {
@@ -692,18 +693,49 @@ impl DerivedRelationBuilder {
self.alias = value;
self
}
+ pub(super) fn projection_names(&mut self, value: Vec<ast::Ident>) -> &mut
Self {
+ self.projection_names = value;
+ self
+ }
fn build(&self) -> Result<ast::TableFactor, BuilderError> {
+ let mut subquery = match self.subquery {
+ Some(ref value) => value.clone(),
+ None => {
+ return
Err(Into::into(UninitializedFieldError::from("subquery")));
+ }
+ };
+ if self
+ .alias
+ .as_ref()
+ .is_none_or(|alias| alias.columns.is_empty())
+ && let ast::SetExpr::Select(select) = subquery.body.as_mut()
+ && select.projection.len() == self.projection_names.len()
+ {
+ for (item, alias) in
select.projection.iter_mut().zip(&self.projection_names)
+ {
+ if let ast::SelectItem::UnnamedExpr(expr) = item {
+ let preserves_name = match expr {
+ ast::Expr::Identifier(name) => name.value ==
alias.value,
+ ast::Expr::CompoundIdentifier(names) => {
+ names.last().is_some_and(|name| name.value ==
alias.value)
+ }
+ _ => false,
+ };
+ if !preserves_name {
+ *item = ast::SelectItem::ExprWithAlias {
+ expr: expr.clone(),
+ alias: alias.clone(),
+ };
+ }
+ }
+ }
+ }
Ok(ast::TableFactor::Derived {
lateral: match self.lateral {
Some(ref value) => *value,
None => return
Err(Into::into(UninitializedFieldError::from("lateral"))),
},
- subquery: match self.subquery {
- Some(ref value) => value.clone(),
- None => {
- return
Err(Into::into(UninitializedFieldError::from("subquery")));
- }
- },
+ subquery,
alias: self.alias.clone(),
sample: None,
})
@@ -713,6 +745,7 @@ impl DerivedRelationBuilder {
lateral: Default::default(),
subquery: Default::default(),
alias: Default::default(),
+ projection_names: Default::default(),
}
}
}
diff --git a/datafusion/sql/src/unparser/plan.rs
b/datafusion/sql/src/unparser/plan.rs
index 69babae4e0..70a802629c 100644
--- a/datafusion/sql/src/unparser/plan.rs
+++ b/datafusion/sql/src/unparser/plan.rs
@@ -593,7 +593,18 @@ impl Unparser<'_> {
alias: Option<ast::TableAlias>,
lateral: bool,
) -> Result<()> {
+ let preserve_names = matches!(plan, LogicalPlan::Projection(_))
+ && alias.as_ref().is_none_or(|alias| alias.columns.is_empty());
let mut derived_builder = DerivedRelationBuilder::default();
+ if preserve_names {
+ derived_builder.projection_names(
+ plan.schema()
+ .fields()
+ .iter()
+ .map(|field| self.column_alias_to_sql(field.name()))
+ .collect::<Result<Vec<_>>>()?,
+ );
+ }
derived_builder.lateral(lateral).alias(alias).subquery({
let inner_statement = self.plan_to_sql(plan)?;
if let ast::Statement::Query(inner_query) = inner_statement {
@@ -2693,18 +2704,9 @@ impl Unparser<'_> {
Expr::Alias(Alias { expr, name, .. }) => {
let inner = self.expr_to_sql(expr)?;
- // Determine the alias name to use
- let col_name = if let Some(rewritten_name) =
- self.dialect.col_alias_overrides(name)?
- {
- rewritten_name.to_string()
- } else {
- name.to_string()
- };
-
Ok(ast::SelectItem::ExprWithAlias {
expr: inner,
- alias: self.new_ident_quoted_if_needs(col_name),
+ alias: self.column_alias_to_sql(name)?,
})
}
_ => {
@@ -2715,6 +2717,14 @@ impl Unparser<'_> {
}
}
+ fn column_alias_to_sql(&self, name: &str) -> Result<Ident> {
+ let name = self
+ .dialect
+ .col_alias_overrides(name)?
+ .unwrap_or_else(|| name.to_string());
+ Ok(self.new_ident_quoted_if_needs(name))
+ }
+
fn sorts_to_sql(&self, sort_exprs: &[SortExpr]) -> Result<OrderByKind> {
Ok(OrderByKind::Expressions(
sort_exprs
diff --git a/datafusion/sql/tests/cases/plan_to_sql.rs
b/datafusion/sql/tests/cases/plan_to_sql.rs
index 191425416c..de22321327 100644
--- a/datafusion/sql/tests/cases/plan_to_sql.rs
+++ b/datafusion/sql/tests/cases/plan_to_sql.rs
@@ -17,6 +17,7 @@
use arrow::datatypes::{DataType, Field, Schema};
+use datafusion_common::tree_node::{Transformed, TransformedResult};
use datafusion_common::{
Column, DFSchema, DFSchemaRef, DataFusionError, Result, TableReference,
assert_contains,
@@ -111,6 +112,25 @@ fn roundtrip_expr(table: TableReference, sql: &str) ->
Result<String> {
Ok(ast.to_string())
}
+fn remove_column_self_aliases(plan: LogicalPlan) -> Result<LogicalPlan> {
+ plan.transform_up_with_subqueries(|plan| {
+ plan.map_expressions(|expr| {
+ if let Expr::Alias(alias) = &expr
+ && alias.relation.is_none()
+ && alias.metadata.is_none()
+ && let Expr::Column(column) = alias.expr.as_ref()
+ && column.relation.is_none()
+ && column.name == alias.name
+ {
+ Ok(Transformed::yes(*alias.expr.clone()))
+ } else {
+ Ok(Transformed::no(expr))
+ }
+ })
+ })
+ .data()
+}
+
#[test]
fn roundtrip_statement() -> Result<()> {
let tests: Vec<&str> = vec![
@@ -254,7 +274,11 @@ fn roundtrip_statement() -> Result<()> {
.sql_statement_to_plan(roundtrip_statement.clone())
.unwrap();
- assert_eq!(plan, plan_roundtrip);
+ // Explicit output names can add unqualified self-aliases without
changing the plan's meaning.
+ assert_eq!(
+ remove_column_self_aliases(plan)?,
+ remove_column_self_aliases(plan_roundtrip)?,
+ );
}
Ok(())
@@ -392,6 +416,37 @@ fn roundtrip_statement_with_dialect_4() -> Result<(),
DataFusionError> {
Ok(())
}
+#[test]
+fn unparse_preserves_derived_aggregate_output_name() -> Result<()> {
+ let schema = Schema::new(vec![Field::new("j1_id", DataType::Int32,
false)]);
+ let aggregate = sum(col("j1.j1_id"));
+ let output =
Expr::Column(Column::from_name(aggregate.schema_name().to_string()));
+ let plan = table_scan(Some("j1"), &schema, None)?
+ .aggregate(Vec::<Expr>::new(), vec![aggregate])?
+ .project(vec![output.clone().alias("visible"), output.clone()])?
+ .project(vec![output])?
+ .build()?;
+
+ let sql = Unparser::new(&UnparserPostgreSqlDialect {})
+ .plan_to_sql(&plan)?
+ .to_string();
+ println!("UNPARSED_SQL={sql}");
+ assert_snapshot!(
+ sql,
+ @r#"SELECT "sum(j1.j1_id)" FROM (SELECT sum("j1"."j1_id") AS
"visible", sum("j1"."j1_id") AS "sum(j1.j1_id)" FROM "j1") AS
"derived_projection""#
+ );
+
+ let sql = Unparser::new(&BigQueryDialect {})
+ .plan_to_sql(&plan)?
+ .to_string();
+ println!("BIGQUERY_SQL={sql}");
+ assert_snapshot!(
+ sql,
+ @r#"SELECT `sum_40j1_46j1_id_41` FROM (SELECT sum(`j1`.`j1_id`) AS
`visible`, sum(`j1`.`j1_id`) AS `sum_40j1_46j1_id_41` FROM `j1`)"#
+ );
+ Ok(())
+}
+
#[test]
fn roundtrip_rebases_derived_projection_references() -> Result<(),
DataFusionError> {
roundtrip_statement_with_dialect_helper!(
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]