Skip to main content

query/optimizer/
enforce_sorting.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Sorting enforcement that runs after GreptimeDB's custom physical rules.
16//!
17//! DataFusion 55 moved the standalone `EnforceSorting` phases into
18//! `EnsureRequirements`. GreptimeDB still needs to rerun those phases after
19//! custom rules modify scan partitioning and distribution.
20
21use std::sync::Arc;
22
23use datafusion::physical_optimizer::PhysicalOptimizerRule;
24use datafusion::physical_optimizer::enforce_sorting::replace_with_order_preserving_variants::{
25    OrderPreservationContext, replace_with_order_preserving_variants,
26};
27use datafusion::physical_optimizer::enforce_sorting::sort_pushdown::{
28    SortPushDown, assign_initial_requirements, pushdown_sorts,
29};
30use datafusion::physical_optimizer::enforce_sorting::{
31    PlanWithCorrespondingCoalescePartitions, PlanWithCorrespondingSort, ensure_sorting,
32    parallelize_sorts, replace_with_partial_sort,
33};
34use datafusion::physical_plan::ExecutionPlan;
35use datafusion_common::Result;
36use datafusion_common::config::ConfigOptions;
37use datafusion_common::tree_node::{Transformed, TransformedResult, TreeNode};
38
39/// Runs the standalone sorting-enforcement pipeline removed in DataFusion 55.
40#[derive(Debug)]
41pub struct EnforceSorting;
42
43impl PhysicalOptimizerRule for EnforceSorting {
44    fn optimize(
45        &self,
46        plan: Arc<dyn ExecutionPlan>,
47        config: &ConfigOptions,
48    ) -> Result<Arc<dyn ExecutionPlan>> {
49        // Phase 1: ensure sorting requirements and remove redundant sorts.
50        let sorting = PlanWithCorrespondingSort::new_default(plan);
51        let sorting = sorting.transform_up(ensure_sorting)?.data;
52
53        // Phase 2: optionally turn CoalescePartitions + Sort into parallel
54        // sorts followed by a SortPreservingMerge.
55        let plan = if config.optimizer.repartition_sorts {
56            let parallel = PlanWithCorrespondingCoalescePartitions::new_default(sorting.plan)
57                .transform_up(parallelize_sorts)
58                .data()?;
59            parallel.plan
60        } else {
61            sorting.plan
62        };
63
64        // Phase 3: use order-preserving executor variants where appropriate.
65        let variants = OrderPreservationContext::new_default(plan);
66        let variants = variants
67            .transform_up(|context| {
68                replace_with_order_preserving_variants(context, false, true, config)
69            })
70            .data()?;
71
72        // Phase 4: push sorts down through order-preserving operators.
73        let mut pushdown = SortPushDown::new_default(variants.plan);
74        assign_initial_requirements(&mut pushdown);
75        let pushed = pushdown_sorts(pushdown)?;
76
77        // Phase 5: exploit an already-satisfied prefix on unbounded inputs.
78        pushed
79            .plan
80            .transform_up(|plan| Ok(Transformed::yes(replace_with_partial_sort(plan)?)))
81            .data()
82    }
83
84    fn name(&self) -> &str {
85        "EnforceSorting"
86    }
87
88    fn schema_check(&self) -> bool {
89        true
90    }
91}