gene-bordegaray commented on code in PR #23484:
URL: https://github.com/apache/datafusion/pull/23484#discussion_r3614687304
##########
datafusion/sqllogictest/test_files/range_partitioning.slt:
##########
@@ -629,7 +630,298 @@ ORDER BY l.range_key;
35 700
##########
-# TEST 18: Sort Merge Join Avoids Repartition for Compatible Range Inputs
+# TEST 18: Right Join on Range Partition Column
+# Compatible Range inputs satisfy the join's partitioning requirements, so no
+# Hash repartitioning is inserted. The left filter keeps its Range partitioning
+# and the unmatched right rows above 150 are preserved.
+##########
+
+query TT
+EXPLAIN SELECT l.value, r.range_key, r.value
+FROM (SELECT range_key, value FROM range_partitioned WHERE value <= 150) l
+RIGHT JOIN range_partitioned r ON l.range_key = r.range_key;
+----
+physical_plan
+01)HashJoinExec: mode=Partitioned, join_type=Right, on=[(range_key@0,
range_key@0)], projection=[value@1, range_key@2, value@3]
+02)--FilterExec: value@1 <= 150
+03)----DataSourceExec: file_groups=<slt:ignore>, projection=[range_key,
value], output_partitioning=Range([range_key@0 ASC], [(10), (20), (30)], 4),
file_type=csv, has_header=false
+04)--DataSourceExec: file_groups=<slt:ignore>, projection=[range_key, value],
output_partitioning=Range([range_key@0 ASC], [(10), (20), (30)], 4),
file_type=csv, has_header=false
+
+query III
+SELECT l.value, r.range_key, r.value
+FROM (SELECT range_key, value FROM range_partitioned WHERE value <= 150) l
+RIGHT JOIN range_partitioned r ON l.range_key = r.range_key
+ORDER BY r.range_key;
+----
+10 1 10
+50 5 50
+100 10 100
+150 15 150
+NULL 20 200
+NULL 25 250
+NULL 30 300
+NULL 35 350
+
+##########
+# TEST 19: Right Semi Join on Range Partition Column
+# Compatible Range inputs avoid Hash repartitioning for RightSemi joins.
+# Only right rows with a match on the filtered left side are returned.
+##########
+
+query TT
+EXPLAIN SELECT r.range_key, r.value
+FROM (SELECT range_key FROM range_partitioned WHERE value <= 150) l
+RIGHT SEMI JOIN range_partitioned r ON l.range_key = r.range_key;
+----
+physical_plan
+01)HashJoinExec: mode=Partitioned, join_type=RightSemi, on=[(range_key@0,
range_key@0)]
+02)--FilterExec: value@1 <= 150, projection=[range_key@0]
+03)----DataSourceExec: file_groups=<slt:ignore>, projection=[range_key,
value], output_partitioning=Range([range_key@0 ASC], [(10), (20), (30)], 4),
file_type=csv, has_header=false
+04)--DataSourceExec: file_groups=<slt:ignore>, projection=[range_key, value],
output_partitioning=Range([range_key@0 ASC], [(10), (20), (30)], 4),
file_type=csv, has_header=false
+
+query II
+SELECT r.range_key, r.value
+FROM (SELECT range_key FROM range_partitioned WHERE value <= 150) l
+RIGHT SEMI JOIN range_partitioned r ON l.range_key = r.range_key
+ORDER BY r.range_key;
+----
+1 10
+5 50
+10 100
+15 150
+
+##########
+# TEST 20: Right Anti Join on Range Partition Column
+# Compatible Range inputs avoid Hash repartitioning for RightAnti joins.
+# Only right rows without a match on the filtered left side are returned.
+##########
+
+query TT
+EXPLAIN SELECT r.range_key, r.value
+FROM (SELECT range_key FROM range_partitioned WHERE value <= 150) l
+RIGHT ANTI JOIN range_partitioned r ON l.range_key = r.range_key;
+----
+physical_plan
+01)HashJoinExec: mode=Partitioned, join_type=RightAnti, on=[(range_key@0,
range_key@0)]
+02)--FilterExec: value@1 <= 150, projection=[range_key@0]
+03)----DataSourceExec: file_groups=<slt:ignore>, projection=[range_key,
value], output_partitioning=Range([range_key@0 ASC], [(10), (20), (30)], 4),
file_type=csv, has_header=false
+04)--DataSourceExec: file_groups=<slt:ignore>, projection=[range_key, value],
output_partitioning=Range([range_key@0 ASC], [(10), (20), (30)], 4),
file_type=csv, has_header=false
+
+query II
+SELECT r.range_key, r.value
+FROM (SELECT range_key FROM range_partitioned WHERE value <= 150) l
+RIGHT ANTI JOIN range_partitioned r ON l.range_key = r.range_key
+ORDER BY r.range_key;
+----
+20 200
+25 250
+30 300
+35 350
+
+##########
+# TEST 21: Incompatible Range Right Join Repartitions
+# The split points of the two inputs differ, so the co-partitioned layout
+# requirement cannot be satisfied and Hash repartitioning repairs both sides
+# of the right join. Results stay correct on the repartitioned path.
+##########
+
+query TT
+EXPLAIN SELECT l.value, r.range_key, r.value
+FROM (SELECT range_key, value FROM range_partitioned WHERE value <= 150) l
+RIGHT JOIN range_partitioned_shifted r ON l.range_key = r.range_key;
+----
+physical_plan
+01)HashJoinExec: mode=Partitioned, join_type=Right, on=[(range_key@0,
range_key@0)], projection=[value@1, range_key@2, value@3]
+02)--RepartitionExec: partitioning=Hash([range_key@0], 4), input_partitions=4
+03)----FilterExec: value@1 <= 150
+04)------DataSourceExec: file_groups=<slt:ignore>, projection=[range_key,
value], output_partitioning=Range([range_key@0 ASC], [(10), (20), (30)], 4),
file_type=csv, has_header=false
+05)--RepartitionExec: partitioning=Hash([range_key@0], 4), input_partitions=4
+06)----DataSourceExec: file_groups=<slt:ignore>, projection=[range_key,
value], output_partitioning=Range([range_key@0 ASC], [(15), (20), (30)], 4),
file_type=csv, has_header=false
+
+query III
+SELECT l.value, r.range_key, r.value
+FROM (SELECT range_key, value FROM range_partitioned WHERE value <= 150) l
+RIGHT JOIN range_partitioned_shifted r ON l.range_key = r.range_key
+ORDER BY r.range_key;
+----
+10 1 10
+50 5 50
+100 10 100
+150 15 150
+NULL 20 200
+NULL 25 250
+NULL 30 300
+NULL 35 350
+
+##########
+# TEST 22: Right Join on Composite Key With Range Subset Rehashes
Review Comment:
please resplace with something like:
```
# TEST 22: Composite-Key Right Join Repartitions
# Range([range_key]) does not satisfy a partitioned join on
# (range_key, non_range_key), so both sides repartition on the full key.
```
##########
datafusion/core/tests/physical_optimizer/enforce_distribution.rs:
##########
@@ -870,6 +870,91 @@ fn
range_inner_hash_join_rehashes_incompatible_range_partitioning() -> Result<()
Ok(())
}
+// Kept as a unit test: `RightMark` join plans cannot be produced from SQL
(they
+// only arise when a statistical swap flips a `LeftMark` join around), so this
+// reuse plan shape has no `range_partitioning.slt` equivalent. The end-to-end
+// mark-join coverage lives in that file's "Mark Join Marker Semantics" case.
+#[test]
+fn range_right_mark_hash_join_reuses_range_partitioning() -> Result<()> {
+ let left = parquet_exec_with_output_partitioning(range_partitioning(
+ "a",
+ [10, 20, 30],
+ SortOptions::default(),
+ )?);
+ let right = parquet_exec_with_output_partitioning(range_partitioning(
+ "a",
+ [10, 20, 30],
+ SortOptions::default(),
+ )?);
+ let join_on = vec![(
+ Arc::new(Column::new_with_schema("a", &left.schema())?) as _,
+ Arc::new(Column::new_with_schema("a", &right.schema())?) as _,
+ )];
+ let join = hash_join_exec(left, right, &join_on, &JoinType::RightMark);
+
+ let plan = TestConfig::default()
+ .with_query_execution_partitions(4)
+ .to_plan(join, &DISTRIB_DISTRIB_SORT);
+
+ assert_plan!(
+ plan,
+ @r"
+ HashJoinExec: mode=Partitioned, join_type=RightMark, on=[(a@0, a@0)]
+ 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
+ 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
+ "
+ );
+
+ Ok(())
+}
+
+// Kept as a unit test: a descending Range input cannot be registered through
the
+// sqllogictest fixtures (they only declare ascending Range layouts), so the
+// sort-direction axis of the co-partition check has no
`range_partitioning.slt`
+// equivalent. Verifies that opposite sort options force Hash repartitioning.
+#[test]
+fn range_right_semi_hash_join_rehashes_incompatible_sort_options() ->
Result<()> {
+ let left = parquet_exec_with_output_partitioning(range_partitioning(
+ "a",
+ [10, 20, 30],
+ SortOptions::default(),
+ )?);
+ // A descending Range requires descending-ordered split points (the
constructor rejects
+ // ascending points under a descending sort), so [30, 20, 10] DESC encodes
the same
+ // partition boundaries as the left's [10, 20, 30] ASC -- only the sort
direction differs,
+ // which is the incompatibility under test.
Review Comment:
ok I think we got confused here. We are specifically tring to isolate
incompatible sort options based on metadata. To isolate this is would be best
to use a single split point and just change the descending metadata rather than
both 👍
##########
datafusion/core/tests/physical_optimizer/enforce_distribution.rs:
##########
@@ -870,6 +870,91 @@ fn
range_inner_hash_join_rehashes_incompatible_range_partitioning() -> Result<()
Ok(())
}
+// Kept as a unit test: `RightMark` join plans cannot be produced from SQL
(they
+// only arise when a statistical swap flips a `LeftMark` join around), so this
+// reuse plan shape has no `range_partitioning.slt` equivalent. The end-to-end
+// mark-join coverage lives in that file's "Mark Join Marker Semantics" case.
+#[test]
+fn range_right_mark_hash_join_reuses_range_partitioning() -> Result<()> {
+ let left = parquet_exec_with_output_partitioning(range_partitioning(
+ "a",
+ [10, 20, 30],
+ SortOptions::default(),
+ )?);
+ let right = parquet_exec_with_output_partitioning(range_partitioning(
+ "a",
+ [10, 20, 30],
+ SortOptions::default(),
+ )?);
+ let join_on = vec![(
+ Arc::new(Column::new_with_schema("a", &left.schema())?) as _,
+ Arc::new(Column::new_with_schema("a", &right.schema())?) as _,
+ )];
+ let join = hash_join_exec(left, right, &join_on, &JoinType::RightMark);
+
+ let plan = TestConfig::default()
+ .with_query_execution_partitions(4)
+ .to_plan(join, &DISTRIB_DISTRIB_SORT);
+
+ assert_plan!(
+ plan,
+ @r"
+ HashJoinExec: mode=Partitioned, join_type=RightMark, on=[(a@0, a@0)]
+ 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
+ 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
+ "
+ );
+
+ Ok(())
+}
+
+// Kept as a unit test: a descending Range input cannot be registered through
the
Review Comment:
ditto
##########
datafusion/core/tests/physical_optimizer/enforce_distribution.rs:
##########
@@ -870,6 +870,91 @@ fn
range_inner_hash_join_rehashes_incompatible_range_partitioning() -> Result<()
Ok(())
}
+// Kept as a unit test: `RightMark` join plans cannot be produced from SQL
(they
+// only arise when a statistical swap flips a `LeftMark` join around), so this
+// reuse plan shape has no `range_partitioning.slt` equivalent. The end-to-end
+// mark-join coverage lives in that file's "Mark Join Marker Semantics" case.
+#[test]
+fn range_right_mark_hash_join_reuses_range_partitioning() -> Result<()> {
+ let left = parquet_exec_with_output_partitioning(range_partitioning(
+ "a",
+ [10, 20, 30],
+ SortOptions::default(),
+ )?);
+ let right = parquet_exec_with_output_partitioning(range_partitioning(
+ "a",
+ [10, 20, 30],
+ SortOptions::default(),
+ )?);
+ let join_on = vec![(
+ Arc::new(Column::new_with_schema("a", &left.schema())?) as _,
+ Arc::new(Column::new_with_schema("a", &right.schema())?) as _,
+ )];
+ let join = hash_join_exec(left, right, &join_on, &JoinType::RightMark);
+
+ let plan = TestConfig::default()
+ .with_query_execution_partitions(4)
+ .to_plan(join, &DISTRIB_DISTRIB_SORT);
+
+ assert_plan!(
+ plan,
+ @r"
+ HashJoinExec: mode=Partitioned, join_type=RightMark, on=[(a@0, a@0)]
+ 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
+ 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
+ "
+ );
+
+ Ok(())
+}
+
+// Kept as a unit test: a descending Range input cannot be registered through
the
+// sqllogictest fixtures (they only declare ascending Range layouts), so the
+// sort-direction axis of the co-partition check has no
`range_partitioning.slt`
+// equivalent. Verifies that opposite sort options force Hash repartitioning.
+#[test]
+fn range_right_semi_hash_join_rehashes_incompatible_sort_options() ->
Result<()> {
+ let left = parquet_exec_with_output_partitioning(range_partitioning(
+ "a",
+ [10, 20, 30],
+ SortOptions::default(),
+ )?);
+ // A descending Range requires descending-ordered split points (the
constructor rejects
+ // ascending points under a descending sort), so [30, 20, 10] DESC encodes
the same
+ // partition boundaries as the left's [10, 20, 30] ASC -- only the sort
direction differs,
+ // which is the incompatibility under test.
Review Comment:
And I don't think this warrants a long comments
##########
datafusion/core/tests/physical_optimizer/enforce_distribution.rs:
##########
@@ -870,6 +870,91 @@ fn
range_inner_hash_join_rehashes_incompatible_range_partitioning() -> Result<()
Ok(())
}
+// Kept as a unit test: `RightMark` join plans cannot be produced from SQL
(they
+// only arise when a statistical swap flips a `LeftMark` join around), so this
Review Comment:
Lets try to keep very PR specific comments out of the code.
Usually phrases like this "Kept as a unit test: ..." does not have any value
outside of the scopoe of this PR becuase a new reaader will not have the
ocntext that you or I have about this test. Rather lets omit the comments and
allow the code to do the talking 👍
##########
datafusion/core/tests/physical_optimizer/enforce_distribution.rs:
##########
@@ -870,6 +870,91 @@ fn
range_inner_hash_join_rehashes_incompatible_range_partitioning() -> Result<()
Ok(())
}
+// Kept as a unit test: `RightMark` join plans cannot be produced from SQL
(they
+// only arise when a statistical swap flips a `LeftMark` join around), so this
+// reuse plan shape has no `range_partitioning.slt` equivalent. The end-to-end
+// mark-join coverage lives in that file's "Mark Join Marker Semantics" case.
+#[test]
+fn range_right_mark_hash_join_reuses_range_partitioning() -> Result<()> {
+ let left = parquet_exec_with_output_partitioning(range_partitioning(
+ "a",
+ [10, 20, 30],
+ SortOptions::default(),
+ )?);
+ let right = parquet_exec_with_output_partitioning(range_partitioning(
+ "a",
+ [10, 20, 30],
+ SortOptions::default(),
+ )?);
+ let join_on = vec![(
+ Arc::new(Column::new_with_schema("a", &left.schema())?) as _,
+ Arc::new(Column::new_with_schema("a", &right.schema())?) as _,
+ )];
+ let join = hash_join_exec(left, right, &join_on, &JoinType::RightMark);
+
+ let plan = TestConfig::default()
+ .with_query_execution_partitions(4)
+ .to_plan(join, &DISTRIB_DISTRIB_SORT);
+
+ assert_plan!(
+ plan,
+ @r"
+ HashJoinExec: mode=Partitioned, join_type=RightMark, on=[(a@0, a@0)]
+ 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
+ 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
+ "
+ );
+
+ Ok(())
+}
+
+// Kept as a unit test: a descending Range input cannot be registered through
the
+// sqllogictest fixtures (they only declare ascending Range layouts), so the
+// sort-direction axis of the co-partition check has no
`range_partitioning.slt`
+// equivalent. Verifies that opposite sort options force Hash repartitioning.
+#[test]
+fn range_right_semi_hash_join_rehashes_incompatible_sort_options() ->
Result<()> {
+ let left = parquet_exec_with_output_partitioning(range_partitioning(
+ "a",
+ [10, 20, 30],
+ SortOptions::default(),
+ )?);
+ // A descending Range requires descending-ordered split points (the
constructor rejects
+ // ascending points under a descending sort), so [30, 20, 10] DESC encodes
the same
+ // partition boundaries as the left's [10, 20, 30] ASC -- only the sort
direction differs,
+ // which is the incompatibility under test.
Review Comment:
this test would still fail if the sort otpions check was removed
##########
datafusion/sqllogictest/test_files/range_partitioning.slt:
##########
@@ -629,7 +630,304 @@ ORDER BY l.range_key;
35 700
##########
-# TEST 18: Sort Merge Join Avoids Repartition for Compatible Range Inputs
+# TEST 18: Right Join on Range Partition Column
+# Right-side partitioned hash joins also opt in to Range satisfying
+# KeyPartitioned requirements. Compatible Range layouts satisfy both the
+# per-child key requirements and the cross-child layout requirement, so no
+# Hash repartitioning is inserted. The filter on the left input keeps its
+# Range partitioning and leaves the right rows above 150 unmatched.
+##########
+
+query TT
+EXPLAIN SELECT l.value, r.range_key, r.value
+FROM (SELECT range_key, value FROM range_partitioned WHERE value <= 150) l
+RIGHT JOIN range_partitioned r ON l.range_key = r.range_key;
+----
+physical_plan
+01)HashJoinExec: mode=Partitioned, join_type=Right, on=[(range_key@0,
range_key@0)], projection=[value@1, range_key@2, value@3]
+02)--FilterExec: value@1 <= 150
+03)----DataSourceExec: file_groups={4 groups:
[[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-0.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-1.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-2.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-3.csv]]},
projection=[range_key, value], output_partitioning=Range([range_key@0 ASC],
[(10), (20), (30)], 4), file_type=csv, has_header=false
+04)--DataSourceExec: file_groups={4 groups:
[[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-0.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-1.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-2.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-3.csv]]},
projection=[range_key, value], output_partitioning=Range([range_key@0 ASC],
[(10), (20), (30)], 4), file_type=csv, has_header=false
+
+query III
+SELECT l.value, r.range_key, r.value
+FROM (SELECT range_key, value FROM range_partitioned WHERE value <= 150) l
+RIGHT JOIN range_partitioned r ON l.range_key = r.range_key
+ORDER BY r.range_key;
+----
+10 1 10
+50 5 50
+100 10 100
+150 15 150
+NULL 20 200
+NULL 25 250
+NULL 30 300
+NULL 35 350
+
+##########
+# TEST 19: Right Semi Join on Range Partition Column
+# Compatible Range inputs avoid Hash repartitioning for RightSemi joins.
+# Only right rows with a match on the filtered left side are returned.
+##########
+
+query TT
+EXPLAIN SELECT r.range_key, r.value
+FROM (SELECT range_key FROM range_partitioned WHERE value <= 150) l
+RIGHT SEMI JOIN range_partitioned r ON l.range_key = r.range_key;
+----
+physical_plan
+01)HashJoinExec: mode=Partitioned, join_type=RightSemi, on=[(range_key@0,
range_key@0)]
+02)--FilterExec: value@1 <= 150, projection=[range_key@0]
+03)----DataSourceExec: file_groups={4 groups:
[[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-0.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-1.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-2.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-3.csv]]},
projection=[range_key, value], output_partitioning=Range([range_key@0 ASC],
[(10), (20), (30)], 4), file_type=csv, has_header=false
+04)--DataSourceExec: file_groups={4 groups:
[[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-0.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-1.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-2.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-3.csv]]},
projection=[range_key, value], output_partitioning=Range([range_key@0 ASC],
[(10), (20), (30)], 4), file_type=csv, has_header=false
+
+query II
+SELECT r.range_key, r.value
+FROM (SELECT range_key FROM range_partitioned WHERE value <= 150) l
+RIGHT SEMI JOIN range_partitioned r ON l.range_key = r.range_key
+ORDER BY r.range_key;
+----
+1 10
+5 50
+10 100
+15 150
+
+##########
+# TEST 20: Right Anti Join on Range Partition Column
+# Compatible Range inputs avoid Hash repartitioning for RightAnti joins.
+# Only right rows without a match on the filtered left side are returned.
+##########
+
+query TT
+EXPLAIN SELECT r.range_key, r.value
+FROM (SELECT range_key FROM range_partitioned WHERE value <= 150) l
+RIGHT ANTI JOIN range_partitioned r ON l.range_key = r.range_key;
+----
+physical_plan
+01)HashJoinExec: mode=Partitioned, join_type=RightAnti, on=[(range_key@0,
range_key@0)]
+02)--FilterExec: value@1 <= 150, projection=[range_key@0]
+03)----DataSourceExec: file_groups={4 groups:
[[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-0.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-1.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-2.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-3.csv]]},
projection=[range_key, value], output_partitioning=Range([range_key@0 ASC],
[(10), (20), (30)], 4), file_type=csv, has_header=false
+04)--DataSourceExec: file_groups={4 groups:
[[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-0.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-1.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-2.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-3.csv]]},
projection=[range_key, value], output_partitioning=Range([range_key@0 ASC],
[(10), (20), (30)], 4), file_type=csv, has_header=false
+
+query II
+SELECT r.range_key, r.value
+FROM (SELECT range_key FROM range_partitioned WHERE value <= 150) l
+RIGHT ANTI JOIN range_partitioned r ON l.range_key = r.range_key
+ORDER BY r.range_key;
+----
+20 200
+25 250
+30 300
+35 350
+
+##########
+# TEST 21: Incompatible Range Right Join Repartitions
+# The split points of the two inputs differ, so the co-partitioned layout
+# requirement cannot be satisfied and Hash repartitioning repairs both sides
+# of the right join. Results stay correct on the repartitioned path.
+##########
+
+query TT
+EXPLAIN SELECT l.value, r.range_key, r.value
+FROM (SELECT range_key, value FROM range_partitioned WHERE value <= 150) l
+RIGHT JOIN range_partitioned_shifted r ON l.range_key = r.range_key;
+----
+physical_plan
+01)HashJoinExec: mode=Partitioned, join_type=Right, on=[(range_key@0,
range_key@0)], projection=[value@1, range_key@2, value@3]
+02)--RepartitionExec: partitioning=Hash([range_key@0], 4), input_partitions=4
+03)----FilterExec: value@1 <= 150
+04)------DataSourceExec: file_groups={4 groups:
[[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-0.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-1.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-2.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-3.csv]]},
projection=[range_key, value], output_partitioning=Range([range_key@0 ASC],
[(10), (20), (30)], 4), file_type=csv, has_header=false
+05)--RepartitionExec: partitioning=Hash([range_key@0], 4), input_partitions=4
+06)----DataSourceExec: file_groups={4 groups:
[[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned_shifted/part-0.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned_shifted/part-1.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned_shifted/part-2.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned_shifted/part-3.csv]]},
projection=[range_key, value], output_partitioning=Range([range_key@0 ASC],
[(15), (20), (30)], 4), file_type=csv, has_header=false
+
+query III
+SELECT l.value, r.range_key, r.value
+FROM (SELECT range_key, value FROM range_partitioned WHERE value <= 150) l
+RIGHT JOIN range_partitioned_shifted r ON l.range_key = r.range_key
+ORDER BY r.range_key;
+----
+10 1 10
+50 5 50
+100 10 100
+150 15 150
+NULL 20 200
+NULL 25 250
+NULL 30 300
+NULL 35 350
+
+##########
+# TEST 22: Right Join on Composite Key With Range Subset Rehashes
+# The join key (range_key, non_range_key) is a strict superset of the exposed
+# Range([range_key]) partitioning column. Unlike aggregates (TEST 3), subset
+# satisfaction is disabled for partitioned joins -- allowing Range([range_key])
+# to satisfy a wider key requirement could let matching rows land in different
+# partitions -- so planning still Hash repartitions both sides on the full
+# join key, even though subset_repartition_threshold is met.
+##########
+
+statement ok
+set datafusion.optimizer.subset_repartition_threshold = 4;
+
+query TT
+EXPLAIN SELECT l.range_key, l.non_range_key, l.value, r.value
+FROM range_partitioned l
+RIGHT JOIN range_partitioned r ON l.range_key = r.range_key AND
l.non_range_key = r.non_range_key;
+----
+physical_plan
+01)HashJoinExec: mode=Partitioned, join_type=Right, on=[(range_key@0,
range_key@0), (non_range_key@1, non_range_key@1)], projection=[range_key@0,
non_range_key@1, value@2, value@5]
+02)--RepartitionExec: partitioning=Hash([range_key@0, non_range_key@1], 4),
input_partitions=4
+03)----DataSourceExec: file_groups={4 groups:
[[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-0.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-1.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-2.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-3.csv]]},
projection=[range_key, non_range_key, value],
output_partitioning=Range([range_key@0 ASC], [(10), (20), (30)], 4),
file_type=csv, has_header=false
+04)--RepartitionExec: partitioning=Hash([range_key@0, non_range_key@1], 4),
input_partitions=4
+05)----DataSourceExec: file_groups={4 groups:
[[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-0.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-1.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-2.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-3.csv]]},
projection=[range_key, non_range_key, value],
output_partitioning=Range([range_key@0 ASC], [(10), (20), (30)], 4),
file_type=csv, has_header=false
+
+query IIII
+SELECT l.range_key, l.non_range_key, l.value, r.value
+FROM range_partitioned l
+RIGHT JOIN range_partitioned r ON l.range_key = r.range_key AND
l.non_range_key = r.non_range_key
+ORDER BY l.range_key;
+----
+1 1 10 10
+5 2 50 50
+10 1 100 100
+15 2 150 150
+20 1 200 200
+25 2 250 250
+30 1 300 300
+35 2 350 350
+
+statement ok
+reset datafusion.optimizer.subset_repartition_threshold;
+
+##########
+# TEST 23: Right Join with Mismatched Range Partition Counts Repartitions
+# Both inputs are range partitioned on range_key, but declare a different
number
+# of partitions (four vs three). The per-child key requirements can be
satisfied
+# by Range, but the co-partitioned layout requirement cannot, so Hash
+# repartitioning repairs both sides of the right join.
+##########
+
+query TT
+EXPLAIN SELECT l.value, r.range_key, r.value
+FROM range_partitioned l
+RIGHT JOIN range_partitioned_narrow r ON l.range_key = r.range_key;
+----
+physical_plan
+01)HashJoinExec: mode=Partitioned, join_type=Right, on=[(range_key@0,
range_key@0)], projection=[value@1, range_key@2, value@3]
+02)--RepartitionExec: partitioning=Hash([range_key@0], 4), input_partitions=4
+03)----DataSourceExec: file_groups={4 groups:
[[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-0.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-1.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-2.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-3.csv]]},
projection=[range_key, value], output_partitioning=Range([range_key@0 ASC],
[(10), (20), (30)], 4), file_type=csv, has_header=false
+04)--RepartitionExec: partitioning=Hash([range_key@0], 4), input_partitions=3
+05)----DataSourceExec: file_groups={3 groups:
[[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned_narrow/part-0.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned_narrow/part-1.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned_narrow/part-2.csv]]},
projection=[range_key, value], output_partitioning=Range([range_key@0 ASC],
[(10), (20)], 3), file_type=csv, has_header=false
+
+query III
+SELECT l.value, r.range_key, r.value
+FROM range_partitioned l
+RIGHT JOIN range_partitioned_narrow r ON l.range_key = r.range_key
+ORDER BY r.range_key;
+----
+10 1 10
+50 5 50
+100 10 100
+150 15 150
+200 20 200
+250 25 250
+300 30 300
+350 35 350
+
+##########
+# TEST 24: Right Join on Non-Range Key Repartitions
+# Both inputs expose Range([range_key]), but the join key is non_range_key.
+# Range([range_key]) does not satisfy KeyPartitioned([non_range_key]), so
+# planning inserts Hash repartitioning on the actual join key for the right
join.
+##########
+
+query TT
+EXPLAIN SELECT l.value, r.range_key, r.value
+FROM (SELECT non_range_key, value FROM range_partitioned WHERE range_key < 10)
l
+RIGHT JOIN range_partitioned r ON l.non_range_key = r.non_range_key;
+----
+physical_plan
+01)HashJoinExec: mode=Partitioned, join_type=Right, on=[(non_range_key@0,
non_range_key@1)], projection=[value@1, range_key@2, value@4]
+02)--RepartitionExec: partitioning=Hash([non_range_key@0], 4),
input_partitions=4
+03)----FilterExec: range_key@0 < 10, projection=[non_range_key@1, value@2]
+04)------DataSourceExec: file_groups={4 groups:
[[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-0.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-1.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-2.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-3.csv]]},
projection=[range_key, non_range_key, value],
output_partitioning=Range([range_key@0 ASC], [(10), (20), (30)], 4),
file_type=csv, has_header=false
+05)--RepartitionExec: partitioning=Hash([non_range_key@1], 4),
input_partitions=4
+06)----DataSourceExec: file_groups={4 groups:
[[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-0.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-1.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-2.csv],
[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch_range_partitioning/range_partitioned/part-3.csv]]},
projection=[range_key, non_range_key, value],
output_partitioning=Range([range_key@0 ASC], [(10), (20), (30)], 4),
file_type=csv, has_header=false
+
+query III
+SELECT l.value, r.range_key, r.value
+FROM (SELECT non_range_key, value FROM range_partitioned WHERE range_key < 10)
l
+RIGHT JOIN range_partitioned r ON l.non_range_key = r.non_range_key
+ORDER BY r.range_key;
+----
+10 1 10
+50 5 50
+10 10 100
+50 15 150
+10 20 200
+50 25 250
+10 30 300
+50 35 350
+
+##########
+# TEST 25: Mark Join Marker Semantics over Range Partitioned Inputs
Review Comment:
please replace with something like:
```text
# TEST 25: Mark Join Marker Semantics
# Mark joins preserve matched, unmatched, and NULL-key marker behavior over
# range-partitioned inputs.
```
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]