zhuqi-lucas commented on code in PR #25688:
URL: https://github.com/apache/datafusion/pull/25688#discussion_r4155680202
##########
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:
Thanks @alamb. No, it is not. Removing the PR history from that comment.
##########
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:
Thanks @alamb. The snapshots embed `target_partitions`, which defaults to
the number of CPUs, so they rendered differently on a 12-core dev machine and
4-core CI. I pinned it to 4 in the test helper and re-accepted the snapshots.
Nothing about the pushdown behaviour changed.
##########
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:
Thanks @alamb. It already is: the default analyzer list is
`[OutputRequirements, EnforceDistribution, EnforceSorting]`, so sorting is
enforced in the analyzer phase, right after distribution. This
`PhysicalOptimizerRule` impl is only kept so custom rule lists that register
`EnforceSorting` by position keep working. I will make that clearer in the
docs, since you had to ask.
##########
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:
Thanks @alamb, agreed. Removing.
##########
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 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:
Thanks @alamb, agreed. Removing and leaving the explanation in one place.
--
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]