Skip to content

Commit 099c3ee

Browse files
naveenp2708szehon-ho
authored andcommitted
[SPARK-59688][SQL][FOLLOWUP] Cover a left outer join in the coinciding-partition-keys test
**What changes were proposed in this pull request?** A follow-up test for SPARK-59688. It parameterizes the existing end-to-end test over an inner and a LEFT OUTER join, on the same identity(id) vs flip_low_bit(id) setup. Test only, no production change. **Why are the changes needed?** The end-to-end test in SPARK-59688 covers an inner join, where the wrong result is missing rows (0 of 2). An outer join hits the same coinciding-keys bug with a different symptom: it keeps the left rows and null-extends them, returning (0, 'a', null) instead of (0, 'a', 'x'). Every left row here has a match, so both join types share the expected rows, the plan assertions, and the test body. The assertions fail on a revert of the fix and pass with it. **Does this PR introduce any user-facing change?** No. Test only. **How was this patch tested?** The parameterized test. It fails on a revert of the partitioning.scala fix and passes with it. KeyGroupedPartitioningSuite runs green (195 tests). **Was this patch authored or co-authored using generative AI tooling?** No Closes #59129 from naveenp2708/SPARK-59688-outer-join-test. Authored-by: naveenp2708 <naveenp2708@gmail.com> Signed-off-by: Szehon Ho <szehon.apache@gmail.com>
1 parent 25a9aa7 commit 099c3ee

1 file changed

Lines changed: 25 additions & 16 deletions

File tree

‎sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala‎

Lines changed: 25 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -1410,30 +1410,39 @@ class KeyGroupedPartitioningSuite
14101410
// report the same partition key list, [0, 1] of LongType, while the rows behind a key differ:
14111411
// the identity side's key 0 holds id 0, the transform side's holds id 1. The identity side's
14121412
// raw keys have to be reduced onto `flip_low_bit` before the partitions can be paired up.
1413+
// Read as they stand, an inner join drops the mispaired rows (0 of 2) and a left outer join
1414+
// null-extends them; the reduce fixes both. Every left row here has a match, so the two join
1415+
// types expect the same rows.
14131416
val cols = Array(Column.create("id", LongType), Column.create("data", StringType))
14141417
createTable("t1", cols, Array(identity("id")))
14151418
sql("INSERT INTO testcat.ns.t1 VALUES (0, 'a'), (1, 'b')")
14161419

14171420
createTable("t2", cols, Array(Expressions.apply("flip_low_bit", Expressions.column("id"))))
14181421
sql("INSERT INTO testcat.ns.t2 VALUES (0, 'x'), (1, 'y')")
14191422

1420-
val df = sql(
1421-
"SELECT t1.id, t1.data, t2.data FROM testcat.ns.t1 JOIN testcat.ns.t2 ON t1.id = t2.id")
1423+
val expected = Seq(Row(0L, "a", "x"), Row(1L, "b", "y"))
1424+
Seq("JOIN", "LEFT OUTER JOIN").foreach { joinType =>
1425+
val df = sql(
1426+
s"SELECT t1.id, t1.data, t2.data FROM testcat.ns.t1 $joinType testcat.ns.t2 " +
1427+
"ON t1.id = t2.id")
14221428

1423-
withSQLConf(SQLConf.V2_BUCKETING_ALLOW_COMPATIBLE_TRANSFORMS.key -> "true") {
1424-
checkAnswer(df, Seq(Row(0L, "a", "x"), Row(1L, "b", "y")))
1425-
val plan = stripAQEPlan(df.queryExecution.executedPlan)
1426-
assert(collectShuffles(plan).isEmpty, "storage-partitioned join should not shuffle")
1427-
val groupPartitions = collectGroupPartitions(plan)
1428-
assert(groupPartitions.size == 2,
1429-
"both sides should be regrouped onto the merged partition keys")
1430-
assert(groupPartitions.count(_.reducers.exists(_.exists(_.isDefined))) == 1,
1431-
"and exactly one side should reduce, the identity one onto `flip_low_bit`")
1432-
// The join subtree, not the whole plan: `ValidateRequirements` walks children and a query
1433-
// stage is a leaf, so validating an AQE plan checks nothing.
1434-
val joins = collect(plan) { case smj: SortMergeJoinExec => smj }
1435-
assert(joins.size == 1, s"test setup: one join to validate:\n$plan")
1436-
assert(ValidateRequirements.validate(joins.head), "the plan that leaves must hold up")
1429+
withSQLConf(SQLConf.V2_BUCKETING_ALLOW_COMPATIBLE_TRANSFORMS.key -> "true") {
1430+
checkAnswer(df, expected)
1431+
val plan = stripAQEPlan(df.queryExecution.executedPlan)
1432+
assert(collectShuffles(plan).isEmpty,
1433+
s"$joinType: storage-partitioned join should not shuffle")
1434+
val groupPartitions = collectGroupPartitions(plan)
1435+
assert(groupPartitions.size == 2,
1436+
s"$joinType: both sides should be regrouped onto the merged partition keys")
1437+
assert(groupPartitions.count(_.reducers.exists(_.exists(_.isDefined))) == 1,
1438+
s"$joinType: exactly one side should reduce, the identity one onto `flip_low_bit`")
1439+
// The join subtree, not the whole plan: `ValidateRequirements` walks children and a
1440+
// query stage is a leaf, so validating an AQE plan checks nothing.
1441+
val joins = collect(plan) { case smj: SortMergeJoinExec => smj }
1442+
assert(joins.size == 1, s"$joinType: test setup: one join to validate:\n$plan")
1443+
assert(ValidateRequirements.validate(joins.head),
1444+
s"$joinType: the plan that leaves must hold up")
1445+
}
14371446
}
14381447
}
14391448
}

0 commit comments

Comments
 (0)