From cd7d651af590865d024aaeb84499f299d014e6d7 Mon Sep 17 00:00:00 2001 From: Himanshu-2005-code Date: Mon, 10 Aug 2026 22:45:02 +0530 Subject: [PATCH 1/2] fix: preserve required ordering in limit pushdown --- .../physical-optimizer/src/limit_pushdown.rs | 54 +++++++++++++------ .../physical-optimizer/src/pushdown_sort.rs | 14 +++-- 2 files changed, 47 insertions(+), 21 deletions(-) diff --git a/datafusion/physical-optimizer/src/limit_pushdown.rs b/datafusion/physical-optimizer/src/limit_pushdown.rs index 01a288f7f1632..230cd6fad05ff 100644 --- a/datafusion/physical-optimizer/src/limit_pushdown.rs +++ b/datafusion/physical-optimizer/src/limit_pushdown.rs @@ -70,6 +70,7 @@ use datafusion_common::error::Result; use datafusion_common::stats::Precision; use datafusion_common::tree_node::{Transformed, TreeNodeRecursion}; use datafusion_common::utils::combine_limit; +use datafusion_physical_expr::LexOrdering; use datafusion_physical_plan::coalesce_partitions::CoalescePartitionsExec; use datafusion_physical_plan::empty::EmptyExec; use datafusion_physical_plan::limit::{GlobalLimitExec, LocalLimitExec}; @@ -96,7 +97,7 @@ pub struct GlobalRequirements { fetch: Option, skip: usize, satisfied: bool, - preserve_order: bool, + required_ordering: Option, } impl LimitPushdown { @@ -116,7 +117,7 @@ impl PhysicalOptimizerRule for LimitPushdown { fetch: None, skip: 0, satisfied: false, - preserve_order: false, + required_ordering: None, }; pushdown_limits(plan, global_state) } @@ -134,7 +135,7 @@ struct LimitInfo { input: Arc, fetch: Option, skip: usize, - preserve_order: bool, + required_ordering: Option, } /// This function is the main helper function of the `LimitPushDown` rule. @@ -160,7 +161,7 @@ pub fn pushdown_limit_helper( ); global_state.skip = skip; global_state.fetch = fetch; - global_state.preserve_order = limit_info.preserve_order; + global_state.required_ordering = limit_info.required_ordering; global_state.satisfied = false; if let Some(fetch) = fetch @@ -219,6 +220,7 @@ pub fn pushdown_limit_helper( pushdown_plan, global_state.skip, None, + global_state.required_ordering.clone(), )), global_state, )) @@ -243,8 +245,12 @@ pub fn pushdown_limit_helper( // Execution plans can't (yet) handle skip, so if we have one, // we still need to add a global limit if global_state.skip > 0 { - new_plan = - add_global_limit(new_plan, global_state.skip, global_state.fetch); + new_plan = add_global_limit( + new_plan, + global_state.skip, + global_state.fetch, + global_state.required_ordering.clone(), + ); } global_state.fetch = skip_and_fetch; global_state.skip = 0; @@ -260,6 +266,7 @@ pub fn pushdown_limit_helper( pushdown_plan, global_state.skip, global_fetch, + global_state.required_ordering.clone(), )), global_state, )) @@ -278,7 +285,7 @@ pub fn pushdown_limit_helper( if global_state.satisfied { if let Some(plan_with_fetch) = maybe_fetchable { let plan_with_preserve_order = plan_with_fetch - .with_preserve_order(global_state.preserve_order) + .with_preserve_order(global_state.required_ordering.is_some()) .unwrap_or(plan_with_fetch); Ok((Transformed::yes(plan_with_preserve_order), global_state)) } else { @@ -288,7 +295,7 @@ pub fn pushdown_limit_helper( global_state.satisfied = true; pushdown_plan = if let Some(plan_with_fetch) = maybe_fetchable { let plan_with_preserve_order = plan_with_fetch - .with_preserve_order(global_state.preserve_order) + .with_preserve_order(global_state.required_ordering.is_some()) .unwrap_or(plan_with_fetch); if global_skip > 0 { @@ -296,12 +303,18 @@ pub fn pushdown_limit_helper( plan_with_preserve_order, global_skip, Some(global_fetch), + global_state.required_ordering.clone(), ) } else { plan_with_preserve_order } } else { - add_limit(pushdown_plan, global_skip, global_fetch) + add_limit( + pushdown_plan, + global_skip, + global_fetch, + global_state.required_ordering.clone(), + ) }; Ok((Transformed::yes(pushdown_plan), global_state)) } @@ -417,7 +430,7 @@ fn extract_limit(plan: &Arc) -> Option { input: Arc::clone(global_limit.input()), fetch: global_limit.fetch(), skip: global_limit.skip(), - preserve_order: global_limit.required_ordering().is_some(), + required_ordering: global_limit.required_ordering().clone(), }) } else { plan.downcast_ref::() @@ -425,7 +438,7 @@ fn extract_limit(plan: &Arc) -> Option { input: Arc::clone(local_limit.input()), fetch: Some(local_limit.fetch()), skip: 0, - preserve_order: local_limit.required_ordering().is_some(), + required_ordering: local_limit.required_ordering().clone(), }) } } @@ -435,17 +448,22 @@ fn combines_input_partitions(plan: &Arc) -> bool { plan.is::() || plan.is::() } -/// Adds a limit to the plan, chooses between global and local limits based on -/// skip value and the number of partitions. +/// Adds a limit to the plan, choosing between global and local limits +/// based on the skip value and the number of partitions. fn add_limit( pushdown_plan: Arc, skip: usize, fetch: usize, + required_ordering: Option, ) -> Arc { if skip > 0 || pushdown_plan.output_partitioning().partition_count() == 1 { - add_global_limit(pushdown_plan, skip, Some(fetch)) + add_global_limit(pushdown_plan, skip, Some(fetch), required_ordering) } else { - Arc::new(LocalLimitExec::new(pushdown_plan, fetch + skip)) as _ + let mut limit = LocalLimitExec::new(pushdown_plan, fetch + skip); + + limit.set_required_ordering(required_ordering); + + Arc::new(limit) } } @@ -454,8 +472,10 @@ fn add_global_limit( pushdown_plan: Arc, skip: usize, fetch: Option, + required_ordering: Option, ) -> Arc { - Arc::new(GlobalLimitExec::new(pushdown_plan, skip, fetch)) as _ + let mut limit = GlobalLimitExec::new(pushdown_plan, skip, fetch); + limit.set_required_ordering(required_ordering); + Arc::new(limit) } - // See tests in datafusion/core/tests/physical_optimizer diff --git a/datafusion/physical-optimizer/src/pushdown_sort.rs b/datafusion/physical-optimizer/src/pushdown_sort.rs index 5dfe221ed24c0..cb3c0f916a341 100644 --- a/datafusion/physical-optimizer/src/pushdown_sort.rs +++ b/datafusion/physical-optimizer/src/pushdown_sort.rs @@ -62,12 +62,11 @@ use datafusion_common::tree_node::{ }; use datafusion_physical_plan::SortOrderPushdownResult; use datafusion_physical_plan::buffer::BufferExec; -use datafusion_physical_plan::limit::{GlobalLimitExec, LocalLimitExec}; +use datafusion_physical_plan::limit::GlobalLimitExec; use datafusion_physical_plan::sorts::sort::SortExec; use datafusion_physical_plan::sorts::sort_preserving_merge::SortPreservingMergeExec; use datafusion_physical_plan::{ExecutionPlan, ExecutionPlanProperties}; use std::sync::Arc; - /// A PhysicalOptimizerRule that attempts to push down sort requirements to data sources. /// /// See module-level documentation for details. @@ -111,7 +110,12 @@ impl PhysicalOptimizerRule for PushdownSort { // Use LocalLimitExec (not Global) since input is multi-partition. let inner = if let Some(fetch) = sort_child.fetch() { inner.with_fetch(Some(fetch)).unwrap_or_else(|| { - Arc::new(LocalLimitExec::new(inner, fetch)) + let mut limit = + GlobalLimitExec::new(inner, 0, Some(fetch)); + limit.set_required_ordering(Some( + sort_child.expr().clone(), + )); + Arc::new(limit) }) } else { inner @@ -171,7 +175,9 @@ impl PhysicalOptimizerRule for PushdownSort { // wrapping with GlobalLimitExec. if let Some(fetch) = sort_exec.fetch() { let inner = inner.with_fetch(Some(fetch)).unwrap_or_else(|| { - Arc::new(GlobalLimitExec::new(inner, 0, Some(fetch))) + let mut limit = GlobalLimitExec::new(inner, 0, Some(fetch)); + limit.set_required_ordering(Some(sort_exec.expr().clone())); + Arc::new(limit) }); Ok(Transformed::yes(inner)) } else { From c48afbf028b32d7e1b4b1be9f96683d12253240f Mon Sep 17 00:00:00 2001 From: Himanshu-2005-code Date: Tue, 11 Aug 2026 01:24:19 +0530 Subject: [PATCH 2/2] test: add limit pushdown ordering regression --- .../physical_optimizer/limit_pushdown.rs | 56 +++++++++++++++++++ 1 file changed, 56 insertions(+) diff --git a/datafusion/core/tests/physical_optimizer/limit_pushdown.rs b/datafusion/core/tests/physical_optimizer/limit_pushdown.rs index b8ebc80348134..4c60f35340bd8 100644 --- a/datafusion/core/tests/physical_optimizer/limit_pushdown.rs +++ b/datafusion/core/tests/physical_optimizer/limit_pushdown.rs @@ -162,6 +162,62 @@ fn transforms_streaming_table_exec_into_fetching_version_and_keeps_the_global_li Ok(()) } +#[test] +fn preserves_required_ordering_when_reinserting_global_limit() -> Result<()> { + let schema = create_schema(); + let projection = + projection_exec(Arc::clone(&schema), empty_exec(Arc::clone(&schema)))?; + let ordering = LexOrdering::new([PhysicalSortExpr::new( + col("c1", schema.as_ref())?, + SortOptions::default(), + )]) + .unwrap(); + + let mut limit = + datafusion_physical_plan::limit::GlobalLimitExec::new(projection, 2, Some(5)); + limit.set_required_ordering(Some(ordering.clone())); + + let optimized = + LimitPushdown::new().optimize(Arc::new(limit), &ConfigOptions::new())?; + let limit = optimized + .as_ref() + .downcast_ref::() + .unwrap() + .input() + .as_ref() + .downcast_ref::() + .unwrap(); + + assert_eq!(limit.required_ordering().as_ref(), Some(&ordering)); + Ok(()) +} + +#[test] +fn preserves_required_ordering_when_reinserting_local_limit() -> Result<()> { + let schema = create_schema(); + let projection = + projection_exec(Arc::clone(&schema), empty_exec(Arc::clone(&schema)))?; + let repartition = repartition_exec(projection)?; + let ordering = LexOrdering::new([PhysicalSortExpr::new( + col("c1", schema.as_ref())?, + SortOptions::default(), + )]) + .unwrap(); + + let mut limit = datafusion_physical_plan::limit::LocalLimitExec::new(repartition, 5); + limit.set_required_ordering(Some(ordering.clone())); + + let optimized = + LimitPushdown::new().optimize(Arc::new(limit), &ConfigOptions::new())?; + let limit = optimized + .as_ref() + .downcast_ref::() + .unwrap(); + + assert_eq!(limit.required_ordering().as_ref(), Some(&ordering)); + Ok(()) +} + fn join_on_columns( left_col: &str, right_col: &str,