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 1628c52, Spark 4.1 profile, cache a 1M-row relation in each format and read the planner's size once it has materialized:
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.
Describe the bug
With Comet's cache serializer, each
CometCachedBatchreports its stored, compressed payload assizeInBytes(sizeInBytes = bytes.sizeinArrowCachedBatchSerializer.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 againstspark.sql.autoBroadcastJoinThreshold, both statically and under AQE through the table-cache stage's runtime statistics, and whatRewriteJoincompares to pick a shuffled hash join's build side.Both of Spark's own cache formats report the decoded size there.
DefaultCachedBatchsums the uncompressed column statistics, and the Arrow cache format Spark added in SPARK-57268, which also compresses per buffer, recordsvector.getBufferSizefor 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-22673andSPARK-36120assert Spark's size semantics. It also explains why the port ofSPARK-37742: AQE reads invalid InMemoryRelation stats and mistakenly plans BHJinCometInMemoryCacheSuitelowered 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
mainat 1628c52, Spark 4.1 profile, cache a 1M-row relation in each format and read the planner's size once it has materialized:zstd(default)noneid, id % 1000, CAST(id AS STRING), id * 1.5With the default 10 MB threshold,
SELECT count(*) FROM range(0, 5000000, 1, 4) r JOIN cached_t c ON r.id = c.idover 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
CometInMemoryCacheBenchmarkuses it for the footprint numbers in the in-memory cache guide, so it should keep being measured, just not throughsizeInBytes. Part of #5487, and it should be fixed before the cache is turned on by default in #5634.