Skip to content

fix: Joom 1.1-3 — COUNT(DISTINCT) under memory pressure, plan task binaries, first/last aggregates, join-condition cost - #4

Merged
msafonov merged 13 commits into
branch-1.1from
joom/1.1-fix3
Oct 7, 2026
Merged

msafonov merged 13 commits into
branch-1.1from
joom/1.1-fix3

Conversation

@msafonov

@msafonov msafonov commented Oct 7, 2026 •

Copy link
Copy Markdown
Collaborator

Fixes for the problems found when the Spark Thrift server and dbt models ran on joom-1.1-2 in 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

  • COUNT(DISTINCT) overcounts under memory pressure (silent wrong results). A single distinct aggregate together with other aggregates is planned without Expand: Partial → PartialMerge → [PartialMerge, Partial] → Final. Comet ran the PartialMerge aggregate as a DataFusion Partial, which emits its groups early when a memory reservation is refused; the next stage counted the duplicated distinct values. PartialMerge now runs as PartialReduce, routed to the spilling final streams that emit state. Affects every build with native PartialMerge (since da187f2).
  • Task binaries of hundreds of MiB in dynamic-partition writes. Every native operator kept its protobuf (with the query text) and the serialized block plan as ordinary fields, so any closure that captured a plan node (Spark's SortExec above a Comet subtree in 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 @transient constructor fields: kept on the driver by makeCopy/transform, never serialized.
  • Spark first/last without ordering and min/max over booleans ran through the row-accumulator adapter (~3× the task time of vanilla Spark on gold profiles). A native groups accumulator for first/last (ignore nulls, FILTER, partial/final merge, partial emission, spill) for primitive, decimal, date/timestamp, string and binary types; boolean min/max map to bool_and/bool_or.
  • Cost-based engine choice and join conditions. A sort-merge join with a non-equi condition is priced (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.
  • Tests: a differential fuzz suite for first/last/boolean aggregates (default part in CI, full matrix behind -Dcomet.test.aggFuzz.full=true); regression tests for each fix that fail on joom-1.1-2.

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

Verification

  • Rust cargo test --workspace: 1910 passed, 0 failed; clippy -D warnings and cargo fmt clean.
  • JVM (Spark 3.5, Scala 2.12), whole spark module: 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.py passes.
  • Fuzz suite full matrix: 31 passed.
  • Production models rerun into junk tables on the new build and compared with production:
    • jms_orders: write-stage task binary 1.27 MiB (was 352.7 MiB), 0 failed tasks; 877,274 rows, all deterministic columns equal, remaining differences only in tie-dependent MAX_BY columns (both values are valid candidates).
    • gold real_user_profiles: PASS.
    • gold device_profiles (final build): PASS — 2,376,885,160 rows, key sets and all exact/boolean/tie columns equal in all 8192 hash buckets; adtech_revenue_total differs only by rounding (max 0.0001, within tolerance).
    • fbj_order_type: the band join runs in Spark; 66,774,282 rows, keys and deterministic columns equal; remaining differences only in ROW_NUMBER tie-dependent columns, all of them valid tie candidates (PASS_WITH_TIES).

🤖 Generated with Claude Code

msafonov and others added 8 commits October 6, 2026 12:48
…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>
msafonov and others added 5 commits October 7, 2026 10:30
…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>
@msafonov
msafonov merged commit 66b2d61 into branch-1.1 Oct 7, 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.

1 participant