diff --git a/datafusion/physical-optimizer/src/ensure_requirements/enforce_sorting/sort_pushdown.rs b/datafusion/physical-optimizer/src/ensure_requirements/enforce_sorting/sort_pushdown.rs index d7f556b90d9fe..0bf107170cef6 100644 --- a/datafusion/physical-optimizer/src/ensure_requirements/enforce_sorting/sort_pushdown.rs +++ b/datafusion/physical-optimizer/src/ensure_requirements/enforce_sorting/sort_pushdown.rs @@ -49,7 +49,10 @@ use datafusion_physical_plan::repartition::RepartitionExec; use datafusion_physical_plan::sorts::sort::SortExec; use datafusion_physical_plan::sorts::sort_preserving_merge::SortPreservingMergeExec; use datafusion_physical_plan::tree_node::PlanContext; -use datafusion_physical_plan::{ExecutionPlan, ExecutionPlanProperties}; +use datafusion_physical_plan::{ + ChildrenPropertiesMode, ExecutionPlan, ExecutionPlanProperties, + ReplaceChildrenOptions, +}; /// "Data class" used by sort pushdown (now driven from `EnsureRequirements`) /// to push down [`SortExec`] in the plan. In some cases the total @@ -993,6 +996,7 @@ fn handle_custom_pushdown( // Collect all unique column indices used in the parent-required sorting // expression: + let output_requirement = parent_required.first().clone(); let requirement = parent_required.into_single(); let all_indices: HashSet = requirement .iter() @@ -1056,7 +1060,10 @@ fn handle_custom_pushdown( .data; Ok(PhysicalSortRequirement::new(updated_columns, req.options)) }) - .collect::>>()?; + .collect::>>(); + let Ok(updated_parent_req) = updated_parent_req else { + return Ok(None); + }; // Prepare the result, populating with the updated requirements for children that maintain order let result = maintains_input_order @@ -1069,7 +1076,56 @@ fn handle_custom_pushdown( None } }) - .collect(); + .collect::>(); + + // Row order preservation does not establish an output-to-input column + // mapping. Prove that this candidate ordering survives the operator's + // value transformations before removing the sort above it. + let mut sorted_children = Vec::with_capacity(plan_children.len()); + for (child, required) in plan_children.into_iter().zip(&result) { + let sorted: Arc = if let Some(required) = required { + let ordering = LexOrdering::from(required.first().clone()); + let schema = child.schema(); + if ordering.iter().any(|sort| { + collect_columns(&sort.expr) + .iter() + .any(|column| column.index() >= schema.fields().len()) + || sort.expr.data_type(&schema).is_err() + }) { + return Ok(None); + } + // SortExec::new assumes valid property derivation. Check its + // fallible steps first; invalid candidates leave the outer sort. + let mut properties = child.equivalence_properties().clone(); + if properties + .extract_common_sort_prefix(ordering.clone()) + .is_err() + || properties.reorder(ordering.clone()).is_err() + { + return Ok(None); + } + Arc::new( + SortExec::new(ordering, Arc::clone(child)) + .with_preserve_partitioning(true), + ) + } else { + Arc::clone(child) + }; + sorted_children.push(sorted); + } + let Ok(candidate) = Arc::clone(plan).replace_children( + sorted_children, + ReplaceChildrenOptions::new(ChildrenPropertiesMode::Recompute), + ) else { + return Ok(None); + }; + if !candidate + .equivalence_properties() + .ordering_satisfy_requirement(output_requirement) + .unwrap_or(false) + { + return Ok(None); + } Ok(Some(result)) } else { @@ -1234,7 +1290,7 @@ mod tests { use arrow::datatypes::{DataType, Field, Schema}; use datafusion_expr::Operator; use datafusion_physical_expr::PhysicalExpr; - use datafusion_physical_expr::expressions::{BinaryExpr, col}; + use datafusion_physical_expr::expressions::{BinaryExpr, NegativeExpr, col, lit}; use datafusion_physical_plan::empty::EmptyExec; use datafusion_physical_plan::limit::GlobalLimitExec; @@ -1352,6 +1408,96 @@ mod tests { LexRequirement::new(reqs).unwrap() } + // Exercise the custom-operator fallback directly with ProjectionExec's + // property derivation. The SQL tests cover dispatch through a custom plan. + #[test] + fn custom_pushdown_accepts_renamed_columns() -> Result<()> { + let child = child_schema(); + let input = Arc::new(EmptyExec::new(Arc::clone(&child))); + let plan: Arc = Arc::new(ProjectionExec::try_new( + [ + (col("a", &child)?, "first".to_string()), + (col("b", &child)?, "second".to_string()), + ], + input, + )?); + let output = plan.schema(); + let required = OrderingRequirements::new(lex([ + req("first", &output, ASC), + req("second", &output, DESC), + ])); + let expected = OrderingRequirements::new(lex([ + req("a", &child, ASC), + req("b", &child, DESC), + ])); + + assert_eq!( + handle_custom_pushdown(&plan, required, &[true])?, + Some(vec![Some(expected)]) + ); + Ok(()) + } + + #[test] + fn custom_pushdown_rejects_reordered_columns() -> Result<()> { + let plan: Arc = reordering_projection(); + let required = + OrderingRequirements::new(lex([req("score", &plan.schema(), ASC)])); + + assert!(handle_custom_pushdown(&plan, required, &[true])?.is_none()); + Ok(()) + } + + #[test] + fn custom_pushdown_rejects_changed_values_with_same_schema() -> Result<()> { + let schema = child_schema(); + let input = Arc::new(EmptyExec::new(Arc::clone(&schema))); + let negative = + Arc::new(NegativeExpr::new(col("a", &schema)?)) as Arc; + let plan: Arc = Arc::new(ProjectionExec::try_new( + [ + (negative, "a".to_string()), + (col("b", &schema)?, "b".to_string()), + (col("c", &schema)?, "c".to_string()), + ], + input, + )?); + assert_eq!(plan.schema(), schema); + let required = OrderingRequirements::new(lex([req("a", &schema, ASC)])); + + assert!(handle_custom_pushdown(&plan, required, &[true])?.is_none()); + Ok(()) + } + + #[test] + fn custom_pushdown_rejects_incompatible_child_expression() -> Result<()> { + let schema = child_schema(); + let input = Arc::new(EmptyExec::new(Arc::clone(&schema))); + let flag = Arc::new(BinaryExpr::new( + col("a", &schema)?, + Operator::Gt, + lit(0_i32), + )) as Arc; + let plan: Arc = Arc::new(ProjectionExec::try_new( + [(flag, "flag".to_string())], + input, + )?); + // Valid against Boolean output, but remapping flag@0 to a@0 would + // produce Int32 AND Boolean in the child schema. + let ordering = Arc::new(BinaryExpr::new( + col("flag", &plan.schema())?, + Operator::And, + lit(true), + )); + let required = OrderingRequirements::new(lex([PhysicalSortRequirement::new( + ordering, + Some(ASC), + )])); + + assert!(handle_custom_pushdown(&plan, required, &[true])?.is_none()); + Ok(()) + } + #[test] fn remap_single_hard_requirement_through_reordering_projection() { let projection = reordering_projection(); diff --git a/datafusion/sqllogictest/src/test_context.rs b/datafusion/sqllogictest/src/test_context.rs index 0c8c2c5add363..d9b19c580038c 100644 --- a/datafusion/sqllogictest/src/test_context.rs +++ b/datafusion/sqllogictest/src/test_context.rs @@ -73,6 +73,7 @@ use log::info; use sqlparser::ast; use tempfile::TempDir; +mod custom_sort_pushdown; mod range_partitioning; /// Context for running tests @@ -182,6 +183,11 @@ impl TestContext { let file_name = relative_path.file_name().unwrap().to_str().unwrap(); match file_name { + "sort_pushdown.slt" => { + custom_sort_pushdown::register_custom_sort_pushdown( + test_ctx.session_ctx(), + ); + } "parquet_missing_bounds.slt" => { register_parquet_missing_bounds(&mut test_ctx).await; } diff --git a/datafusion/sqllogictest/src/test_context/custom_sort_pushdown.rs b/datafusion/sqllogictest/src/test_context/custom_sort_pushdown.rs new file mode 100644 index 0000000000000..ee9985f11c69e --- /dev/null +++ b/datafusion/sqllogictest/src/test_context/custom_sort_pushdown.rs @@ -0,0 +1,185 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +use std::fmt::{self, Formatter}; +use std::sync::Arc; + +use arrow::array::record_batch; +use arrow::datatypes::{self as arrow_schema, SchemaRef}; +use async_trait::async_trait; +use datafusion::catalog::{Session, TableFunctionArgs, TableFunctionImpl, TableProvider}; +use datafusion::common::tree_node::TreeNodeRecursion; +use datafusion::common::{DFSchema, Result, ScalarValue, plan_err}; +use datafusion::execution::{SendableRecordBatchStream, SessionState, TaskContext}; +use datafusion::logical_expr::{Expr, TableType}; +use datafusion::physical_expr::PhysicalExpr; +use datafusion::physical_expr::expressions::Column; +use datafusion::physical_plan::projection::ProjectionExec; +use datafusion::physical_plan::test::TestMemoryExec; +use datafusion::physical_plan::{ + ChildrenPropertiesMode, DisplayAs, DisplayFormatType, ExecutionPlan, PlanProperties, + ReplaceChildrenOptions, +}; +use datafusion::prelude::SessionContext; + +/// Expose custom operators through SQL without changing the optimizer rules. +pub(super) fn register_custom_sort_pushdown(ctx: &SessionContext) { + ctx.register_udtf("custom_projection", Arc::new(CustomProjectionFunction)); +} + +#[derive(Debug)] +struct CustomProjectionFunction; + +impl TableFunctionImpl for CustomProjectionFunction { + fn call_with_args(&self, args: TableFunctionArgs) -> Result> { + let first = record_batch!(("a", Int32, [2, 1]), ("b", Int32, [30, 10]))?; + let second = record_batch!(("a", Int32, [3]), ("b", Int32, [20]))?; + let schema = first.schema(); + let source = TestMemoryExec::try_new_exec( + &[vec![first], vec![second]], + Arc::clone(&schema), + None, + )?; + let df_schema = DFSchema::try_from(schema)?; + let state = args + .session() + .as_any() + .downcast_ref::() + .expect("sqllogictests use SessionState"); + // Parse the expressions here so they execute inside the custom plan, + // rather than in SQL's built-in ProjectionExec. + let expressions = args + .exprs() + .iter() + .map(|arg| { + let Expr::Literal(ScalarValue::Utf8(Some(sql)), _) = arg else { + return plan_err!("custom_projection expects SQL expression strings"); + }; + let expr = state.create_logical_expr(sql, &df_schema)?; + let name = expr.schema_name().to_string(); + Ok((state.create_physical_expr(expr, &df_schema)?, name)) + }) + .collect::>>()?; + let custom = CustomProjection(ProjectionExec::try_new(expressions, source)?); + Ok(Arc::new(CustomTable(Arc::new(custom)))) + } +} + +#[derive(Debug)] +struct CustomTable(Arc); + +#[async_trait] +impl TableProvider for CustomTable { + fn schema(&self) -> SchemaRef { + self.0.schema() + } + + fn table_type(&self) -> TableType { + TableType::Base + } + + async fn scan( + &self, + _: &dyn Session, + projection: Option<&[usize]>, + _: &[Expr], + _: Option, + ) -> Result> { + if let Some(projection) = projection { + let schema = self.schema(); + let expressions = projection.iter().map(|&index| { + let name = schema.field(index).name(); + ( + Arc::new(Column::new(name, index)) as Arc, + name.clone(), + ) + }); + Ok(Arc::new(ProjectionExec::try_new( + expressions, + Arc::clone(&self.0), + )?)) + } else { + Ok(Arc::clone(&self.0)) + } + } +} + +/// A projection with the default custom-plan sort pushdown behavior. Row order +/// is preserved, while column positions and values can change. +#[derive(Debug)] +struct CustomProjection(ProjectionExec); + +impl DisplayAs for CustomProjection { + fn fmt_as(&self, display_type: DisplayFormatType, f: &mut Formatter) -> fmt::Result { + write!(f, "Custom")?; + self.0.fmt_as(display_type, f) + } +} + +impl ExecutionPlan for CustomProjection { + fn name(&self) -> &str { + "CustomProjectionExec" + } + + fn properties(&self) -> &Arc { + self.0.properties() + } + + fn children(&self) -> Vec<&Arc> { + self.0.children() + } + + fn maintains_input_order(&self) -> Vec { + vec![true] + } + + fn replace_children( + self: Arc, + children: Vec>, + _: ReplaceChildrenOptions, + ) -> Result> { + Ok(Arc::new(Self(ProjectionExec::try_new( + self.0.expr().to_vec(), + Arc::clone(&children[0]), + )?))) + } + + fn with_new_children( + self: Arc, + children: Vec>, + ) -> Result> { + self.replace_children( + children, + ReplaceChildrenOptions::new(ChildrenPropertiesMode::Recompute), + ) + } + + fn apply_expressions( + &self, + f: &mut dyn FnMut(&Arc) -> Result, + ) -> Result { + self.0.apply_expressions(f) + } + + fn execute( + &self, + partition: usize, + context: Arc, + ) -> Result { + self.0.execute(partition, context) + } +} diff --git a/datafusion/sqllogictest/test_files/sort_pushdown.slt b/datafusion/sqllogictest/test_files/sort_pushdown.slt index a173e76d6c262..72f17257d1f75 100644 --- a/datafusion/sqllogictest/test_files/sort_pushdown.slt +++ b/datafusion/sqllogictest/test_files/sort_pushdown.slt @@ -15,6 +15,79 @@ # specific language governing permissions and limitations # under the License. +# custom_projection evaluates its SQL expression strings through a custom +# execution plan that preserves row order. Its input has two partitions: +# [(a=2, b=30), (a=1, b=10)] and [(a=3, b=20)]. +statement ok +SET datafusion.execution.target_partitions = 2; + +# Identity and renaming allow the sort to move below the custom operator. +query II +SELECT * FROM custom_projection('a', 'b') ORDER BY a, b; +---- +1 10 +2 30 +3 20 + +query II +SELECT * FROM custom_projection('a AS first', 'b AS second') ORDER BY first, second; +---- +1 10 +2 30 +3 20 + +# Renaming is transparent: keep the sort below the custom operator. +statement ok +SET datafusion.explain.physical_plan_only = true; + +query TT +EXPLAIN SELECT * FROM custom_projection('a AS first', 'b AS second') ORDER BY first, second; +---- +physical_plan +01)SortPreservingMergeExec: [first@0 ASC NULLS LAST, second@1 ASC NULLS LAST] +02)--CustomProjectionExec: expr=[a@0 as first, b@1 as second] +03)----SortExec: expr=[a@0 ASC NULLS LAST, b@1 ASC NULLS LAST], preserve_partitioning=[true] +04)------DataSourceExec: partitions=2, partition_sizes=[1, 1] + +statement ok +SET datafusion.explain.physical_plan_only = false; + +# Reordering output columns requires sorting the custom operator's output. +query II +SELECT * FROM custom_projection('b', 'a') ORDER BY b, a; +---- +10 1 +20 3 +30 2 + +# The output schema is still [a, b], but a contains new values. +query II +SELECT * FROM custom_projection('a + b AS a', 'b') ORDER BY a, b; +---- +11 10 +23 20 +32 30 + +# Negating a reverses its order even though its name and type stay the same. +query II +SELECT * FROM custom_projection('-a AS a', 'b') ORDER BY a, b; +---- +-3 20 +-2 30 +-1 10 + +# Evaluate Boolean sort expressions against the custom operator's output. +query BI +SELECT * FROM custom_projection('b > 10 AS flag', 'a') +ORDER BY flag AND (a > 1), a; +---- +false 1 +true 2 +true 3 + +statement ok +SET datafusion.execution.target_partitions = 4; + #Sort Pushdown for ordered Parquet files statement ok SET datafusion.execution.parquet.pushdown_filters = true;