Repository navigation
fix: Joom 1.1-4 — sort-merge join with a join filter: filter before materializing, FULL OUTER fix, cost rule - #5
Merged
Merged
Conversation
A sort-merge join with a join filter materialized every column of every candidate pair (take over the streamed columns, interleave over the buffered ones), picked the filter columns out of that batch, and for LEFT/RIGHT/FULL joins pushed the whole wide batch, failing pairs included, through the deferred-filtering pipeline (concat, corrected mask, filter). A validity-interval join over large key groups, where almost every pair fails, spent nearly all its time copying rows it then dropped. freeze_streamed_matched now gathers only the filter columns, evaluates the filter, updates the FULL join's buffered filter state from the full mask, and materializes the output columns only for the pairs that can reach the output, reusing the gathered filter columns when every pair is kept. INNER joins keep the passing pairs. Deferred joins keep, within each run of one streamed row's pairs in a freeze, the passing pairs, or the run's last pair when none passed: across the freezes a row spans this keeps all its passing pairs and its overall last pair, which is all get_corrected_filter_mask outputs from it, so the output and its order are unchanged. Semi/anti/mark joins already evaluate the filter on the filter columns alone. Validity-interval bench (100 streamed columns, 11000 buffered rows per key, one passing pair per streamed row, 4.4M pairs): LEFT 471 -> 8.5, FULL 378 -> 12.6, INNER 119 -> 8.0 ns per pair. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Compares Comet's sort-merge join with a join filter against Spark for inner, left, right, full, left semi and left anti joins over generated data with key groups of one row, about a batch and several batches on either side, NULL keys and NULL filter columns, and wide rows with strings, decimals and nested columns. Filters cover a validity interval bounded by one side, a band bounded by the other, one-sided conditions, always true and always false, and almost nothing or almost everything passing, at batch sizes 7, 64 and 1024 and under a memory pool small enough to make the join spill. Results are compared with Spark as multisets, and every case asserts that the join runs as CometSortMergeJoin. EXISTS and IN in a disjunction check the existence join that stays in Spark. The default part runs in CI; -Dcomet.test.smjFuzz.full=true runs the full matrix of shapes, side orders, two join keys, a second seed and groups of thousands of rows. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
get_corrected_filter_mask groups the deferred-filter entries of one streamed row by (batch id, row index) and treated an entry without a row index (a null-joined row) as the end of the row before it. A FULL join stages the null-joined rows of an unmatched buffered key group in freeze_buffered, ahead of the pairs freeze_streamed materializes in the same freeze. When a streamed row's pairs span two freezes and such a group was passed in between, its null-joined rows land between the row's two runs: the first run ended the row, so a row whose first run had no passing pair was emitted null-joined there and again (or with its passing pairs) after the second run. Before eedede1 every freeze pushed batch_size entries, so the deferred output was flushed at the next loop iteration, before the streamed row's remaining pairs could be separated. Keeping only the candidate pairs pushes far fewer entries, the flush comes later and the latent split shows: extra null-joined rows in FULL joins over key groups larger than the batch on both sides. Entries without a row index now neither end a row's group nor reset it, which is all they can be: buffered rows null-joined on the streamed side, or streamed rows that matched no key and form no group. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…ered row A sort-merge join filter like Spark's `o.t > cast(r.eff as timestamp) AND o.t <= cast(r.next_eff as timestamp)` re-evaluated the casts for every candidate pair, and Comet's date to timestamp cast resolves the session time zone per value: over a large key group almost all of the join's time went into casting the same buffered rows once per streamed row. The largest non-volatile subexpressions of the filter that read buffered columns alone are now lifted out (HoistedJoinFilter). A freeze whose pairs all fall in the current key group of their buffered batches evaluates them once over that group's rows of each batch, caches the results on the batch (accounted in its memory reservation) and gathers them per pair; other freezes, e.g. ones holding pairs of an earlier key group or null-joined pairs, evaluate the original filter as before. Also, for large key groups: - a streamed row is paired with a run of buffered rows at once instead of pair by pair; - buffered columns of pairs that form runs of consecutive rows across several batches are copied run by run instead of interleaved row by row; - the FULL join filter state is updated without a per-pair branch. New bench spark-expr/benches/smj_interval_filter.rs: string key, 4 keys x 100 streamed rows (98 payload columns) x 11000 buffered rows with date bounds and the Spark cast, one passing pair per streamed row, 4.4M pairs, batch_size 8192. ns per pair, eedede1 -> this commit: input batches 8192: LEFT 59.6 -> 5.2, INNER 59.6 -> 4.2, FULL 64.5 -> 8.0 input batches 1024: LEFT 61.2 -> 5.8, INNER 60.2 -> 4.6, FULL 64.2 -> 9.6 Unique keys change by +1..4%, the no-filter 1:1 shapes of sort_merge_join.rs by about +3%, inner_1to10 -20%. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
With the join filter tested on each pair of equal keys before the output is built, a native sort-merge join with a condition costs about what Spark does for any condition shape and any number of pairs per key: on the cluster, a left join under a validity interval or a band cost 5.7 to 22 ns per pair natively against 9.6 to 15 in Spark at 15 to 519 output leaves, an inner one 7.4 against 19.3, and one with 100 pairs per key about Spark's price per output row. The smjCondition class, its Line(3000, 80, 0, 850, 23) and the validity-interval exemption (JoinConditionShape) are removed, so a sort-merge join is priced by smj alone with or without a condition. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
New bench spark-expr/benches/smj_streamed_filter.rs, shaped like the fbj_order_type join `cast(o.completed as timestamp) < b.t AND o.t > b.t`: string key, 4 keys x 1000 (200) streamed rows with 16 payload columns x 2000 buffered rows of 3 columns, batch_size 8192. `completed` is a `yyyy-MM-dd HH:mm:ss` string (as in the model) or a date. ns per candidate pair at f5eab19: few_string (1% pass, 8M pairs): LEFT 167, FULL 168, INNER 167 few_date (1% pass, 8M pairs): LEFT 32, FULL 32, INNER 31 all_string (all pass, 1.6M): LEFT 211, FULL 211, INNER 192 Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…amed row The fbj_order_type join `cast(l.completed_dt as timestamp) < r.odt AND l.odt > r.odt` spent about 21 hours in one sort-merge join: completed_dt is a string of the streamed side, so HoistedJoinFilter, which lifted only buffered-only subexpressions, left the string to timestamp cast to run once per candidate pair, over key groups of millions of pairs. HoistedJoinFilter now lifts the largest non-volatile subexpressions that read the columns of either input alone. Streamed ones are evaluated per freeze once per run of a streamed row's pairs (over that row only, so no row the join never pairs is evaluated) and spread to the pairs with a take; when the runs average under two pairs they are evaluated per pair. A filter with streamed-only subexpressions and no buffered-only ones always takes the hoisted path; one with both still falls back to the original filter when the buffered results cannot be cached. Deferred filtering also skips the corrected mask when no pair of the accumulated batch failed the filter: every row is kept as is. smj_streamed_filter bench, ns per candidate pair, f5eab19 -> this: few_string (1% pass, 8M pairs): LEFT 165 -> 5.2, FULL 164 -> 5.5, INNER 164 -> 4.2 few_date (1% pass, 8M pairs): LEFT 34 -> 4.9, FULL 35 -> 5.3, INNER 33 -> 4.2 all_string (all pass, 1.6M): LEFT 206 -> 46, FULL 206 -> 48, INNER 182 -> 29 smj_interval_filter (buffered-side casts) unchanged within noise: group11k_b8192 LEFT 5.0 -> 5.2, FULL 7.9 -> 8.0, INNER 4.0 -> 4.0; group11k_b1024 LEFT 5.5 -> 5.5, FULL 9.4 -> 9.3, INNER 4.5 -> 4.5. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…ressions The validity-interval Rust tests run each case with four equivalent forms of the filter: plain, with buffered-only subexpressions, with streamed-only ones (`t + 0 > eff AND t - 0 <= next_eff AND t + 0 >= t - 0`) and with both, over NULL lookup times, outer and FULL joins, spilling, and key groups across batches on both sides. The hoisting test checks what each form lifts and how the residual filter reads it. CometSmjJoinFilterFuzzSuite gets filters casting a streamed column (string to bigint, and the fbj_order_type shape: a timestamp rendered as a string cast back to timestamp, compared with a buffered value) and casting both sides; they run in the default matrix, at batch 7 and in the large-group sets. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…tors around lifted subexpressions Adds string, int, boolean and date columns to both sides (generated from separate random streams, so existing data is unchanged) with valid, NULL and invalid values, and 26 filter shapes that wrap single-side subexpressions in OR, CASE, IF, NOT, IS [NOT] NULL, IN / NOT IN, COALESCE and null-safe equality. A few run by default; the rest run in the full matrix over every join type, side order, batch size and spill mode. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…th inputs On the cluster (fix4d), left joins with about one passing pair per streamed row under `CASE WHEN l.c = r.c THEN l.t - r.eff ELSE 10 - (r.nxt - l.t) END BETWEEN 1 AND 10` cost 28 to 51 ns per pair natively against 8.7 to 17.6 in Spark, 1.7 to 4.5 times. Cross-input datediff costs about what Spark does, cross-input instr and an OR over a cast of one input 4 to 10 times less, and conditions over one input at most Spark's price. A sort-merge join whose condition holds a CaseWhen or an If referencing columns of both inputs adds smjCrossCondition over every output leaf, at the former smjCondition line Line(3000, 80, 0, 850, 23); any other condition adds nothing. The rule stays disabled by default. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Spark 4 enables ANSI by default, where casting the suite's malformed strings throws in Spark itself instead of returning NULL. 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.
A sort-merge join with a join filter ran several times slower than Spark when a key has many matches: the jms_orders currency joins (LEFT OUTER,
o.t > r.eff AND o.t <= r.next_eff, ~11k rates per currency, ~100 order columns) took 503–968 s per task against 86–317 s on vanilla Spark. Proposed version name:joom-1.1-4.Changes
SortMergeJoinExecbuilt every output column for every candidate pair, then evaluated the filter and dropped almost all of them. The filter is now evaluated on its own columns only, and full rows are materialized only for the pairs that can reach the output (passing pairs, plus the last pair of a streamed row with none, which becomes its null-joined row).cast(effective_date as timestamp)), streamed-side ones once per streamed row (fbj_order_type: a string-to-timestamp cast, which cost ~165 ns per pair).smjConditioncost class and the validity-interval exemption are removed: after the fix a join condition adds nothing to the cost of a sort-merge join (calibration below). The only exception is a condition with a CASE or IF over both inputs (newsmjCrossCondition), which DataFusion evaluates by filtering and merging both branches and stays 1.7–4.5× slower than Spark per pair. The rule stays disabled by default.CometSmjJoinFilterFuzzSuite, a differential Comet-vs-Spark suite for all join types, condition shapes (incl. OR, CASE/IF, NOT, IN, coalesce and casts of invalid strings around lifted subexpressions), group sizes across batch boundaries, NULLs, wide nested rows, small batches and spill (default part in CI, full matrix behind-Dcomet.test.smjFuzz.full=true); Rust regression and property tests; a realistic SMJ benchmark.No upstream default is changed.
ANSI caveat. Lifted subexpressions are evaluated for every row of their input, so under
spark.sql.ansi.enabled=truea fallible one (e.g. a string cast) under OR or in a CASE branch can fail on a row Spark would have skipped: the query fails, results are never wrong. ANSI is not enabled in any of our Spark configs, thrift servers or dbt models; the fuzz suite runs with ANSI off on every Spark version.Measurements
Cluster calibration, ns per candidate pair, LEFT OUTER, 10k matches per key (vanilla Spark / joom-1.1-3 / this PR):
Validity-interval and band conditions cost the same within 5–10% on every engine. Conditions over both inputs:
datediffat parity,instrand OR with a one-side cast 4–10× faster on Comet, CASE 1.7–4.5× slower (charged by the rule).Production models rerun on the final native code and compared with production: fbj_order_type PASS_WITH_TIES, 2.75 task-hours against 7.48 on vanilla Spark on the same inputs; jms_orders all columns equal, 2.54 against 2.96 task-hours; an Iron Head query whose join moves to Comet: identical result, 1.93 against 2.61 task-hours.
🤖 Generated with Claude Code