From 3e3d232fd49789403501626380472a81128f14b0 Mon Sep 17 00:00:00 2001 From: "Lei, HUANG" Date: Fri, 28 Aug 2026 23:14:03 +0800 Subject: [PATCH 1/2] fix: preserve fetch-carrying operators in remove_dist_changing_operators EnforceDistribution::remove_dist_changing_operators strips top-level RepartitionExec / CoalescePartitionsExec / SortPreservingMergeExec, assuming they can be regenerated later if necessary. A fetch on these operators is not a distribution concern but carries global limit semantics (e.g. produced by LimitPushdown folding a GlobalLimitExec into CoalescePartitionsExec: fetch=N), so stripping the operator silently drops the limit. This is latent in a single optimizer pass, but breaks when the physical optimizer pipeline runs twice on the same plan (e.g. GreptimeDB plans a nested subplan separately, then optimizes the enclosing plan again): the second pass removes CoalescePartitionsExec: fetch=N and never re-adds it, so a query like SELECT DISTINCT a FROM t LIMIT 1 can return one row per partition instead of one row in total. Stop the removal loop at any operator carrying a fetch. Regression tests: - EnforceDistribution preserves fetch-carrying CoalescePartitionsExec and SortPreservingMergeExec. - Running the full physical optimizer pipeline a second time over the optimized plan of SELECT DISTINCT a FROM t LIMIT 1 keeps the global limit boundary. Related: apache/datafusion#23800 (orthogonal LimitPushdown fix) Signed-off-by: Lei, HUANG --- .../enforce_distribution.rs | 39 ++++++++++++- .../limited_distinct_aggregation.rs | 55 ++++++++++++++++++- .../src/enforce_distribution.rs | 5 ++ 3 files changed, 97 insertions(+), 2 deletions(-) diff --git a/datafusion/core/tests/physical_optimizer/enforce_distribution.rs b/datafusion/core/tests/physical_optimizer/enforce_distribution.rs index 5df634c70bcbb..231325dc51b89 100644 --- a/datafusion/core/tests/physical_optimizer/enforce_distribution.rs +++ b/datafusion/core/tests/physical_optimizer/enforce_distribution.rs @@ -22,7 +22,7 @@ use std::sync::Arc; use crate::physical_optimizer::test_utils::{ check_integrity, coalesce_partitions_exec, parquet_exec_with_sort, parquet_exec_with_stats, repartition_exec, schema, sort_exec, - sort_exec_with_preserve_partitioning, sort_merge_join_exec, + sort_exec_with_preserve_partitioning, sort_expr, sort_merge_join_exec, sort_preserving_merge_exec, union_exec, }; @@ -3626,6 +3626,43 @@ fn get_schema() -> SchemaRef { Field::new("bank_account", DataType::UInt64, true), ])) } +#[test] +fn keep_fetch_carrying_dist_changing_operators() -> Result<()> { + let config = ConfigOptions::new(); + + // A CoalescePartitionsExec carrying a fetch is not a pure + // distribution-changing operator: it also provides the global limit, so + // EnforceDistribution must not strip it. This can be hit when the + // physical optimizer pipeline runs on an already-optimized plan. + let coalesce = + Arc::new(CoalescePartitionsExec::new(parquet_exec()).with_fetch(Some(1))) + as Arc; + let optimized = EnforceDistribution::new().optimize(coalesce, &config)?; + let coalesce = optimized + .as_any() + .downcast_ref::() + .expect("fetch-carrying CoalescePartitionsExec should be preserved"); + assert_eq!(coalesce.fetch(), Some(1)); + + // Same for SortPreservingMergeExec with a fetch. + let ordering = LexOrdering::new(vec![sort_expr("a", &schema())]).unwrap(); + let spm = Arc::new( + SortPreservingMergeExec::new( + ordering.clone(), + parquet_exec_with_sort(schema(), vec![ordering]), + ) + .with_fetch(Some(1)), + ) as Arc; + let optimized = EnforceDistribution::new().optimize(spm, &config)?; + let spm = optimized + .as_any() + .downcast_ref::() + .expect("fetch-carrying SortPreservingMergeExec should be preserved"); + assert_eq!(spm.fetch(), Some(1)); + + Ok(()) +} + #[test] fn test_replace_order_preserving_variants_with_fetch() -> Result<()> { // Create a base plan diff --git a/datafusion/core/tests/physical_optimizer/limited_distinct_aggregation.rs b/datafusion/core/tests/physical_optimizer/limited_distinct_aggregation.rs index c523b4a752a82..b9151e35aa240 100644 --- a/datafusion/core/tests/physical_optimizer/limited_distinct_aggregation.rs +++ b/datafusion/core/tests/physical_optimizer/limited_distinct_aggregation.rs @@ -25,18 +25,23 @@ use crate::physical_optimizer::test_utils::{ schema, }; +use arrow::array::{BooleanArray, Int32Array, Int64Array}; use arrow::datatypes::DataType; +use arrow::record_batch::RecordBatch; use arrow::{compute::SortOptions, util::pretty::pretty_format_batches}; +use datafusion::datasource::MemTable; use datafusion::prelude::SessionContext; use datafusion_common::Result; +use datafusion_common::config::ConfigOptions; use datafusion_execution::config::SessionConfig; use datafusion_expr::Operator; use datafusion_physical_expr::expressions::{self, cast, col}; use datafusion_physical_expr_common::sort_expr::PhysicalSortExpr; +use datafusion_physical_optimizer::optimizer::PhysicalOptimizer; use datafusion_physical_plan::{ ExecutionPlan, aggregates::{AggregateExec, AggregateMode}, - collect, + collect, displayable, limit::{GlobalLimitExec, LocalLimitExec}, }; @@ -254,6 +259,54 @@ async fn test_distinct_cols_different_than_group_by_cols() -> Result<()> { Ok(()) } +#[tokio::test] +async fn test_global_limit_survives_second_optimizer_pass() -> Result<()> { + // Regression: EnforceDistribution::remove_dist_changing_operators must not + // strip a fetch-carrying CoalescePartitionsExec when the physical + // optimizer pipeline runs on an already-optimized plan (e.g. a nested + // subplan that was planned and optimized separately). + let cfg = SessionConfig::new().with_target_partitions(4); + let ctx = SessionContext::new_with_config(cfg); + let schema = schema(); + let batch = RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(Int64Array::from(vec![1, 2, 3])), + Arc::new(Int64Array::from(vec![4, 5, 6])), + Arc::new(Int64Array::from(vec![7, 8, 9])), + Arc::new(Int32Array::from(vec![1, 2, 3])), + Arc::new(BooleanArray::from(vec![true, false, true])), + ], + )?; + let table = MemTable::try_new(schema, vec![vec![batch]])?; + ctx.register_table("t", Arc::new(table))?; + + let plan = ctx + .sql("SELECT DISTINCT a FROM t LIMIT 1") + .await? + .create_physical_plan() + .await?; + + // Run the physical optimizer pipeline a second time over the + // already-optimized plan. + let mut config = ConfigOptions::new(); + config.execution.target_partitions = 4; + let optimizer = PhysicalOptimizer::new(); + let mut optimized: Arc = plan; + for rule in &optimizer.rules { + optimized = rule.optimize(optimized, &config)?; + } + + let display = displayable(optimized.as_ref()).indent(true).to_string(); + assert!( + display.contains("CoalescePartitionsExec: fetch=1") + || display.contains("GlobalLimitExec: skip=0, fetch=1"), + "global limit must survive a second optimizer pass:\n{display}" + ); + + Ok(()) +} + #[test] fn test_has_order_by() -> Result<()> { let schema = schema(); diff --git a/datafusion/physical-optimizer/src/enforce_distribution.rs b/datafusion/physical-optimizer/src/enforce_distribution.rs index d23a699f715de..6837d73510cc1 100644 --- a/datafusion/physical-optimizer/src/enforce_distribution.rs +++ b/datafusion/physical-optimizer/src/enforce_distribution.rs @@ -981,6 +981,11 @@ fn remove_dist_changing_operators( || is_coalesce_partitions(&distribution_context.plan) || is_sort_preserving_merge(&distribution_context.plan) { + // A fetch carries global limit semantics, not just a distribution + // change, so fetch-carrying operators must be kept. + if distribution_context.plan.fetch().is_some() { + break; + } // All of above operators have a single child. First child is only child. // Remove any distribution changing operators at the beginning: distribution_context = distribution_context.children.swap_remove(0); From 5568299763414095c391538c13403075039e0f36 Mon Sep 17 00:00:00 2001 From: "Lei, HUANG" Date: Sat, 29 Aug 2026 13:13:32 +0800 Subject: [PATCH 2/2] test: use session's physical optimizer rules and config in second pass The regression test test_global_limit_survives_second_optimizer_pass constructed a fresh ConfigOptions::new() and PhysicalOptimizer::new() for the second optimization pass, which could drift from the SessionContext's actual optimizer list and configuration. Reuse ctx.state().physical_optimizers() and ctx.state().config_options() instead, matching how DefaultPhysicalPlanner::optimize_physical_plan runs the pipeline in production. Signed-off-by: Lei, HUANG --- .../limited_distinct_aggregation.rs | 12 ++++-------- 1 file changed, 4 insertions(+), 8 deletions(-) diff --git a/datafusion/core/tests/physical_optimizer/limited_distinct_aggregation.rs b/datafusion/core/tests/physical_optimizer/limited_distinct_aggregation.rs index b9151e35aa240..8213476f039e3 100644 --- a/datafusion/core/tests/physical_optimizer/limited_distinct_aggregation.rs +++ b/datafusion/core/tests/physical_optimizer/limited_distinct_aggregation.rs @@ -32,12 +32,10 @@ use arrow::{compute::SortOptions, util::pretty::pretty_format_batches}; use datafusion::datasource::MemTable; use datafusion::prelude::SessionContext; use datafusion_common::Result; -use datafusion_common::config::ConfigOptions; use datafusion_execution::config::SessionConfig; use datafusion_expr::Operator; use datafusion_physical_expr::expressions::{self, cast, col}; use datafusion_physical_expr_common::sort_expr::PhysicalSortExpr; -use datafusion_physical_optimizer::optimizer::PhysicalOptimizer; use datafusion_physical_plan::{ ExecutionPlan, aggregates::{AggregateExec, AggregateMode}, @@ -288,13 +286,11 @@ async fn test_global_limit_survives_second_optimizer_pass() -> Result<()> { .await?; // Run the physical optimizer pipeline a second time over the - // already-optimized plan. - let mut config = ConfigOptions::new(); - config.execution.target_partitions = 4; - let optimizer = PhysicalOptimizer::new(); + // already-optimized plan, using the session's own rules and config. + let state = ctx.state(); let mut optimized: Arc = plan; - for rule in &optimizer.rules { - optimized = rule.optimize(optimized, &config)?; + for rule in state.physical_optimizers() { + optimized = rule.optimize(optimized, state.config_options())?; } let display = displayable(optimized.as_ref()).indent(true).to_string();