Skip to content

fix: Joom 1.1-2 — production fixes for shuffle read, offsets, timestamps, native memory, task kill and task binaries - #3

Merged
msafonov merged 14 commits into
branch-1.1from
joom/1.1-fix
Oct 5, 2026
Merged

msafonov merged 14 commits into
branch-1.1from
joom/1.1-fix

Conversation

@msafonov

@msafonov msafonov commented Oct 5, 2026 •

Copy link
Copy Markdown
Collaborator

Fixes for the problems found when the Spark Thrift server and dbt models ran on joom-1.1-1 in production (night of 2026-10-05). Proposed version name: joom-1.1-2.

Fixes

  • Shuffle read releases a fetched buffer twice (TaskCompletionListenerException: refCnt: 0, decrement: 1). The coalesced read went through SequenceInputStream, whose close() 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.
  • Arrow 32-bit offset overflow in explode and batch coalescing (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.
  • Comparison of timestamps labelled with different time zones (Invalid comparison operation: Timestamp(µs) >= Timestamp(µs, "UTC")). Port of fix: label native timestamp_seconds results as UTC timestamps apache/datafusion-comet#6341 (timestamp_seconds returns a UTC-labelled timestamp) plus relabelling both sides of a mixed-label timestamp comparison to UTC.
  • Native sort fails during its own spill (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).
  • Fully ordered final aggregate fails after a spill (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.
  • Native execution ignores task kill, so the Spark task reaper killed whole executors. Kill now reaches the native plan (cancellation token, operator wrapper, JNI callbacks and UDF kernels check it) and the task ends with TaskKilledException.
  • Task binaries 10–100× larger than vanilla. The query text is interned for scan planning data and native shuffle plans, and a reading stage no longer carries the map side's plan.
  • Cost-based engine choice prices array grouping keys and grouping keys evaluated through the JVM codegen dispatcher, so distinct-like aggregates over nested keys go to Spark (star_order_2020). No change on the eight benchmark jobs.
  • Hash aggregate string group keys overflow 32-bit offsets (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.
  • CI: the new suites are listed in both workflows; two new suites made independent of the Spark version.

No upstream default is changed; the cost rule and our other plan rules stay disabled by default.

Verification

  • Every fix has a test that fails on joom-1.1-1 with the production error and passes now.
  • Rust cargo test --workspace: 1900 passed, 0 failed; clippy -D warnings clean; cargo fmt clean.
  • JVM (Spark 3.5, Scala 2.12), 31 suites: 1649 passed, 0 failed, 14 canceled (version-gated).
  • Delta contrib: 265 passed, 0 failed.
  • dev/ci/check-suites.py passes.
  • Production models on the fixed build (images fix1/fix2, junk tables): js2_proposals, js2_cluster_variant_weekly_price, moengage_mail_cube, merchants_lifecycle, ads_manual_product_rank_cohorts (concurrently, thrift-like executors, task reaper on) and star_order_2020 — all succeed with 0 failed tasks and none of the night's errors; data matches production except known non-deterministic tie columns.

🤖 Generated with Claude Code

msafonov and others added 11 commits October 5, 2026 08:46
…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>
msafonov and others added 3 commits October 5, 2026 16:11
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>
@msafonov
msafonov merged commit 79586f4 into branch-1.1 Oct 5, 2026
35 checks passed
@msafonov
msafonov deleted the joom/1.1-fix branch October 6, 2026 10:55
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.

1 participant