From 41afc9cd0ec719fe3c0faa171811a3e5393002d1 Mon Sep 17 00:00:00 2001 From: Bhargava Vadlamani Date: Wed, 3 Jun 2026 14:59:50 -0700 Subject: [PATCH 1/9] implement_native_existence_joins --- native/core/src/execution/planner.rs | 1 + native/proto/src/proto/operator.proto | 1 + spark/src/main/scala/org/apache/spark/sql/comet/operators.scala | 2 ++ 3 files changed, 4 insertions(+) diff --git a/native/core/src/execution/planner.rs b/native/core/src/execution/planner.rs index 7ed2b3331c..91371024a8 100644 --- a/native/core/src/execution/planner.rs +++ b/native/core/src/execution/planner.rs @@ -2199,6 +2199,7 @@ impl PhysicalPlanner { Ok(JoinType::FullOuter) => DFJoinType::Full, Ok(JoinType::LeftSemi) => DFJoinType::LeftSemi, Ok(JoinType::LeftAnti) => DFJoinType::LeftAnti, + Ok(JoinType::Existence) => DFJoinType::LeftMark, Err(_) => { return Err(GeneralError(format!( "Unsupported join type: {join_type:?}" diff --git a/native/proto/src/proto/operator.proto b/native/proto/src/proto/operator.proto index 2fcfe7f25b..b5dd6d1eae 100644 --- a/native/proto/src/proto/operator.proto +++ b/native/proto/src/proto/operator.proto @@ -414,6 +414,7 @@ enum JoinType { FullOuter = 3; LeftSemi = 4; LeftAnti = 5; + Existence = 6; } enum BuildSide { diff --git a/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala b/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala index a8c674ca30..6c99503d30 100644 --- a/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala +++ b/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala @@ -2056,6 +2056,7 @@ trait CometHashJoin { case FullOuter => JoinType.FullOuter case LeftSemi => JoinType.LeftSemi case LeftAnti => JoinType.LeftAnti + case ExistenceJoin(_) => JoinType.Existence case _ => // Spark doesn't support other join types withFallbackReason(join, s"Unsupported join type ${join.joinType}") @@ -2549,6 +2550,7 @@ object CometSortMergeJoinExec extends CometOperatorSerde[SortMergeJoinExec] { case FullOuter => JoinType.FullOuter case LeftSemi => JoinType.LeftSemi case LeftAnti => JoinType.LeftAnti + case ExistenceJoin(_) => JoinType.Existence case _ => // Spark doesn't support other join types withFallbackReason(join, s"Unsupported join type ${join.joinType}") From 2b6caed631647991d9476e7f463d5d1fc541ea08 Mon Sep 17 00:00:00 2001 From: Bhargava Vadlamani Date: Wed, 3 Jun 2026 16:33:35 -0700 Subject: [PATCH 2/9] implement_native_existence_joins --- .../apache/spark/sql/comet/operators.scala | 15 ++++++ .../apache/comet/exec/CometJoinSuite.scala | 51 ++++++++++++++++++- 2 files changed, 65 insertions(+), 1 deletion(-) diff --git a/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala b/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala index 6c99503d30..1a8ae8cb25 100644 --- a/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala +++ b/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala @@ -2328,6 +2328,11 @@ case class CometHashJoinExec( override def withNewChildrenInternal(newLeft: SparkPlan, newRight: SparkPlan): SparkPlan = this.copy(left = newLeft, right = newRight) + override def producedAttributes: AttributeSet = joinType match { + case ExistenceJoin(exists) => AttributeSet(exists) + case _ => AttributeSet.empty + } + override def stringArgs: Iterator[Any] = Iterator(leftKeys, rightKeys, joinType, buildSide, condition, left, right) @@ -2469,6 +2474,11 @@ case class CometBroadcastHashJoinExec( override def withNewChildrenInternal(newLeft: SparkPlan, newRight: SparkPlan): SparkPlan = this.copy(left = newLeft, right = newRight) + override def producedAttributes: AttributeSet = joinType match { + case ExistenceJoin(exists) => AttributeSet(exists) + case _ => AttributeSet.empty + } + override def stringArgs: Iterator[Any] = Iterator(leftKeys, rightKeys, joinType, condition, buildSide, left, right) @@ -2660,6 +2670,11 @@ case class CometSortMergeJoinExec( override def withNewChildrenInternal(newLeft: SparkPlan, newRight: SparkPlan): SparkPlan = this.copy(left = newLeft, right = newRight) + override def producedAttributes: AttributeSet = joinType match { + case ExistenceJoin(exists) => AttributeSet(exists) + case _ => AttributeSet.empty + } + override def stringArgs: Iterator[Any] = Iterator(leftKeys, rightKeys, joinType, condition, left, right) diff --git a/spark/src/test/scala/org/apache/comet/exec/CometJoinSuite.scala b/spark/src/test/scala/org/apache/comet/exec/CometJoinSuite.scala index f01d5e2109..374961325e 100644 --- a/spark/src/test/scala/org/apache/comet/exec/CometJoinSuite.scala +++ b/spark/src/test/scala/org/apache/comet/exec/CometJoinSuite.scala @@ -25,7 +25,7 @@ import org.scalatest.Tag import org.apache.spark.sql.CometTestBase import org.apache.spark.sql.catalyst.TableIdentifier import org.apache.spark.sql.catalyst.analysis.UnresolvedRelation -import org.apache.spark.sql.comet.{CometBroadcastExchangeExec, CometBroadcastHashJoinExec, CometBroadcastNestedLoopJoinExec, CometSortMergeJoinExec} +import org.apache.spark.sql.comet.{CometBroadcastExchangeExec, CometBroadcastHashJoinExec, CometBroadcastNestedLoopJoinExec, CometHashJoinExec, CometSortMergeJoinExec} import org.apache.spark.sql.execution.adaptive.AQEShuffleReadExec import org.apache.spark.sql.internal.SQLConf @@ -945,4 +945,53 @@ class CometJoinSuite extends CometTestBase { } } } + + test("ExistenceJoin via BroadcastHashJoin (EXISTS combined with OR)") { + withSQLConf( + 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( + 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 (EXISTS combined with OR)") { + withSQLConf( + 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)") + checkSparkAnswerAndOperator(df, Seq(classOf[CometSortMergeJoinExec])) + } + } + } + } } From 6e8e53ceb3446b32e2bde4e6af3575400276300d Mon Sep 17 00:00:00 2001 From: Bhargava Vadlamani Date: Wed, 3 Jun 2026 22:55:22 -0700 Subject: [PATCH 3/9] implement_native_existence_joins_fix_plans --- .../q10/extended.txt | 57 ++++++++ .../q35/extended.txt | 57 ++++++++ .../q45/extended.txt | 44 ++++++ .../CometExistenceJoinBenchmark.scala | 126 ++++++++++++++++++ 4 files changed, 284 insertions(+) create mode 100644 spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4-spark3_5/q10/extended.txt create mode 100644 spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4-spark3_5/q35/extended.txt create mode 100644 spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4-spark3_5/q45/extended.txt create mode 100644 spark/src/test/scala/org/apache/spark/sql/benchmark/CometExistenceJoinBenchmark.scala diff --git a/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4-spark3_5/q10/extended.txt b/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4-spark3_5/q10/extended.txt new file mode 100644 index 0000000000..3d8ef408c2 --- /dev/null +++ b/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4-spark3_5/q10/extended.txt @@ -0,0 +1,57 @@ +CometNativeColumnarToRow ++- CometTakeOrderedAndProject + +- CometHashAggregate + +- CometExchange + +- CometHashAggregate + +- CometProject + +- CometBroadcastHashJoin + :- CometProject + : +- CometBroadcastHashJoin + : :- CometProject + : : +- CometFilter + : : +- CometBroadcastHashJoin + : : :- CometBroadcastHashJoin + : : : :- CometBroadcastHashJoin + : : : : :- CometFilter + : : : : : +- CometNativeScan parquet spark_catalog.default.customer + : : : : +- CometBroadcastExchange + : : : : +- CometProject + : : : : +- CometBroadcastHashJoin + : : : : :- CometNativeScan parquet spark_catalog.default.store_sales + : : : : : +- CometSubqueryBroadcast + : : : : : +- CometBroadcastExchange + : : : : : +- CometProject + : : : : : +- CometFilter + : : : : : +- CometNativeScan parquet spark_catalog.default.date_dim + : : : : +- CometBroadcastExchange + : : : : +- CometProject + : : : : +- CometFilter + : : : : +- CometNativeScan parquet spark_catalog.default.date_dim + : : : +- CometBroadcastExchange + : : : +- CometProject + : : : +- CometBroadcastHashJoin + : : : :- CometNativeScan parquet spark_catalog.default.web_sales + : : : : +- ReusedSubquery + : : : +- CometBroadcastExchange + : : : +- CometProject + : : : +- CometFilter + : : : +- CometNativeScan parquet spark_catalog.default.date_dim + : : +- CometBroadcastExchange + : : +- CometProject + : : +- CometBroadcastHashJoin + : : :- CometNativeScan parquet spark_catalog.default.catalog_sales + : : : +- ReusedSubquery + : : +- CometBroadcastExchange + : : +- CometProject + : : +- CometFilter + : : +- CometNativeScan parquet spark_catalog.default.date_dim + : +- CometBroadcastExchange + : +- CometProject + : +- CometFilter + : +- CometNativeScan parquet spark_catalog.default.customer_address + +- CometBroadcastExchange + +- CometProject + +- CometFilter + +- CometNativeScan parquet spark_catalog.default.customer_demographics + +Comet accelerated 51 out of 54 eligible operators (94%). Final plan contains 1 transitions between Spark and Comet. \ No newline at end of file diff --git a/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4-spark3_5/q35/extended.txt b/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4-spark3_5/q35/extended.txt new file mode 100644 index 0000000000..3d8ef408c2 --- /dev/null +++ b/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4-spark3_5/q35/extended.txt @@ -0,0 +1,57 @@ +CometNativeColumnarToRow ++- CometTakeOrderedAndProject + +- CometHashAggregate + +- CometExchange + +- CometHashAggregate + +- CometProject + +- CometBroadcastHashJoin + :- CometProject + : +- CometBroadcastHashJoin + : :- CometProject + : : +- CometFilter + : : +- CometBroadcastHashJoin + : : :- CometBroadcastHashJoin + : : : :- CometBroadcastHashJoin + : : : : :- CometFilter + : : : : : +- CometNativeScan parquet spark_catalog.default.customer + : : : : +- CometBroadcastExchange + : : : : +- CometProject + : : : : +- CometBroadcastHashJoin + : : : : :- CometNativeScan parquet spark_catalog.default.store_sales + : : : : : +- CometSubqueryBroadcast + : : : : : +- CometBroadcastExchange + : : : : : +- CometProject + : : : : : +- CometFilter + : : : : : +- CometNativeScan parquet spark_catalog.default.date_dim + : : : : +- CometBroadcastExchange + : : : : +- CometProject + : : : : +- CometFilter + : : : : +- CometNativeScan parquet spark_catalog.default.date_dim + : : : +- CometBroadcastExchange + : : : +- CometProject + : : : +- CometBroadcastHashJoin + : : : :- CometNativeScan parquet spark_catalog.default.web_sales + : : : : +- ReusedSubquery + : : : +- CometBroadcastExchange + : : : +- CometProject + : : : +- CometFilter + : : : +- CometNativeScan parquet spark_catalog.default.date_dim + : : +- CometBroadcastExchange + : : +- CometProject + : : +- CometBroadcastHashJoin + : : :- CometNativeScan parquet spark_catalog.default.catalog_sales + : : : +- ReusedSubquery + : : +- CometBroadcastExchange + : : +- CometProject + : : +- CometFilter + : : +- CometNativeScan parquet spark_catalog.default.date_dim + : +- CometBroadcastExchange + : +- CometProject + : +- CometFilter + : +- CometNativeScan parquet spark_catalog.default.customer_address + +- CometBroadcastExchange + +- CometProject + +- CometFilter + +- CometNativeScan parquet spark_catalog.default.customer_demographics + +Comet accelerated 51 out of 54 eligible operators (94%). Final plan contains 1 transitions between Spark and Comet. \ No newline at end of file diff --git a/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4-spark3_5/q45/extended.txt b/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4-spark3_5/q45/extended.txt new file mode 100644 index 0000000000..7f5a5b390d --- /dev/null +++ b/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4-spark3_5/q45/extended.txt @@ -0,0 +1,44 @@ +CometNativeColumnarToRow ++- CometTakeOrderedAndProject + +- CometHashAggregate + +- CometExchange + +- CometHashAggregate + +- CometProject + +- CometFilter + +- CometBroadcastHashJoin + :- CometProject + : +- CometBroadcastHashJoin + : :- CometProject + : : +- CometBroadcastHashJoin + : : :- CometProject + : : : +- CometBroadcastHashJoin + : : : :- CometProject + : : : : +- CometBroadcastHashJoin + : : : : :- CometFilter + : : : : : +- CometNativeScan parquet spark_catalog.default.web_sales + : : : : : +- CometSubqueryBroadcast + : : : : : +- CometBroadcastExchange + : : : : : +- CometProject + : : : : : +- CometFilter + : : : : : +- CometNativeScan parquet spark_catalog.default.date_dim + : : : : +- CometBroadcastExchange + : : : : +- CometFilter + : : : : +- CometNativeScan parquet spark_catalog.default.customer + : : : +- CometBroadcastExchange + : : : +- CometProject + : : : +- CometFilter + : : : +- CometNativeScan parquet spark_catalog.default.customer_address + : : +- CometBroadcastExchange + : : +- CometProject + : : +- CometFilter + : : +- CometNativeScan parquet spark_catalog.default.date_dim + : +- CometBroadcastExchange + : +- CometProject + : +- CometFilter + : +- CometNativeScan parquet spark_catalog.default.item + +- CometBroadcastExchange + +- CometProject + +- CometFilter + +- CometNativeScan parquet spark_catalog.default.item + +Comet accelerated 40 out of 41 eligible operators (97%). Final plan contains 1 transitions between Spark and Comet. \ No newline at end of file diff --git a/spark/src/test/scala/org/apache/spark/sql/benchmark/CometExistenceJoinBenchmark.scala b/spark/src/test/scala/org/apache/spark/sql/benchmark/CometExistenceJoinBenchmark.scala new file mode 100644 index 0000000000..e9cb2f9de6 --- /dev/null +++ b/spark/src/test/scala/org/apache/spark/sql/benchmark/CometExistenceJoinBenchmark.scala @@ -0,0 +1,126 @@ +/* + * 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. + */ + +package org.apache.spark.sql.benchmark + +import org.apache.spark.SparkConf +import org.apache.spark.sql.SparkSession +import org.apache.spark.sql.internal.SQLConf + +import org.apache.comet.{CometConf, CometSparkSessionExtensions} + +/** + * Benchmark to measure performance of Comet's ExistenceJoin support across the three join + * physical operators (BHJ, SHJ, SMJ). To run this benchmark: + * {{{ + * SPARK_GENERATE_BENCHMARK_FILES=1 make benchmark-org.apache.spark.sql.benchmark.CometExistenceJoinBenchmark + * }}} + * Results will be written to "spark/benchmarks/CometExistenceJoinBenchmark-**results.txt". + */ +object CometExistenceJoinBenchmark extends CometBenchmarkBase { + + override def getSparkSession: SparkSession = { + val conf = new SparkConf() + .setAppName("CometExistenceJoinBenchmark") + .set("spark.master", "local[5]") + .setIfMissing("spark.driver.memory", "3g") + .setIfMissing("spark.executor.memory", "3g") + .set( + "spark.shuffle.manager", + "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager") + + val sparkSession = SparkSession.builder + .config(conf) + .withExtensions(new CometSparkSessionExtensions) + .getOrCreate() + + sparkSession.conf.set(CometConf.COMET_ENABLED.key, "false") + sparkSession.conf.set(CometConf.COMET_EXEC_ENABLED.key, "false") + sparkSession.conf.set(SQLConf.ANSI_ENABLED.key, "false") + sparkSession.conf.set("spark.sql.shuffle.partitions", "2") + + sparkSession + } + + override def runCometBenchmark(mainArgs: Array[String]): Unit = { + val probeRows = 1024 * 1024 + val buildRows = 10000 + + withTempPath { dir => + withTempTable("probe", "build") { + spark + .range(probeRows) + .selectExpr("id AS k", "CASE WHEN id % 3 = 0 THEN 'US' ELSE 'EU' END AS region") + .write + .parquet(s"${dir.getAbsolutePath}/probe") + spark + .range(buildRows) + .selectExpr("id * 7 AS k") + .write + .parquet(s"${dir.getAbsolutePath}/build") + + spark.read.parquet(s"${dir.getAbsolutePath}/probe").createOrReplaceTempView("probe") + spark.read.parquet(s"${dir.getAbsolutePath}/build").createOrReplaceTempView("build") + + val query = + "SELECT count(*) FROM probe p " + + "WHERE p.region = 'US' OR EXISTS (SELECT 1 FROM build b WHERE b.k = p.k)" + + withSQLConf( + CometConf.COMET_ENABLED.key -> "true", + CometConf.COMET_EXEC_ENABLED.key -> "true") { + spark.sql(query).explain() + } + + runBenchmark("ExistenceJoin - BroadcastHashJoin") { + runExpressionBenchmark( + "exists OR predicate (BHJ)", + probeRows, + query, + Map( + SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "10MB", + SQLConf.ADAPTIVE_AUTO_BROADCASTJOIN_THRESHOLD.key -> "10MB")) + } + + runBenchmark("ExistenceJoin - ShuffledHashJoin") { + runExpressionBenchmark( + "exists OR predicate (SHJ)", + probeRows, + query, + Map( + SQLConf.PREFER_SORTMERGEJOIN.key -> "false", + "spark.sql.join.forceApplyShuffledHashJoin" -> "true", + SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1", + SQLConf.ADAPTIVE_AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1")) + } + + runBenchmark("ExistenceJoin - SortMergeJoin") { + runExpressionBenchmark( + "exists OR predicate (SMJ)", + probeRows, + query, + Map( + SQLConf.PREFER_SORTMERGEJOIN.key -> "true", + SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1", + SQLConf.ADAPTIVE_AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1")) + } + } + } + } +} From 8b7f8cab44e701a8cd1435770282e88cf91508bd Mon Sep 17 00:00:00 2001 From: Bhargava Vadlamani Date: Wed, 3 Jun 2026 23:12:38 -0700 Subject: [PATCH 4/9] implement_native_existence_joins_fix_plans --- .../q35/extended.txt | 57 +++++++++++++++++++ 1 file changed, 57 insertions(+) create mode 100644 spark/src/test/resources/tpcds-plan-stability/approved-plans-v2_7-spark3_5/q35/extended.txt diff --git a/spark/src/test/resources/tpcds-plan-stability/approved-plans-v2_7-spark3_5/q35/extended.txt b/spark/src/test/resources/tpcds-plan-stability/approved-plans-v2_7-spark3_5/q35/extended.txt new file mode 100644 index 0000000000..3d8ef408c2 --- /dev/null +++ b/spark/src/test/resources/tpcds-plan-stability/approved-plans-v2_7-spark3_5/q35/extended.txt @@ -0,0 +1,57 @@ +CometNativeColumnarToRow ++- CometTakeOrderedAndProject + +- CometHashAggregate + +- CometExchange + +- CometHashAggregate + +- CometProject + +- CometBroadcastHashJoin + :- CometProject + : +- CometBroadcastHashJoin + : :- CometProject + : : +- CometFilter + : : +- CometBroadcastHashJoin + : : :- CometBroadcastHashJoin + : : : :- CometBroadcastHashJoin + : : : : :- CometFilter + : : : : : +- CometNativeScan parquet spark_catalog.default.customer + : : : : +- CometBroadcastExchange + : : : : +- CometProject + : : : : +- CometBroadcastHashJoin + : : : : :- CometNativeScan parquet spark_catalog.default.store_sales + : : : : : +- CometSubqueryBroadcast + : : : : : +- CometBroadcastExchange + : : : : : +- CometProject + : : : : : +- CometFilter + : : : : : +- CometNativeScan parquet spark_catalog.default.date_dim + : : : : +- CometBroadcastExchange + : : : : +- CometProject + : : : : +- CometFilter + : : : : +- CometNativeScan parquet spark_catalog.default.date_dim + : : : +- CometBroadcastExchange + : : : +- CometProject + : : : +- CometBroadcastHashJoin + : : : :- CometNativeScan parquet spark_catalog.default.web_sales + : : : : +- ReusedSubquery + : : : +- CometBroadcastExchange + : : : +- CometProject + : : : +- CometFilter + : : : +- CometNativeScan parquet spark_catalog.default.date_dim + : : +- CometBroadcastExchange + : : +- CometProject + : : +- CometBroadcastHashJoin + : : :- CometNativeScan parquet spark_catalog.default.catalog_sales + : : : +- ReusedSubquery + : : +- CometBroadcastExchange + : : +- CometProject + : : +- CometFilter + : : +- CometNativeScan parquet spark_catalog.default.date_dim + : +- CometBroadcastExchange + : +- CometProject + : +- CometFilter + : +- CometNativeScan parquet spark_catalog.default.customer_address + +- CometBroadcastExchange + +- CometProject + +- CometFilter + +- CometNativeScan parquet spark_catalog.default.customer_demographics + +Comet accelerated 51 out of 54 eligible operators (94%). Final plan contains 1 transitions between Spark and Comet. \ No newline at end of file From e60de74f35f7943348bbf189a9754922a9de31f6 Mon Sep 17 00:00:00 2001 From: Bhargava Vadlamani Date: Thu, 4 Jun 2026 13:29:07 -0700 Subject: [PATCH 5/9] implement_native_existence_joins_fix_plans --- .../sql-tests/join/existence_join.sql | 164 ++++++++++++++++++ 1 file changed, 164 insertions(+) create mode 100644 spark/src/test/resources/sql-tests/join/existence_join.sql diff --git a/spark/src/test/resources/sql-tests/join/existence_join.sql b/spark/src/test/resources/sql-tests/join/existence_join.sql new file mode 100644 index 0000000000..c2dbbbf525 --- /dev/null +++ b/spark/src/test/resources/sql-tests/join/existence_join.sql @@ -0,0 +1,164 @@ +-- 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. +-- Each query runs against the three physical join strategies (BHJ, SHJ, +-- SMJ) via hints, so we exercise CometBroadcastHashJoinExec, +-- CometHashJoinExec, and CometSortMergeJoinExec all carrying joinType = +-- ExistenceJoin. + +-- ============================================================ +-- Setup: covers NULLs, duplicates, empty build side +-- ============================================================ + +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') + +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) + +-- ============================================================ +-- EXISTS with OR: BHJ build-right +-- ============================================================ + +query +SELECT /*+ BROADCAST(ex_right) */ * FROM ex_left l +WHERE l.region = 'US' OR EXISTS (SELECT 1 FROM ex_right r WHERE r.k = l.k) +ORDER BY l.id + +-- ============================================================ +-- EXISTS with OR: SHJ build-right +-- ============================================================ + +query +SELECT /*+ SHUFFLE_HASH(ex_right) */ * FROM ex_left l +WHERE l.region = 'US' OR EXISTS (SELECT 1 FROM ex_right r WHERE r.k = l.k) +ORDER BY l.id + +-- ============================================================ +-- EXISTS with OR: SMJ +-- ============================================================ + +query +SELECT /*+ MERGE(ex_right) */ * FROM ex_left l +WHERE l.region = 'US' OR EXISTS (SELECT 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 /*+ BROADCAST(ex_right_empty) */ * FROM ex_left l +WHERE l.region = 'US' OR EXISTS (SELECT 1 FROM ex_right_empty r WHERE r.k = l.k) +ORDER BY l.id + +query +SELECT /*+ MERGE(ex_right_empty) */ * FROM ex_left l +WHERE l.region = 'US' OR EXISTS (SELECT 1 FROM ex_right_empty r WHERE r.k = l.k) +ORDER BY l.id + +-- ============================================================ +-- Right side has no NULL: NULL-keyed left row reaches the marker +-- evaluation but cannot match (NULL = anything is NULL → false), +-- so its exists tag is false. +-- ============================================================ + +query +SELECT /*+ BROADCAST(ex_right_no_nulls) */ * FROM ex_left l +WHERE l.region = 'US' OR EXISTS (SELECT 1 FROM ex_right_no_nulls r WHERE r.k = l.k) +ORDER BY l.id + +-- ============================================================ +-- NOT EXISTS combined with OR: also lowers to ExistenceJoin +-- (the optimizer flips the marker via NOT in the filter). +-- ============================================================ + +query +SELECT /*+ BROADCAST(ex_right) */ * FROM ex_left l +WHERE l.region = 'US' OR NOT EXISTS (SELECT 1 FROM ex_right r WHERE r.k = l.k) +ORDER BY l.id + +-- ============================================================ +-- Build with duplicate keys: marker is "at least one match", so duplicates +-- on the right must not multiply the output. +-- ============================================================ + +query +SELECT /*+ BROADCAST(ex_right_dups) */ * FROM ex_left l +WHERE l.region = 'US' OR EXISTS (SELECT 1 FROM ex_right_dups r WHERE r.k = l.k) +ORDER BY l.id + +-- ============================================================ +-- Marker used inside a more complex predicate (NOT exists OR ...). +-- ============================================================ + +query +SELECT /*+ BROADCAST(ex_right) */ id, k, region FROM ex_left l +WHERE l.id > 1 + AND (l.region = 'US' OR NOT EXISTS (SELECT 1 FROM ex_right r WHERE r.k = l.k)) +ORDER BY l.id + +-- ============================================================ +-- Multi-column correlation +-- ============================================================ + +statement +CREATE TABLE ex_left_multi(id int, k1 int, k2 int) USING parquet + +statement +INSERT INTO ex_left_multi VALUES (1, 1, 100), (2, 2, 200), (3, 1, 300) + +statement +CREATE TABLE ex_right_multi(k1 int, k2 int) USING parquet + +statement +INSERT INTO ex_right_multi VALUES (1, 100), (2, 999) + +query +SELECT /*+ BROADCAST(ex_right_multi) */ * FROM ex_left_multi l +WHERE l.id > 0 + AND (l.k1 = 1 + OR EXISTS (SELECT 1 FROM ex_right_multi r WHERE r.k1 = l.k1 AND r.k2 = l.k2)) +ORDER BY l.id From d3d400acea2a12bae19677149225a9a26c607bbd Mon Sep 17 00:00:00 2001 From: Bhargava Vadlamani Date: Fri, 10 Jul 2026 00:18:50 -0700 Subject: [PATCH 6/9] init_commit --- .../scala/org/apache/comet/CometConf.scala | 7 +++ .../apache/spark/sql/comet/operators.scala | 20 ++++++- .../sql-tests/join/existence_join.sql | 3 + .../q10/extended.txt | 57 ------------------- .../q35/extended.txt | 57 ------------------- .../q45/extended.txt | 44 -------------- .../q35/extended.txt | 57 ------------------- .../apache/comet/exec/CometJoinSuite.scala | 3 + 8 files changed, 31 insertions(+), 217 deletions(-) delete mode 100644 spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4-spark3_5/q10/extended.txt delete mode 100644 spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4-spark3_5/q35/extended.txt delete mode 100644 spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4-spark3_5/q45/extended.txt delete mode 100644 spark/src/test/resources/tpcds-plan-stability/approved-plans-v2_7-spark3_5/q35/extended.txt diff --git a/spark/src/main/scala/org/apache/comet/CometConf.scala b/spark/src/main/scala/org/apache/comet/CometConf.scala index 5c130d457e..ca7ca0de91 100644 --- a/spark/src/main/scala/org/apache/comet/CometConf.scala +++ b/spark/src/main/scala/org/apache/comet/CometConf.scala @@ -261,6 +261,13 @@ object CometConf extends ShimCometConf { createExecEnabledConfig("broadcastNestedLoopJoin", defaultValue = true) val COMET_EXEC_SORT_MERGE_JOIN_ENABLED: ConfigEntry[Boolean] = createExecEnabledConfig("sortMergeJoin", defaultValue = true) + val COMET_EXEC_EXISTENCE_JOIN_ENABLED: ConfigEntry[Boolean] = + createExecEnabledConfig( + "existenceJoin", + defaultValue = false, + notes = Some( + "This enables native ExistenceJoin support (EXISTS/NOT EXISTS combined with OR). " + + "It is experimental and disabled by default")) val COMET_EXEC_AGGREGATE_ENABLED: ConfigEntry[Boolean] = createExecEnabledConfig("aggregate", defaultValue = true) val COMET_EXEC_COLLECT_LIMIT_ENABLED: ConfigEntry[Boolean] = diff --git a/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala b/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala index 1a8ae8cb25..b3ea6f549a 100644 --- a/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala +++ b/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala @@ -2056,7 +2056,15 @@ trait CometHashJoin { case FullOuter => JoinType.FullOuter case LeftSemi => JoinType.LeftSemi case LeftAnti => JoinType.LeftAnti - case ExistenceJoin(_) => JoinType.Existence + case ExistenceJoin(_) => + if (!CometConf.COMET_EXEC_EXISTENCE_JOIN_ENABLED.get(join.conf)) { + withFallbackReason( + join, + "Existence join support is experimental. Set " + + s"${CometConf.COMET_EXEC_EXISTENCE_JOIN_ENABLED.key}=true to enable.") + return None + } + JoinType.Existence case _ => // Spark doesn't support other join types withFallbackReason(join, s"Unsupported join type ${join.joinType}") @@ -2560,7 +2568,15 @@ object CometSortMergeJoinExec extends CometOperatorSerde[SortMergeJoinExec] { case FullOuter => JoinType.FullOuter case LeftSemi => JoinType.LeftSemi case LeftAnti => JoinType.LeftAnti - case ExistenceJoin(_) => JoinType.Existence + case ExistenceJoin(_) => + if (!CometConf.COMET_EXEC_EXISTENCE_JOIN_ENABLED.get(join.conf)) { + withFallbackReason( + join, + "Existence join support is experimental. Set " + + s"${CometConf.COMET_EXEC_EXISTENCE_JOIN_ENABLED.key}=true to enable.") + return None + } + JoinType.Existence case _ => // Spark doesn't support other join types withFallbackReason(join, s"Unsupported join type ${join.joinType}") diff --git a/spark/src/test/resources/sql-tests/join/existence_join.sql b/spark/src/test/resources/sql-tests/join/existence_join.sql index c2dbbbf525..f6692a198c 100644 --- a/spark/src/test/resources/sql-tests/join/existence_join.sql +++ b/spark/src/test/resources/sql-tests/join/existence_join.sql @@ -22,6 +22,9 @@ -- CometHashJoinExec, and CometSortMergeJoinExec all carrying joinType = -- ExistenceJoin. +-- Native ExistenceJoin support is experimental and disabled by default. +-- Config: spark.comet.exec.existenceJoin.enabled=true + -- ============================================================ -- Setup: covers NULLs, duplicates, empty build side -- ============================================================ diff --git a/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4-spark3_5/q10/extended.txt b/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4-spark3_5/q10/extended.txt deleted file mode 100644 index 3d8ef408c2..0000000000 --- a/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4-spark3_5/q10/extended.txt +++ /dev/null @@ -1,57 +0,0 @@ -CometNativeColumnarToRow -+- CometTakeOrderedAndProject - +- CometHashAggregate - +- CometExchange - +- CometHashAggregate - +- CometProject - +- CometBroadcastHashJoin - :- CometProject - : +- CometBroadcastHashJoin - : :- CometProject - : : +- CometFilter - : : +- CometBroadcastHashJoin - : : :- CometBroadcastHashJoin - : : : :- CometBroadcastHashJoin - : : : : :- CometFilter - : : : : : +- CometNativeScan parquet spark_catalog.default.customer - : : : : +- CometBroadcastExchange - : : : : +- CometProject - : : : : +- CometBroadcastHashJoin - : : : : :- CometNativeScan parquet spark_catalog.default.store_sales - : : : : : +- CometSubqueryBroadcast - : : : : : +- CometBroadcastExchange - : : : : : +- CometProject - : : : : : +- CometFilter - : : : : : +- CometNativeScan parquet spark_catalog.default.date_dim - : : : : +- CometBroadcastExchange - : : : : +- CometProject - : : : : +- CometFilter - : : : : +- CometNativeScan parquet spark_catalog.default.date_dim - : : : +- CometBroadcastExchange - : : : +- CometProject - : : : +- CometBroadcastHashJoin - : : : :- CometNativeScan parquet spark_catalog.default.web_sales - : : : : +- ReusedSubquery - : : : +- CometBroadcastExchange - : : : +- CometProject - : : : +- CometFilter - : : : +- CometNativeScan parquet spark_catalog.default.date_dim - : : +- CometBroadcastExchange - : : +- CometProject - : : +- CometBroadcastHashJoin - : : :- CometNativeScan parquet spark_catalog.default.catalog_sales - : : : +- ReusedSubquery - : : +- CometBroadcastExchange - : : +- CometProject - : : +- CometFilter - : : +- CometNativeScan parquet spark_catalog.default.date_dim - : +- CometBroadcastExchange - : +- CometProject - : +- CometFilter - : +- CometNativeScan parquet spark_catalog.default.customer_address - +- CometBroadcastExchange - +- CometProject - +- CometFilter - +- CometNativeScan parquet spark_catalog.default.customer_demographics - -Comet accelerated 51 out of 54 eligible operators (94%). Final plan contains 1 transitions between Spark and Comet. \ No newline at end of file diff --git a/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4-spark3_5/q35/extended.txt b/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4-spark3_5/q35/extended.txt deleted file mode 100644 index 3d8ef408c2..0000000000 --- a/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4-spark3_5/q35/extended.txt +++ /dev/null @@ -1,57 +0,0 @@ -CometNativeColumnarToRow -+- CometTakeOrderedAndProject - +- CometHashAggregate - +- CometExchange - +- CometHashAggregate - +- CometProject - +- CometBroadcastHashJoin - :- CometProject - : +- CometBroadcastHashJoin - : :- CometProject - : : +- CometFilter - : : +- CometBroadcastHashJoin - : : :- CometBroadcastHashJoin - : : : :- CometBroadcastHashJoin - : : : : :- CometFilter - : : : : : +- CometNativeScan parquet spark_catalog.default.customer - : : : : +- CometBroadcastExchange - : : : : +- CometProject - : : : : +- CometBroadcastHashJoin - : : : : :- CometNativeScan parquet spark_catalog.default.store_sales - : : : : : +- CometSubqueryBroadcast - : : : : : +- CometBroadcastExchange - : : : : : +- CometProject - : : : : : +- CometFilter - : : : : : +- CometNativeScan parquet spark_catalog.default.date_dim - : : : : +- CometBroadcastExchange - : : : : +- CometProject - : : : : +- CometFilter - : : : : +- CometNativeScan parquet spark_catalog.default.date_dim - : : : +- CometBroadcastExchange - : : : +- CometProject - : : : +- CometBroadcastHashJoin - : : : :- CometNativeScan parquet spark_catalog.default.web_sales - : : : : +- ReusedSubquery - : : : +- CometBroadcastExchange - : : : +- CometProject - : : : +- CometFilter - : : : +- CometNativeScan parquet spark_catalog.default.date_dim - : : +- CometBroadcastExchange - : : +- CometProject - : : +- CometBroadcastHashJoin - : : :- CometNativeScan parquet spark_catalog.default.catalog_sales - : : : +- ReusedSubquery - : : +- CometBroadcastExchange - : : +- CometProject - : : +- CometFilter - : : +- CometNativeScan parquet spark_catalog.default.date_dim - : +- CometBroadcastExchange - : +- CometProject - : +- CometFilter - : +- CometNativeScan parquet spark_catalog.default.customer_address - +- CometBroadcastExchange - +- CometProject - +- CometFilter - +- CometNativeScan parquet spark_catalog.default.customer_demographics - -Comet accelerated 51 out of 54 eligible operators (94%). Final plan contains 1 transitions between Spark and Comet. \ No newline at end of file diff --git a/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4-spark3_5/q45/extended.txt b/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4-spark3_5/q45/extended.txt deleted file mode 100644 index 7f5a5b390d..0000000000 --- a/spark/src/test/resources/tpcds-plan-stability/approved-plans-v1_4-spark3_5/q45/extended.txt +++ /dev/null @@ -1,44 +0,0 @@ -CometNativeColumnarToRow -+- CometTakeOrderedAndProject - +- CometHashAggregate - +- CometExchange - +- CometHashAggregate - +- CometProject - +- CometFilter - +- CometBroadcastHashJoin - :- CometProject - : +- CometBroadcastHashJoin - : :- CometProject - : : +- CometBroadcastHashJoin - : : :- CometProject - : : : +- CometBroadcastHashJoin - : : : :- CometProject - : : : : +- CometBroadcastHashJoin - : : : : :- CometFilter - : : : : : +- CometNativeScan parquet spark_catalog.default.web_sales - : : : : : +- CometSubqueryBroadcast - : : : : : +- CometBroadcastExchange - : : : : : +- CometProject - : : : : : +- CometFilter - : : : : : +- CometNativeScan parquet spark_catalog.default.date_dim - : : : : +- CometBroadcastExchange - : : : : +- CometFilter - : : : : +- CometNativeScan parquet spark_catalog.default.customer - : : : +- CometBroadcastExchange - : : : +- CometProject - : : : +- CometFilter - : : : +- CometNativeScan parquet spark_catalog.default.customer_address - : : +- CometBroadcastExchange - : : +- CometProject - : : +- CometFilter - : : +- CometNativeScan parquet spark_catalog.default.date_dim - : +- CometBroadcastExchange - : +- CometProject - : +- CometFilter - : +- CometNativeScan parquet spark_catalog.default.item - +- CometBroadcastExchange - +- CometProject - +- CometFilter - +- CometNativeScan parquet spark_catalog.default.item - -Comet accelerated 40 out of 41 eligible operators (97%). Final plan contains 1 transitions between Spark and Comet. \ No newline at end of file diff --git a/spark/src/test/resources/tpcds-plan-stability/approved-plans-v2_7-spark3_5/q35/extended.txt b/spark/src/test/resources/tpcds-plan-stability/approved-plans-v2_7-spark3_5/q35/extended.txt deleted file mode 100644 index 3d8ef408c2..0000000000 --- a/spark/src/test/resources/tpcds-plan-stability/approved-plans-v2_7-spark3_5/q35/extended.txt +++ /dev/null @@ -1,57 +0,0 @@ -CometNativeColumnarToRow -+- CometTakeOrderedAndProject - +- CometHashAggregate - +- CometExchange - +- CometHashAggregate - +- CometProject - +- CometBroadcastHashJoin - :- CometProject - : +- CometBroadcastHashJoin - : :- CometProject - : : +- CometFilter - : : +- CometBroadcastHashJoin - : : :- CometBroadcastHashJoin - : : : :- CometBroadcastHashJoin - : : : : :- CometFilter - : : : : : +- CometNativeScan parquet spark_catalog.default.customer - : : : : +- CometBroadcastExchange - : : : : +- CometProject - : : : : +- CometBroadcastHashJoin - : : : : :- CometNativeScan parquet spark_catalog.default.store_sales - : : : : : +- CometSubqueryBroadcast - : : : : : +- CometBroadcastExchange - : : : : : +- CometProject - : : : : : +- CometFilter - : : : : : +- CometNativeScan parquet spark_catalog.default.date_dim - : : : : +- CometBroadcastExchange - : : : : +- CometProject - : : : : +- CometFilter - : : : : +- CometNativeScan parquet spark_catalog.default.date_dim - : : : +- CometBroadcastExchange - : : : +- CometProject - : : : +- CometBroadcastHashJoin - : : : :- CometNativeScan parquet spark_catalog.default.web_sales - : : : : +- ReusedSubquery - : : : +- CometBroadcastExchange - : : : +- CometProject - : : : +- CometFilter - : : : +- CometNativeScan parquet spark_catalog.default.date_dim - : : +- CometBroadcastExchange - : : +- CometProject - : : +- CometBroadcastHashJoin - : : :- CometNativeScan parquet spark_catalog.default.catalog_sales - : : : +- ReusedSubquery - : : +- CometBroadcastExchange - : : +- CometProject - : : +- CometFilter - : : +- CometNativeScan parquet spark_catalog.default.date_dim - : +- CometBroadcastExchange - : +- CometProject - : +- CometFilter - : +- CometNativeScan parquet spark_catalog.default.customer_address - +- CometBroadcastExchange - +- CometProject - +- CometFilter - +- CometNativeScan parquet spark_catalog.default.customer_demographics - -Comet accelerated 51 out of 54 eligible operators (94%). Final plan contains 1 transitions between Spark and Comet. \ No newline at end of file diff --git a/spark/src/test/scala/org/apache/comet/exec/CometJoinSuite.scala b/spark/src/test/scala/org/apache/comet/exec/CometJoinSuite.scala index 374961325e..d834a34b43 100644 --- a/spark/src/test/scala/org/apache/comet/exec/CometJoinSuite.scala +++ b/spark/src/test/scala/org/apache/comet/exec/CometJoinSuite.scala @@ -948,6 +948,7 @@ 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") { @@ -964,6 +965,7 @@ class CometJoinSuite extends CometTestBase { 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", @@ -981,6 +983,7 @@ class CometJoinSuite extends CometTestBase { test("ExistenceJoin via SortMergeJoin (EXISTS combined with OR)") { 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") { From 15b054e47125ca60bde8c0c3a12eb193a402e2b6 Mon Sep 17 00:00:00 2001 From: Bhargava Vadlamani Date: Fri, 17 Jul 2026 08:56:58 -0700 Subject: [PATCH 7/9] make_doc_changes --- spark/src/main/scala/org/apache/comet/CometConf.scala | 2 +- .../src/main/scala/org/apache/spark/sql/comet/operators.scala | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/spark/src/main/scala/org/apache/comet/CometConf.scala b/spark/src/main/scala/org/apache/comet/CometConf.scala index ca7ca0de91..faa2504a20 100644 --- a/spark/src/main/scala/org/apache/comet/CometConf.scala +++ b/spark/src/main/scala/org/apache/comet/CometConf.scala @@ -267,7 +267,7 @@ object CometConf extends ShimCometConf { defaultValue = false, notes = Some( "This enables native ExistenceJoin support (EXISTS/NOT EXISTS combined with OR). " + - "It is experimental and disabled by default")) + "This is highly experimental and disabled by default")) val COMET_EXEC_AGGREGATE_ENABLED: ConfigEntry[Boolean] = createExecEnabledConfig("aggregate", defaultValue = true) val COMET_EXEC_COLLECT_LIMIT_ENABLED: ConfigEntry[Boolean] = diff --git a/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala b/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala index b3ea6f549a..624baf645a 100644 --- a/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala +++ b/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala @@ -2060,7 +2060,7 @@ trait CometHashJoin { if (!CometConf.COMET_EXEC_EXISTENCE_JOIN_ENABLED.get(join.conf)) { withFallbackReason( join, - "Existence join support is experimental. Set " + + "Existence join support is highly experimental. Set " + s"${CometConf.COMET_EXEC_EXISTENCE_JOIN_ENABLED.key}=true to enable.") return None } @@ -2572,7 +2572,7 @@ object CometSortMergeJoinExec extends CometOperatorSerde[SortMergeJoinExec] { if (!CometConf.COMET_EXEC_EXISTENCE_JOIN_ENABLED.get(join.conf)) { withFallbackReason( join, - "Existence join support is experimental. Set " + + "Existence join support is highly experimental. Set " + s"${CometConf.COMET_EXEC_EXISTENCE_JOIN_ENABLED.key}=true to enable.") return None } From ea525e274adbaef4715883c138455eacb761f34c Mon Sep 17 00:00:00 2001 From: Bhargava Vadlamani Date: Sun, 19 Jul 2026 22:57:21 -0700 Subject: [PATCH 8/9] make_doc_default_changes --- .../org/apache/spark/sql/comet/operators.scala | 18 ++---------------- 1 file changed, 2 insertions(+), 16 deletions(-) diff --git a/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala b/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala index 624baf645a..abb10fbaaf 100644 --- a/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala +++ b/spark/src/main/scala/org/apache/spark/sql/comet/operators.scala @@ -2056,14 +2056,7 @@ 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)) { - withFallbackReason( - join, - "Existence join support is highly experimental. Set " + - s"${CometConf.COMET_EXEC_EXISTENCE_JOIN_ENABLED.key}=true to enable.") - return None - } + case ExistenceJoin(_) if CometConf.COMET_EXEC_EXISTENCE_JOIN_ENABLED.get(join.conf) => JoinType.Existence case _ => // Spark doesn't support other join types @@ -2568,14 +2561,7 @@ object CometSortMergeJoinExec extends CometOperatorSerde[SortMergeJoinExec] { case FullOuter => JoinType.FullOuter case LeftSemi => JoinType.LeftSemi case LeftAnti => JoinType.LeftAnti - case ExistenceJoin(_) => - if (!CometConf.COMET_EXEC_EXISTENCE_JOIN_ENABLED.get(join.conf)) { - withFallbackReason( - join, - "Existence join support is highly experimental. Set " + - s"${CometConf.COMET_EXEC_EXISTENCE_JOIN_ENABLED.key}=true to enable.") - return None - } + case ExistenceJoin(_) if CometConf.COMET_EXEC_EXISTENCE_JOIN_ENABLED.get(join.conf) => JoinType.Existence case _ => // Spark doesn't support other join types From 6fc8b3288fe74442f717f480e3782d7f8b0a3ba4 Mon Sep 17 00:00:00 2001 From: Bhargava Vadlamani Date: Fri, 24 Jul 2026 11:54:05 -0700 Subject: [PATCH 9/9] address_review_comments_update_benches --- .../sql/benchmark/CometExistenceJoinBenchmark.scala | 9 +++------ 1 file changed, 3 insertions(+), 6 deletions(-) diff --git a/spark/src/test/scala/org/apache/spark/sql/benchmark/CometExistenceJoinBenchmark.scala b/spark/src/test/scala/org/apache/spark/sql/benchmark/CometExistenceJoinBenchmark.scala index e9cb2f9de6..f0a3cbc109 100644 --- a/spark/src/test/scala/org/apache/spark/sql/benchmark/CometExistenceJoinBenchmark.scala +++ b/spark/src/test/scala/org/apache/spark/sql/benchmark/CometExistenceJoinBenchmark.scala @@ -82,18 +82,13 @@ object CometExistenceJoinBenchmark extends CometBenchmarkBase { "SELECT count(*) FROM probe p " + "WHERE p.region = 'US' OR EXISTS (SELECT 1 FROM build b WHERE b.k = p.k)" - withSQLConf( - CometConf.COMET_ENABLED.key -> "true", - CometConf.COMET_EXEC_ENABLED.key -> "true") { - spark.sql(query).explain() - } - runBenchmark("ExistenceJoin - BroadcastHashJoin") { runExpressionBenchmark( "exists OR predicate (BHJ)", probeRows, query, Map( + CometConf.COMET_EXEC_EXISTENCE_JOIN_ENABLED.key -> "true", SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "10MB", SQLConf.ADAPTIVE_AUTO_BROADCASTJOIN_THRESHOLD.key -> "10MB")) } @@ -104,6 +99,7 @@ object CometExistenceJoinBenchmark extends CometBenchmarkBase { probeRows, query, Map( + CometConf.COMET_EXEC_EXISTENCE_JOIN_ENABLED.key -> "true", SQLConf.PREFER_SORTMERGEJOIN.key -> "false", "spark.sql.join.forceApplyShuffledHashJoin" -> "true", SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "-1", @@ -116,6 +112,7 @@ object CometExistenceJoinBenchmark extends CometBenchmarkBase { probeRows, query, Map( + 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"))