Skip to content

feat: Joom 1.1-1 — Delta scan, correctness fixes, wide-row performance, cost-based engine choice - #2

Merged
msafonov merged 72 commits into
branch-1.1from
joom/1.1-dev
Oct 3, 2026
Merged

msafonov merged 72 commits into
branch-1.1from
joom/1.1-dev

Conversation

@msafonov

@msafonov msafonov commented Oct 3, 2026 •

Copy link
Copy Markdown
Collaborator

Joom build of Comet 1.1 on top of upstream branch-1.1 (499ee08). Proposed version name: joom-1.1-1.

What is in it

Disabled by default: all plan rules (cost-based choice, boundary formats, wide-row sort/shuffle fallback) and the new native window operator PartitionAggregateWindowExec (spark.comet.exec.window.partitionAggregate.enabled; when off, upstream WindowAggExec is used). Workloads enable them explicitly.

Active on default settings (behavior differs from upstream):

  • spill of a large in-memory native sort before producing output (spark.comet.exec.sort.spillBeforeOutputThreshold, default off-heap / task slots / 4; 0 disables) — extra disk IO for big sorts;
  • coalescing of small blocks on Comet shuffle read (spark.comet.shuffle.read.coalesce.enabled, default true);
  • native sort of wide rows uses late materialization and coalesces small input batches (no switch; rows with equal keys may come out in a different order — sort stability is not guaranteed anyway);
  • generated columnar-to-row projections are pooled and reused across tasks (no switch);
  • vendored DataFusion aggregate fix: after spilling an empty table the hash table is resized, not try_resized (85bef53).

Verification

  • Benchmarks. 8 production jobs, Comet with the cost-based choice vs vanilla Spark on the same inputs: −15% task hours in total (QL −46…−52%, ih24h −36%, ExternalLinks −20%, FactOrder −13%, mongo −6%; cube +13%, InstallCube −6% on pr30).
  • Data compare of all 8 jobs vs vanilla: no engine differences. The remaining diffs are float/decimal last-digit noise, non-deterministic queries, and inputs rewritten between runs.
  • Tests on default settings:
    • Rust: cargo test --workspace 1880 passed, 0 failed;
    • clippy -D warnings clean;
    • vendored DataFusion sorts::/joins:: 1151 passed;
    • JVM: all CI suite groups pass;
    • 3 S3 suites pass with Docker + AWS_REGION set: they need a MinIO image and a region in the environment, are not run by CI, and minio/minio is gone from Docker Hub;
    • Delta contrib passes with a UTF-8 locale.
  • CI (PR tier): green — Comet suites on Spark 4.1, TPC-H/TPC-DS result checks, Rust tests, lint/compile for Spark 3.4/3.5/4.0. Also builds with Rust 1.99 (deprecations in upstream code silenced or fixed) and Scala 2.13.

🤖 Generated with Claude Code

notimesea and others added 30 commits September 26, 2026 22:38
Defer native scan DPP partition enumeration until execution so AQE can
coalesce a sibling shuffle without executing an adaptive placeholder.

Pass codegen scalar subqueries as runtime inputs and broadcast length-one
Arrow arguments, including the null fast path. Normalize Etc/UTC in native
timestamp truncation to match scan timestamp types.

Scale sort's eager spill reserve with small off-heap task budgets. Spill
whole-partition aggregate window rows while retaining the existing native
accumulators and their Spark semantics.

Add regressions for planning, multi-row scalar inputs, timezone aliases,
reservation sizing, window spilling, ordering and cancellation cleanup.
Restore DataFusion's default sort_spill_reservation_bytes. The 1/32 cap only
lowers the reservation when the per-task off-heap share is under 320 MiB, so
it is a no-op on our executors (about 48 MiB and 136 MiB per task), and the
native sort merge still fails on the repro query with it applied.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Window nodes with an expression that cannot run with bounded memory used
DataFusion's WindowAggExec, which buffers each partition in memory without
spilling. Generalize PartitionAggregateWindowExec so every natively supported
shape spills instead:

- whole-partition first_value/last_value/nth_value (incl. IGNORE NULLS) track
  the selected row while rows are ingested;
- ntile and percent_rank are computed during the replay from the partition
  size counted on the first pass;
- cume_dist and frames ending at UNBOUNDED FOLLOWING (ROWS/RANGE starting at
  CURRENT ROW, N PRECEDING or N FOLLOWING) use a reverse pass over a narrow
  reverse-order copy of the spilled rows; the per-start frame values are
  buffered in a spillable reservation and read back at each row's frame start;
