diff --git a/dev/diffs/3.4.3.diff b/dev/diffs/3.4.3.diff index 8fc6451ef12..ce9970af34b 100644 --- a/dev/diffs/3.4.3.diff +++ b/dev/diffs/3.4.3.diff @@ -37,6 +37,48 @@ index d3544881af1..0c0618a5e80 100644 +diff --git a/project/SparkBuild.scala b/project/SparkBuild.scala +index 1cbd9a61289..74aa1f684c9 100644 +--- a/project/SparkBuild.scala ++++ b/project/SparkBuild.scala +@@ -445,9 +445,11 @@ object SparkBuild extends PomBuild { + + /* Spark SQL Core console settings */ + enable(SQL.settings)(sql) ++ enable(CometTestSettings.settings)(sql) + + /* Hive console settings */ + enable(Hive.settings)(hive) ++ enable(CometTestSettings.settings)(hive) + + enable(SparkConnectCommon.settings)(connectCommon) + enable(SparkConnect.settings)(connect) +@@ -1183,6 +1185,25 @@ object SQL { + ) + } + ++/** ++ * Comet disables itself when its shuffle manager is not registered. When Comet is enabled, make ++ * that shuffle manager the default for every test SparkConf in the modules that have Comet on ++ * their classpath, so that suites which build their own SparkSession or SparkContext run Comet ++ * too, not just the ones that go through SharedSparkSession or TestHive. ++ */ ++object CometTestSettings { ++ private val cometEnabled = ++ sys.env.get("ENABLE_COMET").forall(v => v == "1" || v.equalsIgnoreCase("true")) ++ ++ lazy val settings: Seq[Setting[_]] = ++ if (cometEnabled && !sys.props.contains("spark.shuffle.manager")) { ++ Seq((Test / javaOptions) += ++ "-Dspark.shuffle.manager=org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager") ++ } else { ++ Seq.empty ++ } ++} ++ + object Hive { + + lazy val settings = Seq( diff --git a/sql/core/pom.xml b/sql/core/pom.xml index b386d135da1..46449e3f3f1 100644 --- a/sql/core/pom.xml @@ -1056,6 +1098,50 @@ index 2dabcf01be7..8fcec0d1ce4 100644 } } } +diff --git a/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionExtensionSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionExtensionSuite.scala +index 48ad10992c5..a164e273b76 100644 +--- a/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionExtensionSuite.scala ++++ b/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionExtensionSuite.scala +@@ -165,7 +165,15 @@ class SparkSessionExtensionSuite extends SparkFunSuite with SQLHelper { + } + } + ++ // Comet: Comet's own columnar rules take over these plans, so the rules the tests below ++ // inject never see the Spark operators they act on or check for. ++ private def assumeCometDisabled(): Unit = { ++ assume(!org.apache.spark.sql.SparkSession.isCometEnabled, ++ "Skipped when Comet is enabled: Comet replaces the operators the injected rules act on") ++ } ++ + test("inject adaptive query prep rule") { ++ assumeCometDisabled() + val extensions = create { extensions => + // inject rule that will run during AQE query stage preparation and will add custom tags + // to the plan +@@ -259,6 +267,7 @@ class SparkSessionExtensionSuite extends SparkFunSuite with SQLHelper { + } + + private def testInjectColumnar(enableAQE: Boolean): Unit = { ++ assumeCometDisabled() + def collectPlanSteps(plan: SparkPlan): Seq[Int] = plan match { + case a: AdaptiveSparkPlanExec => + assert(a.toString.startsWith("AdaptiveSparkPlan isFinalPlan=true")) +@@ -314,6 +323,7 @@ class SparkSessionExtensionSuite extends SparkFunSuite with SQLHelper { + } + + test("reset column vectors") { ++ assumeCometDisabled() + val session = SparkSession.builder() + .master("local[1]") + .config(COLUMN_BATCH_SIZE.key, 2) +@@ -482,6 +492,7 @@ class SparkSessionExtensionSuite extends SparkFunSuite with SQLHelper { + } + + test("SPARK-38697: Extend SparkSessionExtensions to inject rules into AQE Optimizer") { ++ assumeCometDisabled() + def executedPlan(df: Dataset[java.lang.Long]): SparkPlan = { + assert(df.queryExecution.executedPlan.isInstanceOf[AdaptiveSparkPlanExec]) + df.queryExecution.executedPlan.asInstanceOf[AdaptiveSparkPlanExec].executedPlan diff --git a/sql/core/src/test/scala/org/apache/spark/sql/StringFunctionsSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/StringFunctionsSuite.scala index 18123a4d6ec..0fe185baa33 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/StringFunctionsSuite.scala @@ -1328,6 +1414,84 @@ index c0ec8a58bd5..4e8bc6ed3c5 100644 // Fail to read ancient datetime values. withSQLConf(SQLConf.PARQUET_REBASE_MODE_IN_READ.key -> EXCEPTION.toString) { +diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/CoalesceShufflePartitionsSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/CoalesceShufflePartitionsSuite.scala +index 24a98dd83f3..4db28a7a1ce 100644 +--- a/sql/core/src/test/scala/org/apache/spark/sql/execution/CoalesceShufflePartitionsSuite.scala ++++ b/sql/core/src/test/scala/org/apache/spark/sql/execution/CoalesceShufflePartitionsSuite.scala +@@ -51,6 +51,11 @@ class CoalesceShufflePartitionsSuite extends SparkFunSuite { + } + } + ++ // Comet: AQE coalesces post-shuffle partitions by map output size, and Comet's shuffle ++ // writes Arrow IPC blocks whose sizes differ from Spark's, so the expected partition ++ // counts below only hold for Spark's shuffle. ++ private def isCometEnabled: Boolean = org.apache.spark.sql.SparkSession.isCometEnabled ++ + val numInputPartitions: Int = 10 + + def withSparkSession( +@@ -163,9 +168,11 @@ class CoalesceShufflePartitionsSuite extends SparkFunSuite { + assert(shuffleReads.isEmpty) + + case None => +- assert(shuffleReads.length === 2) +- shuffleReads.foreach { read => +- assert(read.outputPartitioning.numPartitions === 2) ++ if (!isCometEnabled) { ++ assert(shuffleReads.length === 2) ++ shuffleReads.foreach { read => ++ assert(read.outputPartitioning.numPartitions === 2) ++ } + } + } + } +@@ -214,9 +221,11 @@ class CoalesceShufflePartitionsSuite extends SparkFunSuite { + assert(shuffleReads.isEmpty) + + case None => +- assert(shuffleReads.length === 2) +- shuffleReads.foreach { read => +- assert(read.outputPartitioning.numPartitions === 2) ++ if (!isCometEnabled) { ++ assert(shuffleReads.length === 2) ++ shuffleReads.foreach { read => ++ assert(read.outputPartitioning.numPartitions === 2) ++ } + } + } + } +@@ -265,9 +274,11 @@ class CoalesceShufflePartitionsSuite extends SparkFunSuite { + assert(shuffleReads.isEmpty) + + case None => +- assert(shuffleReads.length === 2) +- shuffleReads.foreach { read => +- assert(read.outputPartitioning.numPartitions === 3) ++ if (!isCometEnabled) { ++ assert(shuffleReads.length === 2) ++ shuffleReads.foreach { read => ++ assert(read.outputPartitioning.numPartitions === 3) ++ } + } + } + } +@@ -412,10 +423,12 @@ class CoalesceShufflePartitionsSuite extends SparkFunSuite { + // aggregate on the other side of the union. + val finalPlan = resultDf.queryExecution.executedPlan + .asInstanceOf[AdaptiveSparkPlanExec].executedPlan +- assert( +- finalPlan.collect { +- case r @ CoalescedShuffleRead() => r +- }.size == 2) ++ if (!isCometEnabled) { ++ assert( ++ finalPlan.collect { ++ case r @ CoalescedShuffleRead() => r ++ }.size == 2) ++ } + } + withSparkSession(test, 100, None) + } diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/DataSourceScanExecRedactionSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/DataSourceScanExecRedactionSuite.scala index 418ca3430bb..eb8267192f8 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/execution/DataSourceScanExecRedactionSuite.scala @@ -2433,6 +2597,29 @@ index 3a0bd35cb70..b28f06a757f 100644 withTempPath { workDir => val workDirPath = workDir.getAbsolutePath val input = spark.range(5).toDF("id") +diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/joins/BroadcastJoinSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/joins/BroadcastJoinSuite.scala +index 6333808b420..81b2704300c 100644 +--- a/sql/core/src/test/scala/org/apache/spark/sql/execution/joins/BroadcastJoinSuite.scala ++++ b/sql/core/src/test/scala/org/apache/spark/sql/execution/joins/BroadcastJoinSuite.scala +@@ -21,7 +21,7 @@ import scala.reflect.ClassTag + + import org.apache.spark.AccumulatorSuite + import org.apache.spark.internal.config.EXECUTOR_MEMORY +-import org.apache.spark.sql.{Dataset, QueryTest, Row, SparkSession} ++import org.apache.spark.sql.{Dataset, IgnoreComet, QueryTest, Row, SparkSession} + import org.apache.spark.sql.catalyst.expressions.{AttributeReference, BitwiseAnd, BitwiseOr, Cast, Expression, Literal, ShiftLeft} + import org.apache.spark.sql.catalyst.optimizer.{BuildLeft, BuildRight, BuildSide} + import org.apache.spark.sql.catalyst.plans.Inner +@@ -486,7 +486,8 @@ abstract class BroadcastJoinSuiteBase extends QueryTest with SQLTestUtils + } + } + +- test("broadcast join where streamed side's output partitioning is PartitioningCollection") { ++ test("broadcast join where streamed side's output partitioning is PartitioningCollection", ++ IgnoreComet("Comet replaces the join and shuffle operators this test inspects")) { + withSQLConf(SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "500") { + val t1 = (0 until 100).map(i => (i % 5, i % 13)).toDF("i1", "j1") + val t2 = (0 until 100).map(i => (i % 5, i % 14)).toDF("i2", "j2") diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/metric/SQLMetricsSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/metric/SQLMetricsSuite.scala index 26e61c6b58d..cb09d7e116a 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/execution/metric/SQLMetricsSuite.scala diff --git a/dev/diffs/3.5.9.diff b/dev/diffs/3.5.9.diff index 462554a1047..932f558fa56 100644 --- a/dev/diffs/3.5.9.diff +++ b/dev/diffs/3.5.9.diff @@ -37,6 +37,48 @@ index 49eca0f7555..b9e5d69a6c6 100644 +diff --git a/project/SparkBuild.scala b/project/SparkBuild.scala +index daa3ebac992..2b6138c4280 100644 +--- a/project/SparkBuild.scala ++++ b/project/SparkBuild.scala +@@ -459,9 +459,11 @@ object SparkBuild extends PomBuild { + + /* Spark SQL Core console settings */ + enable(SQL.settings)(sql) ++ enable(CometTestSettings.settings)(sql) + + /* Hive console settings */ + enable(Hive.settings)(hive) ++ enable(CometTestSettings.settings)(hive) + + enable(SparkConnectCommon.settings)(connectCommon) + enable(SparkConnect.settings)(connect) +@@ -1198,6 +1200,25 @@ object SQL { + ) + } + ++/** ++ * Comet disables itself when its shuffle manager is not registered. When Comet is enabled, make ++ * that shuffle manager the default for every test SparkConf in the modules that have Comet on ++ * their classpath, so that suites which build their own SparkSession or SparkContext run Comet ++ * too, not just the ones that go through SharedSparkSession or TestHive. ++ */ ++object CometTestSettings { ++ private val cometEnabled = ++ sys.env.get("ENABLE_COMET").forall(v => v == "1" || v.equalsIgnoreCase("true")) ++ ++ lazy val settings: Seq[Setting[_]] = ++ if (cometEnabled && !sys.props.contains("spark.shuffle.manager")) { ++ Seq((Test / javaOptions) += ++ "-Dspark.shuffle.manager=org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager") ++ } else { ++ Seq.empty ++ } ++} ++ + object Hive { + + lazy val settings = Seq( diff --git a/sql/core/pom.xml b/sql/core/pom.xml index d7f34db599a..5f676071ec0 100644 --- a/sql/core/pom.xml @@ -1064,6 +1106,50 @@ index 71af1fd69c3..81a04c93c9c 100644 } } } +diff --git a/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionExtensionSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionExtensionSuite.scala +index 8b4ac474f87..5218e9125c6 100644 +--- a/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionExtensionSuite.scala ++++ b/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionExtensionSuite.scala +@@ -167,7 +167,15 @@ class SparkSessionExtensionSuite extends SparkFunSuite with SQLHelper with Adapt + } + } + ++ // Comet: Comet's own columnar rules take over these plans, so the rules the tests below ++ // inject never see the Spark operators they act on or check for. ++ private def assumeCometDisabled(): Unit = { ++ assume(!org.apache.spark.sql.SparkSession.isCometEnabled, ++ "Skipped when Comet is enabled: Comet replaces the operators the injected rules act on") ++ } ++ + test("inject adaptive query prep rule") { ++ assumeCometDisabled() + val extensions = create { extensions => + // inject rule that will run during AQE query stage preparation and will add custom tags + // to the plan +@@ -261,6 +269,7 @@ class SparkSessionExtensionSuite extends SparkFunSuite with SQLHelper with Adapt + } + + private def testInjectColumnar(enableAQE: Boolean): Unit = { ++ assumeCometDisabled() + def collectPlanSteps(plan: SparkPlan): Seq[Int] = plan match { + case a: AdaptiveSparkPlanExec => + assert(a.toString.startsWith("AdaptiveSparkPlan isFinalPlan=true")) +@@ -316,6 +325,7 @@ class SparkSessionExtensionSuite extends SparkFunSuite with SQLHelper with Adapt + } + + test("reset column vectors") { ++ assumeCometDisabled() + val session = SparkSession.builder() + .master("local[1]") + .config(COLUMN_BATCH_SIZE.key, 2) +@@ -484,6 +494,7 @@ class SparkSessionExtensionSuite extends SparkFunSuite with SQLHelper with Adapt + } + + test("SPARK-38697: Extend SparkSessionExtensions to inject rules into AQE Optimizer") { ++ assumeCometDisabled() + def executedPlan(df: Dataset[java.lang.Long]): SparkPlan = { + assert(df.queryExecution.executedPlan.isInstanceOf[AdaptiveSparkPlanExec]) + df.queryExecution.executedPlan.asInstanceOf[AdaptiveSparkPlanExec].executedPlan diff --git a/sql/core/src/test/scala/org/apache/spark/sql/SubquerySuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/SubquerySuite.scala index 04702201f82..4d38d8d6e51 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/SubquerySuite.scala @@ -1300,6 +1386,84 @@ index ae1c0a86a14..1d3b914fd64 100644 // Fail to read ancient datetime values. withSQLConf(SQLConf.PARQUET_REBASE_MODE_IN_READ.key -> EXCEPTION.toString) { +diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/CoalesceShufflePartitionsSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/CoalesceShufflePartitionsSuite.scala +index e11191da6a9..005fd4fa657 100644 +--- a/sql/core/src/test/scala/org/apache/spark/sql/execution/CoalesceShufflePartitionsSuite.scala ++++ b/sql/core/src/test/scala/org/apache/spark/sql/execution/CoalesceShufflePartitionsSuite.scala +@@ -51,6 +51,11 @@ class CoalesceShufflePartitionsSuite extends SparkFunSuite { + } + } + ++ // Comet: AQE coalesces post-shuffle partitions by map output size, and Comet's shuffle ++ // writes Arrow IPC blocks whose sizes differ from Spark's, so the expected partition ++ // counts below only hold for Spark's shuffle. ++ private def isCometEnabled: Boolean = org.apache.spark.sql.SparkSession.isCometEnabled ++ + val numInputPartitions: Int = 10 + + def withSparkSession( +@@ -163,9 +168,11 @@ class CoalesceShufflePartitionsSuite extends SparkFunSuite { + assert(shuffleReads.isEmpty) + + case None => +- assert(shuffleReads.length === 2) +- shuffleReads.foreach { read => +- assert(read.outputPartitioning.numPartitions === 2) ++ if (!isCometEnabled) { ++ assert(shuffleReads.length === 2) ++ shuffleReads.foreach { read => ++ assert(read.outputPartitioning.numPartitions === 2) ++ } + } + } + } +@@ -214,9 +221,11 @@ class CoalesceShufflePartitionsSuite extends SparkFunSuite { + assert(shuffleReads.isEmpty) + + case None => +- assert(shuffleReads.length === 2) +- shuffleReads.foreach { read => +- assert(read.outputPartitioning.numPartitions === 2) ++ if (!isCometEnabled) { ++ assert(shuffleReads.length === 2) ++ shuffleReads.foreach { read => ++ assert(read.outputPartitioning.numPartitions === 2) ++ } + } + } + } +@@ -265,9 +274,11 @@ class CoalesceShufflePartitionsSuite extends SparkFunSuite { + assert(shuffleReads.isEmpty) + + case None => +- assert(shuffleReads.length === 2) +- shuffleReads.foreach { read => +- assert(read.outputPartitioning.numPartitions === 3) ++ if (!isCometEnabled) { ++ assert(shuffleReads.length === 2) ++ shuffleReads.foreach { read => ++ assert(read.outputPartitioning.numPartitions === 3) ++ } + } + } + } +@@ -473,10 +484,12 @@ class CoalesceShufflePartitionsSuite extends SparkFunSuite { + // aggregate on the other side of the union. + val finalPlan = resultDf.queryExecution.executedPlan + .asInstanceOf[AdaptiveSparkPlanExec].executedPlan +- assert( +- finalPlan.collect { +- case r @ CoalescedShuffleRead() => r +- }.size == 2) ++ if (!isCometEnabled) { ++ assert( ++ finalPlan.collect { ++ case r @ CoalescedShuffleRead() => r ++ }.size == 2) ++ } + } + withSparkSession(test, 100, None) + } diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/DataSourceScanExecRedactionSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/DataSourceScanExecRedactionSuite.scala index 418ca3430bb..eb8267192f8 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/execution/DataSourceScanExecRedactionSuite.scala @@ -2429,6 +2593,29 @@ index b8f3ea3c6f3..bbd44221288 100644 withTempPath { workDir => val workDirPath = workDir.getAbsolutePath val input = spark.range(5).toDF("id") +diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/joins/BroadcastJoinSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/joins/BroadcastJoinSuite.scala +index 5479be86e9f..07c81e4b830 100644 +--- a/sql/core/src/test/scala/org/apache/spark/sql/execution/joins/BroadcastJoinSuite.scala ++++ b/sql/core/src/test/scala/org/apache/spark/sql/execution/joins/BroadcastJoinSuite.scala +@@ -21,7 +21,7 @@ import scala.reflect.ClassTag + + import org.apache.spark.AccumulatorSuite + import org.apache.spark.internal.config.EXECUTOR_MEMORY +-import org.apache.spark.sql.{Dataset, QueryTest, Row, SparkSession} ++import org.apache.spark.sql.{Dataset, IgnoreComet, QueryTest, Row, SparkSession} + import org.apache.spark.sql.catalyst.expressions.{AttributeReference, BitwiseAnd, BitwiseOr, Cast, Expression, Literal, ShiftLeft} + import org.apache.spark.sql.catalyst.optimizer.{BuildLeft, BuildRight, BuildSide} + import org.apache.spark.sql.catalyst.plans.Inner +@@ -487,7 +487,8 @@ abstract class BroadcastJoinSuiteBase extends QueryTest with SQLTestUtils + } + } + +- test("broadcast join where streamed side's output partitioning is PartitioningCollection") { ++ test("broadcast join where streamed side's output partitioning is PartitioningCollection", ++ IgnoreComet("Comet replaces the join and shuffle operators this test inspects")) { + withSQLConf(SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "500") { + val t1 = (0 until 100).map(i => (i % 5, i % 13)).toDF("i1", "j1") + val t2 = (0 until 100).map(i => (i % 5, i % 14)).toDF("i2", "j2") diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/metric/SQLMetricsSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/metric/SQLMetricsSuite.scala index 5cdbdc27b32..307fba16578 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/execution/metric/SQLMetricsSuite.scala diff --git a/dev/diffs/4.0.4.diff b/dev/diffs/4.0.4.diff index 38ec78543aa..a4e645e5e5b 100644 --- a/dev/diffs/4.0.4.diff +++ b/dev/diffs/4.0.4.diff @@ -77,6 +77,48 @@ index 046f357ff79..7c6fd199e41 100644 org.apache.datasketches +diff --git a/project/SparkBuild.scala b/project/SparkBuild.scala +index 9c2312ed256..b151f6a630c 100644 +--- a/project/SparkBuild.scala ++++ b/project/SparkBuild.scala +@@ -414,9 +414,11 @@ object SparkBuild extends PomBuild { + + /* Spark SQL Core settings */ + enable(SQL.settings)(sql) ++ enable(CometTestSettings.settings)(sql) + + /* Hive console settings */ + enable(Hive.settings)(hive) ++ enable(CometTestSettings.settings)(hive) + + enable(HiveThriftServer.settings)(hiveThriftServer) + +@@ -1183,6 +1185,25 @@ object SQL { + } + } + ++/** ++ * Comet disables itself when its shuffle manager is not registered. When Comet is enabled, make ++ * that shuffle manager the default for every test SparkConf in the modules that have Comet on ++ * their classpath, so that suites which build their own SparkSession or SparkContext run Comet ++ * too, not just the ones that go through SharedSparkSession or TestHive. ++ */ ++object CometTestSettings { ++ private val cometEnabled = ++ sys.env.get("ENABLE_COMET").forall(v => v == "1" || v.equalsIgnoreCase("true")) ++ ++ lazy val settings: Seq[Setting[_]] = ++ if (cometEnabled && !sys.props.contains("spark.shuffle.manager")) { ++ Seq((Test / javaOptions) += ++ "-Dspark.shuffle.manager=org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager") ++ } else { ++ Seq.empty ++ } ++} ++ + object Hive { + + lazy val settings = Seq( diff --git a/sql/core/pom.xml b/sql/core/pom.xml index b7b87c013e0..6fcf06859ab 100644 --- a/sql/core/pom.xml @@ -1201,6 +1243,50 @@ index 575a4ae69d1..129d9f27232 100644 } } } +diff --git a/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionExtensionSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionExtensionSuite.scala +index c1c041509c3..d068241d585 100644 +--- a/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionExtensionSuite.scala ++++ b/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionExtensionSuite.scala +@@ -179,7 +179,15 @@ class SparkSessionExtensionSuite extends SparkFunSuite with SQLHelper with Adapt + } + } + ++ // Comet: Comet's own columnar rules take over these plans, so the rules the tests below ++ // inject never see the Spark operators they act on or check for. ++ private def assumeCometDisabled(): Unit = { ++ assume(!org.apache.spark.sql.classic.SparkSession.isCometEnabled, ++ "Skipped when Comet is enabled: Comet replaces the operators the injected rules act on") ++ } ++ + test("inject adaptive query prep rule") { ++ assumeCometDisabled() + val extensions = create { extensions => + // inject rule that will run during AQE query stage preparation and will add custom tags + // to the plan +@@ -273,6 +281,7 @@ class SparkSessionExtensionSuite extends SparkFunSuite with SQLHelper with Adapt + } + + private def testInjectColumnar(enableAQE: Boolean): Unit = { ++ assumeCometDisabled() + def collectPlanSteps(plan: SparkPlan): Seq[Int] = plan match { + case a: AdaptiveSparkPlanExec => + assert(a.toString.startsWith("AdaptiveSparkPlan isFinalPlan=true")) +@@ -328,6 +337,7 @@ class SparkSessionExtensionSuite extends SparkFunSuite with SQLHelper with Adapt + } + + test("reset column vectors") { ++ assumeCometDisabled() + val session = SparkSession.builder() + .master("local[1]") + .config(COLUMN_BATCH_SIZE.key, 2) +@@ -496,6 +506,7 @@ class SparkSessionExtensionSuite extends SparkFunSuite with SQLHelper with Adapt + } + + test("SPARK-38697: Extend SparkSessionExtensions to inject rules into AQE Optimizer") { ++ assumeCometDisabled() + def executedPlan(df: Dataset[java.lang.Long]): SparkPlan = { + assert(df.queryExecution.executedPlan.isInstanceOf[AdaptiveSparkPlanExec]) + stripAQEPlan(df.queryExecution.executedPlan) diff --git a/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionJobTaggingAndCancellationSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionJobTaggingAndCancellationSuite.scala index 5ba69c8f9d9..ac1256afe88 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionJobTaggingAndCancellationSuite.scala @@ -1663,6 +1749,84 @@ index 04d33ecd3d5..450df347297 100644 // Fail to read ancient datetime values. withSQLConf(SQLConf.PARQUET_REBASE_MODE_IN_READ.key -> EXCEPTION.toString) { +diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/CoalesceShufflePartitionsSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/CoalesceShufflePartitionsSuite.scala +index da43b0cfc58..ca1f6f47a9f 100644 +--- a/sql/core/src/test/scala/org/apache/spark/sql/execution/CoalesceShufflePartitionsSuite.scala ++++ b/sql/core/src/test/scala/org/apache/spark/sql/execution/CoalesceShufflePartitionsSuite.scala +@@ -54,6 +54,11 @@ class CoalesceShufflePartitionsSuite extends SparkFunSuite with SQLConfHelper + } + } + ++ // Comet: AQE coalesces post-shuffle partitions by map output size, and Comet's shuffle ++ // writes Arrow IPC blocks whose sizes differ from Spark's, so the expected partition ++ // counts below only hold for Spark's shuffle. ++ private def isCometEnabled: Boolean = org.apache.spark.sql.classic.SparkSession.isCometEnabled ++ + val numInputPartitions: Int = 10 + + def withSparkSession( +@@ -164,9 +169,11 @@ class CoalesceShufflePartitionsSuite extends SparkFunSuite with SQLConfHelper + assert(shuffleReads.isEmpty) + + case None => +- assert(shuffleReads.length === 2) +- shuffleReads.foreach { read => +- assert(read.outputPartitioning.numPartitions === 2) ++ if (!isCometEnabled) { ++ assert(shuffleReads.length === 2) ++ shuffleReads.foreach { read => ++ assert(read.outputPartitioning.numPartitions === 2) ++ } + } + } + } +@@ -214,9 +221,11 @@ class CoalesceShufflePartitionsSuite extends SparkFunSuite with SQLConfHelper + assert(shuffleReads.isEmpty) + + case None => +- assert(shuffleReads.length === 2) +- shuffleReads.foreach { read => +- assert(read.outputPartitioning.numPartitions === 2) ++ if (!isCometEnabled) { ++ assert(shuffleReads.length === 2) ++ shuffleReads.foreach { read => ++ assert(read.outputPartitioning.numPartitions === 2) ++ } + } + } + } +@@ -264,9 +273,11 @@ class CoalesceShufflePartitionsSuite extends SparkFunSuite with SQLConfHelper + assert(shuffleReads.isEmpty) + + case None => +- assert(shuffleReads.length === 2) +- shuffleReads.foreach { read => +- assert(read.outputPartitioning.numPartitions === 3) ++ if (!isCometEnabled) { ++ assert(shuffleReads.length === 2) ++ shuffleReads.foreach { read => ++ assert(read.outputPartitioning.numPartitions === 3) ++ } + } + } + } +@@ -468,10 +479,12 @@ class CoalesceShufflePartitionsSuite extends SparkFunSuite with SQLConfHelper + // Shuffle partition coalescing of the join is performed independent of the non-grouping + // aggregate on the other side of the union. + val finalPlan = stripAQEPlan(resultDf.queryExecution.executedPlan) +- assert( +- finalPlan.collect { +- case r @ CoalescedShuffleRead() => r +- }.size == 2) ++ if (!isCometEnabled) { ++ assert( ++ finalPlan.collect { ++ case r @ CoalescedShuffleRead() => r ++ }.size == 2) ++ } + } + withSparkSession(test, 100, None) + } diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/DataSourceScanExecRedactionSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/DataSourceScanExecRedactionSuite.scala index 418ca3430bb..eb8267192f8 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/execution/DataSourceScanExecRedactionSuite.scala @@ -3154,6 +3318,29 @@ index b8f3ea3c6f3..bbd44221288 100644 withTempPath { workDir => val workDirPath = workDir.getAbsolutePath val input = spark.range(5).toDF("id") +diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/joins/BroadcastJoinSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/joins/BroadcastJoinSuite.scala +index 69dd04e07d5..781018ecc66 100644 +--- a/sql/core/src/test/scala/org/apache/spark/sql/execution/joins/BroadcastJoinSuite.scala ++++ b/sql/core/src/test/scala/org/apache/spark/sql/execution/joins/BroadcastJoinSuite.scala +@@ -21,7 +21,7 @@ import scala.reflect.ClassTag + + import org.apache.spark.AccumulatorSuite + import org.apache.spark.internal.config.EXECUTOR_MEMORY +-import org.apache.spark.sql.{QueryTest, Row} ++import org.apache.spark.sql.{IgnoreComet, QueryTest, Row} + import org.apache.spark.sql.catalyst.expressions.{AttributeReference, BitwiseAnd, BitwiseOr, Cast, Expression, Literal, ShiftLeft} + import org.apache.spark.sql.catalyst.optimizer.{BuildLeft, BuildRight, BuildSide} + import org.apache.spark.sql.catalyst.plans.Inner +@@ -488,7 +488,8 @@ abstract class BroadcastJoinSuiteBase extends QueryTest with SQLTestUtils + } + } + +- test("broadcast join where streamed side's output partitioning is PartitioningCollection") { ++ test("broadcast join where streamed side's output partitioning is PartitioningCollection", ++ IgnoreComet("Comet replaces the join and shuffle operators this test inspects")) { + withSQLConf(SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "500") { + val t1 = (0 until 100).map(i => (i % 5, i % 13)).toDF("i1", "j1") + val t2 = (0 until 100).map(i => (i % 5, i % 14)).toDF("i2", "j2") diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/metric/SQLMetricsSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/metric/SQLMetricsSuite.scala index 0dd90925d3c..7d53ec845ef 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/execution/metric/SQLMetricsSuite.scala diff --git a/dev/diffs/4.1.3.diff b/dev/diffs/4.1.3.diff index a8d9383ced1..f7b87ee7c8b 100644 --- a/dev/diffs/4.1.3.diff +++ b/dev/diffs/4.1.3.diff @@ -77,6 +77,48 @@ index 370bfa4397f..ccb4398de58 100644 org.apache.datasketches +diff --git a/project/SparkBuild.scala b/project/SparkBuild.scala +index 01aa71047fe..cde29aabc83 100644 +--- a/project/SparkBuild.scala ++++ b/project/SparkBuild.scala +@@ -443,9 +443,11 @@ object SparkBuild extends PomBuild { + + /* Spark SQL Core settings */ + enable(SQL.settings)(sql) ++ enable(CometTestSettings.settings)(sql) + + /* Hive console settings */ + enable(Hive.settings)(hive) ++ enable(CometTestSettings.settings)(hive) + + enable(HiveThriftServer.settings)(hiveThriftServer) + +@@ -1347,6 +1349,25 @@ object SQL { + } + } + ++/** ++ * Comet disables itself when its shuffle manager is not registered. When Comet is enabled, make ++ * that shuffle manager the default for every test SparkConf in the modules that have Comet on ++ * their classpath, so that suites which build their own SparkSession or SparkContext run Comet ++ * too, not just the ones that go through SharedSparkSession or TestHive. ++ */ ++object CometTestSettings { ++ private val cometEnabled = ++ sys.env.get("ENABLE_COMET").forall(v => v == "1" || v.equalsIgnoreCase("true")) ++ ++ lazy val settings: Seq[Setting[_]] = ++ if (cometEnabled && !sys.props.contains("spark.shuffle.manager")) { ++ Seq((Test / javaOptions) += ++ "-Dspark.shuffle.manager=org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager") ++ } else { ++ Seq.empty ++ } ++} ++ + object Hive { + + lazy val settings = Seq( diff --git a/sql/core/pom.xml b/sql/core/pom.xml index febc7d747bc..63c8e4e22c0 100644 --- a/sql/core/pom.xml @@ -1219,6 +1261,20 @@ index 885512d4d19..09b1ccaed71 100644 withSQLConf( SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "1") { +diff --git a/sql/core/src/test/scala/org/apache/spark/sql/MapStatusEndToEndSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/MapStatusEndToEndSuite.scala +index abcd346c327..4f305ffc2ab 100644 +--- a/sql/core/src/test/scala/org/apache/spark/sql/MapStatusEndToEndSuite.scala ++++ b/sql/core/src/test/scala/org/apache/spark/sql/MapStatusEndToEndSuite.scala +@@ -37,7 +37,8 @@ class MapStatusEndToEndSuite extends SparkFunSuite with SQLTestUtils { + SparkSession.clearDefaultSession() + } + +- test("Propagate checksum from executor to driver") { ++ test("Propagate checksum from executor to driver", ++ IgnoreComet("https://github.com/apache/datafusion-comet/issues/6414")) { + assert(spark.sparkContext.conf.get(SQLConf.LEAF_NODE_DEFAULT_PARALLELISM.key) == "5") + assert(spark.conf.get(SQLConf.LEAF_NODE_DEFAULT_PARALLELISM.key) == "5") + assert(spark.sparkContext.conf.get(SQLConf.CLASSIC_SHUFFLE_DEPENDENCY_FILE_CLEANUP_ENABLED.key) diff --git a/sql/core/src/test/scala/org/apache/spark/sql/PlanStabilitySuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/PlanStabilitySuite.scala index e4b5e10f7c3..c6efde09c8a 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/PlanStabilitySuite.scala @@ -1321,6 +1377,50 @@ index 23f0144dcec..40d536bb23a 100644 } } } +diff --git a/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionExtensionSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionExtensionSuite.scala +index 66826a9ca76..a330d142eb1 100644 +--- a/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionExtensionSuite.scala ++++ b/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionExtensionSuite.scala +@@ -196,7 +196,15 @@ class SparkSessionExtensionSuite extends SparkFunSuite with SQLHelper with Adapt + } + } + ++ // Comet: Comet's own columnar rules take over these plans, so the rules the tests below ++ // inject never see the Spark operators they act on or check for. ++ private def assumeCometDisabled(): Unit = { ++ assume(!org.apache.spark.sql.classic.SparkSession.isCometEnabled, ++ "Skipped when Comet is enabled: Comet replaces the operators the injected rules act on") ++ } ++ + test("inject adaptive query prep rule") { ++ assumeCometDisabled() + val extensions = create { extensions => + // inject rule that will run during AQE query stage preparation and will add custom tags + // to the plan +@@ -290,6 +298,7 @@ class SparkSessionExtensionSuite extends SparkFunSuite with SQLHelper with Adapt + } + + private def testInjectColumnar(enableAQE: Boolean): Unit = { ++ assumeCometDisabled() + def collectPlanSteps(plan: SparkPlan): Seq[Int] = plan match { + case a: AdaptiveSparkPlanExec => + assert(a.toString.startsWith("AdaptiveSparkPlan isFinalPlan=true")) +@@ -345,6 +354,7 @@ class SparkSessionExtensionSuite extends SparkFunSuite with SQLHelper with Adapt + } + + test("reset column vectors") { ++ assumeCometDisabled() + val session = SparkSession.builder() + .master("local[1]") + .config(COLUMN_BATCH_SIZE.key, 2) +@@ -513,6 +523,7 @@ class SparkSessionExtensionSuite extends SparkFunSuite with SQLHelper with Adapt + } + + test("SPARK-38697: Extend SparkSessionExtensions to inject rules into AQE Optimizer") { ++ assumeCometDisabled() + def executedPlan(df: Dataset[java.lang.Long]): SparkPlan = { + assert(df.queryExecution.executedPlan.isInstanceOf[AdaptiveSparkPlanExec]) + stripAQEPlan(df.queryExecution.executedPlan) diff --git a/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionJobTaggingAndCancellationSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionJobTaggingAndCancellationSuite.scala index d7b2511eac2..d5f5b940b94 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionJobTaggingAndCancellationSuite.scala @@ -1783,6 +1883,84 @@ index fcecaf25d4c..e5a511022cc 100644 // Fail to read ancient datetime values. withSQLConf(SQLConf.PARQUET_REBASE_MODE_IN_READ.key -> EXCEPTION.toString) { +diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/CoalesceShufflePartitionsSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/CoalesceShufflePartitionsSuite.scala +index 28762f01d7a..6cdc3953bf3 100644 +--- a/sql/core/src/test/scala/org/apache/spark/sql/execution/CoalesceShufflePartitionsSuite.scala ++++ b/sql/core/src/test/scala/org/apache/spark/sql/execution/CoalesceShufflePartitionsSuite.scala +@@ -56,6 +56,11 @@ class CoalesceShufflePartitionsSuite extends SparkFunSuite with SQLConfHelper + } + } + ++ // Comet: AQE coalesces post-shuffle partitions by map output size, and Comet's shuffle ++ // writes Arrow IPC blocks whose sizes differ from Spark's, so the expected partition ++ // counts below only hold for Spark's shuffle. ++ private def isCometEnabled: Boolean = org.apache.spark.sql.classic.SparkSession.isCometEnabled ++ + val numInputPartitions: Int = 10 + + def withSparkSession( +@@ -166,9 +171,11 @@ class CoalesceShufflePartitionsSuite extends SparkFunSuite with SQLConfHelper + assert(shuffleReads.isEmpty) + + case None => +- assert(shuffleReads.length === 2) +- shuffleReads.foreach { read => +- assert(read.outputPartitioning.numPartitions === 2) ++ if (!isCometEnabled) { ++ assert(shuffleReads.length === 2) ++ shuffleReads.foreach { read => ++ assert(read.outputPartitioning.numPartitions === 2) ++ } + } + } + } +@@ -221,9 +228,11 @@ class CoalesceShufflePartitionsSuite extends SparkFunSuite with SQLConfHelper + assert(shuffleReads.isEmpty) + + case None => +- assert(shuffleReads.length === 2) +- shuffleReads.foreach { read => +- assert(read.outputPartitioning.numPartitions === 2) ++ if (!isCometEnabled) { ++ assert(shuffleReads.length === 2) ++ shuffleReads.foreach { read => ++ assert(read.outputPartitioning.numPartitions === 2) ++ } + } + } + } +@@ -271,9 +280,11 @@ class CoalesceShufflePartitionsSuite extends SparkFunSuite with SQLConfHelper + assert(shuffleReads.isEmpty) + + case None => +- assert(shuffleReads.length === 2) +- shuffleReads.foreach { read => +- assert(read.outputPartitioning.numPartitions === 3) ++ if (!isCometEnabled) { ++ assert(shuffleReads.length === 2) ++ shuffleReads.foreach { read => ++ assert(read.outputPartitioning.numPartitions === 3) ++ } + } + } + } +@@ -475,10 +486,12 @@ class CoalesceShufflePartitionsSuite extends SparkFunSuite with SQLConfHelper + // Shuffle partition coalescing of the join is performed independent of the non-grouping + // aggregate on the other side of the union. + val finalPlan = stripAQEPlan(resultDf.queryExecution.executedPlan) +- assert( +- finalPlan.collect { +- case r @ CoalescedShuffleRead() => r +- }.size == 2) ++ if (!isCometEnabled) { ++ assert( ++ finalPlan.collect { ++ case r @ CoalescedShuffleRead() => r ++ }.size == 2) ++ } + } + withSparkSession(test, 100, None) + } diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/DataSourceScanExecRedactionSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/DataSourceScanExecRedactionSuite.scala index 418ca3430bb..eb8267192f8 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/execution/DataSourceScanExecRedactionSuite.scala @@ -3348,6 +3526,29 @@ index b8f3ea3c6f3..bbd44221288 100644 withTempPath { workDir => val workDirPath = workDir.getAbsolutePath val input = spark.range(5).toDF("id") +diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/joins/BroadcastJoinSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/joins/BroadcastJoinSuite.scala +index 9bd858608cb..2682ba53513 100644 +--- a/sql/core/src/test/scala/org/apache/spark/sql/execution/joins/BroadcastJoinSuite.scala ++++ b/sql/core/src/test/scala/org/apache/spark/sql/execution/joins/BroadcastJoinSuite.scala +@@ -21,7 +21,7 @@ import scala.reflect.ClassTag + + import org.apache.spark.AccumulatorSuite + import org.apache.spark.internal.config.EXECUTOR_MEMORY +-import org.apache.spark.sql.{QueryTest, Row} ++import org.apache.spark.sql.{IgnoreComet, QueryTest, Row} + import org.apache.spark.sql.catalyst.expressions.{AttributeReference, BitwiseAnd, BitwiseOr, Cast, Expression, Literal, ShiftLeft} + import org.apache.spark.sql.catalyst.optimizer.{BuildLeft, BuildRight, BuildSide} + import org.apache.spark.sql.catalyst.plans.Inner +@@ -488,7 +488,8 @@ abstract class BroadcastJoinSuiteBase extends QueryTest with SQLTestUtils + } + } + +- test("broadcast join where streamed side's output partitioning is PartitioningCollection") { ++ test("broadcast join where streamed side's output partitioning is PartitioningCollection", ++ IgnoreComet("Comet replaces the join and shuffle operators this test inspects")) { + withSQLConf(SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "500") { + val t1 = (0 until 100).map(i => (i % 5, i % 13)).toDF("i1", "j1") + val t2 = (0 until 100).map(i => (i % 5, i % 14)).toDF("i2", "j2") diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/metric/SQLMetricsSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/metric/SQLMetricsSuite.scala index f2e9121d566..2c9f517034f 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/execution/metric/SQLMetricsSuite.scala diff --git a/dev/diffs/4.2.0.diff b/dev/diffs/4.2.0.diff index b129e60847d..5b677814edf 100644 --- a/dev/diffs/4.2.0.diff +++ b/dev/diffs/4.2.0.diff @@ -77,6 +77,48 @@ index 46558134f41..862c9a6eb9e 100644 org.apache.datasketches +diff --git a/project/SparkBuild.scala b/project/SparkBuild.scala +index 7c5e49f3a69..28f83381fa0 100644 +--- a/project/SparkBuild.scala ++++ b/project/SparkBuild.scala +@@ -455,9 +455,11 @@ object SparkBuild extends PomBuild { + + /* Spark SQL Core settings */ + enable(SQL.settings)(sql) ++ enable(CometTestSettings.settings)(sql) + + /* Hive console settings */ + enable(Hive.settings)(hive) ++ enable(CometTestSettings.settings)(hive) + + enable(HiveThriftServer.settings)(hiveThriftServer) + +@@ -1426,6 +1428,25 @@ object SQL { + } + } + ++/** ++ * Comet disables itself when its shuffle manager is not registered. When Comet is enabled, make ++ * that shuffle manager the default for every test SparkConf in the modules that have Comet on ++ * their classpath, so that suites which build their own SparkSession or SparkContext run Comet ++ * too, not just the ones that go through SharedSparkSession or TestHive. ++ */ ++object CometTestSettings { ++ private val cometEnabled = ++ sys.env.get("ENABLE_COMET").forall(v => v == "1" || v.equalsIgnoreCase("true")) ++ ++ lazy val settings: Seq[Setting[_]] = ++ if (cometEnabled && !sys.props.contains("spark.shuffle.manager")) { ++ Seq((Test / javaOptions) += ++ "-Dspark.shuffle.manager=org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager") ++ } else { ++ Seq.empty ++ } ++} ++ + object Hive { + + lazy val settings = Seq( diff --git a/sql/core/pom.xml b/sql/core/pom.xml index 82810f181ac..21a83831188 100644 --- a/sql/core/pom.xml @@ -1248,6 +1290,20 @@ index 7f695e90df8..57fea6f8324 100644 withSQLConf( SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "1") { +diff --git a/sql/core/src/test/scala/org/apache/spark/sql/MapStatusEndToEndSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/MapStatusEndToEndSuite.scala +index 0708b2c9392..d6ff2e1555c 100644 +--- a/sql/core/src/test/scala/org/apache/spark/sql/MapStatusEndToEndSuite.scala ++++ b/sql/core/src/test/scala/org/apache/spark/sql/MapStatusEndToEndSuite.scala +@@ -36,7 +36,8 @@ class MapStatusEndToEndSuite extends QueryTest { + SparkSession.clearDefaultSession() + } + +- test("Propagate checksum from executor to driver") { ++ test("Propagate checksum from executor to driver", ++ IgnoreComet("https://github.com/apache/datafusion-comet/issues/6414")) { + assert(spark.sparkContext.conf.get(SQLConf.LEAF_NODE_DEFAULT_PARALLELISM.key) == "5") + assert(spark.conf.get(SQLConf.LEAF_NODE_DEFAULT_PARALLELISM.key) == "5") + assert(spark.sparkContext.conf.get(SQLConf.CLASSIC_SHUFFLE_DEPENDENCY_FILE_CLEANUP_ENABLED.key) diff --git a/sql/core/src/test/scala/org/apache/spark/sql/PlanStabilitySuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/PlanStabilitySuite.scala index 6cd49948630..b6fcb716a2a 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/PlanStabilitySuite.scala @@ -1402,6 +1458,50 @@ index 395cb67f441..33ac6ed19af 100644 } } } +diff --git a/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionExtensionSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionExtensionSuite.scala +index bfcf583a705..6c177bff0e3 100644 +--- a/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionExtensionSuite.scala ++++ b/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionExtensionSuite.scala +@@ -217,7 +217,15 @@ class SparkSessionExtensionSuite extends PlanTest with SQLHelper with AdaptiveSp + } + } + ++ // Comet: Comet's own columnar rules take over these plans, so the rules the tests below ++ // inject never see the Spark operators they act on or check for. ++ private def assumeCometDisabled(): Unit = { ++ assume(!org.apache.spark.sql.classic.SparkSession.isCometEnabled, ++ "Skipped when Comet is enabled: Comet replaces the operators the injected rules act on") ++ } ++ + test("inject adaptive query prep rule") { ++ assumeCometDisabled() + val extensions = create { extensions => + // inject rule that will run during AQE query stage preparation and will add custom tags + // to the plan +@@ -311,6 +319,7 @@ class SparkSessionExtensionSuite extends PlanTest with SQLHelper with AdaptiveSp + } + + private def testInjectColumnar(enableAQE: Boolean): Unit = { ++ assumeCometDisabled() + def collectPlanSteps(plan: SparkPlan): Seq[Int] = plan match { + case a: AdaptiveSparkPlanExec => + assert(a.toString.startsWith("AdaptiveSparkPlan isFinalPlan=true")) +@@ -366,6 +375,7 @@ class SparkSessionExtensionSuite extends PlanTest with SQLHelper with AdaptiveSp + } + + test("reset column vectors") { ++ assumeCometDisabled() + val session = SparkSession.builder() + .master("local[1]") + .config(COLUMN_BATCH_SIZE.key, 2) +@@ -553,6 +563,7 @@ class SparkSessionExtensionSuite extends PlanTest with SQLHelper with AdaptiveSp + } + + test("SPARK-38697: Extend SparkSessionExtensions to inject rules into AQE Optimizer") { ++ assumeCometDisabled() + def executedPlan(df: Dataset[java.lang.Long]): SparkPlan = { + assert(df.queryExecution.executedPlan.isInstanceOf[AdaptiveSparkPlanExec]) + stripAQEPlan(df.queryExecution.executedPlan) diff --git a/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionJobTaggingAndCancellationSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionJobTaggingAndCancellationSuite.scala index d7b2511eac2..d5f5b940b94 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/SparkSessionJobTaggingAndCancellationSuite.scala @@ -1850,6 +1950,84 @@ index ce549da03b4..47adf7cf94d 100644 // Fail to read ancient datetime values. withSQLConf(SQLConf.PARQUET_REBASE_MODE_IN_READ.key -> EXCEPTION.toString) { +diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/CoalesceShufflePartitionsSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/CoalesceShufflePartitionsSuite.scala +index 3ef22ccb77e..54f67b45387 100644 +--- a/sql/core/src/test/scala/org/apache/spark/sql/execution/CoalesceShufflePartitionsSuite.scala ++++ b/sql/core/src/test/scala/org/apache/spark/sql/execution/CoalesceShufflePartitionsSuite.scala +@@ -56,6 +56,11 @@ class CoalesceShufflePartitionsSuite extends SparkFunSuite with SQLConfHelper + } + } + ++ // Comet: AQE coalesces post-shuffle partitions by map output size, and Comet's shuffle ++ // writes Arrow IPC blocks whose sizes differ from Spark's, so the expected partition ++ // counts below only hold for Spark's shuffle. ++ private def isCometEnabled: Boolean = org.apache.spark.sql.classic.SparkSession.isCometEnabled ++ + val numInputPartitions: Int = 10 + + def withSparkSession( +@@ -166,9 +171,11 @@ class CoalesceShufflePartitionsSuite extends SparkFunSuite with SQLConfHelper + assert(shuffleReads.isEmpty) + + case None => +- assert(shuffleReads.length === 2) +- shuffleReads.foreach { read => +- assert(read.outputPartitioning.numPartitions === 2) ++ if (!isCometEnabled) { ++ assert(shuffleReads.length === 2) ++ shuffleReads.foreach { read => ++ assert(read.outputPartitioning.numPartitions === 2) ++ } + } + } + } +@@ -221,9 +228,11 @@ class CoalesceShufflePartitionsSuite extends SparkFunSuite with SQLConfHelper + assert(shuffleReads.isEmpty) + + case None => +- assert(shuffleReads.length === 2) +- shuffleReads.foreach { read => +- assert(read.outputPartitioning.numPartitions === 2) ++ if (!isCometEnabled) { ++ assert(shuffleReads.length === 2) ++ shuffleReads.foreach { read => ++ assert(read.outputPartitioning.numPartitions === 2) ++ } + } + } + } +@@ -271,9 +280,11 @@ class CoalesceShufflePartitionsSuite extends SparkFunSuite with SQLConfHelper + assert(shuffleReads.isEmpty) + + case None => +- assert(shuffleReads.length === 2) +- shuffleReads.foreach { read => +- assert(read.outputPartitioning.numPartitions === 3) ++ if (!isCometEnabled) { ++ assert(shuffleReads.length === 2) ++ shuffleReads.foreach { read => ++ assert(read.outputPartitioning.numPartitions === 3) ++ } + } + } + } +@@ -530,10 +541,12 @@ class CoalesceShufflePartitionsSuite extends SparkFunSuite with SQLConfHelper + // Shuffle partition coalescing of the join is performed independent of the non-grouping + // aggregate on the other side of the union. + val finalPlan = stripAQEPlan(resultDf.queryExecution.executedPlan) +- assert( +- finalPlan.collect { +- case r @ CoalescedShuffleRead() => r +- }.size == 2) ++ if (!isCometEnabled) { ++ assert( ++ finalPlan.collect { ++ case r @ CoalescedShuffleRead() => r ++ }.size == 2) ++ } + } + withSparkSession(test, 100, None) + } diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/DataSourceScanExecRedactionSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/DataSourceScanExecRedactionSuite.scala index 60bf3a21e96..7693107f107 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/execution/DataSourceScanExecRedactionSuite.scala @@ -3429,6 +3607,29 @@ index b8f3ea3c6f3..bbd44221288 100644 withTempPath { workDir => val workDirPath = workDir.getAbsolutePath val input = spark.range(5).toDF("id") +diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/joins/BroadcastJoinSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/joins/BroadcastJoinSuite.scala +index 52746720eba..15df5e78440 100644 +--- a/sql/core/src/test/scala/org/apache/spark/sql/execution/joins/BroadcastJoinSuite.scala ++++ b/sql/core/src/test/scala/org/apache/spark/sql/execution/joins/BroadcastJoinSuite.scala +@@ -21,7 +21,7 @@ import scala.reflect.ClassTag + + import org.apache.spark.AccumulatorSuite + import org.apache.spark.internal.config.EXECUTOR_MEMORY +-import org.apache.spark.sql.{QueryTest, Row} ++import org.apache.spark.sql.{IgnoreComet, QueryTest, Row} + import org.apache.spark.sql.catalyst.expressions.{AttributeReference, BitwiseAnd, BitwiseOr, Cast, Expression, Literal, ShiftLeft} + import org.apache.spark.sql.catalyst.optimizer.{BuildLeft, BuildRight, BuildSide} + import org.apache.spark.sql.catalyst.plans.Inner +@@ -487,7 +487,8 @@ abstract class BroadcastJoinSuiteBase extends QueryTest + } + } + +- test("broadcast join where streamed side's output partitioning is PartitioningCollection") { ++ test("broadcast join where streamed side's output partitioning is PartitioningCollection", ++ IgnoreComet("Comet replaces the join and shuffle operators this test inspects")) { + withSQLConf(SQLConf.AUTO_BROADCASTJOIN_THRESHOLD.key -> "500") { + val t1 = (0 until 100).map(i => (i % 5, i % 13)).toDF("i1", "j1") + val t2 = (0 until 100).map(i => (i % 5, i % 14)).toDF("i2", "j2") diff --git a/sql/core/src/test/scala/org/apache/spark/sql/execution/metric/MetricsFailureInjectionSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/execution/metric/MetricsFailureInjectionSuite.scala index 2ae2ea6339f..30c8959a057 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/execution/metric/MetricsFailureInjectionSuite.scala diff --git a/dev/diffs/iceberg/1.11.0.diff b/dev/diffs/iceberg/1.11.0.diff index ece8b5bbf98..14133360d12 100644 --- a/dev/diffs/iceberg/1.11.0.diff +++ b/dev/diffs/iceberg/1.11.0.diff @@ -74,6 +74,28 @@ index b5d6415763..0763a723d3 100644 } protected void createTableWithDeleteGranularity( +diff --git a/spark/v4.1/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestPartitionedWritesToWapBranch.java b/spark/v4.1/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestPartitionedWritesToWapBranch.java +index af065451ab..ec7cdc2edb 100644 +--- a/spark/v4.1/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestPartitionedWritesToWapBranch.java ++++ b/spark/v4.1/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestPartitionedWritesToWapBranch.java +@@ -69,6 +69,17 @@ public class TestPartitionedWritesToWapBranch extends PartitionedWritesTestBase + .config("spark.sql.shuffle.partitions", "4") + .config("spark.sql.hive.metastorePartitionPruningFallbackOnException", "true") + .config("spark.sql.legacy.respectNullabilityInTextDatasetConversion", "true") ++ .config("spark.plugins", "org.apache.spark.CometPlugin") ++ .config( ++ "spark.shuffle.manager", ++ "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager") ++ .config("spark.comet.explainFallback.enabled", "true") ++ .config("spark.comet.scan.icebergNative.enabled", "true") ++ .config("spark.comet.write.iceberg.splitOperator.enabled", "true") ++ .config("spark.comet.iceberg.write.enabled", "true") ++ .config("spark.comet.exec.localTableScan.enabled", "true") ++ .config("spark.memory.offHeap.enabled", "true") ++ .config("spark.memory.offHeap.size", "10g") + .config(TestBase.DISABLE_UI) + .enableHiveSupport() + .getOrCreate(); diff --git a/spark/v4.1/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestSystemFunctionPushDownInRowLevelOperations.java b/spark/v4.1/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestSystemFunctionPushDownInRowLevelOperations.java index 934220e5d3..92132d4b8d 100644 --- a/spark/v4.1/spark-extensions/src/test/java/org/apache/iceberg/spark/extensions/TestSystemFunctionPushDownInRowLevelOperations.java diff --git a/docs/source/contributor-guide/spark-sql-tests.md b/docs/source/contributor-guide/spark-sql-tests.md index 04543caf9d1..c4c39764a37 100644 --- a/docs/source/contributor-guide/spark-sql-tests.md +++ b/docs/source/contributor-guide/spark-sql-tests.md @@ -29,6 +29,10 @@ Here is an overview of the changes that we need to make to Spark: - Modify SparkSession to load the Comet extension - Modify TestHive to load Comet - Modify SQLTestUtilsBase to load Comet when `ENABLE_COMET` environment variable exists +- Modify `project/SparkBuild.scala` so that, when Comet is enabled, the `sql` and `hive` test JVMs use + `CometShuffleManager` as the default shuffle manager. Comet disables itself when its shuffle manager is not + registered, so without this default the suites that build their own `SparkSession` or `SparkContext` would + run without Comet Here are the steps involved in running the Spark SQL tests with Comet, using Spark 3.4.3 for this example.