Skip to content

Observed metrics (Dataset.observe) on a cached DataFrame are lost with Comet's native cache scan #6420

Description

@andygrove

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:

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 c4dd525, with Comet's serializer installed and spark.comet.exec.inMemoryCache.enabled=true:

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.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

bugSomething isn't workingpriority:mediumFunctional bugs, performance regressions, broken featuresrequires-triage

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions