gmhelmold commented on code in PR #23484:
URL: https://github.com/apache/datafusion/pull/23484#discussion_r3598404723
##########
datafusion/sqllogictest/test_files/range_partitioning.slt:
##########
@@ -634,6 +635,260 @@ ORDER BY l.range_key;
30 600
35 700
+##########
+# 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: Mark Join Marker Semantics over Range Partitioned Inputs
Review Comment:
Done — moved to the end of the partitioned-hash-join block (now the last
test there, before the SMJ/SHJ cases that came in with the main merge).
##########
datafusion/sqllogictest/test_files/range_partitioning.slt:
##########
@@ -634,6 +635,260 @@ ORDER BY l.range_key;
30 600
35 700
+##########
+# 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: Mark Join Marker Semantics over Range Partitioned Inputs
+# SQL IN subqueries decorrelate to LeftMark joins, which stay on the Hash
+# repartition path (RightMark only arises when a statistical swap turns a
+# LeftMark join around; RightMark plan coverage lives in the
+# enforce_distribution tests). These queries pin the marker semantics --
+# matched, unmatched, and NULL join keys -- for mark joins over range
+# partitioned inputs.
+##########
+
+query TT
+EXPLAIN SELECT r.range_key, r.value
+FROM range_partitioned r
+WHERE r.non_range_key = 2 OR r.range_key IN (
+ SELECT range_key FROM range_partitioned WHERE value <= 150);
+----
+physical_plan
+01)FilterExec: non_range_key@1 = 2 OR mark@3, projection=[range_key@0, value@2]
+02)--HashJoinExec: mode=Partitioned, join_type=LeftMark, on=[(range_key@0,
range_key@0)]
+03)----RepartitionExec: partitioning=Hash([range_key@0], 4), input_partitions=4
Review Comment:
Good question — no, this one stays. This PR is the implementation of #23453
(it closes it), and the plan pinned here is already generated with the change
applied: the range-satisfaction gate in
`HashJoinExec::input_distribution_requirements` covers `Inner | Right |
RightSemi | RightAnti | RightMark` only, and SQL `IN`-subquery decorrelation
produces a `LeftMark` join — a left-anchored type intentionally outside this
PR's covered set — so its inputs keep the Hash repartition path. And exactly
right on the rationale: since `RightMark` is not reachable from SQL (it only
appears when a statistics-based swap flips a `LeftMark`), its reuse behavior is
pinned by the `enforce_distribution.rs` unit test, while this case pins the
mark-marker semantics end-to-end on the SQL-reachable path. It would only
change if co-partition reuse is later extended to the left-anchored variants
under the #22395 umbrella — happy to file that follow-up if there's interest.
##########
datafusion/sqllogictest/test_files/range_partitioning.slt:
##########
@@ -634,6 +635,260 @@ ORDER BY l.range_key;
30 600
35 700
+##########
Review Comment:
Added TEST 22: a RIGHT JOIN on the composite key `(range_key,
non_range_key)` where `Range([range_key])` covers only a subset — pins that,
unlike aggregates (TEST 3), subset satisfaction stays disabled for partitioned
joins, so both sides Hash-repartition on the full key even with the threshold
met.
--
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]