diff --git a/datafusion/optimizer/src/extract_leaf_expressions.rs b/datafusion/optimizer/src/extract_leaf_expressions.rs index 53d6663240fbe..779d404a7fd32 100644 --- a/datafusion/optimizer/src/extract_leaf_expressions.rs +++ b/datafusion/optimizer/src/extract_leaf_expressions.rs @@ -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}; @@ -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), + } } }; @@ -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) "#) } @@ -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) "#) } @@ -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) "#) } @@ -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) "#) } @@ -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) "#) } @@ -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) "#) } diff --git a/datafusion/optimizer/src/optimize_projections/mod.rs b/datafusion/optimizer/src/optimize_projections/mod.rs index 3bbaf887ca594..523d03c45e0bb 100644 --- a/datafusion/optimizer/src/optimize_projections/mod.rs +++ b/datafusion/optimizer/src/optimize_projections/mod.rs @@ -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; diff --git a/datafusion/optimizer/src/push_down_filter.rs b/datafusion/optimizer/src/push_down_filter.rs index 4738904916d5a..000a8b6da1dd3 100644 --- a/datafusion/optimizer/src/push_down_filter.rs +++ b/datafusion/optimizer/src/push_down_filter.rs @@ -4771,9 +4771,11 @@ mod tests { /// `PushDownFilter` and `PushDownLeafProjections` must not undo each other /// for a filter next to a pure extraction projection - /// (). 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. + /// (). 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 = || { @@ -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(()) diff --git a/datafusion/sqllogictest/test_files/unnest.slt b/datafusion/sqllogictest/test_files/unnest.slt index 8e6013328f98f..083e48056a3ed 100644 --- a/datafusion/sqllogictest/test_files/unnest.slt +++ b/datafusion/sqllogictest/test_files/unnest.slt @@ -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] diff --git a/docs/source/library-user-guide/query-optimizer.md b/docs/source/library-user-guide/query-optimizer.md index 73c3ed6ed11dc..03fb6f2e33fe8 100644 --- a/docs/source/library-user-guide/query-optimizer.md +++ b/docs/source/library-user-guide/query-optimizer.md @@ -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