Repository navigation
fix: Joom 1.1-2 — production fixes for shuffle read, offsets, timestamps, native memory, task kill and task binaries - #3
Merged
Conversation
…etween blocks
The coalesced and direct shuffle reads joined the fetched blocks with a
SequenceInputStream. Its close() walks every remaining block, calling
ShuffleBlockFetcherIterator.next() for each. When the task ends early,
the fetcher's completion listener has already released the buffers in
its results queue without removing them, so next() hands out a released
buffer and closing its stream releases it again ("refCnt: 0,
decrement: 1" from the reader's completion listener). With blocks still
in flight, close() instead waits for them, and after an interrupt the
channel runs that close() in the interrupting thread, so a killed task
could not be stopped and the task reaper killed the executor.
The blocks are now read by a stream that closes each block before
fetching the next, never fetches on close, and checks for a kill at
every block boundary.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…m overflowing
Explode repeats every column it does not unnest with take, bounding a
build only by its row count, and never splits one input row. A row
whose nested strings are large, repeated once per element of a long
list, filled more than i32::MAX bytes of one offset buffer, e.g. a
proposal's status_history[].merchant_variant_prices[].variant_id
repeated for each of its target variants ("struct column 2
("merchant_variant_prices") failed: ... offset overflow").
Explode now also bounds each build by what it adds to every 32-bit
offset buffer of the repeated columns, at any depth, and unnests a row
that alone would overflow one in parts, narrowing its lists to the
elements each part outputs. The check costs one pass over the offset
nodes per input batch; rows are only measured when the longest output
length times a whole column could overflow.
The shuffle read coalescer and the window's output coalescing joined
batches by row count alone, so concatenating them could overflow the
same way. They now stop joining before any offset buffer would.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…differently Native timestamp_seconds returned Timestamp(Microsecond, None), the Arrow type of TimestampNTZ, for Spark's TimestampType, so comparing it with a scan column or a timestamp literal failed with "Invalid comparison operation: Timestamp(µs) >= Timestamp(µs, "UTC")", and session-zone dependent expressions read it as wall-clock time. Port apache#6341, which labels its results as UTC. The Etc/UTC alias fix in date_trunc only covers that one label: a native date_trunc in any other session zone, or any other producer that labels its result with a zone other than UTC, still cannot be compared with UTC-labelled timestamps. Comparisons now relabel both sides as UTC when their zone labels differ. Spark only compares timestamps of one type and a TimestampType value is a UTC instant whatever its label, so this keeps every value. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…rter's memory A spill merges the batches the sorter buffered inside a SpillWorkspace that holds the sorter's reservation, twice the batches' size. The merge's builder keeps one sorted batch per run and its row cursors encode each batch's sort key, so when the key is most of the row the merge needs a little more than twice the data. Since d7d47f3 counts the encoded rows a merge keeps, that remainder is asked of the pool, and with the task's Spark share used up the sort failed with "Additional allocation failed for ExternalSorter[0] ... Failed to acquire N bytes plus 0 bytes overcommitted", although the spill was about to release everything it holds. merchants_lifecycle (window over merchant_id, created_date, product_id) failed this way on all 10 attempts. Pristine 55.1.0 passes the same inputs because it does not count the reused rows: 0 of 284 failing vs 130 of 284 at the parent commit in a release grid. Everything the spill merge's children reserve already exists (sorted chunks, encoded rows, computed keys), so the workspace of a spill now grows its first parent with the infallible grow. In Comet the shortfall is carried as overcommit and repaid when the spill frees the workspace; the merge of spill files keeps the fallible workspace it sizes its passes with. Tests: a vendored sort of key-dominated rows in a pool it fills, and a Comet sort with a fixed Spark share and production's merge headroom ratio, which failed with the production message and now peaks within 2% of the share. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…ling A final aggregate that spills replays its runs, merged, through a fully ordered OrderedFinalAggregateStream, which cannot spill and resized its reservation with try_resize after every batch. In a Comet task the native shuffle writer reading the aggregate buffers its output until Spark refuses it, so it soon holds the share the aggregate released, and Spark cannot make it spill; the spill merge feeding the replay also seats as many runs as the pool grants. The replay's next resize was refused and the task failed with "Additional allocation failed for FinalHashAggregateStream[0] ... ShuffleRepartitioner[0] consumed 4.2 GB" (538 times on 2026-10-05, among them ads_manual_product_rank_cohorts). DataFusion 55.1.0 does the same. A fully ordered table holds only the groups the last batch may continue, all of it already allocated, so record it with the infallible resize. Comet carries the shortfall as overcommit and refuses the next fallible request, here the shuffle writer's, until it is repaid, so the writer spills and the memory comes back. A partially ordered aggregate that cannot spill because temporary files are disabled keeps the error. DataFusion's test_sort_reservation_fails_during_spill expected the replay's error in a 500-byte pool; it now checks the query completes and releases everything. Tests: a vendored final aggregate read by a stand-in for the shuffle writer in pools of 4 to 12 MiB, and a Comet final aggregate under ShuffleWriterExec with a fixed Spark share, with and without a writer buffer limit; both failed with the production message and now stay within 5% of the share. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…cost-based engine choice A native hash aggregate over grouping keys that hold arrays costs far more than the agg line predicts: star_order_2020's distinct over 18 key leaves, 7 of them in arrays, took 4.0 us per row in each native phase against 0.7 to 1.0 us in Spark, while the model priced Comet at 0.56 and kept it native. Two classes are added on top of agg: aggArrayKey per leaf of a key holding an array (half per phase), fitted on star_order_2020 and calib29's aggkeys, and codegenDispatch once per leaf of a key computed through the JVM codegen dispatcher, such as the transform normalizing doubles inside an array. Aggregates without such keys keep their prices. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…tion Spark kills a task by marking its TaskContext and, only if asked to, interrupting its thread. Neither reaches a thread inside Native.executePlan: a plan without JVM input runs on a Tokio task while the task thread parks in blocking_recv until the next batch, which a plan rooted at a native shuffle writer only sends once its whole input is written, and a JVM-fed plan is polled on the task thread, where an operator that keeps polling a child that never returns Pending (an aggregate over a join multiplying its probe rows) runs without a break. JVM code called back from native (Comet's ScalaUDF kernels, the row to Arrow conversion) checked for a kill only between batches, or not at all. With spark.task.reaper.enabled the reaper then killed the whole executor after killTimeout: 43 thrift-server executors on 2026-10-05, 21 of them with the task thread inside executePlan under a native shuffle write and 20 more in JVM code it had called back, 19 of them ScalaUDF kernels. Each native plan now has a cancellation, set by cancelPlan from a watcher that polls the task's kill flag every 50 ms. It ends the task thread's wait on the producer task and on native I/O, aborts the producer task, makes every operator of the plan end its stream with an error at its next batch (a transparent wrapper added to each native operator), and stops a shuffle writer's spill or final write at its next batch. The UDF kernels and the row to Arrow reader check for a kill every 256 rows. A cancelled plan's failure is reported as a TaskKilledException. releasePlan aborts the producer task and waits up to a second for it to drop its stream, so that its memory is returned before the task ends. getShufflePartitionOffsets now downcasts through the wrapper. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…ges that read it The native shuffle writer's spec (child plan, metric tree, input RDDs and scan planning data) lived on CometShuffleDependency.nativeShuffleSpec. The dependency is serialized with its own map stage, but also with every stage that reads the shuffle, through CometShuffledBatchRDD.dependency. The spec's input RDDs pull in their upstream shuffle dependencies and specs in turn, so each reduce stage's task binary carried the plan and lineage of every native map stage upstream of it and grew with the depth of the query. Carry the spec on the map-side CometNativeShuffleInputRDD, which only the map stage serializes, and hand it to the writer through the input iterator. The dependency keeps it only on the driver. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…plans convertBlock pools QueryContext SQL text on the root of each native block, but three paths still shipped the whole query text once per expression: - a scan's common planning data (NativeScanCommon data filters, and the Delta contrib's DeltaSparkScan) is serialized from the scan's own operator and injected under the block root on executors; - the native shuffle writer's child plan is the un-interned nativeOp; - a Delta scan sits inside a ContribScan envelope the interner cannot open. With a long query these copies dominated the task binary: a union of 115 filtered scans with a 42 KB query text weighed 92 MB, a four-stage query with 13 KB of text 32 MB. Commons are now interned against the pool of the plan they are injected into (adding missing texts to it), the shuffle writer's child plan is the block's interned plan with its pool hoisted to the ShuffleWriter root, and PlanDataInjector.internScan lets the Delta injector intern its payload in the block plan. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Serializes every stage's task binary the way DAGScheduler does for a union of many filtered scans and a multi-stage query with a long query text, with and without AQE, and bounds each against the largest vanilla Spark stage plus one copy of the query text per native block. Also checks that a reduce stage does not carry the plan of the stage it reads. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…k version Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…-bit offsets A hash aggregate interns the bytes of every distinct Utf8 or Binary group key into one buffer addressed by i32 offsets, and only empties it when the memory pool refuses the table. With a large enough pool a task with more than 2 GiB of distinct key bytes fails with "offset overflow, buffer size > 2147483647" (multi-column keys) or panics in ArrowBytesMap (one key column). Stage 308 of ads_manual_product_rank_cohorts (SELECT DISTINCT date, product_id, country over 5.8 billion feed rows hash-partitioned by date into 16 tasks, ~167 million groups each) failed this way 24 times on 2026-10-05; a retry passed only when the pool spilled first. DataFusion 55.1.0 and its main branch do the same (apache/datafusion#23694, #24704). Each group values implementation now tells whether a batch's keys still fit the 32-bit offsets of what it holds: the byte builders by their buffer length plus the batch's value bytes, the single-column map by the bytes it interned, and the row-format ones, whose decoded offsets never exceed their encoded size, by that size plus the batch's largest offset extent at any depth. When a batch would not fit, the table takes out every group so far before interning it: a partial aggregate emits them as it does under memory pressure, a final or single aggregate sorts and spills them for the merge replay, and a partial-reduce aggregate emits them. Tests: SELECT DISTINCT over 3 GiB of distinct 1 MiB keys in a partial and a final aggregate on (Int32, Utf8) and a partial on Utf8 alone failed with the production error (and the ArrowBytesMap panic) and now return every group. 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.
Fixes for the problems found when the Spark Thrift server and dbt models ran on
joom-1.1-1in production (night of 2026-10-05). Proposed version name:joom-1.1-2.Fixes
TaskCompletionListenerException: refCnt: 0, decrement: 1). The coalesced read went throughSequenceInputStream, whoseclose()fetched the remaining blocks after the fetcher's cleanup and released their buffers again. Each block is now closed before the next is requested,close()fetches nothing, and a killed task stops between blocks.offset overflow: data exceeds the capacity of the offset type). Explode bounded its output only by rows, so repeated nested strings of one row exceeded 2 GiB. Explode, the shuffle read coalescer and the window output coalescing now bound batches by offset extents and split a row that alone would overflow.Invalid comparison operation: Timestamp(µs) >= Timestamp(µs, "UTC")). Port of fix: label native timestamp_seconds results as UTC timestamps apache/datafusion-comet#6341 (timestamp_secondsreturns a UTC-labelled timestamp) plus relabelling both sides of a mixed-label timestamp comparison to UTC.Additional allocation failed for ExternalSorter). Growth past the spill workspace for memory that is already allocated is recorded instead of failing (regression from our d7d47f3).Additional allocation failed for FinalHashAggregateStream, upstream DataFusion behaviour). The aggregate records what it already holds; the native shuffle writer in the same task then gets the refusal and spills.TaskKilledException.offset overflow, buffer size > 2147483647, upstream DataFusion issue Grouping operations on large datasets can overflow i32 offsets apache/datafusion#23694). Before interning a batch whose keys would not fit, the table takes its groups out: a partial aggregate emits them, a final one spills them.No upstream default is changed; the cost rule and our other plan rules stay disabled by default.
Verification
joom-1.1-1with the production error and passes now.cargo test --workspace: 1900 passed, 0 failed; clippy-D warningsclean;cargo fmtclean.dev/ci/check-suites.pypasses.🤖 Generated with Claude Code