alamb commented on code in PR #25688:
URL: https://github.com/apache/datafusion/pull/25688#discussion_r4154971062


##########
datafusion/core/tests/physical_optimizer/ensure_requirements.rs:
##########
@@ -1560,3 +1564,222 @@ fn 
test_sort_pushed_below_limit_with_skip_keeps_skip_rows() -> Result<()> {
     ");
     Ok(())
 }
+
+// ========================================================================
+// Decomposition tests: `EnsureRequirements` was split into three rules

Review Comment:
   Is the history of what happened important to document in comments (that will 
live after this PR)?



##########
datafusion/core/tests/physical_optimizer/filter_pushdown.rs:
##########
@@ -134,7 +134,8 @@ fn test_pushdown_volatile_functions_not_allowed() {
       output:
         Ok:
           - FilterExec: a@0 = random()
-          -   DataSourceExec: file_groups={1 group: [[test.parquet]]}, 
projection=[a, b, c], file_type=test, pushdown_supported=true
+          -   RepartitionExec: partitioning=RoundRobinBatch(4), 
input_partitions=1

Review Comment:
   why did these tests change?



##########
datafusion/physical-optimizer/src/optimizer.rs:
##########
@@ -115,25 +118,12 @@ impl PhysicalOptimizer {
             // window's declared ordering without pattern-matching a SortExec)
             // and before ProjectionPushdown (which embeds projections into 
FilterExec).
             Arc::new(WindowTopN::new()),
-            // Ensures each input plan satisfies the distribution and ordering
-            // requirements declared by 
`ExecutionPlan::required_input_distribution`
-            // and `ExecutionPlan::required_input_ordering`.
-            //
-            // If the requirements are already satisfied, this rule leaves the 
plan
-            // unchanged. For example, it does not add sorting when the input 
is a
-            // file scan whose existing order already satisfies the required 
ordering.
-            // Otherwise, this rule inserts the necessary repartitioning and 
sorting
-            // operators.
-            //
-            // This used to be implemented as two separate rules: 
`EnforceDistribution`
-            // and `EnforceSorting`. It is now a single idempotent rule that 
decides
-            // distribution and sorting together in one bottom-up pass, so the
-            // `pushdown_sorts` step no longer breaks distribution invariants 
set
-            // earlier in the pipeline. See the module-level doc on
-            // [`EnsureRequirements`](crate::ensure_requirements) for the 
per-phase
-            // breakdown, and 
<https://github.com/apache/datafusion/issues/21973>
-            // for the original failure mode.
-            Arc::new(EnsureRequirements::new()),
+            // All enforcement (distribution and ordering) runs first in the

Review Comment:
   the details about enforcement seem irrelevant at this poing



##########
datafusion/physical-optimizer/src/analyzer.rs:
##########
@@ -0,0 +1,153 @@
+// 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.
+
+//! Physical analyzer
+
+use std::sync::Arc;
+
+use crate::ensure_requirements::{EnforceDistribution, EnforceSorting};
+use crate::output_requirements::OutputRequirements;
+
+// Re-export from this module for convenience.
+pub use datafusion_session::PhysicalAnalyzerRule;
+
+/// A rule-based physical analyzer.
+///
+/// Analyzer rules make the plan *valid*: they enforce the invariants every
+/// operator declares (distribution, ordering) rather than making the plan
+/// faster, mirroring the logical layer's `Analyzer`/`Optimizer` split. The
+/// default planner runs the whole analyzer phase before the
+/// [`PhysicalOptimizer`](crate::optimizer::PhysicalOptimizer) phase, so every
+/// optimizer rule can assume it receives a valid plan; see
+/// [`PhysicalAnalyzerRule`] for details.
+#[derive(Clone, Debug)]
+pub struct PhysicalAnalyzer {
+    /// All rules to apply
+    pub rules: Vec<Arc<dyn PhysicalAnalyzerRule + Send + Sync>>,
+}
+
+impl Default for PhysicalAnalyzer {
+    fn default() -> Self {
+        Self::new()
+    }
+}
+
+impl PhysicalAnalyzer {
+    /// Create a new analyzer using the recommended list of rules
+    pub fn new() -> Self {
+        // All enforcement runs here, first, as the analyzer phase: it makes 
the

Review Comment:
   this seems redundant to the documentation on `PhysicalAnalyzer` and 
`PhysicalAnalyzerRule`



##########
datafusion/sqllogictest/test_files/window_topn.slt:
##########
@@ -1575,11 +1575,9 @@ logical_plan
 physical_plan
 01)ProjectionExec: expr=[c1@0 as c1, c2@1 as c2, row_number() PARTITION BY 
[t.c1] ORDER BY [t.c2 DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND 
CURRENT ROW@2 as rn]
 02)--BoundedWindowAggExec: wdw=[row_number() PARTITION BY [t.c1] ORDER BY 
[t.c2 DESC NULLS FIRST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW: 
Field { "row_number() PARTITION BY [t.c1] ORDER BY [t.c2 DESC NULLS FIRST] 
RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW": UInt64 }, frame: RANGE 
BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW], mode=[Sorted]
-03)----SortPreservingMergeExec: [c1@0 ASC NULLS LAST, c2@1 DESC]
-04)------SortExec: expr=[c1@0 ASC NULLS LAST, c2@1 DESC], 
preserve_partitioning=[true]
-05)--------PartitionedTopKExec: fn=row_number, fetch=1, partition=[c1@0], 
order=[c2@1 DESC]
-06)----------RepartitionExec: partitioning=Hash([c1@0], 5), input_partitions=1
-07)------------DataSourceExec: partitions=1, partition_sizes=[1]
+03)----PartitionedTopKExec: fn=row_number, fetch=1, partition=[c1@0], 
order=[c2@1 DESC]

