From 2502691ec3d8371da9f2ad367031d6fb673a9027 Mon Sep 17 00:00:00 2001 From: Preston Zhao <195942578+canyang25@users.noreply.github.com> Date: Mon, 5 Oct 2026 03:40:59 +0000 Subject: [PATCH 1/4] fix: requalify LATERAL filters through derived-table aliases A correlated filter pulled through a SubqueryAlias kept the inner table qualifier. When that alias differed from the table name, the join condition named a column the rewritten subquery no longer had. Fixes #25837 Assisted by Cursor (Grok). --- datafusion/optimizer/src/decorrelate.rs | 65 ++++++++++++++++- .../sqllogictest/test_files/lateral_join.slt | 71 +++++++++++++++++++ 2 files changed, 134 insertions(+), 2 deletions(-) diff --git a/datafusion/optimizer/src/decorrelate.rs b/datafusion/optimizer/src/decorrelate.rs index e4577d6eb0c4c..f0edfab5ab581 100644 --- a/datafusion/optimizer/src/decorrelate.rs +++ b/datafusion/optimizer/src/decorrelate.rs @@ -27,7 +27,8 @@ use datafusion_common::tree_node::{ Transformed, TransformedResult, TreeNode, TreeNodeRecursion, TreeNodeRewriter, }; use datafusion_common::{ - Column, DFSchemaRef, HashMap, Result, ScalarValue, assert_or_internal_err, plan_err, + Column, DFSchemaRef, HashMap, Result, ScalarValue, TableReference, + assert_or_internal_err, plan_err, }; use datafusion_expr::expr::{Alias, GroupingSet}; use datafusion_expr::logical_plan::{Join, JoinType}; @@ -91,6 +92,12 @@ pub struct PullUpCorrelatedExpr { /// The list is cleared when the pull up passes a node that can put a NULL /// back into such a column: an outer join, a union or a grouping set. pub correlated_filters: Vec, + /// `join_filters.len()` at each `SubqueryAlias` on the way down. + /// + /// Filters appended while that alias is being rewritten belong to its + /// subtree. On the way up, only those filters are requalified to the + /// alias, so a sibling filter that names the same column is left alone. + alias_filter_base: Vec, } impl Default for PullUpCorrelatedExpr { @@ -113,6 +120,7 @@ impl PullUpCorrelatedExpr { pull_up_having_expr: None, pulled_up_scalar_agg: false, correlated_filters: Vec::new(), + alias_filter_base: Vec::new(), } } @@ -237,6 +245,17 @@ impl TreeNodeRewriter for PullUpCorrelatedExpr { _ => Ok(self.stop_pull_up(plan)), } } + // Remember how many join filters exist before this alias's + // children are rewritten. `f_up` requalifies only the filters + // those children add. See #25837. + LogicalPlan::SubqueryAlias(_) => { + self.alias_filter_base.push(self.join_filters.len()); + if plan.contains_outer_reference() { + Ok(self.stop_pull_up(plan)) + } else { + Ok(Transformed::no(plan)) + } + } // A correlated filter below these nodes can move above them. The // node itself must not hold an outer reference, because only a // Filter gives its outer references to the join. @@ -263,7 +282,6 @@ impl TreeNodeRewriter for PullUpCorrelatedExpr { | LogicalPlan::Join(_) | LogicalPlan::AsOfJoin(_) | LogicalPlan::Repartition(_) - | LogicalPlan::SubqueryAlias(_) | LogicalPlan::TableScan(_) | LogicalPlan::EmptyRelation(_) | LogicalPlan::Values(_) => { @@ -521,6 +539,19 @@ impl TreeNodeRewriter for PullUpCorrelatedExpr { &self.correlated_subquery_cols_map, &mut local_correlated_cols, ); + // A filter pulled up from below this alias still names the + // inner table (`l.id`). The alias's output names that column + // `l2.id`, and the join condition has to match. + // https://github.com/apache/datafusion/issues/25837 + let base = self + .alias_filter_base + .pop() + .expect("SubqueryAlias f_down pushes alias_filter_base"); + requalify_join_filter_columns( + &mut self.join_filters[base..], + &local_correlated_cols, + &alias.alias, + )?; let mut new_correlated_cols = BTreeSet::new(); for col in local_correlated_cols.iter() { new_correlated_cols @@ -828,6 +859,36 @@ fn can_pullup_over_aggregation(expr: &Expr) -> bool { } } +/// Rewrite columns in `filters` that are listed in `cols` so they use `alias`. +/// +/// `cols` are the correlated columns as named below the alias. Outer +/// references are not in that set, so they stay as they are. The caller +/// passes only the filters added under this alias. +fn requalify_join_filter_columns( + filters: &mut [Expr], + cols: &BTreeSet, + alias: &TableReference, +) -> Result<()> { + if cols.is_empty() || filters.is_empty() { + return Ok(()); + } + for filter in filters { + *filter = filter + .clone() + .transform(|expr| { + if let Expr::Column(col) = &expr + && cols.contains(col) + { + let new_col = Column::new(Some(alias.clone()), col.name.clone()); + return Ok(Transformed::yes(Expr::Column(new_col))); + } + Ok(Transformed::no(expr)) + }) + .data()?; + } + Ok(()) +} + fn collect_local_correlated_cols( plan: &LogicalPlan, all_cols_map: &HashMap>, diff --git a/datafusion/sqllogictest/test_files/lateral_join.slt b/datafusion/sqllogictest/test_files/lateral_join.slt index 0e0ceb9cfa4fa..e7786ddb067d4 100644 --- a/datafusion/sqllogictest/test_files/lateral_join.slt +++ b/datafusion/sqllogictest/test_files/lateral_join.slt @@ -935,6 +935,77 @@ physical_plan 02)--DataSourceExec: partitions=1, partition_sizes=[1] 03)--DataSourceExec: partitions=1, partition_sizes=[1] +########################################################### +# Section 8: Derived-table alias differs from the inner table +# +# https://github.com/apache/datafusion/issues/25837 +# The correlated filter is pulled through a SubqueryAlias whose name is +# not the inner table's name. The join condition must use the alias. +########################################################### + +statement ok +CREATE TABLE o(k INT) AS VALUES (1), (2); + +statement ok +CREATE TABLE l(id INT, v INT) AS VALUES (1, 10), (2, 20); + +statement ok +CREATE TABLE m(id INT) AS VALUES (1), (2); + +query II +SELECT o.k, s.v +FROM o, LATERAL ( + SELECT l2.v FROM (SELECT * FROM l WHERE l.id = o.k) AS l2 +) AS s +ORDER BY o.k; +---- +1 10 +2 20 + +# Projection inside the derived table still uses the inner qualifier. +query II +SELECT o.k, s.v +FROM o, LATERAL ( + SELECT x.v FROM (SELECT l.v, l.id FROM l WHERE l.id = o.k) AS x +) AS s +ORDER BY o.k; +---- +1 10 +2 20 + +# Matching alias and table name already worked; keep it that way. +query II +SELECT o.k, s.v +FROM o, LATERAL ( + SELECT l.v FROM (SELECT * FROM l WHERE l.id = o.k) AS l +) AS s +ORDER BY o.k; +---- +1 10 +2 20 + +# The derived table is one input of a join inside the LATERAL subquery. +query II +SELECT o.k, s.v +FROM o, LATERAL ( + SELECT l2.v + FROM (SELECT * FROM l WHERE l.id = o.k) AS l2 + JOIN m ON m.id = l2.id +) AS s +ORDER BY o.k; +---- +1 10 +2 20 + +statement ok +DROP TABLE o; + +statement ok +DROP TABLE l; + +statement ok +DROP TABLE m; + ########################################################### # Cleanup ########################################################### From f2dd4eef53f62eec65460d96ef70dd30d9ea0340 Mon Sep 17 00:00:00 2001 From: Preston Zhao <195942578+canyang25@users.noreply.github.com> Date: Wed, 7 Oct 2026 14:41:36 +0000 Subject: [PATCH 2/4] fix: map correlated columns through SubqueryAlias output names Pull-up can add a correlation column whose name matches another projected column. SubqueryAlias suffixes the duplicate (id:1), so requalifying only the qualifier targets the wrong column. Map both the join filters and the propagated correlation columns through the rebuilt alias's positional output names. Assisted by Cursor (Grok). --- datafusion/optimizer/src/decorrelate.rs | 83 ++++++++++++++----- .../sqllogictest/test_files/lateral_join.slt | 25 ++++++ 2 files changed, 85 insertions(+), 23 deletions(-) diff --git a/datafusion/optimizer/src/decorrelate.rs b/datafusion/optimizer/src/decorrelate.rs index f0edfab5ab581..3d15a7a2fb971 100644 --- a/datafusion/optimizer/src/decorrelate.rs +++ b/datafusion/optimizer/src/decorrelate.rs @@ -27,8 +27,8 @@ use datafusion_common::tree_node::{ Transformed, TransformedResult, TreeNode, TreeNodeRecursion, TreeNodeRewriter, }; use datafusion_common::{ - Column, DFSchemaRef, HashMap, Result, ScalarValue, TableReference, - assert_or_internal_err, plan_err, + Column, DFSchema, DFSchemaRef, HashMap, Result, ScalarValue, assert_or_internal_err, + internal_err, plan_err, }; use datafusion_expr::expr::{Alias, GroupingSet}; use datafusion_expr::logical_plan::{Join, JoinType}; @@ -543,31 +543,43 @@ impl TreeNodeRewriter for PullUpCorrelatedExpr { // inner table (`l.id`). The alias's output names that column // `l2.id`, and the join condition has to match. // https://github.com/apache/datafusion/issues/25837 - let base = self - .alias_filter_base - .pop() - .expect("SubqueryAlias f_down pushes alias_filter_base"); + // + // Pull-up can add a column whose name another projected + // column already has. `SubqueryAlias::try_new` suffixes that + // duplicate (`id` becomes `id:1`), so the output name is the + // field at the same position, not the inner name. + let Some(base) = self.alias_filter_base.pop() else { + return internal_err!( + "SubqueryAlias f_down pushes alias_filter_base" + ); + }; + // Schema `try_new` qualifies. Read it before rebuilding: the + // projection `try_new` inserts already uses the suffixed names. + let input_schema = Arc::clone(alias.input.schema()); + let new_plan = + if input_schema.fields().len() != alias.schema.fields().len() { + LogicalPlanBuilder::from((*alias.input).clone()) + .alias(alias.alias.clone())? + .build()? + } else { + plan.clone() + }; + let output_schema = Arc::clone(new_plan.schema()); requalify_join_filter_columns( &mut self.join_filters[base..], &local_correlated_cols, - &alias.alias, + &input_schema, + &output_schema, )?; let mut new_correlated_cols = BTreeSet::new(); - for col in local_correlated_cols.iter() { - new_correlated_cols - .insert(Column::new(Some(alias.alias.clone()), col.name.clone())); + for col in &local_correlated_cols { + new_correlated_cols.insert(alias_output_column( + col, + &input_schema, + &output_schema, + )?); } - let new_plan = if alias.input.schema().fields().len() - != alias.schema.fields().len() - { - LogicalPlanBuilder::from((*alias.input).clone()) - .alias(alias.alias.clone())? - .build()? - } else { - plan.clone() - }; - self.correlated_subquery_cols_map .insert(new_plan.clone(), new_correlated_cols); if let Some(input_map) = self.collected_count_expr_map.get(&*alias.input) @@ -859,15 +871,21 @@ fn can_pullup_over_aggregation(expr: &Expr) -> bool { } } -/// Rewrite columns in `filters` that are listed in `cols` so they use `alias`. +/// Rewrite columns in `filters` that are listed in `cols` so they use the +/// alias's output column at the same position. /// /// `cols` are the correlated columns as named below the alias. Outer /// references are not in that set, so they stay as they are. The caller /// passes only the filters added under this alias. +/// +/// `input_schema` is the schema passed to `SubqueryAlias::try_new`. +/// `output_schema` is the schema it built. Fields line up by position, +/// including the dedup suffixes `try_new` applies. fn requalify_join_filter_columns( filters: &mut [Expr], cols: &BTreeSet, - alias: &TableReference, + input_schema: &DFSchema, + output_schema: &DFSchema, ) -> Result<()> { if cols.is_empty() || filters.is_empty() { return Ok(()); @@ -879,7 +897,7 @@ fn requalify_join_filter_columns( if let Expr::Column(col) = &expr && cols.contains(col) { - let new_col = Column::new(Some(alias.clone()), col.name.clone()); + let new_col = alias_output_column(col, input_schema, output_schema)?; return Ok(Transformed::yes(Expr::Column(new_col))); } Ok(Transformed::no(expr)) @@ -889,6 +907,25 @@ fn requalify_join_filter_columns( Ok(()) } +/// The alias output column that corresponds to `col` in `input_schema`. +/// +/// `SubqueryAlias::try_new` keeps field order. When two input fields share a +/// name, the later one's output name is suffixed (`id`, then `id:1`). +fn alias_output_column( + col: &Column, + input_schema: &DFSchema, + output_schema: &DFSchema, +) -> Result { + let idx = input_schema.index_of_column(col)?; + assert_or_internal_err!( + idx < output_schema.fields().len(), + "SubqueryAlias output schema has {} fields, input column {col} is at {idx}", + output_schema.fields().len(), + ); + let (qualifier, field) = output_schema.qualified_field(idx); + Ok(Column::new(qualifier.cloned(), field.name())) +} + fn collect_local_correlated_cols( plan: &LogicalPlan, all_cols_map: &HashMap>, diff --git a/datafusion/sqllogictest/test_files/lateral_join.slt b/datafusion/sqllogictest/test_files/lateral_join.slt index e7786ddb067d4..8c7ba895aaf9f 100644 --- a/datafusion/sqllogictest/test_files/lateral_join.slt +++ b/datafusion/sqllogictest/test_files/lateral_join.slt @@ -997,6 +997,31 @@ ORDER BY o.k; 1 10 2 20 +# Pull-up appends l.id next to another projected id. SubqueryAlias +# keeps the first id and suffixes the pulled-up column to id:1. The +# join has to follow that output name; decoy.id is 100, so using the +# unsuffixed id matches nothing. +statement ok +CREATE TABLE decoy(id INT) AS VALUES (100); + +query II +SELECT o.k, s.v +FROM o, LATERAL ( + SELECT l2.v + FROM ( + SELECT decoy.id, l.v + FROM l, decoy + WHERE l.id = o.k + ) AS l2 +) AS s +ORDER BY o.k; +---- +1 10 +2 20 + +statement ok +DROP TABLE decoy; + statement ok DROP TABLE o; From e8be00684072f8265c5b477825cb5670b0501ea7 Mon Sep 17 00:00:00 2001 From: Preston Zhao <195942578+canyang25@users.noreply.github.com> Date: Wed, 7 Oct 2026 15:17:59 +0000 Subject: [PATCH 3/4] fix: map outer alias columns through SubqueryAlias output names The lateral join and EXISTS/IN rewrites wrapped the pulled-up subquery with SubqueryAlias::try_new, then requalified filters by keeping the inner column name. When try_new suffixes a duplicate (id, id:1, id:2), that name is an earlier column. Both rewrites now use the same positional input-to-output mapping as alias_output_column, and a correlated column missing from the pre-alias schema is an error. Assisted by Cursor (Grok). --- datafusion/optimizer/src/decorrelate.rs | 12 +++- .../optimizer/src/decorrelate_lateral_join.rs | 34 ++++++---- .../src/decorrelate_predicate_subquery.rs | 14 ++++- .../optimizer/src/scalar_subquery_to_join.rs | 14 ++++- datafusion/optimizer/src/utils.rs | 36 ++++++++--- .../sqllogictest/test_files/lateral_join.slt | 62 +++++++++++++++++++ .../sqllogictest/test_files/subquery.slt | 28 +++++++++ 7 files changed, 170 insertions(+), 30 deletions(-) diff --git a/datafusion/optimizer/src/decorrelate.rs b/datafusion/optimizer/src/decorrelate.rs index 3d15a7a2fb971..43af2b0ba47c8 100644 --- a/datafusion/optimizer/src/decorrelate.rs +++ b/datafusion/optimizer/src/decorrelate.rs @@ -553,8 +553,10 @@ impl TreeNodeRewriter for PullUpCorrelatedExpr { "SubqueryAlias f_down pushes alias_filter_base" ); }; - // Schema `try_new` qualifies. Read it before rebuilding: the - // projection `try_new` inserts already uses the suffixed names. + // The projection and the duplicate-name suffixes come from + // `SubqueryAlias::try_new`. Read this schema before rebuilding: + // the projection that constructor inserts already uses the + // suffixed names. let input_schema = Arc::clone(alias.input.schema()); let new_plan = if input_schema.fields().len() != alias.schema.fields().len() { @@ -911,7 +913,11 @@ fn requalify_join_filter_columns( /// /// `SubqueryAlias::try_new` keeps field order. When two input fields share a /// name, the later one's output name is suffixed (`id`, then `id:1`). -fn alias_output_column( +/// +/// Returns an error if `col` is not in `input_schema`, or if that position +/// is past the end of `output_schema`. Callers must not skip the column: +/// leaving it out drops the predicate that referenced it. +pub(crate) fn alias_output_column( col: &Column, input_schema: &DFSchema, output_schema: &DFSchema, diff --git a/datafusion/optimizer/src/decorrelate_lateral_join.rs b/datafusion/optimizer/src/decorrelate_lateral_join.rs index a8df5e69e3f33..5926bfabaf090 100644 --- a/datafusion/optimizer/src/decorrelate_lateral_join.rs +++ b/datafusion/optimizer/src/decorrelate_lateral_join.rs @@ -19,7 +19,9 @@ use std::sync::Arc; -use crate::decorrelate::{PullUpCorrelatedExpr, UN_MATCHED_ROW_INDICATOR}; +use crate::decorrelate::{ + PullUpCorrelatedExpr, UN_MATCHED_ROW_INDICATOR, alias_output_column, +}; use crate::optimizer::ApplyOrder; use crate::utils::evaluates_to_null; use crate::{OptimizerConfig, OptimizerRule}; @@ -137,16 +139,19 @@ fn rewrite_internal(join: Join) -> Result> { // rewrite column references in both the correlation and ON-clause filters. let (right_plan, correlation_filter, original_join_filter) = if let Some(ref alias) = alias { - let inner_schema = Arc::clone(rewritten_subquery.schema()); + let input_schema = Arc::clone(rewritten_subquery.schema()); let right = LogicalPlan::SubqueryAlias(SubqueryAlias::try_new( Arc::new(rewritten_subquery), alias.clone(), )?); + // `try_new` may suffix a duplicate (`id`, then `id:1`). Map from + // the schema it consumed to the schema it produced. + let output_schema = Arc::clone(right.schema()); let corr = correlation_filter - .map(|f| requalify_filter(f, &inner_schema, alias)) + .map(|f| requalify_filter(f, &input_schema, &output_schema)) .transpose()?; let on = original_join_filter - .map(|f| requalify_filter(f, &inner_schema, alias)) + .map(|f| requalify_filter(f, &input_schema, &output_schema)) .transpose()?; (right, corr, on) } else { @@ -347,24 +352,31 @@ fn extract_lateral_subquery( } /// Rewrite column references in a join filter expression so that columns -/// belonging to the inner (right) side use the SubqueryAlias qualifier. +/// belonging to the inner (right) side use the `SubqueryAlias` output column +/// at the same position. /// /// The `PullUpCorrelatedExpr` pass extracts join filters with the inner /// columns qualified by their original table names (e.g., `t2.t1_id`). /// When the inner plan is wrapped in a `SubqueryAlias("sub")`, those -/// columns are re-qualified as `sub.t1_id`. This function applies the -/// same requalification to the filter so it matches the aliased schema. +/// columns are re-qualified from `input_schema` (the schema passed to +/// `SubqueryAlias::try_new`) to `output_schema` (the schema it built). +/// A duplicate name is suffixed there (`id`, then `id:1`), so keeping +/// `col.name` would point at the earlier column. +/// +/// Columns that are not in `input_schema` are outer references and stay as +/// they are. A column that is in `input_schema` but has no output field is +/// an error: skipping it would drop that predicate. fn requalify_filter( filter: Expr, - inner_schema: &DFSchema, - alias: &TableReference, + input_schema: &DFSchema, + output_schema: &DFSchema, ) -> Result { filter .transform(|expr| { if let Expr::Column(col) = &expr - && inner_schema.has_column(col) + && input_schema.has_column(col) { - let new_col = Column::new(Some(alias.clone()), col.name.clone()); + let new_col = alias_output_column(col, input_schema, output_schema)?; return Ok(Transformed::yes(Expr::Column(new_col))); } Ok(Transformed::no(expr)) diff --git a/datafusion/optimizer/src/decorrelate_predicate_subquery.rs b/datafusion/optimizer/src/decorrelate_predicate_subquery.rs index 88a929654bea8..bf890ad912af2 100644 --- a/datafusion/optimizer/src/decorrelate_predicate_subquery.rs +++ b/datafusion/optimizer/src/decorrelate_predicate_subquery.rs @@ -742,6 +742,7 @@ fn build_join( return Ok(None); } + let input_schema = Arc::clone(new_plan.schema()); let sub_query_alias = LogicalPlanBuilder::from(new_plan) .alias(alias.to_string())? .build()?; @@ -752,9 +753,16 @@ fn build_join( .for_each(|cols| all_correlated_cols.extend(cols.clone())); // alias the join filter - let join_filter_opt = conjunction(pull_up.join_filters) - .map_or(Ok(None), |filter| { - replace_qualified_name(filter, &all_correlated_cols, alias).map(Some) + let output_schema = Arc::clone(sub_query_alias.schema()); + let join_filter_opt = + conjunction(pull_up.join_filters).map_or(Ok(None), |filter| { + replace_qualified_name( + filter, + &all_correlated_cols, + &input_schema, + &output_schema, + ) + .map(Some) })?; // The two sides of the `IN` predicate: the value from the outer plan and diff --git a/datafusion/optimizer/src/scalar_subquery_to_join.rs b/datafusion/optimizer/src/scalar_subquery_to_join.rs index c6374660ef1f9..17eaf5791b84b 100644 --- a/datafusion/optimizer/src/scalar_subquery_to_join.rs +++ b/datafusion/optimizer/src/scalar_subquery_to_join.rs @@ -361,6 +361,7 @@ fn build_join( .collected_count_expr_map .get(&decorrelated_subquery) .cloned(); + let input_schema = Arc::clone(decorrelated_subquery.schema()); let aliased_subquery = LogicalPlanBuilder::from(decorrelated_subquery) .alias(subquery_alias.to_string())? .build()?; @@ -372,11 +373,18 @@ fn build_join( .cloned() .collect(); - // Correlated columns now live in the decorrelated subquery's output, - // so re-qualify them with the subquery alias. + // Correlated columns now live in the decorrelated subquery's output. + // `try_new` may suffix a duplicate name, so use that output column. + let output_schema = Arc::clone(aliased_subquery.schema()); let join_filter_opt = conjunction(pull_up.join_filters).map_or(Ok(None), |filter| { - replace_qualified_name(filter, &all_correlated_cols, subquery_alias).map(Some) + replace_qualified_name( + filter, + &all_correlated_cols, + &input_schema, + &output_schema, + ) + .map(Some) })?; // When pull-up did not extract any usable join keys (a correlated subquery diff --git a/datafusion/optimizer/src/utils.rs b/datafusion/optimizer/src/utils.rs index 80045955773ca..969c685c40c48 100644 --- a/datafusion/optimizer/src/utils.rs +++ b/datafusion/optimizer/src/utils.rs @@ -24,7 +24,9 @@ use arrow::array::{Array, RecordBatch, new_null_array}; use arrow::datatypes::{DataType, Field, Schema}; use datafusion_common::TableReference; use datafusion_common::cast::as_boolean_array; -use datafusion_common::tree_node::{TransformedResult, TreeNode, TreeNodeRecursion}; +use datafusion_common::tree_node::{ + Transformed, TransformedResult, TreeNode, TreeNodeRecursion, +}; use datafusion_common::{Column, DFSchema, Result, ScalarValue}; use datafusion_expr::execution_props::ExecutionProps; use datafusion_expr::expr::{Exists, InSubquery, SetComparison}; @@ -159,19 +161,33 @@ pub(crate) fn has_all_column_refs( == column_refs.len() } +/// Requalify correlated columns in `expr` to the `SubqueryAlias` output +/// column at the same position in `output_schema`. +/// +/// `input_schema` is the schema passed to `SubqueryAlias::try_new`. +/// `output_schema` is the schema it built, including dedup suffixes. +/// Only columns listed in `cols` are rewritten. A listed column that is +/// not in `input_schema` is an error: skipping it would drop that predicate. pub(crate) fn replace_qualified_name( expr: Expr, cols: &BTreeSet, - subquery_alias: &str, + input_schema: &DFSchema, + output_schema: &DFSchema, ) -> Result { - let alias_cols: Vec = cols - .iter() - .map(|col| Column::new(Some(subquery_alias), &col.name)) - .collect(); - let replace_map: HashMap<&Column, &Column> = - cols.iter().zip(alias_cols.iter()).collect(); - - replace_col(expr, &replace_map) + expr.transform(|expr| { + if let Expr::Column(col) = &expr + && cols.contains(col) + { + let new_col = crate::decorrelate::alias_output_column( + col, + input_schema, + output_schema, + )?; + return Ok(Transformed::yes(Expr::Column(new_col))); + } + Ok(Transformed::no(expr)) + }) + .data() } /// Column reference to avoid copying string around diff --git a/datafusion/sqllogictest/test_files/lateral_join.slt b/datafusion/sqllogictest/test_files/lateral_join.slt index 8c7ba895aaf9f..90f8cf45413a8 100644 --- a/datafusion/sqllogictest/test_files/lateral_join.slt +++ b/datafusion/sqllogictest/test_files/lateral_join.slt @@ -1004,6 +1004,9 @@ ORDER BY o.k; statement ok CREATE TABLE decoy(id INT) AS VALUES (100); +statement ok +CREATE TABLE other(id INT) AS VALUES (7); + query II SELECT o.k, s.v FROM o, LATERAL ( @@ -1019,9 +1022,68 @@ ORDER BY o.k; 1 10 2 20 +# The lateral alias itself is what SubqueryAlias suffixes. There is no +# inner alias. decoy.id keeps `id`; the pulled-up l.id becomes `id:1`. +query II +SELECT o.k, s.v +FROM o, LATERAL ( + SELECT decoy.id, l.v + FROM l, decoy + WHERE l.id = o.k +) s +ORDER BY o.k; +---- +1 10 +2 20 + +# Same collision on a LEFT JOIN. A filter on decoy.id matches nothing, +# so the outer rows would come back with NULL v. +query II +SELECT o.k, s.v +FROM o +LEFT JOIN LATERAL ( + SELECT decoy.id, l.v + FROM l, decoy + WHERE l.id = o.k +) s ON true +ORDER BY o.k; +---- +1 10 +2 20 + +# decoy.id, other.id, then pulled-up l.id. The filter has to use id:2. +query II +SELECT o.k, s.v +FROM o, LATERAL ( + SELECT decoy.id, other.id, l.v + FROM l, decoy, other + WHERE l.id = o.k +) s +ORDER BY o.k; +---- +1 10 +2 20 + +# l.v is already named id:1, so the pulled-up l.id becomes id:2. +# s."id:1" is l.v. +query II +SELECT o.k, s."id:1" +FROM o, LATERAL ( + SELECT decoy.id, l.v AS "id:1" + FROM l, decoy + WHERE l.id = o.k +) s +ORDER BY o.k; +---- +1 10 +2 20 + statement ok DROP TABLE decoy; +statement ok +DROP TABLE other; + statement ok DROP TABLE o; diff --git a/datafusion/sqllogictest/test_files/subquery.slt b/datafusion/sqllogictest/test_files/subquery.slt index c1fd8564976ac..04202d67caf28 100644 --- a/datafusion/sqllogictest/test_files/subquery.slt +++ b/datafusion/sqllogictest/test_files/subquery.slt @@ -3379,3 +3379,31 @@ DROP TABLE unnest_outer; statement ok DROP TABLE unnest_inner; + +# EXISTS pull-up adds l.id beside decoy.id. The subquery alias suffixes +# the duplicate to id:1. Filtering on the unsuffixed id matches nothing. +statement ok +CREATE TABLE o(k INT) AS VALUES (1), (2); + +statement ok +CREATE TABLE l(id INT, v INT) AS VALUES (1, 10), (2, 20); + +statement ok +CREATE TABLE decoy(id INT) AS VALUES (100); + +query I +SELECT o.k FROM o WHERE EXISTS ( + SELECT decoy.id, l.v FROM l, decoy WHERE l.id = o.k +) ORDER BY o.k; +---- +1 +2 + +statement ok +DROP TABLE o; + +statement ok +DROP TABLE l; + +statement ok +DROP TABLE decoy; From e8554d2b027e3b243a9c03f93c5a3ee98caed398 Mon Sep 17 00:00:00 2001 From: Preston Zhao <195942578+canyang25@users.noreply.github.com> Date: Wed, 7 Oct 2026 16:05:28 +0000 Subject: [PATCH 4/4] fix: keep LATERAL ON clauses on alias output names requalify_filter mapped both the pulled-up correlation predicate and the user's ON clause through the pre-alias schema. ON is already bound to the alias output. When an inner column shares that qualifier, the map sends ON s.id to a later id:1. Pull-up only appends columns, and those are the only names unique_field_aliases suffixes, so the planner's ON names stay valid. EXISTS, IN, and scalar decorrelation only positionally rewrite correlation columns, not a user predicate that already uses output names. Assisted by Cursor (Grok). --- .../optimizer/src/decorrelate_lateral_join.rs | 18 +++--- .../sqllogictest/test_files/lateral_join.slt | 51 +++++++++++++++ .../sqllogictest/test_files/subquery.slt | 64 +++++++++++++++++++ 3 files changed, 125 insertions(+), 8 deletions(-) diff --git a/datafusion/optimizer/src/decorrelate_lateral_join.rs b/datafusion/optimizer/src/decorrelate_lateral_join.rs index 5926bfabaf090..0f78964090709 100644 --- a/datafusion/optimizer/src/decorrelate_lateral_join.rs +++ b/datafusion/optimizer/src/decorrelate_lateral_join.rs @@ -135,8 +135,15 @@ fn rewrite_internal(join: Join) -> Result> { .cloned(); // Re-wrap in SubqueryAlias if the original had one, preserving the alias name. - // The SubqueryAlias re-qualifies all columns with the alias, so we must also - // rewrite column references in both the correlation and ON-clause filters. + // The correlation filter still names the inner tables. `try_new` may suffix + // a duplicate, so those columns are mapped by position. + // + // The user's ON clause already names the alias output. Pull-up only appends + // columns, and `unique_field_aliases` only suffixes those new columns, so + // the names ON uses do not change. Mapping ON through the pre-alias schema + // is wrong when an inner column shares the alias qualifier (`l AS s` inside + // `LATERAL (...) s`): `s.id` would move from the first output `id` to a + // later `id:1`. let (right_plan, correlation_filter, original_join_filter) = if let Some(ref alias) = alias { let input_schema = Arc::clone(rewritten_subquery.schema()); @@ -144,16 +151,11 @@ fn rewrite_internal(join: Join) -> Result> { Arc::new(rewritten_subquery), alias.clone(), )?); - // `try_new` may suffix a duplicate (`id`, then `id:1`). Map from - // the schema it consumed to the schema it produced. let output_schema = Arc::clone(right.schema()); let corr = correlation_filter .map(|f| requalify_filter(f, &input_schema, &output_schema)) .transpose()?; - let on = original_join_filter - .map(|f| requalify_filter(f, &input_schema, &output_schema)) - .transpose()?; - (right, corr, on) + (right, corr, original_join_filter) } else { (rewritten_subquery, correlation_filter, original_join_filter) }; diff --git a/datafusion/sqllogictest/test_files/lateral_join.slt b/datafusion/sqllogictest/test_files/lateral_join.slt index 90f8cf45413a8..c9ffa520c27ef 100644 --- a/datafusion/sqllogictest/test_files/lateral_join.slt +++ b/datafusion/sqllogictest/test_files/lateral_join.slt @@ -1093,6 +1093,57 @@ DROP TABLE l; statement ok DROP TABLE m; +# The user's ON clause names the alias output (`s.id` is decoy.id, the +# first id). The inner table is also `s`, so a positional map of that +# ON predicate follows the later field and filters `l` instead. +statement ok +CREATE TABLE o(k INT) AS VALUES (1), (100); + +statement ok +CREATE TABLE l(id INT, v INT) AS VALUES (1, 10), (100, 99); + +statement ok +CREATE TABLE decoy(id INT) AS VALUES (100); + +query IIII +SELECT o.k, s.id, s."id:1", s.v +FROM o +JOIN LATERAL ( + SELECT decoy.id, s.id, s.v + FROM l AS s, decoy + WHERE s.id = o.k +) s +ON s.id = 100 +ORDER BY o.k, s.v; +---- +1 100 1 10 +100 100 100 99 + +# `s.id` in ON is decoy.id, which is 100, so both correlated rows pass. +# Remapping that ON onto l.id turns the k=1 row into NULLs. +query IIII +SELECT o.k, s.id, s."id:1", s.v +FROM o +LEFT JOIN LATERAL ( + SELECT decoy.id, s.id, s.v + FROM l AS s, decoy + WHERE s.id = o.k +) s +ON s.id = 100 +ORDER BY o.k, s.v; +---- +1 100 1 10 +100 100 100 99 + +statement ok +DROP TABLE o; + +statement ok +DROP TABLE l; + +statement ok +DROP TABLE decoy; + ########################################################### # Cleanup ########################################################### diff --git a/datafusion/sqllogictest/test_files/subquery.slt b/datafusion/sqllogictest/test_files/subquery.slt index 04202d67caf28..a972909bf5d0c 100644 --- a/datafusion/sqllogictest/test_files/subquery.slt +++ b/datafusion/sqllogictest/test_files/subquery.slt @@ -3407,3 +3407,67 @@ DROP TABLE l; statement ok DROP TABLE decoy; + +# The IN value is decoy_hit.id (1). The correlation is l_miss.id (99), +# which matches none of 1, 2, 3, so IN is empty and NOT IN / NOT EXISTS +# keep every outer row. +statement ok +CREATE TABLE o(k INT) AS VALUES (1), (2), (3); + +statement ok +CREATE TABLE l_miss(id INT) AS VALUES (99); + +statement ok +CREATE TABLE decoy_hit(id INT) AS VALUES (1); + +query I +SELECT o.k FROM o WHERE o.k IN ( + SELECT decoy_hit.id FROM l_miss, decoy_hit WHERE l_miss.id = o.k +) ORDER BY o.k; +---- + +query I +SELECT o.k FROM o WHERE o.k NOT IN ( + SELECT decoy_hit.id FROM l_miss, decoy_hit WHERE l_miss.id = o.k +) ORDER BY o.k; +---- +1 +2 +3 + +query I +SELECT o.k FROM o WHERE NOT EXISTS ( + SELECT decoy_hit.id FROM l_miss, decoy_hit WHERE l_miss.id = o.k +) ORDER BY o.k; +---- +1 +2 +3 + +statement ok +DROP TABLE l_miss; + +statement ok +DROP TABLE decoy_hit; + +statement ok +CREATE TABLE l(id INT, v INT) AS VALUES (1, 10), (2, 20); + +statement ok +CREATE TABLE decoy(id INT) AS VALUES (100); + +query I +SELECT o.k FROM o WHERE NOT EXISTS ( + SELECT decoy.id, l.v FROM l, decoy WHERE l.id = o.k +) ORDER BY o.k; +---- +3 + +statement ok +DROP TABLE o; + +statement ok +DROP TABLE l; + +statement ok +DROP TABLE decoy;