- mixed nodes evaluate their streaming expressions in a BoundedWindowAggExec
  below the operator and restore the column order with a projection.

WindowAggExec remains only for expressions without a spilling implementation.
Existing whole-partition aggregate behaviour and metrics are unchanged.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Cover whole-partition first/last/nth_value (with IGNORE NULLS), ntile,
percent_rank and cume_dist with ties and partitions smaller than the bucket
count, a mixed window node, and ROWS/RANGE frames ending at UNBOUNDED
FOLLOWING, asserting that the window stays native.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…rinks

The sort fails with ResourcesExhausted in the in-memory merge's cursor
reservation once more tasks become active and Spark lowers the task's
share below what the sorter already holds.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Adds the crates.io source of datafusion-physical-plan 55.1.0 under
native/vendor and points [patch.crates-io] at it, so the next commit can
carry a small, reviewable patch to its external sort.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
When the external sorter spills, DataFusion 55.1 frees the merge
pre-reservation and lets the sorted runs release their memory to the pool,
then has the merge's cursors, encoded rows and batch buffers request it
again under the non-spillable ExternalSorterMerge consumer. In Comet the
pool hands those bytes back to Spark. Once more tasks are active, Spark's
per-task share is below what the sorter held and it grants nothing, so the
sort fails with ResourcesExhausted although it holds enough memory and has
nothing left to spill.

Patch the vendored datafusion-physical-plan so a spill merges inside the
reservations the sorter already holds: a SpillWorkspace pool takes over the
sorter's and the merge's reservations for the duration of the spill, its
children (runs, cursors, rows, merge buffers) share them, and only growth
beyond them goes to the execution pool under its usual limits. The
workspace is released when the spill ends. This extends the approach of
apache/datafusion#24740, which retains only sort_spill_reservation_bytes
(capped to 1/32 of the task budget in Comet, too small for the merge).

