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 103f95f2e53..7af36cac667 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 @@ -73,6 +73,12 @@ case class CometInMemoryTableScanExec( // for instance -- then reads the wrong column. override def output: Seq[Attribute] = originalPlan.output + // Described by the Spark scan it replaces: the table's name when it has one, the attributes it + // reads and any pruning predicates. The default would print every constructor field, among them + // the CachedRDDBuilder with the whole cached plan, physical and logical, inline and with its raw + // newlines, which breaks the tree of every plan that reads the cache. + override def stringArgs: Iterator[Any] = Iterator(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 ed3d518d3ce..d82f6c021bc 100644 --- a/spark/src/test/scala/org/apache/comet/exec/CometInMemoryCacheSuite.scala +++ b/spark/src/test/scala/org/apache/comet/exec/CometInMemoryCacheSuite.scala @@ -324,6 +324,43 @@ class CometInMemoryCacheSuite extends CometTestBase { } } + test("CometInMemoryTableScan is described by the Spark scan it replaces") { + // With every constructor field printed, the node dumped its CachedRDDBuilder, the whole cached + // plan both physical and logical, into the middle of every plan that read the cache, raw + // newlines and all, in the tree string and in EXPLAIN FORMATTED alike. + withSQLConf( + SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false", + CometConf.COMET_EXEC_IN_MEMORY_CACHE_ENABLED.key -> "true") { + spark.catalog.clearCache() + spark + .range(1000) + .selectExpr("id AS key", "id % 8 AS value") + .createOrReplaceTempView("explain_cache") + spark.catalog.cacheTable("explain_cache") + try { + val df = spark.sql("SELECT key FROM explain_cache WHERE value = 3") + df.collect() + val plan = df.queryExecution.executedPlan + val scans = collect(plan) { case s: CometInMemoryTableScanExec => s } + assert(scans.size == 1) + val line = scans.head.simpleString(SQLConf.get.maxToStringFields) + assert( + line.startsWith("CometInMemoryTableScan Scan In-memory table explain_cache ["), + line) + assert(line.contains("= 3)"), s"the pruning predicates should be shown: $line") + Seq( + plan.treeString, + df.queryExecution.explainString(org.apache.spark.sql.execution.FormattedMode)) + .foreach { text => + assert(!text.contains("CachedRDDBuilder"), text) + assert(!text.contains(classOf[ArrowCachedBatchSerializer].getName), text) + } + } finally { + spark.catalog.clearCache() + } + } + } + test("Comet in-memory cache disabled keeps SparkToColumnar fallback path") { withSQLConf( SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false",