Skip to content

Comet's in-memory cache reports its compressed size to the planner, which changes join planning #6411

Description

@andygrove

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.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

bugSomething isn't workingpriority:mediumFunctional bugs, performance regressions, broken featuresrequires-triage

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions