Repository navigation
feat: add lz4 as an opt-in codec for Comet's in-memory cache - #6859
Conversation
spark.comet.exec.inMemoryCache.compression.codec accepts lz4 alongside zstd and none. Buffers are written as Arrow's standard LZ4_FRAME with lz4-java, the JNI library behind spark.io.compression.codec, rather than Arrow's own pure-Java LZ4 codec, and LZ4_FRAME batches are read with the same codec. lz4 decompresses several times faster than zstd, which closes most of the gap between Spark operators reading Comet's cache format and reading Spark's own. It stays opt-in because it takes around three times zstd's footprint, and on sequential longs more than Spark's own format.
sunchao
left a comment
There was a problem hiding this comment.
Reviewed all seven changed files against base b94d52f43c9edaf341f5c2cc00fa2e351ba9222b at head 9b203d0b4dbcbf4c5276070757a40368b6abdd2f. The PR is non-draft. No introduced P1/P2 issues found within this review.
Summary
- Prior state and problem: Comet’s cache offered
zstdandnone. Wide reads paid substantial decompression cost, while Arrow’s existing LZ4 implementation was unsuitable for the intended performance improvement. - Design approach: Add opt-in
lz4using Spark’s bundled lz4-java library. Encode standard ArrowLZ4_FRAMEbuffers directly from off-heap memory and select the reader from each batch’s recorded codec. - Correctness: Checked framing, block boundaries, stored blocks, length validation, buffer ownership and exception cleanup. The exact codec source passed independent-reader and round-trip checks. Spark’s cache serializer interfaces and compression behavior were compared against supported-version sources. No qualifying correctness issue was identified.
- Compatibility analysis: Spark sources confirm lz4-java
1.8.0on 3.4–4.1 and theat.yawk.lz41.11.0fork on 4.2. Both library versions passed local codec validation. Comparison withbranch-1.1ate9efd9f764ee0a59b7898ff028d6985d4a7a28e1confirms the intended addition oflz4, withzstdremaining the default. Broader cache enablement, fused-reader and statistics changes already exist in the supplied base. - Key design decisions: Keeping LZ4 opt-in reflects its larger cache footprint. Lazy JNI initialization and Java fallback avoid requiring native LZ4 availability. The decoder deliberately accepts the fixed frame format this writer produces.
- Implementation sketch: Extend the existing configuration and codec dispatch, implement independent blocks up to 4 MiB, extend cache tests, and add codec-aware benchmark cases and footprint reporting.
- Performance: Direct buffer compression avoids the stock codec’s heap-copy and stream path. Projection still decompresses only selected buffers. The reported benchmark shows a full-width read falling from 248 ms to 110 ms, with footprint rising from 51 MiB to 145 MiB. These author-reported measurements were not independently rerun. No evidence-backed performance regression was identified on existing
zstdornonepaths. - Design: The change stays within the existing cache serialization boundary and preserves codec-independent planning and statistics. Benchmark setup remains outside timed reads and distinguishes cached copies by serializer and codec.
- Abstraction & complexity: Reusing
AbstractCompressionCodecpreserves Arrow’s prefix, empty-buffer and uncompressed-buffer handling. The shared codec and fixed frame layout keep the custom implementation narrowly scoped. No concrete simplification meeting the reporting bar was identified. - Behavioral changes worth calling out:
spark.comet.exec.inMemoryCache.compression.codec=lz4now affects newly cached data. Previously cached batches retain their recorded codec. Missing JNI support produces a warning and uses Java. Temporary codec buffers use the existing Arrow allocator and accounting model. - Suggested improvements: No additional change meeting the P1/P2 reporting bar was identified.
Routed skills: review-comet-pr, review-comet-memory-pr and review-comet-ffi-pr. Read AGENTS.md, the relevant contributor guidance, and the snapshot and live discussions. There were no existing reviews, issue comments, inline comments or review threads to duplicate.
Exact-head CI: All 33 inspected checks matched the reviewed SHA. Eight passed, eleven were running and fourteen were skipped. Native build, Rust tests, Java lint and strict Scala warnings were still running. Spark SQL, macOS and benchmark checks were skipped. No failures were reported, but no completed integration-test verdict was available.
Validation: Compiled the unchanged codec source and ran a standalone JDK 21 harness with LZ4 1.8.0, 1.11.0, and a forced Java fallback. Each passed 30 size/data-pattern cases, six existing malformed-frame cases and sixteen concurrent round-trips. Independent Arrow and lz4-java readers matched the input, and allocator usage returned to zero. No full Maven/native build, CometInMemoryCacheSuite, Spark SQL matrix or performance benchmark was run. Project files remain unchanged.
Which issue does this PR close?
Part of #5485.
Rationale for this change
#5485 tracks Spark operators reading Comet's cache format more slowly than Spark's own format. Measuring each codec shows that what remains of that gap is almost all decompression: with
none, every read of the benchmark's six-column relation matches or beats Spark's format. The only compressing codec Comet offers iszstd, because Arrow's own LZ4 codec is commons-compress's pure-Java implementation, which is far slower thanzstd. lz4-java, the JNI library behindspark.io.compression.codec, is on the classpath of every Spark version Comet supports (org.lz41.8.0 on 3.4 to 4.1, theat.yawk.lz4fork 1.11.0 on 4.2) and decompresses several times faster thanzstd.From
CometInMemoryCacheBenchmarkon an Apple M3 Ultra (JDK 17, Spark 4.1, release build), over 5M rows of three long and three string columns. Each cell is the read's time relative to Spark's own cache format, so below 1 is faster:zstd, 1 / 3 / 6 columnslz4, 1 / 3 / 6 columnsReading every column of relations of 100 / 200 / 1500 nullable
bigintcolumns through the row reader takes 2.08 / 2.10 / 2.60 times as long as from Spark's format withzstd, and 1.26 / 1.37 / 1.98 times withlz4.zstdlz4noneSpark's own format holds the same relation in 217 MiB. These are from one run; two earlier exploratory runs of the same comparison agreed within a few percent.
It is opt-in rather than the new default because of footprint.
lz4takes about three timeszstd's memory, and on sequential longs, which Spark's format delta-encodes, more than Spark's own format: the wide relations above take about 74 MiB withlz4, against 26 to 27 MiB in Spark's format and withzstd.lz4combined with delta-encoded longs (#5869) is the candidate to revisit the default with.At 1500 columns even
nonetook 1.25 times as long as Spark's format in an exploratory run, so the rest of that gap is per-column decoding rather than the codec, and stays with #5485.What changes are included in this PR?
Lz4FrameCompressionCodec: Arrow'sLZ4_FRAMEcodec on lz4-java. Each buffer is one standard LZ4 frame (independent blocks of up to 4 MiB, no checksums) behind Arrow's uncompressed-length prefix, compressed and decompressed with lz4-java's block API directly on the off-heap buffers. It asks for lz4-java's JNI implementation explicitly, becauseLZ4Factory.fastestInstance()silently settles for the Java one when lz4-java is not loaded by the system class loader (as undermvn exec:java), and falls back to Java, with a warning, only when the native library cannot load. Comet turns off Arrow's bounds checks, so the reader checks every length a frame records before reading through it.spark.comet.exec.inMemoryCache.compression.codecacceptslz4, and the read path decodesLZ4_FRAMEbatches with this codec rather than Arrow's commons-compress one.lz4, with when to choose it and its footprint caveat, and Performance and Limitations say what it changes for reads against Spark's format.CometInMemoryCacheBenchmarkmeasureslz4in the codec, Spark-operator, wide-relation and AQE cases, and prints the footprint of each cached copy.How are these changes tested?
New and extended tests in
CometInMemoryCacheSuite:lz4: full and projected reads through the native scan and through Spark's reader, a row count, and batch pruning.LZ4FrameInputStream, for a one-block buffer and for a three-block buffer whose first block is incompressible and stored verbatim.SparkExceptionnaming LZ4 and release everything they allocated.lz4.The suite passes on Spark 4.1 and on Spark 4.2, which ships the
at.yawk.lz4fork, and the Spark 3.5 strict-warnings build and scalafix check pass.