Repository navigation
feat: add a spill compression codec config, defaulting to lz4 - #6698
Conversation
Add `spark.comet.exec.allowedDataFusionConfigs`, a list of `spark.comet.datafusion.*` keys that are passed to native execution even when the testing flag `spark.comet.exec.respectDataFusionConfigs` is false. The default list holds `spark.comet.datafusion.execution.spill_compression`, so compressing DataFusion's spill files no longer needs the testing flag. Spills stay uncompressed unless that key is set. Document spill compression on the memory tuning page and add a nested and wide data section to the tuning guide. Closes apache#6467.
268dd6b to
0247317
Compare
| Rows with many columns, long strings, or nested columns such as arrays of structs make every batch | ||
| larger, so queries over such data spill more. For these workloads: | ||
|
|
||
| - Compress spill files, which Comet writes uncompressed by default. See |
There was a problem hiding this comment.
Native shuffle already compresses its spills with the shuffle codec, and memory.md says so. Could this bullet say something like "Compress spill files from native sorts, aggregations, and joins, which DataFusion writes uncompressed by default"? Otherwise readers with wide or nested data may expect it to affect shuffle spill too.
There was a problem hiding this comment.
Reworded. The bullet now names native sorts, aggregations, and joins, and reflects the new lz4 default.
| assert(entries.get(spillCompression) == "zstd") | ||
| assert(!entries.containsKey(spillReservation)) | ||
|
|
||
| withSQLConf(CometConf.COMET_ALLOWED_DATAFUSION_CONFIGS.key -> spillReservation) { |
There was a problem hiding this comment.
Since the doc says setting the list replaces the default, could the inner block also assert !entries.containsKey(spillCompression)? As written, it would still pass if the list were appended to the default.
There was a problem hiding this comment.
The allowlist is gone, so this test was replaced by one that compares the spilled_bytes of a spilling native sort with the default codec and with none.
| .booleanConf | ||
| .createWithDefault(false) | ||
|
|
||
| val COMET_ALLOWED_DATAFUSION_CONFIGS: ConfigEntry[Seq[String]] = |
There was a problem hiding this comment.
Thanks for chasing this down, the q21 numbers make a great case for exposing spill compression. I'm wondering about the shape though. Because this list lives in the tuning category, the versioning policy covers it, and per config_conventions.md so does the documented spark.comet.datafusion.execution.spill_compression key. Its accepted values
and default come from DataFusion's SpillCompression, so a DataFusion upgrade could change a key we've promised to keep stable. Users can also list any spark.comet.datafusion.* key here, as the new test does with sort_spill_reservation_bytes. That makes it a supported per-key way around respectDataFusionConfigs. Would a dedicated Comet key work
better, for example spark.comet.exec.spill.compression with checkValues(Set("uncompressed", "lz4_frame", "zstd"))? It could be translated in prepare_datafusion_session_context the same way COMET_PARQUET_ROW_FILTER_PUSHDOWN_ENABLED is, and added to the resolved list in serializeCometSQLConfs. That would also catch a typo like lz4 on the driver.
Today it panics in SessionConfig::set_str on every task.
There was a problem hiding this comment.
Done in 16fb61d, as andygrove also suggested. spark.comet.exec.spill.compression.codec is translated in prepare_datafusion_session_context the same way COMET_PARQUET_ROW_FILTER_PUSHDOWN_ENABLED is, and its resolved value is in the list in serializeCometSQLConfs. I used lz4, zstd and none rather than DataFusion's value names, so the accepted values don't depend on DataFusion's SpillCompression, and checkValues rejects a bad value with a clear message instead of a panic.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Spill compression required enabling the testing-only DataFusion configuration pass-through. This PR exposes selected keys for production tuning.
- Design approach: Add
spark.comet.exec.allowedDataFusionConfigs, defaulting to the spill-compression key, and filter configuration forwarding through that list. - Correctness / compatibility analysis: No introduced P1/P2 issues found within this review. Reviewed the full six-file diff against
38d941f8207241ff6663c51496a5df2c3722f856, surrounding code, and existing discussions. Checked configuration and spill semantics against Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3, and 4.2.0 sources and DataFusion 55.1.0. No substantiated existing P1/P2 blockers remain. - Key design decisions: Unset compression remains
uncompressed. A custom allowlist replaces the default.respectDataFusionConfigs=truestill forwards every DataFusion key. Filtering occurs during plan creation and adds no per-row processing. - Implementation sketch:
CometConfdefines the list,serializeCometSQLConfsapplies it, and the existing JNI path configures DataFusion. The change reuses existing spill implementations and leaves buffer ownership and memory accounting unchanged. - Behavioral changes worth calling out: Compared with
branch-1.1at7b7eec69282a7abc33a0e8e687ca594c72f140ec, explicitly configured allowlisted keys now take effect without the testing flag. This is intentional and documented. Comet shuffle compression remains separate. - Suggested improvements: None meeting the P1/P2 evidence bar beyond the existing discussion.
Reviewed full SHA: d56c0b5f27e7c79e791de14dab94ce661d123f07. The PR remained open and non-draft. Routed skills: review-comet-pr, review-comet-ffi-pr, and review-comet-memory-pr.
Exact-head CI: 25 checks passed and 15 were skipped, with no failures or pending checks. Passing checks include native tests, Linux Spark 4.1 execution suites, and TPC-H/TPC-DS result verification. The execution report confirms the new forwarding test passed. Upstream Spark SQL, Iceberg, and macOS suites were skipped.
Validation: A disposable DataFusion 55.1.0 probe round-tripped 8,000 sliced, nullable nested rows through all three codecs with identical results. It also confirmed the documented invalid-lz4 configuration failure. The probe used cached dependency libraries and verified upstream sources, not a full Comet integration build. No full local build, forced-spill Spark integration run, or performance benchmark was performed. Project files were unchanged.
andygrove
left a comment
There was a problem hiding this comment.
Spill compression still defaults to uncompressed here, so the q21 slowdown in #6467 stays the out-of-the-box behavior. Users only get the fix if they find the tuning guide and set a raw DataFusion key, using DataFusion's value names (lz4_frame, where our shuffle codec and Spark both say lz4). I'd rather Comet own this setting and always pass it to DataFusion.
Could we replace spark.comet.exec.allowedDataFusionConfigs with a Comet config for the spill codec, as @kazuyukitanimura suggested? For example spark.comet.exec.spill.compression.codec, accepting lz4, zstd and none, with lz4 as the default. Spark compresses its own sort and aggregation spills with lz4 by default (spark.shuffle.spill.compress=true with spark.io.compression.codec=lz4), so that default matches Spark, and your q21 numbers show what it buys. The resolved value would go in the list at the end of serializeCometSQLConfs, and prepare_datafusion_session_context would map it to datafusion.execution.spill_compression before the spark.comet.datafusion.* pass-through, the same way COMET_PARQUET_ROW_FILTER_PUSHDOWN_ENABLED is handled. The testing pass-through would still win, so the existing spill tests that set zstd keep working. checkValues would also reject a bad value with a clear message. With this PR, a value like lz4 panics in SessionConfig::set_str in every native plan.
Reusing spark.comet.shuffle.compression.codec would also work if you'd rather not add a key. I lean toward a separate one because the shuffle codec also accepts snappy, which Arrow IPC can't write, and spark.comet.shuffle.compression.zstd.level wouldn't apply because DataFusion's spill writer always uses Arrow's default zstd level.
Since the default would change, could you also share numbers for the flat TPC-H SF1000 run with the new default? It would also be good to have a test where a native sort spills under a small pool and writes fewer spilled_bytes with the default codec than with none. That would show DataFusion actually compressing, not just the key crossing JNI.
|
|
||
| ## Compressing Spill Files | ||
|
|
||
| Native sorts, aggregations, and sort-merge joins spill through DataFusion, which writes spill files |
There was a problem hiding this comment.
DataFusion 55's NestedLoopJoinExec also spills through the same SpillManager when its build side doesn't fit, and Comet's native broadcast nested loop join is on by default. Could this list nested loop joins too?
There was a problem hiding this comment.
Added. When DataFusion 55.1's NestedLoopJoinExec falls back to spilling, it writes both sides through SpillManager with the session's spill_compression, so the new codec covers it.
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Configuring native spill compression required the testing-only DataFusion pass-through flag.
- Design approach: Add
spark.comet.exec.allowedDataFusionConfigs, defaulting to the spill-compression key, and filter forwarding through that list. - Correctness / compatibility analysis: No introduced P1/P2 issues found within this review. Reviewed the full six-file diff against
38d941f8207241ff6663c51496a5df2c3722f856, surrounding code, and existing discussions. Checked relevant configuration and spill semantics against Spark 3.4.3, 3.5.9, 4.0.4, 4.1.3, and 4.2.0 sources and DataFusion 55.1.0. - Key design decisions: Unset compression remains
uncompressed. A custom allowlist replaces the default.respectDataFusionConfigs=trueretains unrestricted forwarding. Filtering runs during plan creation, with no added per-row processing. - Implementation sketch:
CometConfdefines the list,serializeCometSQLConfsfilters explicitly configured keys, and existing JNI code configures DataFusion. Spill implementations, buffer ownership, and memory accounting remain unchanged. - Behavioral changes worth calling out: Compared with
branch-1.1at7b7eec69282a7abc33a0e8e687ca594c72f140ec, configured allowlisted keys now take effect without the testing flag. This is intentional and documented. Shuffle compression remains separate. - Suggested improvements: No additional P1/P2 changes identified. The existing request-changes discussion about a dedicated codec setting and compressed default remains open. The documented invalid-
lz4failure reproduces and is caught at JNI. I did not substantiate an existing introduced P1/P2 regression.
Reviewed full SHA: d56c0b5f27e7c79e791de14dab94ce661d123f07. The PR remains open and non-draft. Routed skills: review-comet-pr, review-comet-ffi-pr, and review-comet-memory-pr.
Exact-head CI: 25 checks passed, 15 skipped, none failed or pending. The downloaded execution report confirms the new forwarding test passed among 1,308 successful tests. Native tests and TPC-H/TPC-DS result checks passed. Upstream Spark SQL, Iceberg, and macOS suites were skipped.
Validation: Rebuilt disposable probes using cached DataFusion 55.1.0 and Arrow 59.3.0 libraries. All three codecs preserved 8,000 sliced nested rows. A sort under a 512 KiB pool spilled 13 times and returned identical 20,000-row results with each codec. Reported spill bytes were 3,679,424 uncompressed, 955,736 with lz4_frame, and 593,440 with zstd. Relevant upstream sources were verified. These were dependency-level probes, not a full Comet integration build. No local Spark integration run or performance benchmark was performed. Project files remain unchanged.
Review state: Comment. This assessment does not resolve the existing API/default discussion.
Replace `spark.comet.exec.allowedDataFusionConfigs` with `spark.comet.exec.spill.compression.codec`, which sets the codec for the files that native sorts, aggregations, sort-merge joins, and nested loop joins spill through DataFusion. It accepts `lz4`, `zstd`, and `none`, and defaults to `lz4`, the codec Spark uses for its own spill files by default, so native spills are now compressed by default. The JVM sends the resolved value, and the native side maps it to DataFusion's `SpillCompression` before the `spark.comet.datafusion.*` testing pass-through, which still wins. Add a test that a spilling native sort writes fewer bytes with the default codec than with `none`, and list nested loop joins in the spill compression docs.
|
Thanks @andygrove, reworked in 16fb61d along these lines. |
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Native spills were uncompressed by default, and configuring compression required the testing-only DataFusion pass-through.
- Design approach: Add
spark.comet.exec.spill.compression.codec, acceptinglz4,zstd, andnone, withlz4as the default. - Correctness / compatibility analysis: No introduced P1/P2 issues found within this review. Reviewed the entire seven-file diff against
38d941f8207241ff6663c51496a5df2c3722f856, surrounding code, and existing discussions. Verified relevant Spark configuration and serialization sources across 3.4.3, 3.5.9, 4.0.4, 4.1.3, and 4.2.0. The default matches Spark’s spill-compression default. - Key design decisions: Comet owns and validates the public values. The testing pass-through retains precedence. Shuffle compression remains separate. Configuration mapping runs once per native plan, and compression reuses DataFusion’s existing implementation without introducing another abstraction.
- Implementation sketch:
CometConfdefines the setting,serializeCometSQLConfssends its resolved value, andprepare_datafusion_session_contextmaps it toSpillCompression. - Behavioral changes worth calling out: Compared with
branch-1.1ate9efd9f764ee0a59b7898ff028d6985d4a7a28e1, native spills now uselz4by default. This intended change is documented, andnonerestores uncompressed spills. Compression adds codec work while reducing disk bytes on the tested data. Runtime performance was not benchmarked. - Suggested improvements: None meeting the P1/P2 evidence bar. Earlier API/default concerns are addressed, and no substantiated existing P1/P2 blocker remains.
Reviewed full SHA: 16fb61d17c85d4d1a09dc63af5f5b83b23ffae27. The PR remains open and non-draft. Routed skills: review-comet-pr, review-comet-ffi-pr, and review-comet-memory-pr.
Exact-head CI: 18 checks passed, 14 skipped, and 7 remain in progress. No failures were reported. Native compilation passed. Rust tests, four Linux Spark 4.1 test groups, and TPC-H/TPC-DS verification remain pending. Upstream Spark SQL, Iceberg, and macOS suites were skipped.
Validation: Recompiled disposable probes using verified DataFusion 55.1.0 and Arrow 59.3.0 sources and cached libraries. A 20,000-row sort under a 512 KiB pool spilled 13 times with identical results across all codecs. Spill bytes were 3,679,424 with none, 955,736 with default lz4, and 593,440 with zstd. All codecs also preserved 8,000 sliced nested rows. Exact-head native mapping, invalid-value rejection, and override checks passed. These were dependency-level probes. No full local Comet build or Spark integration run was performed, and the requested flat SF1000 benchmark remains unavailable. Project files were unchanged.
andygrove
left a comment
There was a problem hiding this comment.
Thanks @comphead, this is what I was hoping for. The codec is Comet's own, lz4 is the default, a bad value is rejected by checkValues, and the testing pass-through still wins. I checked that last part by running the spilling sort with both keys set, and the pass-through won whichever way I set them.
I'm no longer asking for the flat SF1000 run, because I measured the codec cost directly. I wrote and read back TPC-H SF1 lineitem in 8192-row batches with the same Arrow IPC options DataFusion 55.1 uses for its spill files. lz4 made them 2.5x smaller and cost about 1.1 s of extra CPU per GB spilled. zstd made them 4.3x smaller for about 2 s per GB. A disk slower than roughly 550 MB/s comes out ahead with lz4, and on fast local NVMe it is a small net cost that none removes. That seems a fair price for matching Spark's default.
I'd like two small additions to the tests before this merges. The description says a bad value is rejected with a clear message, but nothing checks it. Could we add a test that serializeCometSQLConfs throws for a value like lz4_frame under withSQLConf? It would also pin the entry in the resolved list. Nothing does that today, because the native side falls back to lz4 when the key is missing. I deleted the entry and the new test and both SQLConf serde tests still passed.
Second, zstd has no coverage through the new key, since the existing zstd tests set the DataFusion key. Could the spill test also run with zstd and assert that it writes fewer bytes than the default? On this data I get 183760 bytes for zstd, 340392 for the default and 9041216 for none.
Which issue does this PR close?
Closes #6467.
Rationale for this change
On the nested TPC-H benchmark in #6467, Comet lost q21 to Spark because its spills were uncompressed. DataFusion 55.1 defaults
datafusion.execution.spill_compressiontouncompressed, while Spark compresses its own spill files withlz4by default (spark.shuffle.spill.compress=truewithspark.io.compression.codec=lz4). Withlz4_framespill compression, Comet's q21 went from 380 s to 250 s, against 367 s for Spark (Spark 4, Iceberg, SF1000, numbers from the issue).Until now the DataFusion key was reachable only through the testing flag
spark.comet.exec.respectDataFusionConfigs, which forwards everyspark.comet.datafusion.*key, and the versioning policy says a testing key must not be the only way to reach a behavior that production users need.What changes are included in this PR?
spark.comet.exec.spill.compression.codec, in the tuning category. It sets the codec for the files that native sorts, aggregations, sort-merge joins, and nested loop joins spill through DataFusion. It acceptslz4,zstd, andnone, and defaults tolz4, so native spills are now compressed by default, as Spark's are. Any other value is rejected with a clear message instead of panicking inSessionConfig::set_str.CometExecIterator.serializeCometSQLConfssends the resolved value, andprepare_datafusion_session_contextmaps it to DataFusion'sSpillCompressionbefore thespark.comet.datafusion.*testing pass-through, the same wayspark.comet.parquet.rowFilterPushdown.enabledis handled. An explicitspark.comet.datafusion.execution.spill_compressionwithrespectDataFusionConfigs=truestill wins, so the existing spill tests that setzstdare unchanged.spark.comet.shuffle.compression.codec. A separate key avoidssnappy, which Arrow IPC cannot write, and the shuffle zstd level, which DataFusion's spill writer does not apply.How are these changes tested?
A new test in
CometExecSuiteruns a native sort that spills under a small memory pool, once with the default codec and once withnone, and checks that the default writes fewerspilled_bytes. DataFusion counts the bytes it writes to the spill file, so this shows the codec reaching DataFusion's spill writer, not just crossing JNI.