Repository navigation
fix: Joom 1.1-3 — COUNT(DISTINCT) under memory pressure, plan task binaries, first/last aggregates, join-condition cost - #4
Merged
Conversation
…engine choice A native sort-merge join with a join condition builds every pair of rows of equal keys, with all their output columns, before the condition drops them, while Spark tests the condition first. fbj_order_type's left band joins took 23 to 28 us per output row natively against 6 to 8 in Spark, and a local micro-benchmark of left band joins shows the same 3.6 to 4 times ratio, but the model priced only the smj line, which favours Comet, and kept them native. A class smjCondition, Line(3000, 80, 0, 850, 23), is added on top of smj over the output leaves of a sort-merge join that has a condition; joins without one keep their prices. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…jCondition A join on an SCD validity interval, L <= V < U with V from one input and L, U from the other, runs about as fast natively as in Spark: few pairs per key fail the condition. smjCondition is now skipped when the condition is exactly one such interval, recognized by expression structure only: comparisons in either order, BETWEEN, casts and date truncations, COALESCE(U, literal) and U IS NULL OR V < U. U must not be L shifted by a constant through the aliases of the plan below the join; a band such as v BETWEEN ts - 30 days AND ts, a CTE projecting l + 30 days as the upper bound, a one-sided bound, two lower bounds, bounds from different inputs or any extra conjunct keep the charge. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Every CometNativeExec held its native Operator (the protobuf of its whole subtree, with QueryContext SQL text per expression) and the block's serialized plan as ordinary fields. Any task closure that captures a plan node, such as Spark's SortExec above a Comet subtree in a dynamic-partition write, Java-serialized all of them, growing quadratically with plan depth and query text: 352.7 MiB task binaries for jms_orders against 1.15 MiB on vanilla Spark. nativeOp and serializedPlanOpt are now @transient constructor fields. They stay in the product, so makeCopy, copy, transform, withNewChildren, equality and canonicalization keep them on the driver; executors never read them, since CometExecRDD, the native shuffle spec and the native write paths carry the plan bytes they need explicitly. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Spark First/Last without ordering ran through DataFusion's GroupsAccumulatorAdapter (one row Accumulator per group). SparkFirstLast keeps the FirstValue/LastValue state layout (value, is_set) and adds a GroupsAccumulator for boolean, integer, float, decimal, date, timestamp, string and binary inputs, honouring ignore nulls, FILTER, merge, EmitTo::First and convert_to_state. Nested types keep the row accumulator. Spark Min/Max over booleans map to bool_and/bool_or, which have a groups accumulator and the same null semantics. Partial aggregate of 10M rows into 1M groups (bench first_last): first ignore nulls with filter over longs 5.83 s -> 0.27 s, last ignore nulls with filter over strings 7.05 s -> 0.64 s, first over strings 4.41 s -> 0.29 s, max over booleans 4.17 s -> 0.25 s. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…egates Compares Comet with Spark for first/last (with and without IGNORE NULLS and FILTER), any_value, min/max over booleans, bool_and/bool_or/every/some/any and the FILTER shape of Spark's COUNT(DISTINCT) rewrite, over generated data with fixed null patterns per group. Ordered single-split inputs are compared exactly; many-partition and spilling runs check that each value belongs to its group. The default part runs in CI; -Dcomet.test.aggFuzz.full=true runs the full matrix of types, shapes, groupings and a 200k-group dataset. An ignored test records a pre-existing bug: under memory pressure a single COUNT(DISTINCT) overcounts because the PartialMerge aggregate runs as a DataFusion Partial aggregate and emits groups early. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…ng a group twice
Spark plans a single COUNT(DISTINCT x) (with other aggregates) as
Partial(k, x) -> PartialMerge(k, x) -> [PartialMerge, Partial count(x)](k) -> Final.
The PartialMerge stage must hand each (k, x) to the next aggregate once:
that aggregate counts x without de-duplicating it. Comet ran PartialMerge as a
DataFusion Partial aggregate with MergeAsPartial-wrapped accumulators, and
DataFusion's PartialHashAggregateStream emits all held groups and starts over
when its reservation is refused (hash_stream.rs "Step 4: Larger-than-memory
execution (early emit)"), so under memory pressure the same (k, x) left it
several times and COUNT(DISTINCT) overcounted (3x in the new test); SUM/MAX
and other non-distinct aggregates stayed correct.
PartialMerge now runs as DataFusion PartialReduce (states in, states out).
DataFusion 55.1 routes PartialReduce to PartialReduceHashAggregateStream,
which cannot spill, unless the pool reports a finite limit, and then to the
legacy GroupedHashAggregateStream. Comet's pools report an unknown limit, so
the vendored crate now sends PartialReduce under any non-infinite pool to
FinalHashAggregateStream / OrderedFinalAggregateStream, which spill and
replay sorted runs, and emit merged states instead of final values for it.
The mixed {PartialMerge, Partial} stage keeps Partial + MergeAsPartial: a
Final always merges its output. Skip-partial stays disabled for plans with a
PartialMerge.
Tests: a native test feeds a PartialMerge three copies of 24000 (k, x) states
through a 256 KiB fair pool (23679 groups came out more than once before);
CometAggregateSuite checks a grouped and a global single COUNT(DISTINCT)
whose PartialMerge spills against Spark (25716 vs 8572 per group before);
the fuzz suite's known-bug case is a test again and the spill shapes of the
full matrix include distinct_one.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
… one input fbj_order_type's left join took its lower bound from replenishments and its upper bound from orders joined below the sort-merge join, a band join that the validity-interval shape let run natively. Trace each bound to the operator that produces it and charge smjCondition when a join or union lies between them; a LEAD-built upper bound counts as its own source. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…nar format A Comet shuffle with no consumer in the plan, such as the repartition at the root of a dbt model's write (DISTRIBUTE BY), could only keep its native format, which a Spark producer cannot feed. Every operator of its stage was pinned native, so fbj_order_type's band sort-merge join stayed native whatever its smjCondition price. Such a shuffle may now also take the columnar format: its output is Arrow either way. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…ch Spark 4.1 renamed Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
The late-materialized sort interleaved every output batch from all the batches it buffered, ~1600 batches of 8192 rows at 13M rows per task, column by column. It now concatenates each payload column of the buffered batches once, releasing the source column as soon as it is copied, and takes every output batch from it. The memory held grows by at most one column while it is copied. A column is concatenated whole when its reservation grows by the column's size, in up to 16 chunks otherwise, and is left in its batches when not even a chunk fits; an output batch then interleaves the few pieces its rows come from. A spill concatenates in chunks of at most 1/16 of what it buffered, since its reservation cannot refuse. Columns whose concatenation or take is unsupported, or whose 32-bit list offsets would overflow, stay in their batches. An order that reads at most 16 batches per output batch on average (input already nearly sorted) keeps interleaving and releases each batch once its last row is output. SortSpillBenchSuite, 22 leaf columns (string, date, long, double and boolean, mostly null), 4 keys, one task, in memory, end to end ns/row, alternating runs, before -> after: 6M rows 505/465 -> 431/406; 13M rows 1166/556 -> 484/498 (the old gather ran at either speed per JVM). Native sort time at 13M: 12.3/6.9 -> 5.2/5.6 s. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
The keys of a late-materialized sort over several columns were sorted as (row bytes, index) pairs, every comparison a memcmp through a pointer into the rows. They are now sorted as (first 8 bytes big-endian, index) pairs, and only runs of equal prefixes are sorted again by the whole row. The row format already encodes descending order, null order and every type as bytes, and a row shorter than 8 bytes is zero-padded, which can only tie with the rows it is a prefix of, so the order is the same. Keys of the SortSpillBenchSuite rows (a 24-character hex id, two dates, an int; every row ties on its prefix with ~17 others) alone, ns/row, before -> after: 1M 97 -> 81, 6M 132 -> 123, 13M 160 -> 123. The whole sort, end to end, alternating runs: 6M 428/395 -> 380/386, 13M 539/513 -> 518/459. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…ed tasks The sort lines came from calib29's small tasks, where Comet's sort cost 224 + 0.023 * L^2 against Spark's flat 646 ns per row. A local single-core benchmark (sortWithinPartitions minus the same scan and conversion, so Spark's inserts and copies count, not only its sort time) after the gather and key-sort changes measured 3 to 30 flat leaves and 8 and 22 leaves in structs and arrays at 1M, 6M and 13M rows per task. Both engines cost more per row the more rows a task sorts, Comet more: Comet/Spark is 0.4 at 1M rows, 0.75 flat and about 1 nested at 13M, where Comet's per-leaf gather from separate columns outgrows the caches while Spark copies whole rows. The lines are fitted at 13M rows and scaled by 1.6, the ratio of the pr29 prices of small tasks to the benchmark at 1M rows: sort flat Comet 321 + 21.2 L (was 224 + 0.023 L^2), Spark 432 + 24.1 L (646) sort nested Comet 1110 + 3.3 L (244 + 2.69 L + 0.016 L^2), Spark 621 + 25.4 L (770 + 2.34 L) sortSpill flat Comet 22.2 L (265 + 49 L), Spark 93 + 39.3 L (495 + 38.5 L) sortSpill nested Comet 34.6 L (324 + 45.3 L), Spark 66.7 L (502 + 33.5 L) sortSpillFraction stays 0. A narrow native sort now gains about 120 ns per row over Spark instead of 420, less than a columnar shuffle written for it by a Spark producer, so the reduce-side sort of a Spark sort aggregate moves to Spark with its shuffle; the suite asserts that, and that both sorts stay native when the native sort is cheap. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
The local fit on 8 and 22 nested leaves extrapolated to 509 nested leaves sent mongo finance's sort before its window group limit native, where it cost 13.7 us per row and spilled 955 GB against 0.28 us and no spill in Spark (the whole job 39.96 task-hours against 35.80). The nested sort and sortSpill lines go back to calib29, measured on the cluster at 8 to 512 leaves, which keep that sort in Spark; the flat lines stay recalibrated. 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-2in production (night of 2026-10-06), plus a silent data-corruption bug found by a new differential fuzz suite. Proposed version name:joom-1.1-3.Fixes
Partial, which emits its groups early when a memory reservation is refused; the next stage counted the duplicated distinct values. PartialMerge now runs asPartialReduce, routed to the spilling final streams that emit state. Affects every build with native PartialMerge (since da187f2).INSERT OVERWRITE ... PARTITION) Java-serialized them all: 352.7 MiB against 1.15 MiB on vanilla Spark, executors ran out of heap. They are now@transientconstructor fields: kept on the driver bymakeCopy/transform, never serialized.first/lastwithout ordering andmin/maxover booleans ran through the row-accumulator adapter (~3× the task time of vanilla Spark on gold profiles). A native groups accumulator forfirst/last(ignore nulls, FILTER, partial/final merge, partial emission, spill) for primitive, decimal, date/timestamp, string and binary types; booleanmin/maxmap tobool_and/bool_or.smjCondition), except a validity interval whose bounds read one row of one input (SCD dimensions,LEAD-built intervals). Bounds from two joined inputs (overlapping intervals per key) are charged. A Spark producer may now feed a root Comet shuffle with no consumer through the columnar format, so such a join can move to Spark.first/last/boolean aggregates (default part in CI, full matrix behind-Dcomet.test.aggFuzz.full=true); regression tests for each fix that fail onjoom-1.1-2.No upstream default is changed; the cost rule and our other plan rules stay disabled by default.
Verification
cargo test --workspace: 1910 passed, 0 failed; clippy-D warningsandcargo fmtclean.sparkmodule: 3821 passed, 0 failed (3 S3/Iceberg suites need Docker, run by CI). Delta contrib: 265 passed. Spark 4.0 / Scala 2.13 test-compile passes.dev/ci/check-suites.pypasses.MAX_BYcolumns (both values are valid candidates).adtech_revenue_totaldiffers only by rounding (max 0.0001, within tolerance).🤖 Generated with Claude Code