Skip to content

Commit e236f7b

Browse files
committed
feat: allow Range partitioning to satisfy KeyPartitioned requirement for Window/TopK
1 parent 6930807 commit e236f7b

4 files changed

Lines changed: 337 additions & 2 deletions

File tree

datafusion/physical-plan/src/sorts/partitioned_topk.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -311,6 +311,7 @@ impl ExecutionPlan for PartitionedTopKExec {
311311
crate::InputDistributionRequirements::new(vec![Distribution::KeyPartitioned(
312312
partition_exprs,
313313
)])
314+
.allow_range_satisfaction_for_key_partitioning()
314315
}
315316

316317
fn maintains_input_order(&self) -> Vec<bool> {

datafusion/physical-plan/src/windows/bounded_window_agg_exec.rs

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -37,7 +37,8 @@ use crate::windows::{
3737
use crate::{
3838
ColumnStatistics, DisplayAs, DisplayFormatType, Distribution, ExecutionPlan,
3939
ExecutionPlanProperties, InputOrderMode, PlanProperties, RecordBatchStream,
40-
SendableRecordBatchStream, Statistics, WindowExpr, check_if_same_properties,
40+
SendableRecordBatchStream, Statistics, WindowExpr,
41+
check_if_same_properties,
4142
};
4243

4344
use arrow::compute::take_record_batch;
@@ -342,6 +343,7 @@ impl ExecutionPlan for BoundedWindowAggExec {
342343
} else {
343344
vec![Distribution::KeyPartitioned(self.partition_keys().clone())]
344345
})
346+
.allow_range_satisfaction_for_key_partitioning()
345347
}
346348

347349
fn maintains_input_order(&self) -> Vec<bool> {

datafusion/physical-plan/src/windows/window_agg_exec.rs

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -33,7 +33,8 @@ use crate::windows::{
3333
use crate::{
3434
ColumnStatistics, DisplayAs, DisplayFormatType, Distribution, ExecutionPlan,
3535
ExecutionPlanProperties, PhysicalExpr, PlanProperties, RecordBatchStream,
36-
SendableRecordBatchStream, Statistics, WindowExpr, check_if_same_properties,
36+
SendableRecordBatchStream, Statistics, WindowExpr,
37+
check_if_same_properties,
3738
};
3839

3940
use arrow::array::ArrayRef;
@@ -250,6 +251,7 @@ impl ExecutionPlan for WindowAggExec {
250251
} else {
251252
vec![Distribution::KeyPartitioned(self.partition_keys())]
252253
})
254+
.allow_range_satisfaction_for_key_partitioning()
253255
}
254256

255257
fn with_new_children(

0 commit comments

Comments
 (0)