Skip to content

fix: spill a native sort before output only when the JVM reads that output - #13

Merged
msafonov merged 1 commit into
branch-1.1from
joom/1.1-fix-sort-spill-jvm-consumer
Oct 10, 2026
Merged

msafonov merged 1 commit into
branch-1.1from
joom/1.1-fix-sort-spill-jvm-consumer

Conversation

@msafonov

Copy link
Copy Markdown
Collaborator

Problem

spark.comet.exec.sort.spillBeforeOutputThreshold was 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 with UNABLE_TO_ACQUIRE_MEMORY.

The threshold was a session extension, so every ExternalSorter in 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.
  • On the 2026-10-10 night, CometSort operators under native consumers spilled 21.6 TB, about one spill per task:
Consumer Spilled
CometWindowGroupLimitExec 12.6 TB
CometWindowExec 3.7 TB
CometSortMergeJoin 3.7 TB
native CometExchange 2.6 TB

Change

  • SortExec carries a spill_before_output flag. with_fetch, projection swaps and filter pushdown keep it. ExternalSorter applies the threshold only when the flag is set.
  • The planner sets the flag for a Sort reached 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.
  • Window, partition-aggregate and sorted window operators do not use ExternalSorter, so they are unaffected. Late materialization runs inside the same ExternalSorter and follows the flag.

Tests

  • Planner (Rust): jvm_output_sorts marks only sorts reached through streaming operators. A planned SortExec carries the flag for a root sort and for a sort under a Filter, but not for a sort under another sort.
  • Native sort: sort_spills_before_output_above_the_threshold and late_materialized_sort_spills_before_output_above_the_threshold still spill with a JVM consumer. With a native consumer, the same threshold leaves the input in memory with no spill.
  • CometExecSuite: with a 1 byte threshold, sorts read through columnar-to-row (ORDER BY, and ORDER BY with a filter and projection) spill. Sorts under a window and under a sort-merge join do not. All results match Spark.
  • Native suite: cargo test --release -p datafusion-comet --lib gives 660 passed, 0 failed.
  • Scala suites: CometExecSuite, CometWindowExecSuite, CometPartitionAggregateWindowSuite, CometJoinSuite, WideRowSortFallbackSuite, TPC-DS v1.4 and v2.7 plan stability. Result: 513 succeeded, 0 failed, 2 canceled. Both canceled tests are pre-existing version-gated ones.
  • Lint: clippy -D warnings, rustfmt and spotless are clean.

Measurement

A 1/40 sample of sat_product_info_loader's main stage, on the msaf-spark-comet:joom-1.1-6 image with the patched libcomet.so loaded from java.library.path (verified via /proc/self/maps on every executor):

Variant task-h spilled
joom-1.1-6, default threshold 0.824 / 0.839 29.2 GB
joom-1.1-6, threshold 0 0.747 / 0.756 / 0.733 0
this PR, default threshold 0.736 0

Not to be merged on its own: it goes into the next batch release.

🤖 Generated with Claude Code

…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>
@github-actions github-actions Bot added the bug Something isn't working label Oct 10, 2026
msafonov added a commit that referenced this pull request Oct 10, 2026
@msafonov
msafonov merged commit 03040f3 into branch-1.1 Oct 10, 2026
36 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

bug Something isn't working

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant