query/optimizer/
enforce_sorting.rs1use 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#[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 let sorting = PlanWithCorrespondingSort::new_default(plan);
51 let sorting = sorting.transform_up(ensure_sorting)?.data;
52
53 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 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 let mut pushdown = SortPushDown::new_default(variants.plan);
74 assign_initial_requirements(&mut pushdown);
75 let pushed = pushdown_sorts(pushdown)?;
76
77 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}