andygrove commented on code in PR #6421:
URL: https://github.com/apache/datafusion-comet/pull/6421#discussion_r4155983152
##########
spark/src/test/scala/org/apache/comet/exec/CometInMemoryCacheSuite.scala:
##########
@@ -2504,4 +2507,54 @@ class CometInMemoryCacheSuite extends CometTestBase {
spark.catalog.clearCache()
}
}
+
+ test("Comet in-memory cache keeps the observed metrics recorded in a cached
plan") {
+ // Spark collects the metrics of an observe() inside a cached plan only
through an
+ // InMemoryTableScanExec over it, so the scan of such a relation has to
stay Spark's. Replaced,
+ // the metrics come back empty, and on Spark 3.4 Observation.get never
returns. Nested the way
+ // SPARK-35695's test nests it, with a shuffle in the inner cached plan so
that AQE plans it.
+ withAQECache {
+ val df = spark
+ .range(0, 100, 1, 2)
+ .repartition(4)
+ .observe("inner_event", count(lit(1)).as("rows"),
max($"id").as("max_id"))
+ .persist()
+ .observe("outer_event", min($"id").as("min_id"))
Review Comment:
You're right, the test as it was still passed with that arm removed.
59b4958890 adds a top level with no `observe` of its own (`.filter($"id" >
10).persist()`) and keeps the same expected metrics. With the arm removed, the
observed metrics now come back empty and the test fails.
##########
spark/src/main/scala/org/apache/comet/rules/CometExecRule.scala:
##########
@@ -389,14 +393,20 @@ case class CometExecRule(session: SparkSession)
scan,
s"Comet in-memory cache requires
${classOf[ArrowCachedBatchSerializer].getName} " +
s"but this relation was cached with
${serializer.getClass.getName}")
- } else if (nativeCacheEnabled) {
+ } else if (nativeCacheEnabled && !cometCacheFormat) {
val unsupported = scan.relation.output
.filterNot(a =>
ArrowCachedBatchSerializer.supportsType(a.dataType))
.map(a => s"${a.name}: ${a.dataType.simpleString}")
withFallbackReason(
scan,
"Comet in-memory cache does not support the type of these cached
columns, so the " +
s"relation was cached in Spark's default format:
${unsupported.mkString(", ")}")
+ } else if (nativeCacheEnabled) {
Review Comment:
Done in 59b4958890. `recordsObservedMetrics` is now a `lazy val` next to
`nativeScan`, and this branch tests it.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]