Skip to content

fix(query): a multi-key select with a distinct count beside an expression aggregate - #760

Open
singaraiona wants to merge 2 commits into
devfrom
fix/756-multikey-distinct-count
Open

singaraiona wants to merge 2 commits into
devfrom
fix/756-multikey-distinct-count

Conversation

@singaraiona

Copy link
Copy Markdown
Collaborator

A select grouped by two or more keys failed with name when it held a (count (distinct x)) beside an aggregate over an expression, such as (sum (if (== st 'bad) 1 0)). Every narrower form worked.

Fixes #756.

The distinct count is not a DAG aggregate, so with more than one key the select fell to the eval-level composite grouping. That path evaluated a streaming aggregate's argument with plain ray_eval and no column in scope, so only a bare column name worked. It is also serial, which is the 8x slowdown the issue measured on the inner-select workaround.

Three changes, each closing one of the issue's observations:

  • The eval-level grouping evaluates an aggregate's argument over the selected rows with the columns bound, compiled as a row-wise projection through group_key_eval (fix(query): group by a computed key beside distinct aggregates and row expressions #708's helper). A plain scoped eval would still have been wrong: the eval-level if tests whole-vector truthiness and returns a scalar.
  • A plain multi-key by: stays on the parallel DAG path when the only outputs the DAG cannot serve are distinct counts and literal broadcasts. The post-group scatter maps rows to groups through the composite key: rgid_build_multikey hashes the result's key tuples and probes every row in parallel, applying the where: selection as the single-key probe does. The existing per-group kernels then serve the distinct counts, including the global-hash kernel past 50 K groups. STR, GUID, LIST and mapcommon keys keep their previous routes; a row expression beside a composite key keeps the eval-level grouping.
  • The two shapes beside the failing one are routed the same way. An aggregate over an if dropped the whole group onto the legacy engine, 6x slower under a where: selection (109 ms against 14 ms without one): the v2 expression bridge in exec_group_v2_exprs now materializes an input the element-wise compiler declines through the DAG executor, with the selection cleared around the call so the vector comes back at full length. A distinct count alone took the two-pass rewrite, grouping by (keys, value) and then counting per key. When that composite is not dense it is a radix grouping that materializes and orders every tuple: 1.8 GB and 3x the time of the slice dedup at 10 M rows with 45 groups of 3 M values. The rewrite is now gated on cd_two_pass_preferred: it keeps the rewrite when the (keys, value) space packs within the row count, when the groups outnumber 256 K, or when they are fewer than the workers; otherwise the select groups by the keys and counts over the slices. The estimate comes from agg_group_card_estimate, a strided sample of up to 65,536 rows: span for integer keys, distinct codes for SYM keys (codes of a shared domain interleave, so a span would overstate a 5-value key by orders of magnitude), a fraction of a millisecond at any size.

Validation:

  • Full ASan/UBSan suite: 4,111 of 4,111 passed.
  • test/rfl/regress/issue_756.rfl compares every shape against the two-select-and-join oracle: both aggregate orders, an aggregate over arithmetic, where: on whole keys and on part of every morsel, three keys with an integer key, literal broadcasts, desc:/take:, the by-dict form, a distinct count alone with and without where:, the eval-level grouping forced by a row expression, and an if aggregate under where: for one and two keys.
  • Beyond the file: a 300 K-row table with the where: prefilter shape (desc: s take: 5), two distinct counts beside a sum, a distinct over an expression, and a 210 K-group integer composite on the global-hash kernel, all equal to the join oracle.

Benchmark: the issue's table, 10 M rows, release build, 7 cores, median of repeats, milliseconds. "dev" is 24c14e85.

Form dev this PR
The failing select name error 50
The failing select with where: (!= k1 'a) name error 45
Inner select, then group (the issue's workaround) 883 48
Two selects and left-join 134 58
Distinct count alone, two keys 119 45
Distinct count beside (sum id) 875 45
Expression aggregate alone, with where: 109 16
Expression aggregate alone 14 13

The distinct-count gate across group and value cardinalities, 10 M rows, the two-pass rewrite on dev against this PR's choice:

Keys × distinct values rewrite this PR
45 × 3 M (the issue) 125 41
45 × 100 15 15
25 K × 3 M 66 35
25 K × 100 25 25
45 K × 600 K 57 36
2.5 M × 100 50 54
2.5 M × 600 K 75 74
by-dict {a: k1 b: k2}, 45 × 3 M 128 39

Peak memory for the issue's distinct count alone falls from +1.8 GB (the (keys, value) radix grouping) to +78 MB. The single-key fused kernel (ray_cd_fused) is untouched.

singaraiona and others added 2 commits October 9, 2026 12:35
…sion aggregate (#756)

A select grouped by two or more keys failed with `name` when it held a
`(count (distinct x))` beside an aggregate over an expression, such as
`(sum (if (== st 'bad) 1 0))`.  The distinct count is not a DAG aggregate,
so with more than one key the select fell to the eval-level grouping,
which evaluated an aggregate's argument with no column in scope — and
which is serial, 8x slower than the parallel path for the shapes that did
not fail.

The eval-level grouping now evaluates an aggregate's argument over the
selected rows with the columns bound, as a row-wise projection
(group_key_eval), since the eval-level `if` is scalar.

A plain multi-key by: stays on the DAG path when the outputs the DAG
cannot serve are distinct counts and literal broadcasts: the post-group
scatter maps rows to groups through the composite key (rgid_build_multikey,
a parallel tuple probe that applies the selection), and the existing
per-group kernels serve the distinct counts.

Two shapes the issue's workarounds exposed go the same way.  An aggregate
over an `if` dropped the group onto the legacy engine, 6x slower under a
where: selection: the v2 expression bridge now materializes an input the
element-wise compiler declines through the DAG executor.  A distinct count
alone took the two-pass rewrite — group by (keys, value), then count —
which is a radix grouping of every tuple when that composite is not dense
(1.8 GB and 3x the time at 10 M rows with 45 groups of 3 M values); the
rewrite is now gated on a sampled estimate of the group and value counts
(agg_group_card_estimate) and the select otherwise groups by the keys and
counts over the slices.

Closes #756.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
@ser-vasilich

Copy link
Copy Markdown
Collaborator

Reviewed 2d624e2. The fix works for #756. The wins hold, and everything else matched dev and independent oracles, except a wrong result with -0.0 keys.

Ours vs dev, 10M rows, -c 8, ms:

  • The failing select: 106–118 (dev: name error).
  • The inner-select workaround: 2515 → 106.
  • Distinct count alone, 2 keys: 230 → 100, peak memory 1390 → 216 MB.
  • if aggregate with where:: 143 → 26.
  • Multi-key literal + sum: 2355 → 15, a gain the description does not mention.
  • No regression on: single-key distinct counts (fused kernel, with where:, with sum) and multi-key sum/count without distinct.

Correctness: composite scatter results matched dev's eval-level grouping, an encoded single-key oracle, and a count(distinct rowid) == count invariant. Coverage:

  • key types I64/I32/I16/U8/BOOL/SYM/DATE/TIME/TIMESTAMP and F64 without -0.0;
  • nulls in one key component, in all components, and in the distinct column;
  • where: keeping 0 rows, 1 row, sparse, dense, whole keys, or the table tail;
  • 2 and 3 keys, from 45 up to 1.44M groups;
  • splayed and parted tables, keys from two FILE domains, and union-all of two splayed domains;
  • the by-dict form, two distinct counts, distinct over an expression, desc:/take:.

-c 8 equals -c 1. The if bridge returns full-length aligned vectors for every selection shape. TSan is clean, and make test (ASan/UBSan) passes 4111/4111.

BUG: F64/F32 keys holding -0.0 lose rows on the composite path.

(set T (table [p s x] (list (take [0.0 1.5] 8) (take ['a 'b 'c 'd] 8) (til 8))))
(select {from: T by: [p s] n: (count (distinct x)) c: (count x) asc: [p s]})
;; dev: n [2 2 2 2]   this PR: n [0 0 2 2]

The grouping canonicalises -0.0 to +0.0 when it reads a key (group.c:62-89, #407), but key_read_i64 (query.c:6799) bit-casts F64/F32 without that step, and rgid_mk_key_type_ok admits F32/F64.

  • In release builds the literal 0.0 already carries -0.0 bits, so ordinary code hits this.
  • A distinct count alone loses the same rows: n=0 against 250 at 300K rows. So do (as 'F64 [-0.0]) and CSV -0.0 loaded as F32/F64.
  • Single-key by: p gives 0 on dev too, so the bug existed on the single-key path; this PR spreads it to composite keys.
  • Either canonicalise -0.0 in key_read_i64 for F32/F64, or keep F32/F64 keys off this path.

The gate (agg_group_card_estimate / cd_two_pass_preferred). The sample is 64 contiguous 1,024-row blocks, not the strided sample the description says. Clustered or skewed SYM keys are underestimated, so the gate picks scatter where the rewrite is faster. 10M rows, -c 8, ms / peak-live MB:

shape dev (rewrite) this PR (scatter)
2 keys, skewed SYM (90% one value, ~1M groups) 163–189 / 574 258–330 / 670–695
2 keys, clustered SYM (4 rows per value, 5M groups) 316–328 / 1345 531–670 / 595
1 key + where:, skewed SYM, 857K groups 99–107 / 342 405–408 / 529
1 key + where:, clustered SYM, 2.5M groups 196–199 / 1217 240–289 / 467

Single-key where: count-distinct now goes through this gate, and the skewed case is ~4x slower than dev.

Other gate notes:

  • "Fewer groups than workers" removes the fix on big machines. The issue's own 45-group query at -c 46 picks the rewrite: 199 ms / 1375 MB, against 70 ms / 512 MB with scatter forced. It still picks scatter at -c 23.
  • Overestimates keep the rewrite where scatter wins. These are not regressions against dev:
    • TIMESTAMP keys, 45 groups (the span overflows): 262 / 1405 vs 132 / 216;
    • sorted integers, 750K groups (above the 256K cap): 271 / 1405 vs 165 / 422.
  • Planning cost. agg_sym_distinct_sampled allocates and memsets a 2 MB set per SYM key, about 36 µs each. A distinct count alone on a 10-row table goes from 6.6–11 µs to 43–48 µs with one SYM key and 79–82 µs with two. Single-key queries pay it too whenever ray_cd_fused declines (small input, or where:).

Behaviour changes users will see:

  • Literals beside a distinct count become LIST columns. one: 'z lit: 1 f: 2.5 beside n: (count (distinct x)) with two keys returns LIST columns of atoms; dev returns typed SYM/I64/F64 (STR for a string). Single-key already behaved this way.
  • Column order changes for composite keys. DAG aggregates now come first: [k1 k2 n s] → [k1 k2 s n], and [k1 k2 one s] → [k1 k2 s one] with no distinct at all. issue_756.rfl:15 pins the new order.
  • Row order without asc: changes. At 600K rows dev emits first-seen order and the PR packed-key order; ties under desc: c take: break differently.
  • By-dict aliases now depend on the data. The rewrite already drops the aliases (dev always returns [k1 k2 n]). Now the same query returns [a b n] when the gate picks scatter and [k1 k2 n] when it picks the rewrite; single-key aliases with where: behave the same way.
  • Window aggregates under where: change when the query falls back to eval level. (sum (deltas x)) / (sum (sums x)) see the whole column on the DAG path (as on dev), but only the filtered rows on the new eval-level path. So adding an unrelated row expression such as r: (+ x 1) changes the aggregate's value.
  • A cancel may be lost (from the code, not reproduced). The nested ray_execute → ray_execute_inner resets pool->cancelled, so a cancel that arrives before the nested call would be dropped.

Tests. I reverted each piece in turn:

  • The composite probe disabled: only issue_756:55 fails. Lines 18–20 pass, because the test's oracle O for n goes through the same probe.
  • The probe ignoring the where: selection: not caught, because it only matters on the >50K-group kernel. A 60K-group where: case catches it.
  • The gate always picking scatter: not caught.
  • The ray_execute fallback in the v2 bridge removed: not caught; the full suite passes.
  • Routing reverted: caught only by the column-order assertion at :15.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Bug: multi-key group-by with a distinct count fails with name beside an expression aggregate, and falls to a serial path 8× slower

2 participants