Review Comment:
   is this plan change ok? it looks like maybe the partitioning doesn't run in 
parallel anymore?



##########
datafusion/physical-optimizer/src/analyzer.rs:
##########
@@ -0,0 +1,153 @@
+// 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.
+
+//! Physical analyzer
+
+use std::sync::Arc;
+
+use crate::ensure_requirements::{EnforceDistribution, EnforceSorting};
+use crate::output_requirements::OutputRequirements;
+
+// Re-export from this module for convenience.
+pub use datafusion_session::PhysicalAnalyzerRule;
+
+/// A rule-based physical analyzer.
+///
+/// Analyzer rules make the plan *valid*: they enforce the invariants every
+/// operator declares (distribution, ordering) rather than making the plan
+/// faster, mirroring the logical layer's `Analyzer`/`Optimizer` split. The
+/// default planner runs the whole analyzer phase before the
+/// [`PhysicalOptimizer`](crate::optimizer::PhysicalOptimizer) phase, so every
+/// optimizer rule can assume it receives a valid plan; see
+/// [`PhysicalAnalyzerRule`] for details.
+#[derive(Clone, Debug)]
+pub struct PhysicalAnalyzer {
+    /// All rules to apply
+    pub rules: Vec<Arc<dyn PhysicalAnalyzerRule + Send + Sync>>,
+}
+
+impl Default for PhysicalAnalyzer {
+    fn default() -> Self {
+        Self::new()
+    }
+}
+
+impl PhysicalAnalyzer {
+    /// Create a new analyzer using the recommended list of rules
+    pub fn new() -> Self {
+        // All enforcement runs here, first, as the analyzer phase: it makes 
the
+        // plan *valid* (distribution first, then ordering) before any 
optimizer
+        // rule sees it, mirroring the logical Analyzer/Optimizer split. The
+        // optimizer rules that change these requirements (`JoinSelection`,
+        // `WindowTopN`, `FilterPushdown`) re-establish validity themselves, so
+        // enforcement is not repeated as an optimizer pass. The sort
+        // *optimizations* (`OptimizeSorts`) are not enforcement and stay in 
the
+        // optimizer phase.
+        let rules: Vec<Arc<dyn PhysicalAnalyzerRule + Send + Sync>> = vec![
+            // Establish the output-requirement boundary first, so enforcement 
can
+            // see it (parallelize top-level scans below it, preserve the 
query's
+            // final ordering). The matching remove pass runs late in the 
optimizer.
+            Arc::new(OutputRequirements::new_add_mode()),
+            Arc::new(EnforceDistribution::new()),

Review Comment:
   why did you split EnforceDistribution and EnforceSorting when they were 
unified before?



##########
datafusion/physical-optimizer/src/ensure_requirements/mod.rs:
##########
@@ -176,6 +177,302 @@ impl EnsureRequirements {
     }
 }
 
+/// Phases 0-2a: make the plan valid with respect to **distribution**
+/// requirements only (normalize interleave, join-key reordering, distribution
+/// enforcement). Split from ordering enforcement so it can be used on its own:
+/// the [`EnforceDistribution`] analyzer rule runs it as the first enforcement
+/// step, and rules that change only distribution (`JoinSelection`,
+/// `FilterPushdown`) call it directly to re-establish the partitioning their
+/// rewrite disturbed, without touching ordering.
+pub fn enforce_distribution_requirements(
+    plan: Arc<dyn ExecutionPlan>,
+    context: &dyn PhysicalOptimizerContext,
+) -> Result<Arc<dyn ExecutionPlan>> {
+    let config = context.config_options();
+    // Phase 0: Normalize `InterleaveExec` back to `UnionExec` (top-down).
+    // Interleaves are distribution artifacts of Phase 2, which re-derives
+    // them from the children's final partitioning. Keeping them would
+    // fail as soon as a child loses the partitioning they depend on.
+    use super::enforce_distribution::replace_interleave_with_union;
+    let plan = plan.transform_down(replace_interleave_with_union).data()?;
+
+    // Phase 1: Join key reordering (top-down, from EnforceDistribution)
+    use super::enforce_distribution::{
+        PlanWithKeyRequirements, adjust_input_keys_ordering,
+    };
+    let top_down_join_key_reordering = 
config.optimizer.top_down_join_key_reordering;
+    let plan = if top_down_join_key_reordering {
+        let ctx = PlanWithKeyRequirements::new_default(plan);
+        ctx.transform_down(adjust_input_keys_ordering).data()?.plan
+    } else {
+        use super::enforce_distribution::reorder_join_keys_to_inputs;
+        plan.transform_up(|p| 
Ok(Transformed::yes(reorder_join_keys_to_inputs(p)?)))
+            .data()?
+    };
+
+    // Phase 2a: Distribution enforcement (bottom-up)
+    use super::enforce_distribution::{
+        DistributionContext, ensure_distribution_with_stats,
+    };
+    let dist_ctx = DistributionContext::new_default(plan);
+    // Share one statistics context across the whole distribution pass so each
+    // subtree's statistics are computed once instead of once per ancestor.
+    // Build it from the session's statistics registry so registered providers
+    // are consulted (an empty registry, the default, is unchanged behavior).
+    // `StatsCache` is keyed by raw node pointer, so reset it after any node
+    // whose plan pointer actually changed: a rewrite can free a cached node
+    // and a later allocation could reuse its address. A node that makes no
+    // change cannot free anything, so the cache safely persists across the
+    // no-op nodes that dominate a deep plan.
+    let stats_ctx = match context.statistics_registry() {
+        Some(registry) => 
StatisticsContext::new_with_registry(registry.clone()),
+        None => StatisticsContext::new(),
+    };
+    let dist_ctx = dist_ctx
+        .transform_up(|ctx| {
+            let before = Arc::clone(&ctx.plan);
+            let result = ensure_distribution_with_stats(ctx, config, 
&stats_ctx)?;
+            if !Arc::ptr_eq(&before, &result.data.plan) {
+                stats_ctx.reset_cache();
+            }
+            Ok(result)
+        })
+        .data()?;
+    Ok(dist_ctx.plan)
+}
+
+/// Phase 2b: enforce **ordering** requirements by inserting `SortExec`s on a
+/// distribution-fixed plan (bottom-up). This is the enforcement half that is
+/// *not* idempotent, so in the default pipeline it runs exactly once via the
+/// [`EnforceSorting`] analyzer rule. Exposed as a free function so the 
combined
+/// [`enforce_requirements`] (used by rules that disturb ordering, and by the
+/// [`EnsureRequirements`] compatibility shim) can reuse it.
+pub fn enforce_sorting_requirements(
+    plan: Arc<dyn ExecutionPlan>,
+) -> Result<Arc<dyn ExecutionPlan>> {
+    use super::enforce_sorting::{PlanWithCorrespondingSort, ensure_sorting};
+    let sort_ctx = PlanWithCorrespondingSort::new_default(plan);
+    let sort_ctx = sort_ctx.transform_up(ensure_sorting)?.data;
+    Ok(sort_ctx.plan)
+}
+
+/// Phases 0-2: full requirement enforcement (distribution via
+/// [`enforce_distribution_requirements`], then sorting via
+/// [`enforce_sorting_requirements`]). The default pipeline uses the
+/// finer-grained [`EnforceDistribution`] / [`EnforceSorting`] rules instead;
+/// this stays for the [`EnsureRequirements`] compatibility shim.
+pub fn enforce_requirements(
+    plan: Arc<dyn ExecutionPlan>,
+    context: &dyn PhysicalOptimizerContext,
+) -> Result<Arc<dyn ExecutionPlan>> {
+    let plan = enforce_distribution_requirements(plan, context)?;
+    enforce_sorting_requirements(plan)
+}
+
+/// Phase 3: sort and distribution *optimizations* that make an already-valid
+/// plan faster (parallelize sorts, order-preserving variants, sort pushdown,
+/// partial sort). Split out of enforcement so it can run in the optimizer 
phase,
+/// after other optimizer rules (such as `WindowTopN`) have produced the
+/// operators it parallelizes. Exposed as the [`OptimizeSorts`] rule.
+pub fn optimize_sorts(
+    plan: Arc<dyn ExecutionPlan>,
+    config: &ConfigOptions,
+) -> Result<Arc<dyn ExecutionPlan>> {
+    // 3a: Parallelize sorts (Coalesce+Sort → SPM+Sort)
+    use super::enforce_sorting::{
+        PlanWithCorrespondingCoalescePartitions, parallelize_sorts,
+        replace_with_partial_sort,
+    };
+    let plan = if config.optimizer.repartition_sorts {
+        let ctx = PlanWithCorrespondingCoalescePartitions::new_default(plan);
+        ctx.transform_up(parallelize_sorts).data()?.plan
+    } else {
+        plan
+    };
+
+    // 3b: Order-preserving variants
+    use super::enforce_sorting::replace_with_order_preserving_variants::{
+        OrderPreservationContext, replace_with_order_preserving_variants,
+    };
+    let ctx = OrderPreservationContext::new_default(plan);
+    let plan = ctx
+        .transform_up(|c| replace_with_order_preserving_variants(c, false, 
true, config))
+        .data()?
+        .plan;
+
+    // 3c: Sort pushdown (distribution-aware)
+    use super::enforce_sorting::sort_pushdown::{
+        SortPushDown, assign_initial_requirements, pushdown_sorts,
+    };
+    let mut sort_pushdown = SortPushDown::new_default(plan);
+    assign_initial_requirements(&mut sort_pushdown);
+    let adjusted = pushdown_sorts(sort_pushdown)?;
+
+    // 3d: Partial sort
+    adjusted
+        .plan
+        .transform_up(|p| Ok(Transformed::yes(replace_with_partial_sort(p)?)))
+        .data()
+}
+
+/// Enforces **distribution** requirements (Phases 0-2a) via
+/// [`enforce_distribution_requirements`]. In the default pipeline it runs as 
the
+/// [`PhysicalAnalyzerRule`] that makes the plan distribution-valid before the
+/// optimizer rules see it. It is idempotent enough to run more than once, and
+/// still implements [`PhysicalOptimizerRule`] so downstream pipelines that 
splice
+/// it in by position keep working; the default optimizer rules that change
+/// distribution (`JoinSelection`, `FilterPushdown`) instead call
+/// [`enforce_distribution_requirements`] directly to re-establish it 
themselves.
+#[derive(Default, Debug)]
+pub struct EnforceDistribution {}
+
+impl EnforceDistribution {
+    #[expect(missing_docs)]
+    pub fn new() -> Self {
+        Self {}
+    }
+}
+
+impl PhysicalOptimizerRule for EnforceDistribution {
+    fn optimize(
+        &self,
+        plan: Arc<dyn ExecutionPlan>,
+        config: &ConfigOptions,
+    ) -> Result<Arc<dyn ExecutionPlan>> {
+        enforce_distribution_requirements(plan, 
&ConfigOnlyContext::new(config))
+    }
+
+    fn optimize_with_context(
+        &self,
+        plan: Arc<dyn ExecutionPlan>,
+        context: &dyn PhysicalOptimizerContext,
+    ) -> Result<Arc<dyn ExecutionPlan>> {
+        enforce_distribution_requirements(plan, context)
+    }
+
+    fn name(&self) -> &str {
+        "EnforceDistribution"
+    }
+
+    fn schema_check(&self) -> bool {
+        true
+    }
+}
+
+impl PhysicalAnalyzerRule for EnforceDistribution {
+    fn analyze(
+        &self,
+        plan: Arc<dyn ExecutionPlan>,
+        config: &ConfigOptions,
+    ) -> Result<Arc<dyn ExecutionPlan>> {
+        enforce_distribution_requirements(plan, 
&ConfigOnlyContext::new(config))
+    }
+
+    fn analyze_with_context(
+        &self,
+        plan: Arc<dyn ExecutionPlan>,
+        context: &dyn PhysicalOptimizerContext,
+    ) -> Result<Arc<dyn ExecutionPlan>> {
+        enforce_distribution_requirements(plan, context)
+    }
+
+    fn name(&self) -> &str {
+        "EnforceDistribution"
+    }
+
+    fn schema_check(&self) -> bool {
+        true
+    }
+}
+
+/// Enforces **ordering** requirements (Phase 2b) via
+/// [`enforce_sorting_requirements`]. Not idempotent, so it runs exactly once 
in
+/// the default pipeline, as the ordering-enforcement [`PhysicalAnalyzerRule`]
+/// (after [`EnforceDistribution`], on the distribution-fixed plan). The later
+/// optimizer rules that disturb ordering (`WindowTopN`) re-establish it
+/// themselves via [`enforce_requirements`].
+#[derive(Default, Debug)]
+pub struct EnforceSorting {}
+
+impl EnforceSorting {
+    #[expect(missing_docs)]
+    pub fn new() -> Self {
+        Self {}
+    }
+}
+
+impl PhysicalOptimizerRule for EnforceSorting {

Review Comment:
   Shouldn't enforcing sorting be a PhysicalAnalyzer rule too? Until that is 
run the sort order is not correct for declared input requirements



##########
datafusion/physical-optimizer/src/optimizer.rs:
##########
@@ -89,9 +89,12 @@ impl PhysicalOptimizer {
         //   queries, and will likely increase the optimization time. Please 
extend
         //   existing rules when possible, rather than adding a new rule.
         let rules: Vec<Arc<dyn PhysicalOptimizerRule + Send + Sync>> = vec![
-            // If there is a output requirement of the query, make sure that
-            // this information is not lost across different rules during 
optimization.
-            Arc::new(OutputRequirements::new_add_mode()),
+            // The output-requirement boundary (`OutputRequirements` add mode) 
and

Review Comment:
   this seems like comments related to this PR (not something that will be 
relevant to future readers)



-- 
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]

Reply via email to