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
81 changes: 44 additions & 37 deletions datafusion/optimizer/src/extract_leaf_expressions.rs
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,19 @@
//! that runs before the extraction projection exists, so row group pruning and
//! source level filtering still happen.
//!
//! # Yields to [`OptimizeProjections`]
//!
//! When the extraction projection cannot move below the input of a projection
//! (for example a `TableScan` or an `Aggregate`), [`PushDownLeafProjections`]
//! leaves the projection unchanged. Splitting it in place would give a
//! recovery projection over an extraction projection on the same input, and
//! `OptimizeProjections` merges those two back into the original projection.
//! The two rules would then undo each other on every optimizer pass. The split
//! has no use there: the source absorbs the leaf expressions of the original
//! projection just as well.
//!
//! [`PushDownFilter`]: crate::push_down_filter::PushDownFilter
//! [`OptimizeProjections`]: crate::optimize_projections::OptimizeProjections

use indexmap::{IndexMap, IndexSet};
use std::collections::{BTreeSet, HashMap};
Expand Down Expand Up @@ -1113,18 +1125,33 @@ fn split_and_push_projection(
// Only pre-existing __datafusion_extracted aliases and columns, no new
// extractions from routing_extract. The original projection is
// already an extraction projection that couldn't be pushed
// further. Return None.
// further. Return None. This return also ends the recursion
// of the `try_push_input` call below, which comes back here
// with a pure extraction projection.
return Ok(None);
}
// Build extraction projection in-place (couldn't push down)
// Build the extraction projection in place, then push it. It is a
// pure extraction projection, which can go through nodes that the
// original projection cannot.
let input_arc = Arc::clone(input);
let extraction = build_extraction_projection_impl(
let extraction = LogicalPlan::Projection(build_extraction_projection_impl(
&extraction_pairs,
columns_needed,
&input_arc,
input_schema.as_ref(),
)?;
LogicalPlan::Projection(extraction)
)?);
match try_push_input(&extraction, alias_generator)? {
Some(pushed) => pushed,
// The extraction projection cannot move below `input`. Leave
// the original projection unchanged: the split would put the
// extraction projection on the same input, and
// `OptimizeProjections` merges the recovery and extraction
// projections back into the original projection. The rules
// would then undo each other in every optimizer pass. The
// source absorbs the leaf expressions of the original
// projection just as well.
None => return Ok(None),
}
}
};

Expand Down Expand Up @@ -1653,13 +1680,10 @@ mod tests {
(same as original)

## After Pushdown
Projection: __datafusion_extracted_1 AS leaf_udf(test.user,Utf8("name"))
Projection: leaf_udf(test.user, Utf8("name")) AS __datafusion_extracted_1, test.user
TableScan: test projection=[user]
(same as after extraction)

## Optimized
Projection: leaf_udf(test.user, Utf8("name"))
TableScan: test projection=[user]
(same as after pushdown)
"#)
}

Expand All @@ -1683,13 +1707,10 @@ mod tests {
(same as original)

## After Pushdown
Projection: __datafusion_extracted_1 IS NOT NULL AS has_name
Projection: leaf_udf(test.user, Utf8("name")) AS __datafusion_extracted_1, test.user
TableScan: test projection=[user]
(same as after extraction)

## Optimized
Projection: leaf_udf(test.user, Utf8("name")) IS NOT NULL AS has_name
TableScan: test projection=[user]
(same as after pushdown)
"#)
}

Expand Down Expand Up @@ -1888,13 +1909,10 @@ mod tests {
(same as original)

## After Pushdown
Projection: __datafusion_extracted_1 AS username
Projection: leaf_udf(test.user, Utf8("name")) AS __datafusion_extracted_1, test.user
TableScan: test projection=[user]
(same as after extraction)

## Optimized
Projection: leaf_udf(test.user, Utf8("name")) AS username
TableScan: test projection=[user]
(same as after pushdown)
"#)
}

Expand Down Expand Up @@ -1954,13 +1972,10 @@ mod tests {
(same as original)

## After Pushdown
Projection: __datafusion_extracted_1 AS leaf_udf(test.user,Utf8("name")), __datafusion_extracted_1 AS name2
Projection: leaf_udf(test.user, Utf8("name")) AS __datafusion_extracted_1, test.user
TableScan: test projection=[user]
(same as after extraction)

## Optimized
Projection: leaf_udf(test.user, Utf8("name")), leaf_udf(test.user, Utf8("name")) AS name2
TableScan: test projection=[user]
(same as after pushdown)
"#)
}

Expand Down Expand Up @@ -2168,13 +2183,10 @@ mod tests {
(same as original)

## After Pushdown
Projection: __datafusion_extracted_1 AS leaf_udf(test.user,Utf8("name"))
Projection: leaf_udf(test.user, Utf8("name")) AS __datafusion_extracted_1, test.user
TableScan: test projection=[user]
(same as after extraction)

## Optimized
Projection: leaf_udf(test.user, Utf8("name"))
TableScan: test projection=[user]
(same as after pushdown)
"#)
}

Expand Down Expand Up @@ -2355,15 +2367,10 @@ mod tests {
(same as original)

## After Pushdown
Projection: __datafusion_extracted_1 IS NOT NULL AS has_name, COUNT(Int32(1))
Projection: leaf_udf(test.user, Utf8("name")) AS __datafusion_extracted_1, test.user, COUNT(Int32(1))
Aggregate: groupBy=[[test.user]], aggr=[[COUNT(Int32(1))]]
TableScan: test projection=[user]
(same as after extraction)

## Optimized
Projection: leaf_udf(test.user, Utf8("name")) IS NOT NULL AS has_name, COUNT(Int32(1))
Aggregate: groupBy=[[test.user]], aggr=[[COUNT(Int32(1))]]
TableScan: test projection=[user]
(same as after pushdown)
"#)
}

