Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

Combine Join Strategies

Combine is the N-ary record-combining operator: every input is declared up front, the where: predicate matches records across inputs, and the cxl: body shapes the output row. This page covers the parts an engine engineer reaches for when reasoning about how a Combine executes — the strategy selection the planner performs from the predicate shape, the heap cost of materializing build sides, how reconciled document boundaries flow through every join path, the runtime mechanics of correlation-key propagation across all four execution paths, and the join-planner statistics catalog that drives build-side selection and grace-hash partitioning.

User-facing view: the User Guide’s “Combine Nodes” page.

Predicate classification

The where: expression is a CXL boolean evaluated for every candidate record pair across inputs. The planner splits a compound and-predicate into three conjunct classes, and the classification is what selects the execution strategy:

  • Equi conjunct — a cross-input equality (a.x == b.y). Drives the hash lookup or the sort-merge join.
  • Range conjunct — a cross-input ordered comparison (a.start <= b.ts and b.ts <= a.end). Handled by the block-band IEJoin path, with or without an equi conjunct on the same input pair (IEJoin for pure range, HashPartitionIEJoin when an equality is present); a pure-range predicate (no equality) with a single range conjunct over ascending-presorted inputs can take SortMerge instead.
  • Residual conjunct — any other CXL predicate (intra-input filter, function call, and so on). Applied as a post-filter after the equi/range match succeeds.

At least one cross-input equality is required for every Combine, except for pure-range predicates, which IEJoin handles without an equi conjunct.

Match selection versus the projection body

Every strategy separates two steps: match selection (the where: predicate, including any residual re-check) chooses which build records pair with a driver, and the cxl: body projects each selected pair into an output row. Under match: first this distinction is load-bearing and is contractually uniform across the four physical join paths — in-memory hash build-probe, grace-hash, sort-merge, and the block-band IEJoin. The block-band path serves both IEJoin plan variants: pure-range (IEJoin) and equi+range (HashPartitionIEJoin), where equality is an added block-pair prune axis on the same machinery, so both are bounded and behave identically at the emit boundary.

  • One candidate order: build arrival. A driver’s candidates are its matching build rows in ascending BuildSeq, the position the build input delivered each row at. first takes the lowest, all emits in that order within a driver, and collect builds its array in that order, on every strategy:

    • the hash table’s collision chains append, so a probe walks a key’s build rows oldest first;
    • grace-hash mints BuildSeq where it receives its build input and carries it beside each row through partitions, spill files, recursive repartition and fallback chunks. A probe checks that its candidates rise in BuildSeq and fails as an internal error otherwise, so a table walked out of order cannot silently pick another row;
    • sort-merge replays its window ascending over inputs declared ascending on the range key, so window order is arrival order;
    • the block-band path (both IEJoin variants) keeps the minimum build input index per driver and sorts its output by (driver order, driver index, build index).

    The pick is therefore a function of the data and the build input’s order, not of the strategy, the predicate’s shape or the memory-derived block, spill or hash layout. Grace-hash’s block-nested-loop fallback, which a reloaded partition takes when it is still over budget and either the assigner is at its MAX_HASH_BITS (12-bit) cap or a split leaves the larger child with more than 80% of the parent’s rows (SKEW_REDUCTION_THRESHOLD), still decides first, collect and on_miss once per build chunk instead of once per driver, and its output order depends on which partitions spilled; both are known defects tracked separately.

  • One verdict per driver. Each candidate’s residual has one of three outcomes (PredicateOutcome in pipeline/combine_verdict.rs): true, not true (false or null), or failed. Every strategy feeds a driver’s outcomes to one DriverScan, which folds them into the driver’s DriverVerdict under its match mode; no strategy decides “did this driver match” from its own container. first is decided by the earliest candidate that is not “not true”: a true one is selected, a failed one is the driver’s only result, and a later candidate is outside the result. all acts on every true and every failed pair as it is found. collect writes no row once a candidate failed. Only a driver with no true and no failed candidate is a miss, and every on_miss path takes the MissToken only DriverScan::finish (or a driver whose keys admit no candidate) can make, so on_miss after a failure does not compile.

    • Hash, grace-hash and sort-merge walk candidates in arrival order and stop a first driver once its scan settles.
    • The block-band path visits candidates in block order, so each driver’s scan keeps the earliest candidate by build input index, true or failed, and the driver block’s finalize reads its verdict. A held candidate is charged in the match state’s held bytes; the fixed per-driver scan slot is charged with the driver block.
    • Under fail_fast a failure outside the result (after a first driver’s deciding candidate) never aborts the run.
  • The body runs once, as a post-match projection. A body that skips the chosen build (a filter that fails, an EvalResult::Skip) or defers on a recoverable error drops only that output row. The engine does not retry a later matching build, and it does not route the driver to on_miss: the driver matched the predicate, so it is not a zero-match driver. The verdict reads predicate outcomes only, so a body result never reaches it, under any match mode.

  • A recoverable failure belongs to one pair. On every strategy a failing residual or body on one (driver, build) pair that is part of the driver’s result writes one combine_output_row failure with that pair’s build row. A failing probe key or on_miss: null_fields body is one failure with no build row. Each pair is evaluated under the driver’s row, so a failure’s detail names the same row on every strategy.

    • The materialized hash loop writes each driver’s failures before it reads the next driver, so it holds none. The streaming probe cannot: its thread has no way to reach the dead-letter output, so it holds every failure until the thread joins and the dispatcher replays them. Under match: all that set grows with a hot key’s fan-out, so each held failure (the driver row, the build row and the failure itself, priced by held_record_bytes) is charged to the Combine’s consumer as it is appended and discharged as it is written, and the thread polls the arbitrator after each driver once 10,000 rows and failures have accumulated. The overshoot past the limit is one 10,000-entry window plus the failures of the driver in flight, never the whole stream’s.
  • The residual covers every conjunct the kernel does not. The block-band kernel verifies the first two range conjuncts itself and evaluates the residual when the predicate has a third range or any conjunct that is not a range (DecomposedPredicate::residual_exceeds_ranges), so no part of where: goes unapplied.

