sunchao commented on code in PR #4587:
URL: https://github.com/apache/datafusion-comet/pull/4587#discussion_r4131690071


##########
spark/src/test/scala/org/apache/comet/exec/CometJoinSuite.scala:
##########
@@ -1572,4 +1572,131 @@ class CometJoinSuite extends CometTestBase {
       }
     }
   }
+
+  test("ExistenceJoin via BroadcastHashJoin (EXISTS combined with OR)") {
+    withSQLConf(
+      CometConf.COMET_EXEC_EXISTENCE_JOIN_ENABLED.key -> "true",
+      SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "10MB",
+      SQLConf.ADAPTIVE_AUTO_BROADCASTJOIN_THRESHOLD.key -> "10MB") {
+      withParquetTable((0 until 100).map(i => (i, if (i % 3 == 0) "US" else 
"EU")), "tbl_a") {
+        withParquetTable((0 until 30).map(i => (i, i + 1)), "tbl_b") {
+          val df = sql("SELECT * FROM tbl_a a " +
+            "WHERE a._2 = 'US' OR EXISTS (SELECT /*+ BROADCAST(b) */ 1 FROM 
tbl_b b WHERE b._1 = a._1)")
+          checkSparkAnswerAndOperator(
+            df,
+            Seq(classOf[CometBroadcastExchangeExec], 
classOf[CometBroadcastHashJoinExec]))
+        }
+      }
+    }
+  }
+
+  test("ExistenceJoin via ShuffledHashJoin (EXISTS combined with OR)") {
+    withSQLConf(
+      CometConf.COMET_EXEC_EXISTENCE_JOIN_ENABLED.key -> "true",
+      SQLConf.PREFER_SORTMERGEJOIN.key -> "false",
+      "spark.sql.join.forceApplyShuffledHashJoin" -> "true",
+      SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1",
+      SQLConf.ADAPTIVE_AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1") {
+      withParquetTable((0 until 100).map(i => (i, if (i % 3 == 0) "US" else 
"EU")), "tbl_a") {
+        withParquetTable((0 until 30).map(i => (i, i + 1)), "tbl_b") {
+          val df = sql(
+            "SELECT * FROM tbl_a a " +
+              "WHERE a._2 = 'US' OR EXISTS (SELECT 1 FROM tbl_b b WHERE b._1 = 
a._1)")
+          checkSparkAnswerAndOperator(df, Seq(classOf[CometHashJoinExec]))
+        }
+      }
+    }
+  }
+
+  test("ExistenceJoin via SortMergeJoin falls back to Spark") {
+    // Existence sort-merge joins are not executed natively; verify the 
fallback.
+    withSQLConf(
+      CometConf.COMET_EXEC_EXISTENCE_JOIN_ENABLED.key -> "true",
+      SQLConf.PREFER_SORTMERGEJOIN.key -> "true",
+      SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1",
+      SQLConf.ADAPTIVE_AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1") {
+      withParquetTable((0 until 100).map(i => (i, if (i % 3 == 0) "US" else 
"EU")), "tbl_a") {
+        withParquetTable((0 until 30).map(i => (i, i + 1)), "tbl_b") {
+          val df = sql(
+            "SELECT * FROM tbl_a a " +
+              "WHERE a._2 = 'US' OR EXISTS (SELECT 1 FROM tbl_b b WHERE b._1 = 
a._1)")
+          checkSparkAnswerAndFallbackReason(df, "Unsupported join type")
+        }
+      }
+    }
+  }
+
+  test("ExistenceJoin with residual condition falls back to Spark") {
+    // A non-equi residual predicate is evaluated over the whole candidate 
batch by DataFusion's
+    // LeftMark join (no first-match short-circuit), so Comet keeps it on 
Spark; verify parity.
+    withSQLConf(
+      CometConf.COMET_EXEC_EXISTENCE_JOIN_ENABLED.key -> "true",
+      SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "10MB",
+      SQLConf.ADAPTIVE_AUTO_BROADCASTJOIN_THRESHOLD.key -> "10MB") {
+      withParquetTable((0 until 100).map(i => (i, if (i % 3 == 0) "US" else 
"EU")), "tbl_a") {
+        withParquetTable((0 until 30).map(i => (i, i + 1)), "tbl_b") {
+          val df = sql(
+            "SELECT * FROM tbl_a a WHERE a._2 = 'US' OR EXISTS " +
+              "(SELECT /*+ BROADCAST(b) */ 1 FROM tbl_b b WHERE b._1 = a._1 
AND b._2 > a._1)")
+          checkSparkAnswerAndFallbackReason(df, "residual (non-equi)")

Review Comment:
   [P2] Make the new fallback-reason assertions pass with the default planning 
configuration. This residual-condition test and the computed-key test at line 
1659 both fail in Spark 4.1 CI: results match Spark and the final plans contain 
Spark `BroadcastHashJoin`, but `checkSparkAnswerAndFallbackReason` finds no 
fallback reasons. These assertions add two execution-suite failures. Preserve 
the guard-specific reasons on the final plan, or test the reasons at a stable 
planning phase while retaining an assertion that execution falls back to Spark.
   
   Evidence: Head-associated CI job 
https://github.com/apache/datafusion-comet/actions/runs/36515604951/job/109239141889
 reports `Expected fallback reason 'residual (non-equi)' but no fallback 
reasons were found` and the equivalent failure for `computed (non-column)`. 
`CometTestBase.checkSparkAnswerAndFallbackReasons` checks the executed 
Comet-enabled plan after result comparison. The relevant guards and assertions 
are unchanged between the requested head and the CI merge revision. With native 
build prerequisites available, isolate these tests using `./mvnw test 
-Pspark-4.1 -Dtest=none -Dsuites="org.apache.comet.exec.CometJoinSuite 
ExistenceJoin"`.



-- 
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]

Reply via email to