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]

Reply via email to