Expand Down
11 changes: 11 additions & 0 deletions datafusion/optimizer/src/optimize_projections/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,17 @@
// under the License.

//! [`OptimizeProjections`] identifies and eliminates unused columns
//!
//! # Precedence over `PushDownLeafProjections`
//!
//! `OptimizeProjections` merges adjacent projections, including a recovery
//! projection over an extraction projection that
//! [`PushDownLeafProjections`] creates. To keep the two rules from undoing each
//! other, `PushDownLeafProjections` only creates that pair when the extraction
//! projection moves below the input. See the module documentation of
//! [`extract_leaf_expressions`](crate::extract_leaf_expressions).
//!
//! [`PushDownLeafProjections`]: crate::extract_leaf_expressions::PushDownLeafProjections

mod required_indices;

Expand Down
20 changes: 13 additions & 7 deletions datafusion/optimizer/src/push_down_filter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4771,9 +4771,11 @@ mod tests {

/// `PushDownFilter` and `PushDownLeafProjections` must not undo each other
/// for a filter next to a pure extraction projection
/// (<https://github.com/apache/datafusion/issues/14540>). Neither rule may
/// change the plan in the last optimizer pass. The source does not absorb
/// filters, so the `Filter` node stays in the plan.
/// (<https://github.com/apache/datafusion/issues/14540>). The source does
/// not absorb filters, so the `Filter` node stays in the plan. Also,
/// `PushDownLeafProjections` must not split a projection that
/// `OptimizeProjections` merges back. No rule may change the plan in the
/// last optimizer pass.
#[test]
fn filter_and_extraction_projection_reach_fixed_point() -> Result<()> {
let scan = || {
Expand All @@ -4794,13 +4796,17 @@ mod tests {
.project(vec![leaf_udf_expr(col("a")), col("b")])?
.build()?;

for plan in [simple, two_filters] {
// No filter: the extraction projection cannot move below the scan.
let no_filter = LogicalPlanBuilder::from(test_table_scan()?)
.project(vec![leaf_udf_expr(col("a"))])?
.build()?;

for plan in [simple, two_filters, no_filter] {
let passes = rules_that_changed_plan_per_pass(plan)?;
let last = passes.last().unwrap();
assert!(
!last.iter().any(|rule| rule == "push_down_filter"
|| rule == "push_down_leaf_projections"),
"pushdown rules changed the plan in the last pass: {passes:?}"
last.is_empty(),
"rules changed the plan in the last pass: {passes:?}"
);
}
Ok(())
Expand Down
2 changes: 1 addition & 1 deletion datafusion/sqllogictest/test_files/unnest.slt
Original file line number Diff line number Diff line change
Expand Up @@ -668,7 +668,7 @@ explain select unnest(unnest(unnest(column3)['c1'])), column3 from recursive_unn
logical_plan
01)Projection: __unnest_placeholder(UNNEST(recursive_unnest_table.column3)[c1],depth=2) AS UNNEST(UNNEST(UNNEST(recursive_unnest_table.column3)[c1])), recursive_unnest_table.column3
02)--Unnest: lists[__unnest_placeholder(UNNEST(recursive_unnest_table.column3)[c1])|depth=2] structs[]
03)----Projection: get_field(__unnest_placeholder(recursive_unnest_table.column3,depth=1), Utf8("c1")) AS __unnest_placeholder(UNNEST(recursive_unnest_table.column3)[c1]), recursive_unnest_table.column3
03)----Projection: get_field(__unnest_placeholder(recursive_unnest_table.column3,depth=1) AS UNNEST(recursive_unnest_table.column3), Utf8("c1")) AS __unnest_placeholder(UNNEST(recursive_unnest_table.column3)[c1]), recursive_unnest_table.column3
04)------Unnest: lists[__unnest_placeholder(recursive_unnest_table.column3)|depth=1] structs[]
05)--------Projection: recursive_unnest_table.column3 AS __unnest_placeholder(recursive_unnest_table.column3), recursive_unnest_table.column3
06)----------TableScan: recursive_unnest_table projection=[column3]
Expand Down
3 changes: 2 additions & 1 deletion docs/source/library-user-guide/query-optimizer.md
Original file line number Diff line number Diff line change
Expand Up @@ -168,9 +168,10 @@ Two rules can want the opposite order for the same pair of adjacent plan nodes.

Do not let two rules compete. Give one rule precedence, make the competing rule yield, and record the decision in the module documentation of both rules.

There is one such decision today:
There are two such decisions today:

- `PushDownFilter` yields to a _pure extraction projection_. A pure extraction projection is a projection whose expressions are only `__datafusion_extracted_N` aliases and pass-through columns. `ExtractLeafExpressions` creates it, and `PushDownLeafProjections` moves it towards the leaves. `PushDownFilter` does not move a filter below such a projection, so the projection stays next to the scan. A Parquet scan then merges the projection into the file projection and reads only the struct leaf. The filter loses nothing, because `PushDownFilter` records the predicate in `TableScan::filters` in the pass that runs before the extraction projection exists. See [issue #14540](https://github.com/apache/datafusion/issues/14540).
- `PushDownLeafProjections` yields to `OptimizeProjections`. When the extraction projection cannot move below the input of a projection, for example a `TableScan` or an `Aggregate`, `PushDownLeafProjections` leaves the projection unchanged. A split in place would give a recovery projection over an extraction projection on the same input, which `OptimizeProjections` merges back. The source absorbs the leaf expressions of the original projection just as well. See [PR #26085](https://github.com/apache/datafusion/pull/26085).

### Expression Naming

Expand Down
Loading