andygrove opened a new issue, #6420:
URL: https://github.com/apache/datafusion-comet/issues/6420
### Describe the bug
With Comet's cache serializer installed and
`spark.comet.exec.inMemoryCache.enabled=true`, metrics from a `Dataset.observe`
placed before `persist()` or `cache()` are silently lost.
`QueryExecution.observedMetrics` comes back empty, `Observation.get` returns an
empty map on Spark 3.5 and later, and on Spark 3.4 `Observation.get` never
returns.
Spark collects observed metrics after a query with
`CollectMetricsExec.collect`, which only looks inside a cached plan through an
`InMemoryTableScanExec`:
```scala
case tableScan: InMemoryTableScanExec =>
CollectMetricsExec.collect(tableScan.relation.cachedPlan)
```
That is the same in every Spark version Comet supports, 3.4 through 4.2.
`CometInMemoryTableScanExec` replaces that node, so the `CollectMetricsExec`
inside the cached plan is never visited. The cached data is not the problem.
With the native scan turned off at runtime, Spark's own scan reads the same
`CometCachedBatch` payloads and the metrics come back.
I found this by running Spark's SQL suites with Comet's serializer
installed, where `DataFrameCallbackSuite`'s `SPARK-35695: get observable
metrics with persist by callback` fails with `0 did not equal 2`.
### Steps to reproduce
On `main` at c4dd52503, with Comet's serializer installed and
`spark.comet.exec.inMemoryCache.enabled=true`:
```scala
val df = spark.range(100)
.observe("my_event", count(lit(1)).as("rows"), max("id").as("max_id"))
.persist()
df.collect()
df.queryExecution.observedMetrics // Map()
val obs = Observation("obs")
val df2 = spark.range(100).observe(obs, count(lit(1)).as("rows")).persist()
df2.collect()
obs.get // Map() on Spark 3.5 and later, never returns on Spark 3.4
```
With `spark.comet.exec.inMemoryCache.enabled=false` set at runtime, the same
code returns `Map(my_event -> [100,99])` and `Map(rows -> 100)`.
### Expected behavior
Observed metrics recorded in a cached plan are collected the same way they
are with Spark's cache scan.
### Additional context
The native cache scan is off by default, so this only affects applications
that turned it on, but it has to be fixed before #5634 turns it on by default.
Part of #5487.
--
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]