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
134 changes: 119 additions & 15 deletions datafusion/optimizer/src/decorrelate.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Expand Down Expand Up @@ -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<Expr>,
/// `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<usize>,
}

impl Default for PullUpCorrelatedExpr {
Expand All @@ -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(),
}
}

Expand Down Expand Up @@ -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.
Expand All @@ -263,7 +282,6 @@ impl TreeNodeRewriter for PullUpCorrelatedExpr {
| LogicalPlan::Join(_)
| LogicalPlan::AsOfJoin(_)
| LogicalPlan::Repartition(_)
| LogicalPlan::SubqueryAlias(_)
| LogicalPlan::TableScan(_)
| LogicalPlan::EmptyRelation(_)
| LogicalPlan::Values(_) => {
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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<Column>,
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<Column> {
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<LogicalPlan, BTreeSet<Column>>,
Expand Down
46 changes: 30 additions & 16 deletions datafusion/optimizer/src/decorrelate_lateral_join.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Expand Down Expand Up @@ -133,22 +135,27 @@ fn rewrite_internal(join: Join) -> Result<Transformed<LogicalPlan>> {
.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)
};
Expand Down Expand Up @@ -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<Expr> {
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))
Expand Down
14 changes: 11 additions & 3 deletions datafusion/optimizer/src/decorrelate_predicate_subquery.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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()?;
Expand All @@ -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
Expand Down
14 changes: 11 additions & 3 deletions datafusion/optimizer/src/scalar_subquery_to_join.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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()?;
Expand All @@ -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
Expand Down
36 changes: 26 additions & 10 deletions datafusion/optimizer/src/utils.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Expand Down Expand Up @@ -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<Column>,
subquery_alias: &str,
input_schema: &DFSchema,
output_schema: &DFSchema,
) -> Result<Expr> {
let alias_cols: Vec<Column> = 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
Expand Down
Loading