feat: Joom 1.1-1 — Delta scan, correctness fixes, wide-row performance, cost-based engine choice - #2
Merged
Merged
Conversation
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>
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>
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>
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>
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.
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
spark.comet.exec.costBasedEngines.enabled), calibrated on synthetic and production workloads, plus two structural special cases (filters over native scans, partial aggregates over native inputs). Also the earlier wide-row sort/shuffle fallback and boundary-format rules.datafusion-physical-plan55.1 vendored undernative/vendor(needed for the sort/SMJ changes; most of the line count).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, upstreamWindowAggExecis used). Workloads enable them explicitly.Active on default settings (behavior differs from upstream):
spark.comet.exec.sort.spillBeforeOutputThreshold, default off-heap / task slots / 4;0disables) — extra disk IO for big sorts;spark.comet.shuffle.read.coalesce.enabled, defaulttrue);try_resized (85bef53).Verification
cargo test --workspace1880 passed, 0 failed;-D warningsclean;sorts::/joins::1151 passed;AWS_REGIONset: they need a MinIO image and a region in the environment, are not run by CI, andminio/miniois gone from Docker Hub;🤖 Generated with Claude Code