Repository navigation
Conversation
fced6bc to
47d0a85
Compare
47d0a85 to
72c3d8c
Compare
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Compressed cache buffers store logical long values directly. Sequential values can occupy less space when compressed as deltas.
- Design approach: Add an opt-in setting for top-level long columns. Keep deltas only when their compressed buffer is over 25% smaller than the ordinary compressed buffer.
- Correctness / compatibility analysis: Encoding and reconstruction preserve wrapping long arithmetic, validity buffers, and logical pruning statistics. Per-batch flags keep reads independent of subsequent configuration changes. Checked Spark’s cache contracts and reader paths in 3.4.3, 3.5.9, 4.0.4, 4.1.3, and 4.2.0, plus Arrow 18.3.0 buffer ownership.
- Key design decisions: Reusing
CachedBatchIpc.Projectiongives Spark row, Spark columnar, and native consumers the same reconstruction path. The per-column flags and existing compression machinery keep the added abstraction small. - Implementation sketch: Capture the write setting, identify eligible data buffers after dictionary decoding, compare compressed representations, retain encoding flags, and reconstruct selected columns before exposing their vectors.
- Behavioral changes worth calling out: The option adds cache-build work, temporary off-heap storage, and a prefix sum during reads. It remains disabled by default and is skipped for uncompressed caches. Against
branch-1.1at7b7eec69282a7abc33a0e8e687ca594c72f140ec, delta encoding is an intended addition. Differences in decoded-size statistics, row reading, and cache-registration safeguards already exist in the supplied base. - Suggested improvements: No introduced P1/P2 issues found within this review. No additional changes requested at that severity.
Reviewed full SHA 870187953cee14a9ae5ed57369db637bd8eb6c64 against 83285bbe1b604e3958bffceca2e4dc9137d98fec, including all nine base-to-head file differences. The NOTICE packaging and sequence-test differences come from newer base-only commits. The PR remains non-draft. The snapshot and live discussion contained no reviews, issue comments, inline comments, or review threads.
Routed skills: review-comet-pr, review-comet-ffi-pr, review-comet-memory-pr, and review-comet-expression-pr.
Exact-head CI: 27 successful checks, 15 skipped, and no failures. Linux Spark 4.1 suites, native tests/build, cross-version lint, and Spark 3.5 strict compilation passed. Downloaded exec artifacts confirm 1,295 passing tests, including cache, delta, disk-persistence, Kryo, and row-reader coverage. Spark SQL matrix and macOS checks were skipped.
Local validation: Compiled the exact-head CachedBatchIpc.scala with real Spark 4.1.3/Arrow 18.3.0 dependencies and small Comet vector adapters. All 84 round-trip scenarios passed, covering nulls, overflow, nested-column indexing, repeated/reordered projections, concurrent reads, logical sizes, and compression-failure cleanup. No full local Maven/native build, Spark SQL matrix, or performance benchmark was run. The author’s reported local SQL-suite result was not independently reproduced. Project files remain unchanged.
andygrove
left a comment
There was a problem hiding this comment.
I merged this with current main and ran CometInMemoryCacheSuite and CometInMemoryCacheKryoSuite on Spark 3.4, 4.1 and 4.2, and everything passes. The only conflict with main is the config table in in-memory-cache.md, where #5634 changed the enabled row. I also compared answers against uncached Spark on a 5M-row relation and on a Parquet scan with a column near Long.MinValue, and they match. So the round trip looks right to me.
My main concern is what the option costs on bigint columns that are not sequential. Details are inline, along with a doc fix and a test helper that drops the new flags.
| deltas.writerIndex(buffer.writerIndex()) | ||
| // Skip a second compression for full-width irregular longs. The size comparison | ||
| // still rejects poorly compressing deltas from narrower distributions. | ||
| if (smallDeltas.toLong * 2 >= batch.getLength) { |
There was a problem hiding this comment.
This check sends every column whose values fit in an int through a second zstd pass, whether or not the deltas end up smaller. I cached a 5M-row relation of three such bigint columns: uniform below 1e6, uniform below 2^31, and a 1 ms step with up to 1 ms of jitter. With the option on, none of the 512 batches kept deltas for any column and the footprint stayed at 46.9 MiB. The build went from 210 ms to 299 ms, median of 7, and every run with the option on was slower than every run without it. Per 10,000-row column, uniform values below 1000 took 328 us to serialize instead of 169 us. Their deltas compress to 1.38 times the size of the values, so they can never pass the bar.
Could the loop also compare how wide the deltas are with how wide the values are, and only try the second pass when the deltas are clearly narrower? For example, summing 64 - numberOfLeadingZeros of each value's zigzag form, and trying only when the deltas come out at least four bits narrower on average. That skips every uniform case I tried and still tries all the sequential, constant-step and cyclic ones.
|
|
||
| Set `spark.comet.exec.inMemoryCache.deltaEncoding.enabled=true` to try delta encoding | ||
| for top-level `bigint` columns in compressed caches. A column uses deltas only when | ||
| its compressed data buffer is over 25% smaller than the plain representation. |
There was a problem hiding this comment.
The code compares the compressed deltas with the compressed values (encoded.writerIndex() < packed.writerIndex() * 3 / 4), not with the plain representation, so this sentence describes a different threshold. The CometConf doc string has the same gap, since it only says the delta buffer is over 25% smaller. Could both say the comparison is between the two compressed forms?
It would also help to give numbers here, the way the codec table under Storage format does, perhaps from a delta arm in runCodecBenchmark. That benchmark's relation is a good case for this option. I measured 54.8 MiB without deltas and 44.2 MiB with them. id, k and v went from 6.3, 0.9 and 6.5 MiB to 1.3, 0.4 and 1.5 MiB, and reads were no slower. The build cost on columns that are not sequential is worth stating too.
| CachedBatchIpc.serialize(batch, codec, allocator, chunkSize)._1 | ||
| chunkSize: Int = 1024 * 1024, | ||
| deltaEncoding: Boolean = false): ChunkedByteBuffer = | ||
| CachedBatchIpc.serialize(batch, codec, allocator, chunkSize, deltaEncoding)._1 |
There was a problem hiding this comment.
With deltaEncoding = true this returns the payload without its flags, and cachedBatch below then builds a batch that has none. A test that round trips through these and load would silently read the deltas back as values. The doc on cachedBatch says it is the batch the writer would have stored, which no longer holds. Could serialize return the flags as well, or the CachedBatch itself?
|
@peterxcli, a heads-up on related work. While digging into #5485 I found that most of what remains of the Spark-reader gap is That is the data this PR targets. The delta pass compresses with whatever codec the batch uses, so the two should compose without further changes once both land, and |
Which issue does this PR close?
Related to #5485. Builds on the merged #5859.
Rationale for this change
Sequential long columns can occupy less cache space when stored as compressed deltas. Cache creation does extra work, and reads reconstruct the values, so this is opt-in and does not promise faster reads.
What changes are included in this PR?
Add
spark.comet.exec.inMemoryCache.deltaEncoding.enabled(default false) to the Arrow IPC cache format. Select top-level long data buffers only when compressed deltas are over 25% smaller. Full-width irregular longs skip the second compression; uncompressed caches skip encoding.The shared projection reader restores values for Spark and native consumers. Nulls and logical pruning statistics remain intact, and changing the setting leaves existing caches readable.
How are these changes tested?
Spark 4.1 and Spark 3.5 cache, Kryo, row-iterator, and utility suites pass; Spark 3.5 strict-warning compilation also passes. Coverage includes nulls, long overflow, projections, all readers, disk persistence, recaching after changing the setting, and compression-failure cleanup. CI configuration checks pass.
Spark 4.1 SQL
sql_core-1passes: 8,335 tests succeeded with no failures (69 canceled, 581 ignored).