Skip to content

feat: add a spill compression codec config, defaulting to lz4 - #6698

Merged
comphead merged 3 commits into
apache:mainfrom
comphead:datafusion-conf-allowlist-6467
Oct 6, 2026
Merged

comphead merged 3 commits into
apache:mainfrom
comphead:datafusion-conf-allowlist-6467

Conversation

@comphead

@comphead comphead commented Oct 5, 2026 •

Copy link
Copy Markdown
Contributor

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_compression to uncompressed, while Spark compresses its own spill files with lz4 by default (spark.shuffle.spill.compress=true with spark.io.compression.codec=lz4). With lz4_frame spill 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 every spark.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?

  • A new config, 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 accepts lz4, zstd, and none, and defaults to lz4, so native spills are now compressed by default, as Spark's are. Any other value is rejected with a clear message instead of panicking in SessionConfig::set_str.
  • CometExecIterator.serializeCometSQLConfs sends the resolved value, and prepare_datafusion_session_context maps it to DataFusion's SpillCompression before the spark.comet.datafusion.* testing pass-through, the same way spark.comet.parquet.rowFilterPushdown.enabled is handled. An explicit spark.comet.datafusion.execution.spill_compression with respectDataFusionConfigs=true still wins, so the existing spill tests that set zstd are unchanged.
  • Comet's shuffle writers are unaffected and keep using spark.comet.shuffle.compression.codec. A separate key avoids snappy, which Arrow IPC cannot write, and the shuffle zstd level, which DataFusion's spill writer does not apply.
  • Docs: a "Compressing Spill Files" section on the memory tuning page, and a "Nested and Wide Data" section in the tuning guide overview that links the settings that matter for such data. The configuration reference picks up the new config from its description.

How are these changes tested?

A new test in CometExecSuite runs a native sort that spills under a small memory pool, once with the default codec and once with none, and checks that the default writes fewer spilled_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.

@github-actions github-actions Bot added the enhancement New feature or request label Oct 5, 2026
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.
@comphead
comphead force-pushed the datafusion-conf-allowlist-6467 branch from 268dd6b to 0247317 Compare October 5, 2026 20:24
@comphead
comphead requested a review from andygrove October 5, 2026 23:43
Comment thread docs/source/user-guide/latest/tuning.md Outdated
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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

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.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

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) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

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.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

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]] =

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

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.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

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 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.

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=true still forwards every DataFusion key. Filtering occurs during plan creation and adds no per-row processing.
  • Implementation sketch: CometConf defines the list, serializeCometSQLConfs applies 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.1 at 7b7eec69282a7abc33a0e8e687ca594c72f140ec, 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 andygrove 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.

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

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.

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?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

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 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.

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=true retains unrestricted forwarding. Filtering runs during plan creation, with no added per-row processing.
  • Implementation sketch: CometConf defines the list, serializeCometSQLConfs filters 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.1 at 7b7eec69282a7abc33a0e8e687ca594c72f140ec, 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-lz4 failure 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.
@comphead comphead changed the title feat: allowlist DataFusion configs, starting with spill compression feat: add a spill compression codec config, defaulting to lz4 Oct 6, 2026
@comphead

comphead commented Oct 6, 2026

Copy link
Copy Markdown
Contributor Author

Thanks @andygrove, reworked in 16fb61d along these lines. spark.comet.exec.allowedDataFusionConfigs is gone. spark.comet.exec.spill.compression.codec accepts lz4, zstd and none, defaults to lz4, and is validated with checkValues. serializeCometSQLConfs sends the resolved value, and prepare_datafusion_session_context maps it to DataFusion's SpillCompression before the testing pass-through, so the CometTaskMetricsSuite tests that set zstd keep working. I went with a separate key for the reasons you gave. The new CometExecSuite test checks that a spilling native sort writes fewer spilled_bytes with the default than with none.

@comphead
comphead requested a review from andygrove October 6, 2026 15:52

@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.

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, accepting lz4, zstd, and none, with lz4 as 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: CometConf defines the setting, serializeCometSQLConfs sends its resolved value, and prepare_datafusion_session_context maps it to SpillCompression.
  • Behavioral changes worth calling out: Compared with branch-1.1 at e9efd9f764ee0a59b7898ff028d6985d4a7a28e1, native spills now use lz4 by default. This intended change is documented, and none restores 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 andygrove 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.

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.

@comphead
comphead added this pull request to the merge queue Oct 6, 2026
Merged via the queue into apache:main with commit 6cd5ea0 Oct 6, 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.

Investigate nested TPC-H q21 slowdown: Comet 37% slower than Spark at SF1000

4 participants