Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
Loading