diff --git a/datafusion/optimizer/src/decorrelate.rs b/datafusion/optimizer/src/decorrelate.rs index e4577d6eb0c4c..43af2b0ba47c8 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, 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}; @@ -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,22 +539,49 @@ 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 + // + // 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" + ); + }; + // 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() { + 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, + &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) @@ -828,6 +873,65 @@ fn can_pullup_over_aggregation(expr: &Expr) -> bool { } } +/// 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, + input_schema: &DFSchema, + output_schema: &DFSchema, +) -> 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 = alias_output_column(col, input_schema, output_schema)?; + return Ok(Transformed::yes(Expr::Column(new_col))); + } + Ok(Transformed::no(expr)) + }) + .data()?; + } + 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`). +/// +/// 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, +) -> 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/optimizer/src/decorrelate_lateral_join.rs b/datafusion/optimizer/src/decorrelate_lateral_join.rs index a8df5e69e3f33..0f78964090709 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}; @@ -133,22 +135,27 @@ 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 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(), )?); + 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)) - .transpose()?; - (right, corr, on) + (right, corr, original_join_filter) } else { (rewritten_subquery, correlation_filter, original_join_filter) }; @@ -347,24 +354,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 0e0ceb9cfa4fa..c9ffa520c27ef 100644 --- a/datafusion/sqllogictest/test_files/lateral_join.slt +++ b/datafusion/sqllogictest/test_files/lateral_join.slt @@ -935,6 +935,215 @@ 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 + +# 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); + +statement ok +CREATE TABLE other(id INT) AS VALUES (7); + +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 + +# 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; + +statement ok +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 c1fd8564976ac..a972909bf5d0c 100644 --- a/datafusion/sqllogictest/test_files/subquery.slt +++ b/datafusion/sqllogictest/test_files/subquery.slt @@ -3379,3 +3379,95 @@ 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; + +# 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;