Skip to content

Commit f2e6d2c

Browse files
Support grouped distribution requirements
1 parent 7ff7278 commit f2e6d2c

37 files changed

Lines changed: 1635 additions & 473 deletions

datafusion/core/tests/physical_optimizer/enforce_distribution.rs

Lines changed: 49 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -265,8 +265,12 @@ impl ExecutionPlan for SinglePartitionMaintainsOrderExec {
265265
vec![&self.input]
266266
}
267267

268-
fn required_input_distribution(&self) -> Vec<Distribution> {
269-
vec![Distribution::SinglePartition]
268+
fn input_distribution_requirements(
269+
&self,
270+
) -> datafusion_physical_plan::InputDistributionRequirements {
271+
datafusion_physical_plan::InputDistributionRequirements::new(vec![
272+
Distribution::SinglePartition,
273+
])
270274
}
271275

272276
fn maintains_input_order(&self) -> Vec<bool> {
@@ -823,6 +827,49 @@ fn range_grouping_set_aggregate_rehashes_with_grouping_id() -> Result<()> {
823827
Ok(())
824828
}
825829

830+
#[test]
831+
fn range_inner_hash_join_rehashes_incompatible_range_partitioning() -> Result<()> {
832+
let left = parquet_exec_with_output_partitioning(range_partitioning(
833+
"a",
834+
[10, 20, 30],
835+
SortOptions::default(),
836+
)?);
837+
let right = projection_exec_with_alias(
838+
parquet_exec_with_output_partitioning(range_partitioning(
839+
"a",
840+
[10, 30, 40],
841+
SortOptions::default(),
842+
)?),
843+
vec![
844+
("a".to_string(), "a1".to_string()),
845+
("b".to_string(), "b1".to_string()),
846+
],
847+
);
848+
let join_on = vec![(
849+
Arc::new(Column::new_with_schema("a", &left.schema())?) as _,
850+
Arc::new(Column::new_with_schema("a1", &right.schema())?) as _,
851+
)];
852+
let join = hash_join_exec(left, right, &join_on, &JoinType::Inner);
853+
854+
let plan = TestConfig::default()
855+
.with_query_execution_partitions(4)
856+
.to_plan(join, &DISTRIB_DISTRIB_SORT);
857+
858+
assert_plan!(
859+
plan,
860+
@r"
861+
HashJoinExec: mode=Partitioned, join_type=Inner, on=[(a@0, a1@0)]
862+
RepartitionExec: partitioning=Hash([a@0], 4), input_partitions=4
863+
DataSourceExec: file_groups={4 groups: [[p0], [p1], [p2], [p3]]}, projection=[a, b, c, d, e], output_partitioning=Range([a@0 ASC], [(10), (20), (30)], 4), file_type=parquet
864+
RepartitionExec: partitioning=Hash([a1@0], 4), input_partitions=4
865+
ProjectionExec: expr=[a@0 as a1, b@1 as b1]
866+
DataSourceExec: file_groups={4 groups: [[p0], [p1], [p2], [p3]]}, projection=[a, b, c, d, e], output_partitioning=Range([a@0 ASC], [(10), (30), (40)], 4), file_type=parquet
867+
"
868+
);
869+
870+
Ok(())
871+
}
872+
826873
#[test]
827874
fn multi_hash_joins() -> Result<()> {
828875
let left = parquet_exec();

datafusion/core/tests/physical_optimizer/ensure_requirements.rs

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -974,8 +974,12 @@ impl ExecutionPlan for MockReqExec {
974974
fn children(&self) -> Vec<&Arc<dyn ExecutionPlan>> {
975975
vec![&self.input]
976976
}
977-
fn required_input_distribution(&self) -> Vec<Distribution> {
978-
vec![self.dist.clone()]
977+
fn input_distribution_requirements(
978+
&self,
979+
) -> datafusion_physical_plan::InputDistributionRequirements {
980+
datafusion_physical_plan::InputDistributionRequirements::new(vec![
981+
self.dist.clone(),
982+
])
979983
}
980984
fn required_input_ordering(&self) -> Vec<Option<OrderingRequirements>> {
981985
vec![

datafusion/core/tests/physical_optimizer/projection_pushdown.rs

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -800,7 +800,9 @@ fn test_output_req_after_projection() -> Result<()> {
800800
if let Distribution::KeyPartitioned(vec) = after_optimize
801801
.downcast_ref::<OutputRequirementExec>()
802802
.unwrap()
803-
.required_input_distribution()[0]
803+
.input_distribution_requirements()
804+
.child_distribution(0)
805+
.unwrap()
804806
.clone()
805807
{
806808
assert!(

datafusion/core/tests/physical_optimizer/sanity_checker.rs

Lines changed: 47 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -19,9 +19,9 @@ use insta::assert_snapshot;
1919
use std::sync::Arc;
2020

2121
use crate::physical_optimizer::test_utils::{
22-
bounded_window_exec, global_limit_exec, local_limit_exec, memory_exec,
23-
projection_exec, repartition_exec, sort_exec, sort_expr, sort_expr_options,
24-
sort_merge_join_exec, sort_preserving_merge_exec, union_exec,
22+
bounded_window_exec, global_limit_exec, hash_join_exec, local_limit_exec,
23+
memory_exec, projection_exec, repartition_exec, sort_exec, sort_expr,
24+
sort_expr_options, sort_merge_join_exec, sort_preserving_merge_exec, union_exec,
2525
};
2626

2727
use arrow::compute::SortOptions;
@@ -30,8 +30,8 @@ use datafusion::datasource::stream::{FileStreamProvider, StreamConfig, StreamTab
3030
use datafusion::prelude::{CsvReadOptions, SessionContext};
3131
use datafusion_common::config::ConfigOptions;
3232
use datafusion_common::{JoinType, Result, ScalarValue};
33-
use datafusion_physical_expr::Partitioning;
3433
use datafusion_physical_expr::expressions::{Literal, col};
34+
use datafusion_physical_expr::{Partitioning, RangePartitioning, SplitPoint};
3535
use datafusion_physical_expr_common::sort_expr::LexOrdering;
3636
use datafusion_physical_optimizer::PhysicalOptimizerRule;
3737
use datafusion_physical_optimizer::sanity_checker::SanityCheckPlan;
@@ -400,6 +400,49 @@ fn assert_sanity_check(plan: &Arc<dyn ExecutionPlan>, is_sane: bool) {
400400
);
401401
}
402402

403+
fn range_partitioned_exec(
404+
schema: &SchemaRef,
405+
key: &str,
406+
split_points: impl IntoIterator<Item = i32>,
407+
) -> Result<Arc<dyn ExecutionPlan>> {
408+
let split_points = split_points
409+
.into_iter()
410+
.map(|value| SplitPoint::new(vec![ScalarValue::Int32(Some(value))]))
411+
.collect();
412+
let partitioning = Partitioning::Range(RangePartitioning::try_new(
413+
[sort_expr(key, schema)].into(),
414+
split_points,
415+
)?);
416+
RepartitionExec::try_new(memory_exec(schema), partitioning)
417+
.map(|exec| Arc::new(exec) as Arc<dyn ExecutionPlan>)
418+
}
419+
420+
#[test]
421+
fn test_partitioned_hash_join_requires_co_partitioned_children() -> Result<()> {
422+
let schema = create_test_schema2();
423+
let join_on = vec![(col("a", &schema)?, col("a", &schema)?)];
424+
425+
let compatible_join = hash_join_exec(
426+
range_partitioned_exec(&schema, "a", [10])?,
427+
range_partitioned_exec(&schema, "a", [10])?,
428+
join_on.clone(),
429+
None,
430+
&JoinType::Inner,
431+
)?;
432+
assert_sanity_check(&compatible_join, true);
433+
434+
let incompatible_join = hash_join_exec(
435+
range_partitioned_exec(&schema, "a", [10])?,
436+
range_partitioned_exec(&schema, "a", [20])?,
437+
join_on,
438+
None,
439+
&JoinType::Inner,
440+
)?;
441+
assert_sanity_check(&incompatible_join, false);
442+
443+
Ok(())
444+
}
445+
403446
#[tokio::test]
404447
/// Tests that plan is valid when the sort requirements are satisfied.
405448
async fn test_bounded_window_agg_sort_requirement() -> Result<()> {

datafusion/core/tests/user_defined/user_defined_plan.rs

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -708,8 +708,12 @@ impl ExecutionPlan for TopKExec {
708708
&self.cache
709709
}
710710

711-
fn required_input_distribution(&self) -> Vec<Distribution> {
712-
vec![Distribution::SinglePartition]
711+
fn input_distribution_requirements(
712+
&self,
713+
) -> datafusion_physical_plan::InputDistributionRequirements {
714+
datafusion_physical_plan::InputDistributionRequirements::new(vec![
715+
Distribution::SinglePartition,
716+
])
713717
}
714718

715719
fn children(&self) -> Vec<&Arc<dyn ExecutionPlan>> {

datafusion/datasource/src/sink.rs

Lines changed: 11 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -31,8 +31,9 @@ use datafusion_physical_expr_common::sort_expr::{LexRequirement, OrderingRequire
3131
use datafusion_physical_plan::metrics::MetricsSet;
3232
use datafusion_physical_plan::stream::RecordBatchStreamAdapter;
3333
use datafusion_physical_plan::{
34-
DisplayAs, DisplayFormatType, ExecutionPlan, ExecutionPlanProperties, Partitioning,
35-
PlanProperties, SendableRecordBatchStream, execute_input_stream,
34+
DisplayAs, DisplayFormatType, ExecutionPlan, ExecutionPlanProperties,
35+
InputDistributionRequirements, Partitioning, PlanProperties,
36+
SendableRecordBatchStream, execute_input_stream,
3637
};
3738

3839
use async_trait::async_trait;
@@ -189,9 +190,16 @@ impl ExecutionPlan for DataSinkExec {
189190
}
190191

191192
fn required_input_distribution(&self) -> Vec<Distribution> {
193+
self.input_distribution_requirements().into_per_child()
194+
}
195+
196+
fn input_distribution_requirements(&self) -> InputDistributionRequirements {
192197
// DataSink is responsible for dynamically partitioning its
193198
// own input at execution time, and so requires a single input partition.
194-
vec![Distribution::SinglePartition; self.children().len()]
199+
InputDistributionRequirements::new(vec![
200+
Distribution::SinglePartition;
201+
self.children().len()
202+
])
195203
}
196204

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

datafusion/physical-expr/src/lib.rs

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -63,7 +63,9 @@ pub use equivalence::{
6363
AcrossPartitions, ConstExpr, EquivalenceProperties, calculate_union,
6464
};
6565
pub use expressions::{DynamicFilterTracker, DynamicFilterTracking};
66-
pub use partitioning::{Distribution, Partitioning, RangePartitioning};
66+
pub use partitioning::{
67+
Distribution, Partitioning, PartitioningSatisfaction, RangePartitioning,
68+
};
6769
pub use physical_expr::{
6870
add_offset_to_expr, add_offset_to_physical_sort_exprs, create_lex_ordering,
6971
create_ordering, create_physical_partitioning, create_physical_sort_expr,

0 commit comments

Comments
 (0)