comphead commented on code in PR #6421:
URL: https://github.com/apache/datafusion-comet/pull/6421#discussion_r4146862913


##########
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:
   Nit: this branch means 'observed metrics' only because the two branches 
above already ruled out the other cases. Would it read better to hoist the 
`recordsObservedMetrics` result into a named `lazy val` next to `nativeScan` 
(so the plan walk stays conditional) and test that here? Then the reason sits 
next to its condition, and a branch added above later cannot change which 
reason gets reported.



##########
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:
   Both cached levels here have their own `observe`, so I think this test would 
still pass if the nested-cache arm in `recordedIn` (`case scan: 
InMemoryTableScanExec`) were dropped. Would it be worth appending one more 
level with no `observe` of its own, for example `.filter($"id" > 
10).persist()`, and keeping the same expected metrics? I haven't run it, but I 
expect the top scan would then be replaced and the inner metrics lost without 
that arm.



-- 
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