comphead commented on code in PR #4587:
URL: https://github.com/apache/datafusion-comet/pull/4587#discussion_r4157863390
##########
spark/src/main/scala/org/apache/spark/sql/comet/operators.scala:
##########
@@ -2574,6 +2574,32 @@ trait CometHashJoin {
case FullOuter => JoinType.FullOuter
case LeftSemi => JoinType.LeftSemi
case LeftAnti => JoinType.LeftAnti
+ case ExistenceJoin(_) if
CometConf.COMET_EXEC_EXISTENCE_JOIN_ENABLED.get(join.conf) =>
+ // Native only for equi-key joins with bare column/literal keys and
no residual.
+ if (join.condition.isDefined) {
+ withFallbackReason(
+ join,
+ "ExistenceJoin with a residual (non-equi) condition is not
supported natively")
+ return None
+ } else if (!(join.leftKeys ++ join.rightKeys).forall {
+ case _: Attribute | _: Literal => true
+ case _ => false
+ }) {
+ withFallbackReason(
+ join,
+ "ExistenceJoin with a computed (non-column) join key is not
supported natively")
+ return None
+ } else {
+ JoinType.Existence
+ }
+ case ExistenceJoin(_) =>
+ // The guard above matched only when the flag is enabled; reaching
here means it is off.
+ // Report a toggle-specific reason so the plan does not read like a
permanent limitation.
+ withFallbackReason(
+ join,
+ "Native ExistenceJoin is disabled; set " +
+ s"${CometConf.COMET_EXEC_EXISTENCE_JOIN_ENABLED.key}=true to
enable it")
+ return None
Review Comment:
These guards would read better as flat early exits in the pre-flight block
above (after the null-aware check, before `val condition`), like the collation
and null-aware guards:
```scala
if (join.joinType.isInstanceOf[ExistenceJoin]) {
if (!CometConf.COMET_EXEC_EXISTENCE_JOIN_ENABLED.get(join.conf)) {
withFallbackReason(join, "<current flag-off reason>")
return None
}
// same shape for join.condition.isDefined and
!joinKeys.forall(_.isInstanceOf[Attribute])
}
```
The match then keeps a single `case ExistenceJoin(_) => JoinType.Existence`.
Each guard keeps its own reason, so the per-guard tests should be unaffected
(expectation, not run). This also removes:
- the second `case ExistenceJoin(_)` arm, which only works because of arm
order, and the comment that explains that order
- the `join.leftKeys ++ join.rightKeys` rebuild, since `joinKeys` is already
in scope
- the `exprToProto` call on the residual that is then thrown away. If the
residual holds an unsupported expression, the early `return None` fires first
and the `residual (non-equi)` reason is never recorded.
`CometSortMergeJoinExec.convert` already checks its join-filter flag before
`val condition`.
Two smaller points:
- I expect `Literal` cannot reach this check. `ExtractEquiJoinKeys` skips
equalities with a reference-free side (`l.references.isEmpty ||
r.references.isEmpty`, same in Spark 3.4.3, 3.5.5 and 4.0.2), so
`_.isInstanceOf[Attribute]` would be enough and the new `Literal` import could
go. This is from reading the Spark source, not run.
- The reason for the guards (batch-level evaluation is eager while Spark
stops at the first match) only lives in test comments. One sentence here would
be more useful than `Native only for equi-key joins ...`.
##########
spark/src/test/scala/org/apache/comet/exec/CometJoinSuite.scala:
##########
@@ -1574,6 +1574,186 @@ 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") {
Review Comment:
Three of these tests re-run queries that `existence_join.sql` already has,
with the same assertion: `ExistenceJoin via SortMergeJoin falls back to Spark`,
`ExistenceJoin with duplicate build keys runs natively (markers not
multiplied)` and `ExistenceJoin with NOT EXISTS combined with OR runs natively`.
- The fixture covers duplicates. `ex_right` holds key 2 twice and `ex_left`
id 2 has region `EU`, so a multiplied marker would add a row, and
`ex_right_dups` repeats key 1 under `SHUFFLE_HASH`.
- `NOT` is applied by a Filter above the same `ExistenceJoin`, and the
duplicate-key and `NOT EXISTS` tests carry the same class pin as the first test.
- The SMJ fallback is decided by the join type alone, and the fixture's
`MERGE` queries assert the same `Unsupported join type` reason.
These three can go. The first two tests are worth keeping because they pin
the join class, which a `.sql` query cannot do.
The residual, computed-key and flag-off tests share the same tables and
AQE-off block and differ only in the correlation predicate (and the flag). They
could be one test with a single #6442 comment, or `.sql` fixtures that set `--
Config: spark.sql.adaptive.enabled=false` and use `query expect_fallback(...)`.
The flag-off case would need its own fixture because `-- Config:` is file-wide.
##########
spark/src/main/scala/org/apache/comet/rules/RewriteJoin.scala:
##########
@@ -67,8 +67,10 @@ object RewriteJoin extends JoinSelectionHelper {
def rewrite(plan: SparkPlan): SparkPlan = plan match {
case smj: SortMergeJoinExec =>
getSmjBuildSide(smj) match {
- case Some(BuildRight) if smj.joinType == LeftSemi =>
+ case Some(BuildRight)
+ if smj.joinType == LeftSemi ||
smj.joinType.isInstanceOf[ExistenceJoin] =>
// LeftSemi https://github.com/apache/datafusion-comet/issues/2667
+ // ExistenceJoin
https://github.com/apache/datafusion-comet/issues/2697
Review Comment:
#2697 is titled "Investigate potential issues with HashJoin and LeftSemi"
and neither it nor #2667 mentions `ExistenceJoin`, so this line reads as if the
issue covered existence joins. Something like `// ExistenceJoin is excluded by
analogy with LeftSemi until the root cause in #2697 is found` would be more
accurate. The forced-SHJ test in `CometJoinSuite` has the same wording in its
comment.
##########
spark/src/test/scala/org/apache/comet/exec/CometJoinSuite.scala:
##########
@@ -1574,6 +1574,186 @@ 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.
+ // AQE is disabled because a declined broadcast join loses its fallback
reason from the
+ // AQE-final plan; see
https://github.com/apache/datafusion-comet/issues/6442.
+ withSQLConf(
+ CometConf.COMET_EXEC_EXISTENCE_JOIN_ENABLED.key -> "true",
+ SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false",
+ 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)")
+ }
+ }
+ }
+ }
+
+ test("ExistenceJoin with computed join key falls back to Spark") {
+ // A computed join key (not a bare column reference) is evaluated eagerly
over the batch by
+ // the native join, so Comet keeps it on Spark; verify parity.
+ // AQE is disabled because a declined broadcast join loses its fallback
reason from the
+ // AQE-final plan; see
https://github.com/apache/datafusion-comet/issues/6442.
+ withSQLConf(
+ CometConf.COMET_EXEC_EXISTENCE_JOIN_ENABLED.key -> "true",
+ SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false",
+ 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 +
1)")
+ checkSparkAnswerAndFallbackReason(df, "computed (non-column)")
+ }
+ }
+ }
+ }
+
+ test("ExistenceJoin with duplicate build keys runs natively (markers not
multiplied)") {
+ 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") {
+ // Build side repeats keys 0..4 many times; the existence marker must
stay "at least one
+ // match" and must not multiply matching probe rows.
+ withParquetTable((0 until 30).map(i => (i % 5, i)), "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 with NOT EXISTS combined with OR runs natively") {
+ 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 NOT EXISTS " +
+ "(SELECT /*+ BROADCAST(b) */ 1 FROM tbl_b b WHERE b._1 = a._1)")
+ checkSparkAnswerAndOperator(
+ df,
+ Seq(classOf[CometBroadcastExchangeExec],
classOf[CometBroadcastHashJoinExec]))
+ }
+ }
+ }
+ }
+
+ test("ExistenceJoin falls back to Spark when the feature flag is disabled") {
+ // With the flag off, the join is declined with a toggle-specific reason
rather than a generic
+ // "Unsupported join type" message that reads like a permanent limitation.
AQE is disabled
+ // because a declined broadcast join loses its fallback reason from the
AQE-final plan; see
+ // https://github.com/apache/datafusion-comet/issues/6442.
+ withSQLConf(
+ CometConf.COMET_EXEC_EXISTENCE_JOIN_ENABLED.key -> "false",
+ SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false",
+ 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)")
+ checkSparkAnswerAndFallbackReason(df, "Native ExistenceJoin is
disabled")
+ }
+ }
+ }
+ }
+
+ test("ExistenceJoin via SortMergeJoin stays on Spark even with
forceShuffledHashJoin") {
+ // RewriteJoin refuses to rewrite an existence SortMergeJoin into a
BuildRight ShuffledHashJoin
+ // (https://github.com/apache/datafusion-comet/issues/2697), so even with
forced SHJ the join
+ // must stay on Spark rather than run as a native hash join. Assert on the
plan (no native join
+ // form) rather than the fallback reason, which is the robust signal for
this rewrite guard.
+ withSQLConf(
+ CometConf.COMET_EXEC_EXISTENCE_JOIN_ENABLED.key -> "true",
+ CometConf.COMET_FORCE_SHJ.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)")
+ val (_, cometPlan) = checkSparkAnswer(df)
+ assert(
+ collect(cometPlan) { case h: CometHashJoinExec => h }.isEmpty &&
+ collect(cometPlan) { case s: CometSortMergeJoinExec => s
}.isEmpty,
+ s"Existence SMJ must not be rewritten to a native hash
join:\n$cometPlan")
Review Comment:
This walks the plan twice for one emptiness check. A single `collectFirst`
keeps the same assertion (`nativeHashJoins` earlier in this suite is another
one-pass helper if only the hash join half is needed):
```suggestion
assert(
collectFirst(cometPlan) {
case j @ (_: CometHashJoinExec | _: CometSortMergeJoinExec) =>
j
}.isEmpty,
s"Existence SMJ must not be rewritten to a native hash
join:\n$cometPlan")
```
##########
spark/src/test/resources/sql-tests/join/existence_join.sql:
##########
@@ -0,0 +1,281 @@
+-- Licensed to the Apache Software Foundation (ASF) under one
+-- or more contributor license agreements. See the NOTICE file
+-- distributed with this work for additional information
+-- regarding copyright ownership. The ASF licenses this file
+-- to you under the Apache License, Version 2.0 (the
+-- "License"); you may not use this file except in compliance
+-- with the License. You may obtain a copy of the License at
+--
+-- http://www.apache.org/licenses/LICENSE-2.0
+--
+-- Unless required by applicable law or agreed to in writing,
+-- software distributed under the License is distributed on an
+-- "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+-- KIND, either express or implied. See the License for the
+-- specific language governing permissions and limitations
+-- under the License.
+
+-- Tests for ExistenceJoin: produced when EXISTS / NOT EXISTS is combined
+-- with another predicate via OR, preventing rewrite to LeftSemi / LeftAnti.
+--
+-- Strategy hints are placed INSIDE the EXISTS subquery (on the subquery's own
+-- relation alias), because an outer hint referencing the subquery table cannot
+-- resolve it and silently falls back to a broadcast hash join. Subquery-local
+-- hints (verified on Spark 4.1.3 with AQE off) select ShuffledHashJoin /
+-- SortMergeJoin.
+--
+-- Native existence support is currently hash-only: BROADCAST and SHUFFLE_HASH
+-- cases exercise CometBroadcastHashJoinExec / CometHashJoinExec, while MERGE
+-- cases fall back to Spark's SortMergeJoin (existence SMJ is not yet native)
and
+-- verify result parity under the Comet-enabled config.
+
+-- Native ExistenceJoin support is enabled by default; this fixture pins the
flag on explicitly.
+-- Config: spark.comet.exec.existenceJoin.enabled=true
+
+-- ============================================================
+-- Setup: NULLs (both sides), duplicates, empty build, all-NULL build
+-- ============================================================
+
+statement
+CREATE TABLE ex_left(id int, k int, region string) USING parquet
+
+statement
+INSERT INTO ex_left VALUES
+ (1, 1, 'US'),
+ (2, 2, 'EU'),
+ (3, NULL, 'US'),
+ (4, 4, 'EU'),
+ (5, 5, 'EU'),
+ (6, NULL, 'EU')
+
+statement
+CREATE TABLE ex_right(id int, k int) USING parquet
+
+statement
+INSERT INTO ex_right VALUES (10, 1), (11, 2), (12, 2), (13, NULL)
+
+statement
+CREATE TABLE ex_right_no_nulls(id int, k int) USING parquet
+
+statement
+INSERT INTO ex_right_no_nulls VALUES (10, 1), (11, 5)
+
+statement
+CREATE TABLE ex_right_empty(id int, k int) USING parquet
+
+statement
+CREATE TABLE ex_right_dups(id int, k int) USING parquet
+
+statement
+INSERT INTO ex_right_dups VALUES (10, 1), (11, 1), (12, 1), (13, 2)
+
+statement
+CREATE TABLE ex_right_all_null(id int, k int) USING parquet
+
+statement
+INSERT INTO ex_right_all_null VALUES (10, NULL), (11, NULL)
+
+statement
+CREATE TABLE ex_left_empty(id int, k int, region string) USING parquet
+
+-- ============================================================
+-- EXISTS with OR across all three strategies (hint in subquery)
+-- ============================================================
+
+query
+SELECT * FROM ex_left l
+WHERE l.region = 'US'
+ OR EXISTS (SELECT /*+ BROADCAST(r) */ 1 FROM ex_right r WHERE r.k = l.k)
+ORDER BY l.id
+
+query
+SELECT * FROM ex_left l
+WHERE l.region = 'US'
+ OR EXISTS (SELECT /*+ SHUFFLE_HASH(r) */ 1 FROM ex_right r WHERE r.k = l.k)
+ORDER BY l.id
+
+query expect_fallback(Unsupported join type)
+SELECT * FROM ex_left l
+WHERE l.region = 'US'
+ OR EXISTS (SELECT /*+ MERGE(r) */ 1 FROM ex_right r WHERE r.k = l.k)
+ORDER BY l.id
+
+-- ============================================================
+-- Empty build: every left row is unmatched, only OR-arm rows survive
+-- ============================================================
+
+query
+SELECT * FROM ex_left l
+WHERE l.region = 'US'
+ OR EXISTS (SELECT /*+ BROADCAST(r) */ 1 FROM ex_right_empty r WHERE r.k =
l.k)
+ORDER BY l.id
+
+query expect_fallback(Unsupported join type)
Review Comment:
There are four `MERGE` queries and all assert `Unsupported join type`, which
`CometSortMergeJoinExec` reports from the join type alone. The data, `NOT` and
duplicate keys cannot change that outcome, so the first one is enough and this
one, the `NOT EXISTS` one and the duplicate-keys one can go.
This one could instead become `SHUFFLE_HASH` over `ex_right_empty`. None of
the new tests covers SHJ with an empty build side, while BHJ is covered a few
lines up.
Optionally the `NOT EXISTS` + `SHUFFLE_HASH` query as well, since it is the
plain `EXISTS` + `SHUFFLE_HASH` join with a `NOT` Filter above it. The banner
of the "Empty left" section says "each strategy" but only `SHUFFLE_HASH` runs.
--
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]