Repository navigation
fix: spill a native sort before output only when the JVM reads that output - #13
Merged
Merged
Conversation
…utput spark.comet.exec.sort.spillBeforeOutputThreshold (568955d) protects a Spark sort, sort-merge join or window that reads a native sort's output through a columnar-to-row transition in the same task: the native sort cannot be made to release memory, so it spills its in-memory input before producing output and holds only the merge buffers. The threshold was a session extension, so every ExternalSorter in the native plan applied it, including sorts whose output never leaves the plan: under a window, a window group limit, a sort-merge join, an aggregate or the native shuffle writer. Those sorts wrote and read back their whole input once per task for nothing, because their consumer reserves memory from the same native pool and the sort releases it as its output is merged. SortExec now carries a spill_before_output flag, kept by with_fetch, projection swaps and filter pushdown, and ExternalSorter applies the threshold only when it is set. The planner sets it for a Sort reached from the root of the native plan through streaming operators only (Projection, Filter, Limit, Sample, Expand, Explode), whose output is what the JVM consumes. Window, partition-aggregate and sorted window operators do not use ExternalSorter and are unaffected. On a 1/40 sample of sat_product_info_loader's main stage (scan, window, sort-merge join, union, sort, native shuffle) with the joom-1.1-6 image and the patched library loaded from java.library.path, the query took 0.736 task-h with no spill, the same as threshold 0 (0.733-0.756), against 0.824-0.839 task-h and 29.2 GB spilled before. On the 2026-10-10 night, CometSort operators under native consumers spilled 21.6 TB, about one spill per task, mostly under window group limit, window, sort-merge join and native shuffle. Tests: the planner marks only sorts reached through streaming operators, and a planned SortExec carries the flag; the native sort with a JVM consumer still spills before output at the threshold, and with a native consumer it keeps its input in memory; in CometExecSuite, with a 1 byte threshold, sorts read through columnar-to-row spill and sorts under a window and a sort-merge join do not, all matching Spark. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Problem
spark.comet.exec.sort.spillBeforeOutputThresholdwas added in 568955d (#2). It protects a Spark sort, sort-merge join or window that reads a native sort's output through a columnar-to-row transition in the same task. Before it, that consumer was refused memory and failed withUNABLE_TO_ACQUIRE_MEMORY.The threshold was a session extension, so every
ExternalSorterin the native plan applied it. That includes sorts whose output never leaves the plan: under a window, a window group limit, a sort-merge join, an aggregate or the native shuffle writer. With the default threshold (offHeap / tasks / 4, about 1 GiB for 34g and 8 cores), such sorts wrote and read back their whole input once per task.Prod examples:
sat_product_info_loader, stage 4: 9422 spills, 1.41 TB.Change
SortExeccarries aspill_before_outputflag.with_fetch, projection swaps and filter pushdown keep it.ExternalSorterapplies the threshold only when the flag is set.Sortreached from the root of the native plan through streaming operators only: Projection, Filter, Limit, Sample, Expand, Explode. The output of such a sort is what the JVM consumes, so the original case stays protected.ExternalSorter, so they are unaffected. Late materialization runs inside the sameExternalSorterand follows the flag.Tests
jvm_output_sortsmarks only sorts reached through streaming operators. A plannedSortExeccarries the flag for a root sort and for a sort under a Filter, but not for a sort under another sort.sort_spills_before_output_above_the_thresholdandlate_materialized_sort_spills_before_output_above_the_thresholdstill spill with a JVM consumer. With a native consumer, the same threshold leaves the input in memory with no spill.cargo test --release -p datafusion-comet --libgives 660 passed, 0 failed.-D warnings, rustfmt and spotless are clean.Measurement
A 1/40 sample of
sat_product_info_loader's main stage, on themsaf-spark-comet:joom-1.1-6image with the patchedlibcomet.soloaded fromjava.library.path(verified via/proc/self/mapson every executor):Not to be merged on its own: it goes into the next batch release.
🤖 Generated with Claude Code