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]