Strategy hint

The strategy config field carries a hint; the planner has final say.

ValueBehavior
auto (default)Planner picks a strategy from the predicate shape. Pure equality: HashBuildProbe, or GraceHash when the build side’s estimated size approaches the memory limit. Equality plus range: HashPartitionIEJoin. Pure range: SortMerge when there is exactly one range conjunct on plain fields of an integer, float, date or datetime axis (not decimal, not mixed integer/float) and both inputs arrive globally sorted ascending on those fields (a declared sort_order whose first key is that field, on an input that is not partitioned per file); otherwise IEJoin. Both IEJoin variants run on the same block-band executor.
grace_hashForce grace hash join (disk-spilling partitioned hash). Applies only to pure-equi predicates; ignored on predicates carrying range conjuncts.

grace_hash is the right hint when build-side inputs may be larger than the memory budget but fit on disk after partitioning. It is a behavioral switch, not only an assertion: the hash-versus-grace choice is made once, at plan time (grace_hash_should_fire in clinker-plan/src/plan/combine.rs), and there is no runtime fallback. An in-memory HashBuildProbe build that outgrows the budget returns CombineError::MemoryLimitExceeded, which the dispatcher maps straight to PipelineError::MemoryBudgetExceeded (E310) where it calls CombineHashTable::build in executor/combine_dispatch.rs; nothing converts the join to grace hash at that point (#1337). strategy: grace_hash pins the spilling strategy regardless of the plan-time size estimate.

The choice of in-memory hash versus grace-hash for a pure-equality Combine is driven by the build-side row-count estimate (see Join-planner statistics below): a build side large enough to risk overrunning the memory limit is what tips the planner from the in-memory hash strategy to the disk-spilling grace-hash strategy.

Memory considerations

Build-side inputs are materialized in memory as hash tables keyed by the equi columns. For each non-driving input, plan for roughly 1.5–2× the raw CSV size in heap. A 50 MB product catalog typically occupies 75–100 MB of hash-table memory — the multiplier covers the per-key Value boxing, the bucket array overhead, and the per-entry chaining structure on top of the raw payload bytes.

This heap cost is the quantity the memory arbitrator charges against pipeline.memory.limit. For an in-memory HashBuildProbe Combine, exceeding the limit during build is an E310 abort; it is not flipped to grace-hash spill at runtime (#1337). See Memory Arbitration & Scheduling for the spill thresholds, the back-pressure knob, and strategy overrides.

Block-band IEJoin: bounded on both input axes and the output

Interactive companion: the range-join explainer runs this path on small inputs: the sort and slice into min/max-tagged blocks, the block-pair pruning, the IEJoin kernel step by step, the nested-loop fallback, and match:/on_miss:, checked against a brute-force join.

The block-band path that serves every range and equi+range Combine is bounded on both input axes and the output axis, so a range join whose inputs — or whose result — exceed the memory budget spills and completes rather than failing:

  • Input axes. Each side is drained into a payload-ordered, spillable sort buffer, then sliced into contiguous, key-sorted, min/max-tagged blocks. Blocks stay resident under a shared budget and spill to their own files past it; the scheduler joins one block-pair at a time, so at most a bounded working set (the resident blocks, one loaded spilled block per side, and the kernel’s O(n) sort arrays) is live. Equi+range adds an equality-hash prune axis on the same block machinery. Both deferred on_miss piles (NULL/non-orderable-keyed drivers and in-block zero-match drivers) drain to their own spillable, driver-index-ordered buffers rather than resident vectors.
  • Numeric range axis. Each range key reduces to a single order-preserving i128 so the IEJoin/block-band kernels compare, sort, and bounds-prune one integer per axis. The reduction is chosen per axis at plan time from the CONCRETE operand types: integers widen directly, finite floats take the orderable-float transform, dates take their day number (num_days_from_ce), datetimes take their canonical i128-nanosecond key (timestamp_seconds × 10⁹ + subsecond_nanos, exact and injective for the full calendar range — no microsecond rounding, no saturation outside 1677–2262), and decimals are placed on a shared 10^18 fixed-point grid (mantissa × 10^(18−scale)) — exact and scale-invariant (2.5 and 2.50 reduce identically). A mixed integer/float axis reduces both sides through the float transform, so the two encodings agree; a decimal value whose scaled magnitude overflows i128 or whose scale exceeds 18 fractional digits is not representable and aborts the run with E326 rather than dropping the row. An operand pairing that cannot reduce to one exact axis — an ambiguous Numeric operand (abs/min/max/clamp, whose runtime value may be an integer compared exactly or a float) or a non-orderable pairing — is rejected at plan time with E327 rather than routed to the axis.
  • SortMerge key encoding. The alternate pure-range strategy (SortMerge, chosen for a single-inequality predicate on pre-sorted inputs) maps each range Value to a monotone key WITHIN its own type: integers, floats (orderable transform), and dates (day number, num_days_from_ce) reduce to an i128 key, and datetimes reduce through the SAME canonical i128-nanosecond key the IEJoin axis and the memcomparable sort bytes use, so a presorted (SortMerge) and an unsorted (IEJoin) datetime run agree row for row. That within-type reduction is correct only for a same-type pair, so strategy selection restricts SortMerge to Integer/Float/Date/DateTime axes; a decimal axis (which it would map to a constant) and a mixed int/float axis (disjoint per-type encodings) route to IEJoin instead, which reduces both sides to one i128 space. The result is identical either way — only the physical strategy differs.
  • Output axis. Emitted rows accumulate in a payload-ordered sort buffer keyed on (driver order, driver index, build index) that spills on its own byte threshold, so the O(N·M) result of a high-fan-out join never has to fit in RAM. The final order is realized by that (possibly external) sort, so the output is byte-identical across memory budgets.
  • Dense block pairs. A surviving block-pair normally materializes its candidate vector before emitting. The kernel is asked for at most max_pairs = (hard limit − reserved pair peak) / 48 pairs (16 bytes per pair, times three for the vector’s growth), and stops the moment the pair’s actual match count would exceed that, without allocating the over-budget vector (iejoin/block.rs, the materialize-vs-fallback decision). A selective pair (matches far below driver × build) takes the materialized path however large its blocks are; only a genuinely dense pair, such as a single hot equality value whose block-pair is a near cross product that finer slicing cannot reduce, is instead streamed through a bounded block-nested-loop that buffers only a small tile of candidate index pairs (at most min(headroom / 16, 4096)), emitting each match immediately through the same spillable output sink. Extra residency is O(tile), so the pair completes within the budget instead of aborting. The nested-loop path is byte-identical to the materialized path (the output re-sorts, so candidate arrival order never shows) and preserves match: first/all/collect, the residual filter, and on_miss exactly.

max_output_rows, when set on the node, is a strategy-agnostic result-size runaway guard enforced at each strategy’s output-emit chokepoint — the spill-backed paths (block-band IEJoin, sort-merge) check the output sort’s cumulative total_rows() before each push, and the RAM-vec paths (hash build-probe, grace-hash) check the output vector’s cumulative length. On a breach the run fails loud with E325 rather than truncating. Because every emitted row on every strategy passes one such check, the cap covers all match modes, the deferred on_miss rows, and both the materialized and nested-loop block-pair paths. It is a result-size guard, orthogonal to the byte-based spill / RSS machinery above (which answers memory pressure by spilling or a typed budget abort, not by capping rows). Recoverable dead-lettered rows never reach the emit chokepoint, so they are not counted; under match: collect the counted rows are the per-driver output rows.

Document boundaries

A Combine forwards reconciled document boundaries to its output on every strategy — the inline hash build-probe, IEJoin, grace-hash, sort-merge, and the streaming-probe path. The boundary semantics are uniform across the strategy matrix, so a downstream operator never has to know which join algorithm ran.

Concretely:

  • A per-document Aggregate downstream of a join flushes per document. A driver source that carries several documents (a glob: over monthly files, say) produces one roll-up per driver document after the join, not one fold spanning all of them.
  • A document that spans both join inputs — the same document carried on the driver and on the build side — opens and closes exactly once downstream. The boundary is reconciled, never double-fired: the join does not emit a separate open/close for the driver-side and build-side appearances of the same document.

This reconciliation is what lets the per-document aggregation model compose with joins without special-casing the operator order.

Correlation-key propagation

Combine declares which correlation-key columns its output rows carry via the required propagate_ck field. The choice shapes both the compile-time output schema and the runtime record builder — those are the two internal surfaces an engine engineer touches when changing CK behavior.

propagate_ck valueCompile-time output schemaRuntime record builder
driverCarries only the driver input’s $ck.<field> columns.Build-side records contribute body fields only; their CK identity is consumed by the match and not copied onto the output row.
allCarries every input’s $ck.<field> columns.Copies build-side CK values onto each output row alongside the body’s emit columns. Use when the build side carries CK fields downstream operators must read.
{ named: [<field>, ...] }Carries the explicit subset, intersected with what is actually present upstream.Copies exactly the named subset. Use to project a multi-field CK down to a single field after a join.

Driver wins on a name collision. If both the driver and a build input declare $ck.<field>, the column appears once on the output schema and the runtime keeps the driver’s value.

propagate_ck is required on every Combine; a pipeline without an explicit value fails to compile.

Match-mode interaction across the strategy paths

The propagation contract holds identically across the hash build-probe, IEJoin, grace-hash, and sort-merge paths — the record builder is shared, so a build-side CK value lands on the output row the same way regardless of which algorithm produced the match. The interaction that does vary is by match mode rather than by strategy:

  • match: first / match: all — each emitted row is one driver × one build pairing, so the propagated $ck.<field> slot holds a single value (the driver’s, or the build’s, per the table above).
  • match: collect — the propagated CK slot is single-valued (it tracks the driver’s correlation-group identity), while the collected array column preserves the full lineage of every build match. The single-valued slot and the array column are distinct: the slot answers “which correlation group does this output row belong to,” the array answers “which build records were gathered.”

See the User Guide’s correlation-keys reference for the per-mode lifecycle narrative; the lifecycle and rollback-narrowing mechanics on the engine side are in Correlation Keys: Lifecycle & Rollback Narrowing.

Join-planner statistics

When the plan carries column statistics, --explain ends with a === Statistics === section listing the planner-wide statistics catalog. These are the figures that drive build-side selection and grace-hash partition-bit choice, so they belong to the join planner. Every figure is tagged with its provenance, so a metadata-derived estimate is distinguishable from a record-exact measurement.

Row counts — [file metadata] vs [exec sketch]

One line per source node, for example:

orders: ≈90 rows [file metadata] (informs combine build/probe + partition bits)
  • A [file metadata] figure is derived at plan time by dividing the input file’s on-disk byte length by an average-record-bytes constant, before any record is read. This is the same row count that drives a Combine’s build-side selection and its grace-hash partition-bit choice. A build side large enough to risk overrunning the memory limit is what tips a pure-equality Combine from the in-memory hash strategy to the disk-spilling grace-hash strategy.
  • An [exec sketch] figure is the exact count a source measured during a run, superseding the plan-time estimate.

Row counts also appear inline on each Combine’s driving and build inputs (est. 90 [file metadata] rows).

Column sketches — distinct counts, heavy hitters, membership filters

Three sketch kinds are populated by operators while records flow. All three are maintained by the grace-hash Combine over its build-side join keys, recorded under the build input’s (node, column):

  • Distinct-count estimate — product_id: 12,431 distinct [exec sketch].
  • Top-k heavy-hitter list with lower-bound counts — product_id: heavy hitters [exec sketch, lower bound]: widget=9,000, gadget=3,200, .... The list is a lower bound on frequency: a value absent from it may still be frequent, so it is only ever used to promote a key, never to exclude one.
  • Membership filter — product_id: membership filter, 119048 bits / 7 probes [exec sketch, sized from estimate]. Sized up front from the build node’s plan-time row-count estimate, built in the single build pass with no per-row buffer, and skipped entirely when no plan-time estimate is available.

Honest nulls and missing sections

A statistic that was never gathered renders as null rather than a fabricated zero. A plan over sources whose sizes cannot be read — a glob/regex multi-file source, a network source, or a missing/unreadable file — adds no Statistics section at all, and (per the membership-filter rule above) skips the membership filter that the plan-time estimate would have sized. Confirm the live shape via clinker metrics collect after the first production run, since the planner has no group-cardinality side-table to consult before the run.