andygrove opened a new issue, #6411:
URL: https://github.com/apache/datafusion-comet/issues/6411
### Describe the bug
With Comet's cache serializer, each `CometCachedBatch` reports its stored,
compressed payload as `sizeInBytes` (`sizeInBytes = bytes.size` in
`ArrowCachedBatchSerializer.encodeBatches`), and the per-column size fields in
its statistics row are the compressed buffer lengths. Once a cached relation
materializes, `InMemoryRelation.computeStats()` returns the sum of those as the
relation's size. That is what join planning compares against
`spark.sql.autoBroadcastJoinThreshold`, both statically and under AQE through
the table-cache stage's runtime statistics, and what `RewriteJoin` compares to
pick a shuffled hash join's build side.
Both of Spark's own cache formats report the decoded size there.
`DefaultCachedBatch` sums the uncompressed column statistics, and the Arrow
cache format Spark added in SPARK-57268, which also compresses per buffer,
records `vector.getBufferSize` for each column, with a comment that an
understated size makes a relation "wrongly eligible for broadcast". So with
Comet's format on, a cached relation looks several times smaller to the planner
than the same relation in Spark's format, and a join that Spark would plan as a
sort-merge join broadcasts the cached side instead. Results are correct. The
risk is memory, since the broadcast builds its hash table from the decoded data.
I found this by running Spark's own cache suites with Comet's serializer
installed, where `InMemoryRelation statistics`, `SPARK-22673` and `SPARK-36120`
assert Spark's size semantics. It also explains why the port of `SPARK-37742:
AQE reads invalid InMemoryRelation stats and mistakenly plans BHJ` in
`CometInMemoryCacheSuite` lowered Spark's threshold from 1048584 to 1024 bytes.
With Spark's threshold it plans exactly the broadcast join that Spark's test
exists to catch.
### Steps to reproduce
On `main` at 1628c5200, Spark 4.1 profile, cache a 1M-row relation in each
format and read the planner's size once it has materialized:
```scala
val df = spark.sql(
"SELECT id, CAST(rand(1) * 1000 AS INT) AS k, CAST(rand(3) AS STRING) AS
s, rand(2) AS d " +
"FROM range(0, 1000000, 1, 4)")
df.cache()
df.count()
df.queryExecution.withCachedData
.collectFirst { case r: InMemoryRelation => r }
.get
.computeStats()
.sizeInBytes
```
| Relation (1M rows) | Spark's format | Comet,
`zstd` (default) | Comet, `none` |
| --------------------------------------------- | -------------: |
----------------------: | ------------: |
| the query above | 42.3 MB |
21.3 MB | 42.8 MB |
| `id, id % 1000, CAST(id AS STRING), id * 1.5` | 33.9 MB |
7.1 MB | 42.4 MB |
With the default 10 MB threshold, `SELECT count(*) FROM range(0, 5000000, 1,
4) r JOIN cached_t c ON r.id = c.id` over the first relation plans a sort-merge
join with Spark's format and a broadcast hash join with Comet's.
### Expected behavior
A cached relation's planner size is its decoded size and does not depend on
the compression codec, so queries over it are planned the way they would be
with Spark's cache format.
### Additional context
The stored payload size is still what the cache occupies in memory, and
`CometInMemoryCacheBenchmark` uses it for the footprint numbers in the
in-memory cache guide, so it should keep being measured, just not through
`sizeInBytes`. Part of #5487, and it should be fixed before the cache is turned
on by default in #5634.
--
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]