diff --git a/.github/workflows/pr_build_linux.yml b/.github/workflows/pr_build_linux.yml index 0b0a69a1a75..9e51cf03e22 100644 --- a/.github/workflows/pr_build_linux.yml +++ b/.github/workflows/pr_build_linux.yml @@ -518,6 +518,7 @@ jobs: org.apache.comet.exec.CometInMemoryCacheKryoSuite org.apache.comet.exec.CometInMemoryCacheKryoUnregisteredSuite org.apache.comet.exec.CometInMemoryCacheKryoClassesToRegisterSuite + org.apache.comet.exec.CometInMemoryCacheLastAttemptMetricSuite org.apache.comet.exec.CometGenerateExecSuite org.apache.comet.exec.CometWindowExecSuite org.apache.comet.exec.CometJoinSuite diff --git a/.github/workflows/pr_build_macos.yml b/.github/workflows/pr_build_macos.yml index 5d91aeb53fe..5d759d0b647 100644 --- a/.github/workflows/pr_build_macos.yml +++ b/.github/workflows/pr_build_macos.yml @@ -222,6 +222,7 @@ jobs: org.apache.comet.exec.CometInMemoryCacheKryoSuite org.apache.comet.exec.CometInMemoryCacheKryoUnregisteredSuite org.apache.comet.exec.CometInMemoryCacheKryoClassesToRegisterSuite + org.apache.comet.exec.CometInMemoryCacheLastAttemptMetricSuite org.apache.comet.exec.CometGenerateExecSuite org.apache.comet.exec.CometWindowExecSuite org.apache.comet.exec.CometJoinSuite diff --git a/spark/src/main/scala/org/apache/spark/sql/comet/CometInMemoryTableScanExec.scala b/spark/src/main/scala/org/apache/spark/sql/comet/CometInMemoryTableScanExec.scala index ac1ca3854b7..207513f0a77 100644 --- a/spark/src/main/scala/org/apache/spark/sql/comet/CometInMemoryTableScanExec.scala +++ b/spark/src/main/scala/org/apache/spark/sql/comet/CometInMemoryTableScanExec.scala @@ -86,6 +86,15 @@ case class CometInMemoryTableScanExec( // ExtendedExplainInfo.executionInnerChildren leaves them out of Comet's own reporting. override def innerChildren: Seq[QueryPlan[_]] = Seq(originalPlan.relation) + // SparkPlanInfo, behind the SQL tab's graph and the event log, draws Spark's own scan (matched + // by class) with its cached plan below, and any other node with its children and subqueries. + // So expose that scan as the one subquery. It never runs, since subqueries run from + // expressions, and other walkers of subqueries stop at it, as at Spark's scan in Spark's own + // plans. Exposing the cached plan instead lets them in: Spark 4.2's last-attempt metrics then + // give up on its Comet shuffle. innerChildren above must not add super.innerChildren, which is + // this list. A lazy val overrides both Spark 3.x's lazy val and Spark 4's def. + @transient override lazy val subqueries: Seq[SparkPlan] = Seq(originalPlan) + // `originalPlan` is a plan-typed field rather than a child, so QueryPlan's canonicalization // walks straight past it: its attributes and predicates keep the expression IDs of whichever // occurrence of the cached relation produced them. Two scans of one cache then compare unequal, diff --git a/spark/src/test/scala/org/apache/comet/exec/CometInMemoryCacheSuite.scala b/spark/src/test/scala/org/apache/comet/exec/CometInMemoryCacheSuite.scala index 3d0caec810b..be8f07a812e 100644 --- a/spark/src/test/scala/org/apache/comet/exec/CometInMemoryCacheSuite.scala +++ b/spark/src/test/scala/org/apache/comet/exec/CometInMemoryCacheSuite.scala @@ -41,7 +41,7 @@ import org.apache.spark.sql.comet.{CometBroadcastHashJoinExec, CometInMemoryTabl import org.apache.spark.sql.comet.execution.arrow.{ArrowCachedBatchSerializer, CometCachedBatchHelper} import org.apache.spark.sql.comet.execution.shuffle.CometCelebornShuffleManager import org.apache.spark.sql.comet.util.Utils -import org.apache.spark.sql.execution.{FormattedMode, SortExec} +import org.apache.spark.sql.execution.{CometSparkPlanInfoHelper, FormattedMode, SortExec, SparkPlanInfo} import org.apache.spark.sql.execution.adaptive.{AdaptiveSparkPlanExec, AQEShuffleReadExec, QueryStageExec, ShuffleQueryStageExec} import org.apache.spark.sql.execution.columnar.{CometInMemoryRelationHelper, InMemoryRelation, InMemoryTableScanExec} import org.apache.spark.sql.execution.exchange.{Exchange, ReusedExchangeExec, ShuffleExchangeLike} @@ -430,6 +430,50 @@ class CometInMemoryCacheSuite extends CometTestBase { } } + test("the SQL tab and event log draw the cached plan below CometInMemoryTableScan") { + // Spark builds both from SparkPlanInfo, which gives its own cache scan the cached plan as a + // child. See https://github.com/apache/datafusion-comet/issues/6463. + def scanInfos(info: SparkPlanInfo): Seq[SparkPlanInfo] = + (if (info.nodeName == "CometInMemoryTableScan") Seq(info) else Nil) ++ + info.children.flatMap(scanInfos) + + Seq("false", "true").foreach { aqe => + withClue(s"AQE $aqe: ") { + withSQLConf( + SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> aqe, + CometConf.COMET_EXEC_IN_MEMORY_CACHE_ENABLED.key -> "true") { + withTempView("plan_info_cache") { + // The shuffle gives the cached plan an adaptive plan of its own when AQE is on. + spark + .range(1000) + .selectExpr("id % 10 AS k") + .groupBy("k") + .count() + .createOrReplaceTempView("plan_info_cache") + spark.catalog.cacheTable("plan_info_cache") + val df = spark.sql("SELECT * FROM plan_info_cache WHERE k > 1") + df.collect() + val plan = df.queryExecution.executedPlan + val scans = collect(plan) { case s: CometInMemoryTableScanExec => s } + assert(scans.size == 1, plan) + val sparkScan = scans.head.originalPlan + val cachedPlanInfo = + CometSparkPlanInfoHelper.fromSparkPlan(sparkScan.relation.cachedPlan) + // Spark's own scan of the cache sits below, and the cached plan below that. + val infos = scanInfos(CometSparkPlanInfoHelper.fromSparkPlan(plan)) + assert( + infos.map(_.children.map(info => (info.nodeName, info.children))) == + Seq(Seq((sparkScan.nodeName, Seq(cachedPlanInfo)))), + plan) + // Other walkers of subqueries stop at Spark's scan, as in Spark's own plans, so they + // do not find the cached plan's shuffle in this query, which has none of its own. + assert(collectWithSubqueries(plan) { case e: ShuffleExchangeLike => e }.isEmpty, plan) + } + } + } + } + } + test("Comet in-memory cache disabled keeps SparkToColumnar fallback path") { withSQLConf( SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false", diff --git a/spark/src/test/scala/org/apache/spark/sql/execution/CometSparkPlanInfoHelper.scala b/spark/src/test/scala/org/apache/spark/sql/execution/CometSparkPlanInfoHelper.scala new file mode 100644 index 00000000000..f8ccd6ef50e --- /dev/null +++ b/spark/src/test/scala/org/apache/spark/sql/execution/CometSparkPlanInfoHelper.scala @@ -0,0 +1,29 @@ +/* + * 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.execution + +/** + * Test-only access to `SparkPlanInfo.fromSparkPlan`, which builds the plan tree behind the SQL + * tab's graph and the plans that the event log records. Its object is `private[execution]`, hence + * this shim. + */ +object CometSparkPlanInfoHelper { + def fromSparkPlan(plan: SparkPlan): SparkPlanInfo = SparkPlanInfo.fromSparkPlan(plan) +} diff --git a/spark/src/test/spark-4.2/org/apache/comet/exec/CometInMemoryCacheLastAttemptMetricSuite.scala b/spark/src/test/spark-4.2/org/apache/comet/exec/CometInMemoryCacheLastAttemptMetricSuite.scala new file mode 100644 index 00000000000..3f30e44aa06 --- /dev/null +++ b/spark/src/test/spark-4.2/org/apache/comet/exec/CometInMemoryCacheLastAttemptMetricSuite.scala @@ -0,0 +1,73 @@ +/* + * 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.comet.exec + +import org.apache.spark.SparkConf +import org.apache.spark.sql.CometTestBase +import org.apache.spark.sql.comet.CometInMemoryTableScanExec +import org.apache.spark.sql.comet.execution.shuffle.CometShuffleExchangeExec +import org.apache.spark.sql.execution.columnar.CometInMemoryRelationHelper +import org.apache.spark.sql.execution.metric.SQLLastAttemptMetrics + +import org.apache.comet.CometConf + +/** Spark 4.2's last-attempt metrics over a relation cached in Comet's format. */ +class CometInMemoryCacheLastAttemptMetricSuite extends CometTestBase { + + import testImplicits._ + + override protected def beforeAll(): Unit = { + CometInMemoryRelationHelper.clearSerializer() + super.beforeAll() + } + + override protected def afterAll(): Unit = { + try { + super.afterAll() + } finally { + CometInMemoryRelationHelper.clearSerializer() + } + } + + override protected def sparkConf: SparkConf = super.sparkConf + .set("spark.plugins", "org.apache.spark.CometPlugin") + .set( + "spark.sql.cache.serializer", + "org.apache.spark.sql.comet.execution.arrow.ArrowCachedBatchSerializer") + + test("a last-attempt metric outside the cache keeps its value") { + // Spark finds the metric's stages by walking the plan and its subqueries, and gives up on a + // shuffle it does not know. A Comet shuffle inside the cached plan must stay out of that walk. + // https://github.com/apache/datafusion-comet/pull/6577#discussion_r4174502176 + withSQLConf(CometConf.COMET_EXEC_IN_MEMORY_CACHE_ENABLED.key -> "true") { + val metric = SQLLastAttemptMetrics.createMetric(spark.sparkContext, "rows") + val cached = spark.range(0, 100, 1, 2).repartition(2).cache() + cached.count() + val df = cached.map { id => metric.add(1); id } + df.collect() + val plan = df.queryExecution.executedPlan + val scans = collect(plan) { case s: CometInMemoryTableScanExec => s } + assert(scans.size == 1, plan) + val cachedPlan = scans.head.originalPlan.relation.cachedPlan + assert(collect(cachedPlan) { case s: CometShuffleExchangeExec => s }.nonEmpty, cachedPlan) + assert(metric.lastAttemptValueForDataset(df) == Some(100L)) + } + } +}