stuhood opened a new issue, #25642: URL: https://github.com/apache/datafusion/issues/25642
### Describe the bug `OptimizeProjections` initializes child `RequiredIndices` with `projection_beneficial: false` when traversing a `LogicalPlan::Extension` node, preventing `add_projection_on_top_if_helpful` from inserting a `Projection` above the extension node's children. When `optimize_projections` handles `LogicalPlan::Extension`: https://github.com/apache/datafusion/blob/c4f72bad34249798182ac7de605eb63ab134f26e/datafusion/optimizer/src/optimize_projections/mod.rs#L351-L358 It constructs child requirements via `RequiredIndices::new_from_indices(necessary_indices)`. In `required_indices.rs`: https://github.com/apache/datafusion/blob/c4f72bad34249798182ac7de605eb63ab134f26e/datafusion/optimizer/src/optimize_projections/required_indices.rs#L61-L67 `RequiredIndices::new_from_indices` initializes `projection_beneficial` to `false`. Later, when rewriting the plan's children in `map_children`: https://github.com/apache/datafusion/blob/c4f72bad34249798182ac7de605eb63ab134f26e/datafusion/optimizer/src/optimize_projections/mod.rs#L476-L487 `add_projection_on_top_if_helpful` is guarded by `if projection_beneficial`. Because `projection_beneficial` is `false`: 1. Any incoming `indices.projection_beneficial()` from ancestors (such as a `Sort` or `Filter` above the extension) is dropped. 2. Even when `extension.node.necessary_children_exprs(...)` specifies that only a subset of child columns is required, `add_projection_on_top_if_helpful` is never called above the extension's child. Leaf `TableScan` children are not impacted because `TableScan` prunes directly against `indices.into_inner()`. However, intermediate operators whose schemas cannot prune themselves—such as `LogicalPlan::Join` (where schema is the concatenation of left and right inputs) or `LogicalPlan::Filter` (where predicates reference extra columns)—cannot eliminate unneeded columns without a `Projection` placed above them. For a child `LogicalPlan::Join`, the join is forced to output all columns (including join keys that are not needed downstream) directly into the extension node. Furthermore, because no `Projection` sits directly above the join, downstream physical optimizations such as embedding `HashJoinExec::projection` via `ProjectionPushdown` cannot take place. ### To Reproduce The following test case (which can be added to `datafusion/optimizer/src/optimize_projections/mod.rs`) demonstrates the issue using an extension node wrapping a `Left Join`: ```rust #[test] fn test_user_defined_logical_plan_above_join() -> Result<()> { let table_scan = test_table_scan()?; let schema = Schema::new(vec![Field::new("c1", DataType::UInt32, false)]); let table2_scan = scan_empty(Some("test2"), &schema, None)?.build()?; // Join test(a, b, c) and test2(c1) on test.a = test2.c1 let join_plan = LogicalPlanBuilder::from(table_scan) .join(table2_scan, JoinType::Left, (vec!["a"], vec!["c1"]), None)? .build()?; let custom_plan = LogicalPlan::Extension(Extension { node: Arc::new(NoOpUserDefined::new( Arc::clone(join_plan.schema()), Arc::new(join_plan), )), }); // Parent only requires test.a and test.b; join key test2.c1 is not needed downstream. let plan = LogicalPlanBuilder::from(custom_plan) .project(vec![col("test.a"), col("test.b")])? .build()?; let optimized_plan = optimize(plan)?; // Observed plan: No Projection is placed above Left Join. // Left Join outputs test2.c1 into NoOpUserDefined even though only test.a and test.b are required. assert_snapshot!( optimized_plan, @r" Projection: test.a, test.b NoOpUserDefined Left Join: test.a = test2.c1 TableScan: test projection=[a, b] TableScan: test2 projection=[c1] " ); Ok(()) } ``` Similarly, if a `Sort` is placed above the extension node (`Sort -> Projection -> Extension -> Join`), `Sort` sets `projection_beneficial = true`, but the flag is dropped by `LogicalPlan::Extension`, placing a `Projection` above `NoOpUserDefined` instead of above `Left Join`: ``` Sort: test.a ASC NULLS LAST Projection: test.a, test.b NoOpUserDefined Left Join: test.a = test2.c1 TableScan: test projection=[a, b] TableScan: test2 projection=[c1] ``` ### Expected behavior When `extension.node.necessary_children_exprs(...)` prunes child columns (or when parent requirements have `projection_beneficial == true`), `OptimizeProjections` should mark child requirements with `projection_beneficial = true`. With `projection_beneficial = true`, `add_projection_on_top_if_helpful` will insert a `Projection` above the child: ``` Projection: test.a, test.b NoOpUserDefined Projection: test.a, test.b Left Join: test.a = test2.c1 TableScan: test projection=[a, b] TableScan: test2 projection=[c1] ``` This allows operators like `Left Join` to eliminate unneeded join keys from their output batches before passing data to `NoOpUserDefined`, and allows physical planning to embed the projection into `HashJoinExec`. ### Additional context Other operators in `OptimizeProjections` (such as `Sort`, `Filter`, `Repartition`, `Union`, and `Join`) unconditionally set `with_projection_beneficial()` for their children because narrower input batches improve execution performance and memory usage. Additionally, `UserDefinedLogicalNodeCore` currently does not provide an interface for extension nodes to declare whether their inputs benefit from projection pruning (unlike `supports_limit_pushdown`). -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
