andygrove opened a new pull request, #6538:
URL: https://github.com/apache/datafusion-comet/pull/6538

   ## Which issue does this PR close?
   
   Part of #5485. This is option 1 from that issue.
   
   ## Rationale for this change
   
   `spark.sql.cache.serializer` is static, so a relation cached in Comet's 
format stays in that format for the life of the application. If a session then 
turns Comet or its native execution off, Spark's `InMemoryTableScanExec` reads 
it, which is slower than reading Spark's own format by the factors #5485 
measured. Option 1 in #5485 is to report that through the existing fallback 
reasons, so a user can see why such a scan is slow and how to keep Spark's 
format instead.
   
   `CometExecRule` already records a reason on the cache scan when the native 
cache scan is turned off, when a relation is not stored in Comet's format, and 
when its cached plan records observed metrics. It returns before looking at the 
plan when Comet or its native execution is off, though, so those plans recorded 
nothing. The review of #5634, which turns the cache on by default, asked for 
this case to be covered first.
   
   ## What changes are included in this PR?
   
   - When Comet is disabled for a plan, or `spark.comet.exec.enabled=false`, 
`CometExecRule` records a fallback reason on each `InMemoryTableScanExec` that 
reads Comet's format. The reason says why Spark reads it, and that 
`spark.comet.exec.inMemoryCache.enabled=false` at startup keeps caches in 
Spark's format.
   - `CometExecRule.readsCometCacheFormat` names the check the rule already 
made for whether a scan reads Comet's format, so both paths share it.
   
   The plan is unchanged. The reason is a tag on the scan, shown by 
`ExtendedExplainInfo`, through `spark.sql.extendedExplainProviders` on Spark 
4.0 and later or through `getFallbackReasons`.
   
   ## How are these changes tested?
   
   A new test in `CometInMemoryCacheSuite` caches one relation in Comet's 
format and one that Comet's serializer delegates to Spark's format. With 
`spark.comet.enabled=false` and with `spark.comet.exec.enabled=false`, each 
with AQE off and on, it checks that a query over the first records the reason 
and one over the second does not. The test fails without the change.
   
   `CometInMemoryCacheSuite`, `CometInMemoryCachePruningSuite` and 
`CometInMemoryCacheKryoSuite` pass locally on Spark 4.1, and 
`CometInMemoryCacheSuite` on Spark 3.4, which has no table-cache query stages. 
The scalafix check passes on Spark 3.4.
   


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