Tests cover the shrinking share at several points, smaller output batches,
many small input batches, a key-only row, a fixed share, and both sorts of
a sort-merge join, and check that the pool and Spark are back to zero.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
MultiLevelMergeBuilder frees the memory reserved for a merge pass when the
pass ends (StreamAttachedReservation) and when it retries a pass with less
read-ahead, so the sort_spill_reservation_bytes headroom that sort() hands
over protects only the first pass. Later passes must win the memory back from
the pool, which fails once another consumer or a shrinking Spark share has
taken it (apache/datafusion#25804, finding 2).

Merge the spill files inside a SpillWorkspace built from the merge
reservation: every pass's reservations are its children, what one pass
releases stays reserved for the next, and the workspace is closed when the
final pass is chosen, so it keeps only what that pass holds. This follows
apache/datafusion#24740, which keeps the headroom across passes with a
MergeMemoryPool.

The test hands every byte the sort releases to another consumer once the
merge starts, for a pass that ends and for a pass that falls back to a
smaller read-ahead.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
RowCursorStream keeps two Rows buffers per input stream for reuse, but only a
live cursor's reservation covers one. Once the cursor is dropped its buffer
stays allocated in ReusableRows with nothing reserved for it, until the merge
ends (apache/datafusion#25804, finding 4).

Following apache/datafusion#25372, ReusableRows now holds one reservation
that covers every buffer it keeps, resized when a buffer is refilled, and the
cursor gets an empty reservation so the bytes are not counted twice. A
finished stream drops the buffers no cursor still holds and shrinks the
reservation.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
FieldCursorStream::convert_batch reserves the buffers of the sort key it
evaluates. When the key is a column of the batch, those are the batch's own
buffers, which BatchBuilder::push_batch charges as part of the batch, so a
single-column merge counted its key twice (apache/datafusion#25804,
finding 7).

Reserve nothing for a key that is a column of the batch: the batch is held
and charged by the merge's BatchBuilder at least as long as the cursor over
it. A key the sort expression computes is still reserved in full.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The final pass of a sort's spill merge seats as many runs as the pool grants
(fan-in is unlimited by default), so it grows until the pool refuses and holds
nearly the whole task share while its output is read. The operators reading
that output then find nothing left, and in Comet an infallible grow() is
recorded as overcommit instead (apache/datafusion#25804, finding 9). In a
Comet test with a fixed 32 MiB share the final pass held 30.7 MiB.

As apache/datafusion#25383 does for aggregate spill merges, admit the final
pass only if the pool could grant as much again as its buffers need,
checked through SpillWorkspace::can_grow, which counts unused workspace and
asks the pool only for the rest, then gives it back. Try less read-ahead
first. Otherwise merge only enough runs in this pass that the final one would
fit twice in what the merge holds, and put the rest back, so the merge goes
multi-pass. The smallest merge, two runs, still runs without the spare and
without read-ahead. All of it stays reserved as before, and only a sort's
merge (one in a SpillWorkspace) is affected.

Tests: a vendored sort whose final pass must leave room for a consumer, and a
Comet sort with a fixed Spark share whose final merge must hold at most half
of it, with the pool's peak within the share.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
With rows of a few KiB and Comet's batch size, one spill batch can need more
than a merge stream may reserve (twice its size): every run is a single batch,
and a merge pass that consumed a split run writes batches of up to half the
batch size, which the next pass could not seat at all. The merge then failed
with ResourcesExhausted although the batch could be split.

For a sort's merge, re-spill such a run in halves instead of failing, as is
already done when two runs do not fit. The re-spill now reads the run without
read-ahead, which held up to three batches while two were reserved, and
reserves what it holds, the batch read plus one half encoded for the new file,
so a batch too wide to reserve twice can still be split. It still fails for a
single row, or when even that cannot be reserved.

The test sorts ~4.5 KiB rows in a fixed 2 MiB Spark share with a batch size
whose full batch exceeds the share, and checks the output, many spills and a
multi-pass merge, a pool peak within the share, and that memory returns to 0.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
consume_and_spill_append freed the sorter's reservation before writing the
sorted batches it had taken, so they stayed in memory unaccounted while the
spill file was written (apache/datafusion#25804, finding 3).

Take the reservation instead and drop it when the batches have been written
or the write fails, as the reservation part of apache/datafusion#24923 does
(backport apache/datafusion#25806). The async spill-writing API of #24923 is
not ported.

The test records the pool's reservation on every spill write and checks it
covers the buffered input the spill workspace still holds plus the sorted
batch.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
sort_batch_stream sorts a batch into batch_size chunks with take, which
shares dictionary values and view data buffers between the chunks, and then
summed get_record_batch_memory_size over the chunks, charging the shared
buffers once per chunk (apache/datafusion#25804, finding 6). A batch that fit
its reservation could not be sorted in it.

Port apache/datafusion#25800: count the chunks with one
RecordBatchMemoryCounter, charging each shared buffer to the last chunk that
holds it, and release each chunk's share as it is output. ReservationStream,
used only here, is removed with its tests as in the PR. get_sliced_size now
counts a view data buffer listed more than once, in one array or across the
batch's arrays, once, as the PR does; 55.1.0 adds those buffers by hand, per
array.

Tests: one Utf8View batch and one dictionary batch are each sorted into four
chunks in a pool holding just the batch's reservation, and the PR's
get_sliced_size test for a buffer listed three times in two columns.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
reserve_memory_for_batch_and_maybe_spill charged every input batch its full
buffer capacity, so each zero-copy slice of one parent batch, as
AggregateExec emits for EmitTo::All, was charged the whole parent. Sorting 64
slices of a 512 KiB batch in a 2 MiB pool spilled 21 times
(apache/datafusion#25804, finding 5).

As apache/datafusion#22862 does for the hash join build side, count the
buffered batches with one RecordBatchMemoryCounter, so a shared buffer is
reserved in full with the first batch holding it; the sliced-size half of the
estimate is unchanged. The counter restarts whenever the buffered batches
are sorted. The runs of an in-memory merge split the reservation the same
way, a shared buffer going with its first run, and coalescing realigns to
that total. sort_batch_stream no longer asserts that its reservation is the
batch's full estimate.

The test sorts the 64 slices, concatenated and merged as runs, without a
spill, and checks the reservation is the parent once plus the slices' rows.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
When nothing was spilled, sort() freed the sort_spill_reservation_bytes
headroom to the pool and then merged the buffered batches with a new empty
reservation, so the merge's buffers had to win that memory back. Once the
pool had given it to another consumer or Spark had lowered the task's share,
the merge failed with ResourcesExhausted, and the final merge cannot spill
(apache/datafusion#25804, finding 1). The spill path and the concat path
during a spill already merge in a SpillWorkspace since 880575b; this was the
remaining case, open on main as well.

Merge the buffered batches in a SpillWorkspace built from the merge's and
the sorter's reservations, as a spill does, but keep only up to the headroom
unused: the merge's cursors and buffers use it first, and what the runs
release beyond it goes back to the pool while the output is read. Growth is
charged to the merge's consumer, as before. SpillWorkspace::keep_at_most
limits the retention; close() is keep_at_most(0). A single batch or a
concatenated sort has no merge and still returns the headroom.

The test gives the pool left after sort() to another consumer and hands it
everything the sort releases, and checks the merge completes.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
A spill merge pass admitted each run at twice its largest batch per
read-ahead slot and then merged against an unbounded pool. That leaves out
the batch the merge holds while spawn_buffered refills the read-ahead (one
even with read-ahead 1), the rows the cursor encodes the batch's sort key
into, which a row cursor keeps two buffers of, and the batch the output
builder keeps when a cursor moves to the next one
(apache/datafusion#25804, finding 8; apache/datafusion#23760).

For a sort's merge (one in a SpillWorkspace) reserve per run the read-ahead,
the merge's batch and its rows (one buffer for a single primitive or
string/binary key, two for a row cursor), each estimated at the run's
largest batch, plus the largest batch once per pass. The final-pass bound
uses the same sizes. Other merges keep DataFusion's estimate. With
read-ahead 2 a single-column run reserves as before, plus the one crossing
batch per pass; a row-cursor run reserves one batch more, and with
read-ahead 1 each run reserves one or two batches more. apache/datafusion#25565
(output construction headroom) is not ported.

The test checks the admission for one- and two-column runs, and that it
falls back to read-ahead 1 when read-ahead 2 no longer fits.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Revert d153b3a to restore the per-task initial reservation cap and the off_heap_limit argument used by the sort spill regressions. Native library tests: 512 passed, 5 ignored.

Co-Authored-By: Codex <noreply@openai.com>
Avoid reading session SQLConf while validating the static default. Clamp both JVM shuffle writers to the execution batch limit. CometConfSuite: 17 tests passed.

Co-Authored-By: Codex <noreply@openai.com>
Opt in with COMET_DEBUG_BATCH_MEMORY=1; count shared Arrow buffers once. Native tests: 512 passed, 5 ignored.

Co-Authored-By: Codex <noreply@openai.com>
Compare 8192-row and 512-row output batches under the same fixed memory share. Prove the final large batch survives after the reservation is released, while small output bounds the unreserved handoff. Regression: 1 passed.

Co-Authored-By: Codex <noreply@openai.com>
Compare Binary and BinaryView take/merge pipelines with identical materialized outputs, including narrow-data overhead.

Co-Authored-By: Codex <noreply@openai.com>
Select BinaryView for wide non-key binary payloads inside full external sort and materialize the original schema at output. Preserve narrow, ordered and TopK paths. Cover slices, nulls, metadata, bounded spill (including zstd), reservation release and JVM round trips.

Co-Authored-By: Codex <noreply@openai.com>
Compare dictionary-to-Binary and dictionary-to-BinaryView boundaries and equivalent materialized sort outputs, with narrow controls and nullable repeated payloads.

Co-Authored-By: Codex <noreply@openai.com>
Serialize ShuffleScan for newly converted exchanges when direct read is enabled, instead of relying on later AQE replacement of Scan inputs retained by native parents. Cover JVM/native shuffle, AQE on/off, both direct-read settings, and nullable wide binary results in initial and executed plans.

Co-Authored-By: Codex <noreply@openai.com>
Co-Authored-By: Codex <noreply@openai.com>
Borrow plain Arrow binary spans only during the non-codegen UnsafeProjection copy; retain the ordinary path for other vector representations.

Co-Authored-By: Codex <noreply@openai.com>
Co-Authored-By: Codex <noreply@openai.com>
…e row conversion

When the JVM columnar shuffle rejects a dictionary because it is not
efficient and every non-null key refers to a distinct value in insertion
order, return the dictionary values buffer as the plain array instead of
copying it through cast. Adds a wide Binary row_columnar benchmark.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
msafonov and others added 14 commits October 1, 2026 22:01
Port apache/datafusion#24573 (d1fe988942) to the vendored DataFusion 55.1.
A LEFT or RIGHT sort-merge join with a join filter stages its filtered
output in a second coalescer, but the final flush at end of input emitted
its batch directly, ahead of rows still buffered there. The join advertises
maintains_input_order, so an aggregate grouping on the streamed key in the
same native plan ran in sorted mode, closed groups early and emitted the
same group twice with partial aggregates. The final batch now goes through
the coalescer like every other filtered batch.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Every map task writes one block per reduce partition, so with many
partitions a reducer reads hundreds of blocks of a few rows each and
passes each one on as its own batch. On wide rows the per-batch cost of
every column (FFI export and import, every native operator's per-batch
work) then dominates the read.

The local block store reader now reads all fetched blocks as one stream
and decodes them natively into a ShuffleReadCoalescer, exporting one
batch of at least spark.comet.batchSize rows to the JVM. ShuffleScan
does the same for native consumers reading the shuffle directly,
reconciling each block with the declared schema before joining it.
Controlled by spark.comet.shuffle.read.coalesce.enabled (default true).
The Celeborn JVM reader keeps decoding block by block.

Coalescing hands a final aggregate several partial states per batch,
so the bloom filter aggregate now merges every row of a batch instead
of asserting there is one.

Wide shuffle read benchmark, 512 leaf columns, 64 maps x 250 partitions
(about 8 rows per block), reduce task time: shuffle then project
28.7 s -> 6.2 s, shuffle then sort 9.4 s -> 8.2 s.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Comet's native and columnar shuffles write and read every leaf column of
every row as Arrow, so their cost grows with rows times leaf columns,
while Spark's shuffle moves whole rows. A shuffle whose payload outside
the partitioning key has at least
spark.comet.shuffle.wideRowFallback.minLeafColumns leaf columns (0, the
default, disables the rule) now stays a Spark shuffle.

Leaves are counted by LeafColumns, for reuse by the sort rule: a struct
has the leaves of its fields, an array those of its element, a map those
of its key and value, and any other type one.

The rule is one more reason in shuffleSupported, so the operators
reading the shuffle stay in Spark and a native producer converts its
batches to rows once before the write. It is also checked where boundary
formats ask whether a columnar shuffle is available, so no later rule
turns the shuffle back into a Comet one. It depends on the schema alone,
so AQE re-plans decide the same.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
The non-codegen columnar-to-row path built its UnsafeProjection in every
task. Janino compilation hits Spark's code cache, but generating and
formatting the source does not, and that grows with the number of
columns: in the wide shuffle read benchmark it took about 40% of reduce
task samples at 512 leaf columns.

Generated projections are now pooled per executor, keyed by their bound
references. A task takes one from the pool, or generates one, and hands
it back when it completes, so a projection is never shared between two
users at once.

Wide shuffle read benchmark, 512 leaf columns, 64 maps x 250 partitions,
reduce task time: shuffle 6.4 s -> 3.1 s, shuffle then project
6.2 s -> 3.5 s, shuffle then sort 8.2 s -> 3.5 s.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
WideRowSortFallback moved a sort to Spark for a variable-width column
outside its key, or for rows above 1 KB on average with a key under a
fifth of the row, the row size taken as the larger of the stage
statistics and the schema estimate so that AQE re-plans kept the
decision. It now counts leaf columns like the shuffle rule: a sort read
by a Spark operator moves to Spark when its input has at least
spark.comet.exec.sort.wideRowFallback.minLeafColumns (default 50) leaf
columns outside the sort key, counted by LeafColumns, with the columns
the sort key references left out as WideRowShuffleFallback leaves out
the partitioning key. The decision reads only the schema, so every plan
of a query makes the same one without statistics.

The variable-width condition, the row size condition and their
settings minAvgRowBytes, maxKeyFraction and variableWidthTypes.enabled
are removed. The rule still runs only when
spark.comet.exec.sort.wideRowFallback.enabled is set, and a sort read
by a native operator stays native.

spark.comet.shuffle.wideRowFallback.minLeafColumns now defaults to 50
instead of 0, which disabled the rule.

Tests cover flat, struct, array and map payloads at 49 and 50 leaves,
sort key columns left out of the count, Spark and native consumers, and
a wide sort with the shuffle rule at its default; the sort suite turns
the shuffle rule off elsewhere so that its sorts read Comet shuffles.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Spark cannot make a native operator release memory: the native memory
consumer's spill() returns 0. A native sort whose input stayed in
memory kept its whole reservation until its last output batch, and late
materialization, which estimates its input at 1.1 to 1.3 times instead
of 2, keeps larger sorts in memory. A Spark sort, sort-merge join or
window reading that output through a columnar-to-row transition in the
same task was then refused memory and failed with
UNABLE_TO_ACQUIRE_MEMORY.

ExternalSorter::sort() now sorts and spills the buffered batches before
producing output when the input fit in memory but its reservation
exceeds spark.comet.exec.sort.spillBeforeOutputThreshold, and produces
its output by reading the spilled run back, so it holds only the merge
buffers while its output is consumed. Batches are written as they are
sorted, as late materialization already did, instead of being kept
until the pool refuses them. Both sort paths, with and without late
materialization, are covered; TopK and fetch do not use ExternalSorter.
Neither apache/datafusion nor apache/datafusion-comet has an equivalent
setting.

The threshold is resolved on the JVM, like the other configs native
code parses: when unset it is a quarter of a task's share of
spark.memory.offHeap.size, spark.memory.offHeap.size /
(spark.executor.cores / spark.task.cpus) / 4, which is 384 MiB for 12g
and 8 cores, and 0, which disables it, in on-heap mode. The native
session passes it to the sort as a SessionConfig extension,
SpillBeforeOutputThreshold.

Native tests sort 900k rows with a 1 KB payload, and 450k rows with a
1 KB key, in input batches of 8192 and 3 rows, under a fair unified
pool with a 1.5 GiB Spark share. Without the threshold the reservation
stays at 888-910 MiB until the last output batch with no spill; with
384 MiB the sort spills once and holds nothing while producing output,
without raising its peak; with a threshold above the reservation
nothing changes. Output order and every row are checked.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…nces

The native shuffle writer's gather cost scaled with partitions times
buffered batches times leaf columns instead of with rows. Every
partition's interleave received all buffered batches, and arrow's
interleave compares each batch's data type with the first one at every
struct level, recursing into the whole subtree whenever the two types
are equal but separate instances. A scan hands the writer new type
instances with every batch, so on mongo financeOrderCosts (509 leaves,
4800 partitions, about 1000-row input batches) partition interleaving
took 32.4 of the stage's 40.9 task hours.

The writer now takes its input batches' column types from its own
schema when they are equal to it, or differ only in field names or
metadata, rebuilding struct, list, fixed-size list and map arrays around
the same buffers; a batch that does not match keeps its own types. Each
output chunk interleaves only the batches its rows come from, renumbered
densely.

wide_write_bench, one map task of 220,000 rows, zstd(1), writer time
before -> after: 508 leaves in structs at 4800 partitions 11.6 -> 6.5 s
with shared type instances, 25.7 -> 7.3 s with three, 33.5 -> 7.4 s with
one per batch; 513 flat leaves 8.4 -> 5.9 s. At 250 partitions and for
8 and 64 leaves times are unchanged.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
CostBasedEngineChoice weighed every native operator at -1, every Spark
operator at 0 and every conversion at 1, whatever the rows. It now
minimizes an estimated time in ns: each priced operator, shuffle and
conversion costs its rows times a price per row from EngineCostTable,
k0*L + k1*L*L natively and c0 + k*L in Spark, where L counts the leaf
columns of the output outside the key (the sort order of a sort, the
ordering a sort-merge join, window, window group limit or expand
requires, the partitioning of a shuffle). The coefficients depend on
the class (shuffleWrite, shuffleRead, sort for sorts and the buffering
operators, rowLocal for filters, projects, broadcast hash joins, unions,
hash aggregates, coalesces and limits, c2r for conversions) and on the
form, nested when at least half of the leaves are inside structs,
arrays or maps. A native shuffle write is scaled by
1 + 0.08 * max(0, partitions / 250 - 1); a Comet columnar shuffle over
Spark rows costs a native shuffle plus one c2r, an extrapolation.

A conversion from a native operator to a Spark one inside a stage costs
oomRiskPenalty (2000 ns per row) more when a native sort, sort-merge
join, hash join or hash aggregate is at or below it and a Spark sort,
window, sort-merge join, sort aggregate or object hash aggregate at or
above it. Below a native operator the stage is native and above a Spark
one it is Spark, so the solver sees both sides there and can move the
whole stage to one engine.

Rows are the runtime statistics of a materialized query stage, else the
row count of the logical plan, else the largest estimate among the
children, else 1, which compares the engines per row.

BoundaryFormats takes a Pricing, conversion counts by default, so that
the formats the solver priced are the ones applied; ChooseBoundaryFormats
uses the cost model when the cost-based choice is enabled.

spark.comet.exec.costBasedEngines.costTable overrides any line or scalar
with entries such as shuffleWrite.flat.comet=50.3,0.221 or
oomRiskPenalty=2000, and spark.comet.exec.costBasedEngines.log.enabled
(or spark.comet.explain.fallback.enabled) logs every decided operator,
shuffle and conversion with its class, form, leaf columns, rows and
costs. conversionWeight is removed; cometOperatorWeight,
sparkOperatorWeight and cometOperatorWeights now price only operators
outside the table, such as shuffled hash joins. WideRowSortFallback and
WideRowShuffleFallback do not run while the cost-based choice is
enabled; it stays disabled by default.

Tests cover a narrow schema staying native, a wide sort and shuffle
under a Spark consumer moving to Spark, a wide sort under a native
consumer kept or moved with its stage by the table, a sort-merge join
with one wide input, native sorts under a Spark window, sort aggregate
and sort-merge join moved by the memory risk and kept without it, the
wide-row rules not running, the table overrides and their errors, the
formulas and the form.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…ry plan

EngineCostModel no longer estimates rows: every operator, shuffle and
conversion costs its price per row from EngineCostTable, so the choice
depends only on the schema and the shape of the plan. An expand costs its
price once per projection; every other operator once. The memory risk
penalty (oomRiskPenalty and the sets of operators holding memory) is
removed.

The flat and nested lines are blended by the fraction f of leaves inside
structs, arrays or maps, (1 - f) * flat + f * nested, instead of switching
at half. A Comet line gains a constant c0, c0 + k0*L + k1*L*L: 250 ns per
row for shuffleWrite, 0 elsewhere. Shuffle widths include the partitioning
key. The Comet sort is 15*L + 0.028*L*L flat and 15*L + 0.009*L*L nested.
Window group limits and expands move to rowLocal; a project counts only
the expressions it computes, and a filter the leaves its predicate
references plus filterPassThroughPerLeaf (0.5 ns) per output leaf in both
engines. A new agg class prices hash, object hash and sort aggregates over
the grouping key leaves plus one per aggregate function, provisionally at
400 + 60*L in Spark and 240 + 36*L natively; a sort aggregate also costs a
Spark sort of the same width. shuffleReadPerByte.comet and
shuffleReadPerByte.spark (0 by default) price shuffle reads per byte of
the estimated UnsafeRow size, times cometShuffleBytesRatio (0.5) for
Comet.

Operators the rule reverts are tagged ENGINE_CHOICE_SPARK_TAG instead of
KEEP_ON_SPARK_TAG. AQE's per-stage conversion still keeps them in Spark,
but CometRule converts a whole plan (the plan without AQE, the initial
plan and every re-optimization) with CometExecRule(wholePlan = true),
which converts them again, so the rule decides each re-optimized plan from
its own shape, including the subtrees AQE carries over from the previous
plan above a materialized stage. Tags of the other rules are honored as
before.

spark.comet.exec.costBasedEngines.costTable takes c0,k0,k1 or k0,k1 for
a Comet line, the agg class and the new scalars; the log lists classes,
leaf columns and nested fraction instead of form and rows.

Tests cover prices with no rows, the blend with no step at half, projects,
filters, aggregates, expands and shuffles with their keys and c0, the
per-byte terms, the sorts formerly moved by the memory risk now moved only
by price, the whole-plan conversion of tagged operators, and an AQE query
whose sort-merge join becomes a broadcast hash join, where the final
aggregate reverted under the sort-merge join runs natively after the
re-optimization.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
spark.comet.exec.costBasedEngines.enabled now defaults to true and
spark.comet.shuffle.wideRowFallback.minLeafColumns to 0, so the cost-based
choice is the only plan rule enabled by default. The suites of the
boundary-format and wide-row rules run with it disabled, and the wide-row
shuffle suite sets its old threshold of 50 explicitly.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
EngineCostTable takes the prices measured on pr29 (calib29, calib29c,
calib29d) with L counting every leaf of the row. shuffleWrite,
shuffleRead, sort, agg and c2r get new lines, and new classes replace
the stand-ins: r2c for row-to-columnar conversions, smj and bhj instead
of sort and rowLocal, predicate, projectPassThrough and expression for
filters and projects, window, wglPartial and wglFinal, expand (free in
Comet and in Spark with codegen) and generate, and an optional sortSpill
applied to a fraction sortSpillFraction of rows, none by default.
rowLocal now prices unions, coalesces and limits at nothing.

An aggregate costs its grouping keys, half per phase of a two-phase
aggregate, plus the price of the class of each function (aggDeclarative,
aggCollectList, aggCollectSet, aggPercentile, aggPercentileApprox,
aggOther) and aggObjectHash for an object hash aggregate. A window costs
the row_number line over its input plus windowAggregate, windowOffset or
windowRank per function. In Spark, operators beyond
spark.sql.codegen.maxFields take aggDeclarativeNoCodegen,
expandNoCodegen and generateNoCodegen, and a project over a scan passes
its columns for free and computes at expressionOverScan.

The partition factors of the native write and of the Comet read grow
with the leaves (0.04 + 0.00036 * L and 0.06 + 0.00025 * L); Comet's
columnar shuffle has its own write slope (0.001 * L), a constant of
400 ns and one r2c. The filter pass-through is per engine (Comet 1.5,
Spark 0). Shuffles
and Comet sorts pay per byte of the estimated row beyond 12 per leaf. The
quadratic term of a Comet line is capped at k1 * L * min(L, 600). Arrays
still count the leaves of their element once: the plan has no average
length. The cost table override accepts the new classes and scalars,
and <class>.<engine> for both forms.

Rows are not estimated, so the model cannot see what a filter drops or an
aggregate reduces. Two flags of the table, true by default, keep such
operators native in the solver whatever their prices:
keepFiltersOverNativeScans keeps a native filter over a native scan and
the native projects over it, and keepPartialAggregatesOverNativeInputs
keeps a native partial aggregate directly over a native scan, filter or
project, so the conversion to rows sits above them.

Tests cover the new formulas and overrides, filters, projects over a scan
and after a shuffle, aggregate and window functions, expands and
generates with and without codegen, shuffle formats, partitions and
bytes, a selective filter and a partial aggregate over a wide native
scan kept native whatever the prices, and a
cube over wide rows whose reduce side runs in Spark with one conversion.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
The clippy::clone_on_ref_ptr lint denies `.clone()` on an Arc, which failed
`cargo clippy --all-targets --workspace -- -D warnings`.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
spark.comet.exec.costBasedEngines.enabled defaults to false again, so no
plan rule of this fork is enabled by default and Comet's defaults match
upstream. Workloads that want the cost-based choice enable it explicitly.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…spills

The test needs the native child sort to spill. Each of its 5 tasks sorts
4000 rows of ~264 bytes under the fair_unified pool, whose limit with
fraction 0.002 is 4294967 / 3 consumers = 1431655 bytes per consumer.
Upstream reserves twice each 1024-row batch (540680 bytes), so its third
batch exceeds the limit and the sort spills. Late materialization of wide
rows reserves the input once plus the keys and 48 bytes per row
(327684 bytes per batch, 1305360 for the task), which fits, so the sort
finishes in memory and reports no spill.

Fraction 0.001 halves the limit to 715827 bytes, below the input itself,
so the sort spills whatever the per-row overhead of its reservation. The
assertions are unchanged and pass, including the stage disk spill total.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
@msafonov msafonov changed the title Joom 1.1-1: Delta scan, correctness fixes, wide-row perf, cost-based engine choice feat: Joom 1.1-1 — Delta scan, correctness fixes, wide-row performance, cost-based engine choice Oct 3, 2026
msafonov and others added 3 commits October 3, 2026 20:18
Window nodes with an expression that cannot stream run in DataFusion's
WindowAggExec again by default, as upstream. The spilling
PartitionAggregateWindowExec is used only with
spark.comet.exec.window.partitionAggregate.enabled=true.

CometExecIterator sends the resolved flag to the native side, which marks
the DataFusion session with a PartitionAggregateWindowEnabled extension;
the planner tries PartitionAggregateWindowExec only when it is present.

CometWindowExecSuite runs on the default, and
CometPartitionAggregateWindowSuite runs the same tests with the flag on.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
CometBatchRowProjection took its projections lazily and registered the
completion listener that returns them to the pool on first use. When a
Python UDF writer thread reads the rows, that is after the PythonRunner
listener, and listeners run in reverse order, so the projection went back
to the pool before the runner joined the writer thread, and another task
could take it while the writer still wrote to its row buffer (the class of
SPARK-33277).

The listener is now registered when CometBatchRowProjection is created in
the task thread, before any reader exists, and returns whatever was taken
later. A projection taken after the task completed is not pooled.

A generated projection's row buffer grows to its largest row and never
shrinks, so a projection whose buffer exceeds 1 MiB is dropped instead of
pooled, and a schema keeps at most as many projections as there are
available processors (at most 64) instead of 64.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Add CometShuffleReadCoalesceSuite and CometShuffleExternalSorterSpillSuite
to the shuffle group and WideRowShuffleFallbackSuite next to
WideRowSortFallbackSuite in both PR build workflows. The contrib/delta-spark
suites only run with -Pdelta, so check-suites.py ignores them.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
@github-actions github-actions Bot added the enhancement New feature or request label Oct 3, 2026
msafonov and others added 3 commits October 3, 2026 20:33
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Annotate the default cost lines map so Scala 2.13 infers a Map, rename a
test helper that clashes with SQLTestUtils.withTable on Spark 4, drop a
core::f64 import that shadows the primitive's constants, and allow the
fetch_update deprecation until the MSRV reaches try_update.

Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
@msafonov
msafonov merged commit b9c24e4 into branch-1.1 Oct 3, 2026
35 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants