From 60c672debb6fe08bdc462ee2fd399d6c455ca39b Mon Sep 17 00:00:00 2001 From: linfeng Date: Thu, 23 Jul 2026 23:46:37 +0800 Subject: [PATCH 1/2] fix: preserve column qualifiers in CTAS with explicit schema --- datafusion/sql/src/statement.rs | 10 +++++----- datafusion/sqllogictest/test_files/ddl.slt | 8 ++++++++ 2 files changed, 13 insertions(+), 5 deletions(-) diff --git a/datafusion/sql/src/statement.rs b/datafusion/sql/src/statement.rs index ae7579c8c4dfb..3cfbb45688984 100644 --- a/datafusion/sql/src/statement.rs +++ b/datafusion/sql/src/statement.rs @@ -53,7 +53,7 @@ use datafusion_expr::{ LogicalPlan, LogicalPlanBuilder, OperateFunctionArg, PlanType, Prepare, ResetVariable, SetVariable, SortExpr, Statement as PlanStatement, ToStringifiedPlan, TransactionAccessMode, TransactionConclusion, TransactionEnd, - TransactionIsolationLevel, TransactionStart, Volatility, WriteOp, cast, col, + TransactionIsolationLevel, TransactionStart, Volatility, WriteOp, cast, }; use sqlparser::ast::{ self, BeginTransactionKind, CheckConstraint, ForeignKeyConstraint, IndexColumn, @@ -555,14 +555,14 @@ impl SqlToRel<'_, S> { input_schema.fields().len() ); } - let input_fields = input_schema.fields(); + let input_columns = input_schema.columns(); let project_exprs = schema .fields() .iter() - .zip(input_fields) - .map(|(field, input_field)| { + .zip(input_columns) + .map(|(field, input_column)| { cast( - col(input_field.name()), + Expr::Column(input_column), field.data_type().clone(), ) .alias(field.name()) diff --git a/datafusion/sqllogictest/test_files/ddl.slt b/datafusion/sqllogictest/test_files/ddl.slt index e1a48ce5e8e3c..d4aebb730459c 100644 --- a/datafusion/sqllogictest/test_files/ddl.slt +++ b/datafusion/sqllogictest/test_files/ddl.slt @@ -979,6 +979,14 @@ CREATE TABLE dup_src AS VALUES(1, 2); statement error DataFusion error: Schema error: Schema contains duplicate unqualified field name column1 CREATE TABLE dup_ctas AS SELECT * FROM dup_src LEFT JOIN dup_src y ON dup_src.column1 = y.column2; +statement ok +CREATE TABLE dup_ctas_with_schema(left_c1 bigint, right_c1 bigint) AS +SELECT dup_src.column1, y.column1 +FROM dup_src CROSS JOIN dup_src y; + +statement ok +DROP TABLE IF EXISTS dup_ctas_with_schema; + statement error DataFusion error: Schema error: Schema contains duplicate unqualified field name column1 CREATE VIEW dup_view AS SELECT * FROM dup_src LEFT JOIN dup_src y ON dup_src.column1 = y.column2; From fae3a8bc74803e4678d6f959c7cbaf599d2bdb97 Mon Sep 17 00:00:00 2001 From: linfeng <33561138+lyne7-sc@users.noreply.github.com> Date: Fri, 24 Jul 2026 14:39:17 +0800 Subject: [PATCH 2/2] verify ctas explicit schema column mapping --- datafusion/sqllogictest/test_files/ddl.slt | 12 +++++++++--- 1 file changed, 9 insertions(+), 3 deletions(-) diff --git a/datafusion/sqllogictest/test_files/ddl.slt b/datafusion/sqllogictest/test_files/ddl.slt index d4aebb730459c..672bab553330d 100644 --- a/datafusion/sqllogictest/test_files/ddl.slt +++ b/datafusion/sqllogictest/test_files/ddl.slt @@ -981,11 +981,17 @@ CREATE TABLE dup_ctas AS SELECT * FROM dup_src LEFT JOIN dup_src y ON dup_src.co statement ok CREATE TABLE dup_ctas_with_schema(left_c1 bigint, right_c1 bigint) AS -SELECT dup_src.column1, y.column1 -FROM dup_src CROSS JOIN dup_src y; +SELECT dup_src.column1, right_src.column1 +FROM dup_src +CROSS JOIN (SELECT column2 AS column1 FROM dup_src) right_src; + +query II +SELECT left_c1, right_c1 FROM dup_ctas_with_schema; +---- +1 2 statement ok -DROP TABLE IF EXISTS dup_ctas_with_schema; +DROP TABLE dup_ctas_with_schema; statement error DataFusion error: Schema error: Schema contains duplicate unqualified field name column1 CREATE VIEW dup_view AS SELECT * FROM dup_src LEFT JOIN dup_src y ON dup_src.column1 = y.column2;