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]
