From 5b3fc2ee51102f7b8c551611efd594ef5f6829e2 Mon Sep 17 00:00:00 2001 From: comphead Date: Sat, 3 Oct 2026 11:43:28 -0700 Subject: [PATCH 1/3] fix: draw the cached plan below CometInMemoryTableScan in the SQL tab and event log Spark builds the SQL tab's graph and the plans in the event log with SparkPlanInfo.fromSparkPlan. It gives its own InMemoryTableScanExec the cached plan as a child, but recognizes that scan by its class, and for any other node takes the children and the subqueries. Expose the cached plan as the one subquery of CometInMemoryTableScanExec. Spark runs subqueries from a plan's expressions and only walks this list, so the cached plan does not become part of the query that reads the cache. Closes #6463. --- .../comet/CometInMemoryTableScanExec.scala | 9 ++++ .../comet/exec/CometInMemoryCacheSuite.scala | 42 ++++++++++++++++++- .../execution/CometSparkPlanInfoHelper.scala | 29 +++++++++++++ 3 files changed, 79 insertions(+), 1 deletion(-) create mode 100644 spark/src/test/scala/org/apache/spark/sql/execution/CometSparkPlanInfoHelper.scala 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..95b4cedb47c 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, which the SQL tab's graph and the event log's plans are built from, gives + // Spark's own scan its cached plan as a child, but recognizes that scan by its class. For any + // other node it takes the children and the subqueries, so expose the cached plan as the one + // subquery. A child would make the cached plan part of the query that reads the cache, but + // Spark only walks this list: subqueries run from a plan's expressions. Its other walkers, such + // as collectWithSubqueries, follow it into the cached plan too. A lazy val, because Spark 3.x + // declares subqueries as one, and a lazy val overrides Spark 4's def as well. + @transient override lazy val subqueries: Seq[SparkPlan] = Seq(originalPlan.relation.cachedPlan) + // `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..63c308e0c4a 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,46 @@ 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 => + withSQLConf( + SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> aqe, + CometConf.COMET_EXEC_IN_MEMORY_CACHE_ENABLED.key -> "true") { + spark.catalog.clearCache() + // 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") + try { + 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, s"AQE $aqe: $plan") + val cachedPlan = scans.head.originalPlan.relation.cachedPlan + + val infos = scanInfos(CometSparkPlanInfoHelper.fromSparkPlan(plan)) + assert(infos.size == 1, s"AQE $aqe: $plan") + assert( + infos.head.children == Seq(CometSparkPlanInfoHelper.fromSparkPlan(cachedPlan)), + s"AQE $aqe: ${infos.head.children.map(_.simpleString)}") + } finally { + spark.catalog.clearCache() + } + } + } + } + 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) +} From 668512594a5324c7b61c0a82c2ec9cdbff9ab97b Mon Sep 17 00:00:00 2001 From: comphead Date: Sat, 3 Oct 2026 13:39:58 -0700 Subject: [PATCH 2/3] Simplify the SQL tab test and the subqueries comment --- .../comet/CometInMemoryTableScanExec.scala | 13 +++-- .../comet/exec/CometInMemoryCacheSuite.scala | 50 +++++++++---------- 2 files changed, 29 insertions(+), 34 deletions(-) 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 95b4cedb47c..909edb6e548 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,13 +86,12 @@ case class CometInMemoryTableScanExec( // ExtendedExplainInfo.executionInnerChildren leaves them out of Comet's own reporting. override def innerChildren: Seq[QueryPlan[_]] = Seq(originalPlan.relation) - // SparkPlanInfo, which the SQL tab's graph and the event log's plans are built from, gives - // Spark's own scan its cached plan as a child, but recognizes that scan by its class. For any - // other node it takes the children and the subqueries, so expose the cached plan as the one - // subquery. A child would make the cached plan part of the query that reads the cache, but - // Spark only walks this list: subqueries run from a plan's expressions. Its other walkers, such - // as collectWithSubqueries, follow it into the cached plan too. A lazy val, because Spark 3.x - // declares subqueries as one, and a lazy val overrides Spark 4's def as well. + // SparkPlanInfo, behind the SQL tab's graph and the event log, gives Spark's own scan (matched + // by class) its cached plan as a child, and any other node its children and subqueries. Expose + // the cached plan as a subquery: unlike a child it is only walked, never run, since subqueries + // run from expressions. Walkers such as collectWithSubqueries reach it too. innerChildren above + // leaves out super.innerChildren, which is this list, or EXPLAIN would draw the cached plan + // twice. 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.relation.cachedPlan) // `originalPlan` is a plan-typed field rather than a child, so QueryPlan's canonicalization 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 63c308e0c4a..65c1eff2a38 100644 --- a/spark/src/test/scala/org/apache/comet/exec/CometInMemoryCacheSuite.scala +++ b/spark/src/test/scala/org/apache/comet/exec/CometInMemoryCacheSuite.scala @@ -438,33 +438,29 @@ class CometInMemoryCacheSuite extends CometTestBase { info.children.flatMap(scanInfos) Seq("false", "true").foreach { aqe => - withSQLConf( - SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> aqe, - CometConf.COMET_EXEC_IN_MEMORY_CACHE_ENABLED.key -> "true") { - spark.catalog.clearCache() - // 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") - try { - 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, s"AQE $aqe: $plan") - val cachedPlan = scans.head.originalPlan.relation.cachedPlan - - val infos = scanInfos(CometSparkPlanInfoHelper.fromSparkPlan(plan)) - assert(infos.size == 1, s"AQE $aqe: $plan") - assert( - infos.head.children == Seq(CometSparkPlanInfoHelper.fromSparkPlan(cachedPlan)), - s"AQE $aqe: ${infos.head.children.map(_.simpleString)}") - } finally { - spark.catalog.clearCache() + 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 cachedPlanInfo = + CometSparkPlanInfoHelper.fromSparkPlan(scans.head.originalPlan.relation.cachedPlan) + val infos = scanInfos(CometSparkPlanInfoHelper.fromSparkPlan(plan)) + assert(infos.map(_.children) == Seq(Seq(cachedPlanInfo)), plan) + } } } } From c190f6e3ac0252987d3d77dc7bdf15313ce7147c Mon Sep 17 00:00:00 2001 From: comphead Date: Sun, 4 Oct 2026 09:36:30 -0700 Subject: [PATCH 3/3] Expose Spark's own scan, not the cached plan, as the subquery Spark 4.2's SQLLastAttemptAccumulator walks a plan's subqueries to find a metric's stages and gives up on a shuffle it does not know. With the cached plan as the subquery it reached the cached plan's Comet shuffle, so a last-attempt metric used outside the cache returned None. SparkPlanInfo draws the cached plan below Spark's own InMemoryTableScanExec, matched by class, so expose that scan instead. Every other walker of subqueries stops at it, as in Spark's own plans. Adds a check on every version that collectWithSubqueries does not reach the cached plan's shuffle, and a Spark 4.2 suite for the outside-cache metric. --- .github/workflows/pr_build_linux.yml | 1 + .github/workflows/pr_build_macos.yml | 1 + .../comet/CometInMemoryTableScanExec.scala | 15 ++-- .../comet/exec/CometInMemoryCacheSuite.scala | 12 ++- ...tInMemoryCacheLastAttemptMetricSuite.scala | 73 +++++++++++++++++++ 5 files changed, 93 insertions(+), 9 deletions(-) create mode 100644 spark/src/test/spark-4.2/org/apache/comet/exec/CometInMemoryCacheLastAttemptMetricSuite.scala 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 909edb6e548..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,13 +86,14 @@ 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, gives Spark's own scan (matched - // by class) its cached plan as a child, and any other node its children and subqueries. Expose - // the cached plan as a subquery: unlike a child it is only walked, never run, since subqueries - // run from expressions. Walkers such as collectWithSubqueries reach it too. innerChildren above - // leaves out super.innerChildren, which is this list, or EXPLAIN would draw the cached plan - // twice. 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.relation.cachedPlan) + // 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 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 65c1eff2a38..be8f07a812e 100644 --- a/spark/src/test/scala/org/apache/comet/exec/CometInMemoryCacheSuite.scala +++ b/spark/src/test/scala/org/apache/comet/exec/CometInMemoryCacheSuite.scala @@ -456,10 +456,18 @@ class CometInMemoryCacheSuite extends CometTestBase { 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(scans.head.originalPlan.relation.cachedPlan) + 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) == Seq(Seq(cachedPlanInfo)), 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) } } } 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)) + } + } +}