Skip to content

feat: add lz4 as an opt-in codec for Comet's in-memory cache - #6859

Merged
andygrove merged 1 commit into
apache:mainfrom
andygrove:andygrove/datafusion-comet-issue-brainstorm
Oct 11, 2026
Merged

andygrove merged 1 commit into
apache:mainfrom
andygrove:andygrove/datafusion-comet-issue-brainstorm

Conversation

@andygrove

Copy link
Copy Markdown
Member

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 is zstd, because Arrow's own LZ4 codec is commons-compress's pure-Java implementation, which is far slower than zstd. lz4-java, the JNI library behind spark.io.compression.codec, is on the classpath of every Spark version Comet supports (org.lz4 1.8.0 on 3.4 to 4.1, the at.yawk.lz4 fork 1.11.0 on 4.2) and decompresses several times faster than zstd.

From CometInMemoryCacheBenchmark on 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:

Read zstd, 1 / 3 / 6 columns lz4, 1 / 3 / 6 columns
Native scan, Comet operators above (AQE) 0.69 / 1.17 / 0.91 0.63 / 0.70 / 0.39
Native scan, Spark operator above (AQE) 0.69 / 1.14 / 1.10 0.62 / 0.69 / 0.59
Spark's scan, fused reader 0.67 / 1.19 / 0.95 0.58 / 0.73 / 0.51
Spark's scan, row reader 1.15 / 1.71 / 1.48 1.08 / 1.25 / 1.02

Reading every column of relations of 100 / 200 / 1500 nullable bigint columns through the row reader takes 2.08 / 2.10 / 2.60 times as long as from Spark's format with zstd, and 1.26 / 1.37 / 1.98 times with lz4.

Codec Materialize Footprint Read 6 of 6
zstd 1232 ms 51 MiB 248 ms
lz4 1115 ms 145 MiB 110 ms
none 857 ms 315 MiB 54 ms

Spark'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. lz4 takes about three times zstd'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 with lz4, against 26 to 27 MiB in Spark's format and with zstd. lz4 combined with delta-encoded longs (#5869) is the candidate to revisit the default with.

At 1500 columns even none took 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's LZ4_FRAME codec 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, because LZ4Factory.fastestInstance() silently settles for the Java one when lz4-java is not loaded by the system class loader (as under mvn 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.codec accepts lz4, and the read path decodes LZ4_FRAME batches with this codec rather than Arrow's commons-compress one.
  • User guide: the codec table gains lz4, with when to choose it and its footprint caveat, and Performance and Limitations say what it changes for reads against Spark's format.
  • CometInMemoryCacheBenchmark measures lz4 in 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:

  • The codec round-trip test also covers lz4: full and projected reads through the native scan and through Spark's reader, a row count, and batch pruning.
  • What the codec writes decodes with lz4-java's own LZ4FrameInputStream, for a one-block buffer and for a three-block buffer whose first block is incompressible and stored verbatim.
  • Corrupt frames (not a frame, cut off inside the header or before the end mark, a block that runs past the end, decoding to less or more than the recorded length) fail with a SparkException naming LZ4 and release everything they allocated.
  • The test that a column failing to decode releases its vectors also runs under lz4.

The suite passes on Spark 4.1 and on Spark 4.2, which ships the at.yawk.lz4 fork, and the Spark 3.5 strict-warnings build and scalafix check pass.

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.
@github-actions github-actions Bot added the enhancement New feature or request label Oct 10, 2026

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 zstd and none. Wide reads paid substantial decompression cost, while Arrow’s existing LZ4 implementation was unsuitable for the intended performance improvement.
  • Design approach: Add opt-in lz4 using Spark’s bundled lz4-java library. Encode standard Arrow LZ4_FRAME buffers 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.0 on 3.4–4.1 and the at.yawk.lz4 1.11.0 fork on 4.2. Both library versions passed local codec validation. Comparison with branch-1.1 at e9efd9f764ee0a59b7898ff028d6985d4a7a28e1 confirms the intended addition of lz4, with zstd remaining 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 zstd or none paths.
  • 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 AbstractCompressionCodec preserves 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=lz4 now 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.

@andygrove
andygrove added this pull request to the merge queue Oct 11, 2026
Merged via the queue into apache:main with commit afa477f Oct 11, 2026
41 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

enhancement New feature or request

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants