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

   ## Which issue does this PR close?
   
   Closes #6411.
   
   ## Rationale for this change
   
   Spark's planner reads the size of a materialized cached relation from its 
batches' `sizeInBytes`, and Comet's cache format reported each batch's 
compressed payload there. Both of Spark's own cache formats report the decoded 
size, so with the default `zstd` codec a relation cached in Comet's format 
could look several times smaller than the same relation in Spark's format and 
be broadcast where Spark would have shuffled it. The issue has the measurements.
   
   ## What changes are included in this PR?
   
   `CachedBatchIpc.serialize` now measures each column on the plain record 
batch, before compression, so the size field in the statistics row is the 
column's decoded Arrow size. That is what `getBufferSize` reports for a vector, 
and what Spark's `ArrowCachedBatchSerializer` records. `CometCachedBatch` no 
longer carries a `sizeInBytes` of its own and inherits 
`SimpleMetricsCachedBatch`'s, which sums those fields, the same way Spark's 
`ArrowCachedBatch` does.
   
   `CometInMemoryCacheBenchmark` read its footprint numbers from 
`computeStats`, which no longer changes with the codec, so it now sums the 
stored payloads instead. The footprint figures in the guide keep their meaning. 
The guide also says which size the planner sees.
   
   ## How are these changes tested?
   
   The port of Spark's `SPARK-37742` test in `CometInMemoryCacheSuite` goes 
back to Spark's own broadcast threshold of 1048584 bytes. It had been lowered 
to 1024 because of this bug. With Spark's threshold, `main` plans the extra 
broadcast hash join that the test exists to catch, and this branch does not.
   
   The per-column size test now compares each statistics field with the column 
decoded back out of the payload, and checks that `sizeInBytes` is their sum. A 
new test caches the same relation under `zstd` and `none` and checks that the 
relation's size is the same under both, while the `zstd` payload is less than 
half of it. Both fail on `main`.
   
   `CometInMemoryCacheSuite`, `CometInMemoryCacheKryoSuite` and 
`CometInMemoryCachePruningSuite` pass locally on Spark 3.4, 3.5, 4.0, 4.1 and 
4.2.
   


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