Skip to content

Commit de0572d

Browse files
allow range to satisfy key distribution generally
1 parent ca02890 commit de0572d

11 files changed

Lines changed: 182 additions & 537 deletions

File tree

datafusion/physical-expr/src/partitioning.rs

Lines changed: 160 additions & 273 deletions
Large diffs are not rendered by default.

datafusion/physical-optimizer/src/ensure_requirements/enforce_distribution.rs

Lines changed: 11 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -35,7 +35,7 @@ use std::sync::Arc;
3535
use crate::output_requirements::OutputRequirementExec;
3636
use crate::utils::{
3737
add_sort_above_with_check, is_coalesce_partitions, is_repartition,
38-
is_sort_preserving_merge, range_partitioning_satisfies_key_partitioning,
38+
is_sort_preserving_merge,
3939
};
4040

4141
use arrow::compute::SortOptions;
@@ -698,18 +698,13 @@ fn add_roundrobin_on_top(
698698
}
699699
}
700700

701-
// TODO: remove this temporary bridge once [`Partitioning::Range`]
702-
// generally satisfies [`Distribution::KeyPartitioned`] through
703-
// [`Partitioning::satisfaction`].
704-
// <https://github.com/apache/datafusion/issues/23266>.
705-
//
706-
// Partial aggregates do not require key partitioning, but they preserve their
707-
// input partitioning for the final aggregate. Until Range satisfies
708-
// KeyPartitioned generally, this check keeps preserve_file_partitions from
709-
// inserting RoundRobin between a reusable Range input and the partial aggregate.
710-
fn partial_aggregate_preserves_reusable_partitioning(
701+
// Partial aggregates require unspecified input distribution, but their output
702+
// may already satisfy the final aggregate's key distribution because partial
703+
// aggregation preserves/projects input partitioning. Keep that reusable output
704+
// partitioning intact when preserve_file_partitions would otherwise insert
705+
// RoundRobin below the partial aggregate.
706+
fn partial_aggregate_output_satisfies_final_partitioning(
711707
plan: &Arc<dyn ExecutionPlan>,
712-
child: &Arc<dyn ExecutionPlan>,
713708
allow_subset_satisfy_partitioning: bool,
714709
) -> bool {
715710
let Some(aggregate) = plan.downcast_ref::<AggregateExec>() else {
@@ -722,24 +717,15 @@ fn partial_aggregate_preserves_reusable_partitioning(
722717
return false;
723718
}
724719

725-
let group_exprs = aggregate.group_expr().input_exprs();
726-
let output_partitioning = child.output_partitioning();
727-
let eq_properties = child.equivalence_properties();
728-
let key_distribution = Distribution::KeyPartitioned(group_exprs.clone());
720+
let key_distribution = Distribution::KeyPartitioned(aggregate.output_group_expr());
729721

730-
output_partitioning
722+
plan.output_partitioning()
731723
.satisfaction(
732724
&key_distribution,
733-
eq_properties,
725+
plan.equivalence_properties(),
734726
allow_subset_satisfy_partitioning,
735727
)
736728
.is_satisfied()
737-
|| range_partitioning_satisfies_key_partitioning(
738-
output_partitioning,
739-
&group_exprs,
740-
eq_properties,
741-
allow_subset_satisfy_partitioning,
742-
)
743729
}
744730

745731
/// Adds a [`SortPreservingMergeExec`] or a [`CoalescePartitionsExec`] operator
@@ -1308,9 +1294,8 @@ pub fn ensure_distribution(
13081294

13091295
let preserve_partial_aggregate_partitioning =
13101296
preserve_file_partition_threshold_met
1311-
&& partial_aggregate_preserves_reusable_partitioning(
1297+
&& partial_aggregate_output_satisfies_final_partitioning(
13121298
&plan,
1313-
&child.plan,
13141299
allow_subset_satisfy_partitioning,
13151300
);
13161301

datafusion/physical-optimizer/src/utils.rs

Lines changed: 1 addition & 58 deletions
Original file line numberDiff line numberDiff line change
@@ -18,10 +18,7 @@
1818
use std::sync::Arc;
1919

2020
use datafusion_common::Result;
21-
use datafusion_physical_expr::{
22-
Distribution, EquivalenceProperties, LexOrdering, LexRequirement, Partitioning,
23-
PhysicalExpr, physical_exprs_equal,
24-
};
21+
use datafusion_physical_expr::{Distribution, LexOrdering, LexRequirement};
2522
use datafusion_physical_plan::coalesce_partitions::CoalescePartitionsExec;
2623
use datafusion_physical_plan::limit::{GlobalLimitExec, LocalLimitExec};
2724
use datafusion_physical_plan::repartition::RepartitionExec;
@@ -161,60 +158,6 @@ pub fn is_repartition(plan: &Arc<dyn ExecutionPlan>) -> bool {
161158
plan.is::<RepartitionExec>()
162159
}
163160

164-
/// TODO: remove once Range generally satisfies KeyPartitioned requirements
165-
/// through Partitioning::satisfaction.
166-
/// See <https://github.com/apache/datafusion/issues/23266>.
167-
///
168-
/// Checks whether range partitioning satisfies a key partitioning requirement.
169-
/// This is intentionally separate from general partitioning satisfaction while
170-
/// range reuse is rolled out operator by operator.
171-
pub(crate) fn range_partitioning_satisfies_key_partitioning(
172-
partitioning: &Partitioning,
173-
required_exprs: &[Arc<dyn PhysicalExpr>],
174-
eq_properties: &EquivalenceProperties,
175-
allow_subset: bool,
176-
) -> bool {
177-
match partitioning {
178-
Partitioning::Range(range) => {
179-
let partition_exprs = range
180-
.ordering()
181-
.iter()
182-
.map(|sort_expr| Arc::clone(&sort_expr.expr))
183-
.collect::<Vec<_>>();
184-
185-
if partition_exprs.is_empty() || required_exprs.is_empty() {
186-
return false;
187-
}
188-
189-
let eq_group = eq_properties.eq_group();
190-
let normalized_partition_exprs = partition_exprs
191-
.iter()
192-
.map(|expr| eq_group.normalize_expr(Arc::clone(expr)))
193-
.collect::<Vec<_>>();
194-
let normalized_required_exprs = required_exprs
195-
.iter()
196-
.map(|expr| eq_group.normalize_expr(Arc::clone(expr)))
197-
.collect::<Vec<_>>();
198-
199-
if physical_exprs_equal(
200-
&normalized_required_exprs,
201-
&normalized_partition_exprs,
202-
) {
203-
return true;
204-
}
205-
206-
allow_subset
207-
&& normalized_partition_exprs.len() < normalized_required_exprs.len()
208-
&& normalized_partition_exprs.iter().all(|partition_expr| {
209-
normalized_required_exprs
210-
.iter()
211-
.any(|required_expr| partition_expr.eq(required_expr))
212-
})
213-
}
214-
_ => false,
215-
}
216-
}
217-
218161
/// Checks whether the given operator is a limit;
219162
/// i.e. either a [`LocalLimitExec`] or a [`GlobalLimitExec`].
220163
pub fn is_limit(plan: &Arc<dyn ExecutionPlan>) -> bool {

datafusion/physical-plan/src/aggregates/mod.rs

Lines changed: 2 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -1934,7 +1934,7 @@ impl ExecutionPlan for AggregateExec {
19341934
}
19351935

19361936
fn input_distribution_requirements(&self) -> InputDistributionRequirements {
1937-
let requirements = InputDistributionRequirements::new(match &self.mode {
1937+
InputDistributionRequirements::new(match &self.mode {
19381938
AggregateMode::Partial | AggregateMode::PartialReduce => {
19391939
vec![Distribution::UnspecifiedDistribution]
19401940
}
@@ -1944,15 +1944,7 @@ impl ExecutionPlan for AggregateExec {
19441944
AggregateMode::Final | AggregateMode::Single => {
19451945
vec![Distribution::SinglePartition]
19461946
}
1947-
});
1948-
match &self.mode {
1949-
AggregateMode::FinalPartitioned | AggregateMode::SinglePartitioned
1950-
if !self.group_by.has_grouping_set() =>
1951-
{
1952-
requirements.allow_range_satisfaction_for_key_partitioning()
1953-
}
1954-
_ => requirements,
1955-
}
1947+
})
19561948
}
19571949

19581950
fn required_input_ordering(&self) -> Vec<Option<OrderingRequirements>> {

0 commit comments

Comments
 (0)