From ec8d6fb3fc1717658aa22455521e5508819d2dc1 Mon Sep 17 00:00:00 2001 From: Tam Nguyen Duc <1218621+tamnd@users.noreply.github.com> Date: Tue, 25 Aug 2026 13:56:26 +0700 Subject: [PATCH 1/4] exec: fold the group tables up a tree Every worker sees keys from all over the scan, so each partial ends up about as wide as the answer and folding eight of them into the first is seven full width merges one after another, at the point in the query where nothing else is running. On the hundred thousand group bench that is 8 to 18 ms of a 65 ms query. Same seven merges, three rounds, each round on the pool the scan has just finished with. An odd table carries to the next round rather than joining a pair, so every merge in a round is the same size, and this thread takes one pair itself the way worker zero does on the scan. --- crates/zu-exec/src/sink.rs | 83 +++++++++++++++++++++++++++++++------- 1 file changed, 68 insertions(+), 15 deletions(-) diff --git a/crates/zu-exec/src/sink.rs b/crates/zu-exec/src/sink.rs index 94086a9a..34fdbda7 100644 --- a/crates/zu-exec/src/sink.rs +++ b/crates/zu-exec/src/sink.rs @@ -9,6 +9,7 @@ use std::cmp::Ordering; use std::collections::BTreeSet; +use std::sync::Mutex; use zu_common::gqlstatus::codes; use zu_common::{Result, ZuError}; @@ -1091,14 +1092,7 @@ pub(crate) fn finish_agg( // Fold every other worker into the first non-empty table rather // than into a fresh one, so the biggest partial is usually the one // nobody has to rehash. - let mut merged: Option = None; - for p in partials { - let Some(t) = p.groups else { continue }; - match &mut merged { - None => merged = Some(t), - Some(m) => m.merge_from(&t)?, - } - } + let merged = fold_tables(partials)?; // Asked here rather than per row on the way in: the table holds one // copy of each distinct key and every row that did not create a // group has the bytes of the one that did. @@ -1118,15 +1112,74 @@ pub(crate) fn finish_agg( /// merge only has to find the keys one table holds that another one /// does not. pub(crate) fn finish_distinct(partials: Vec) -> Result { - let mut merged: Option = None; - for p in partials { - let Some(t) = p.groups else { continue }; - match &mut merged { - None => merged = Some(t), - Some(m) => m.merge_from(&t)?, + Ok(fold_tables(partials)?.map_or(0, |m| m.groups() as i64)) +} + +/// Folds the workers' group tables into one. +/// +/// Pairwise up a tree rather than one after another into the first, +/// because the scheduler hands every worker morsels from all over the +/// scan, so each partial ends up about as wide as the answer and every +/// fold is a full width merge. Eight of them in a row is seven of those +/// at the one point in the query where nothing else is running, which +/// on the hundred thousand group bench is as much wall clock as the +/// eight worker scan that produced them. The tree does the same seven +/// merges in three rounds and runs each round on the pool the scan just +/// finished with, so what was seven merges deep is three. +/// +/// An odd table carries to the next round rather than being merged into +/// a pair, which keeps every merge in a round the same size. +fn fold_tables(partials: Vec) -> Result> { + let mut tables: Vec = partials.into_iter().filter_map(|p| p.groups).collect(); + while tables.len() > 1 { + let carry = (tables.len() % 2 == 1).then(|| tables.pop().expect("an odd count has a last")); + let mut pairs: Vec<(GroupTable, GroupTable)> = Vec::with_capacity(tables.len() / 2); + let mut it = tables.into_iter(); + while let (Some(a), Some(b)) = (it.next(), it.next()) { + pairs.push((a, b)); + } + // One pair for this thread, the same way the scan keeps worker + // zero, so a round of one pair never goes near the pool. + let mine = pairs.pop().expect("a round starts with two tables"); + let slots: Vec>>> = + pairs.iter().map(|_| Mutex::new(None)).collect(); + let here = { + let jobs: Vec> = pairs + .into_iter() + .zip(&slots) + .map(|((mut a, b), slot)| { + Box::new(move || { + *slot.lock().unwrap() = Some(a.merge_from(&b).map(|()| a)); + }) as Box + }) + .collect(); + let pending = crate::pool::submit(jobs); + let (mut a, b) = mine; + let here = a.merge_from(&b).map(|()| a); + pending.wait(); + here + }; + let mut next = Vec::with_capacity(slots.len() + 2); + let mut first_err = None; + for res in std::iter::once(here).chain(slots.into_iter().map(|slot| { + slot.into_inner().unwrap().unwrap_or_else(|| { + Err(ZuError::InvalidArgument( + "aggregation merge panicked".into(), + )) + }) + })) { + match res { + Ok(t) => next.push(t), + Err(e) => first_err = first_err.or(Some(e)), + } + } + if let Some(e) = first_err { + return Err(e); } + next.extend(carry); + tables = next; } - Ok(merged.map_or(0, |m| m.groups() as i64)) + Ok(tables.pop()) } /// Stitches row batches back into scan order and applies the posts. From ac5b6acce5e0b67b4e2e2e6fde4167a77ac9af5f Mon Sep 17 00:00:00 2001 From: Tam Nguyen Duc <1218621+tamnd@users.noreply.github.com> Date: Tue, 25 Aug 2026 14:51:04 +0700 Subject: [PATCH 2/4] exec: build the answer rows on the pool The rows of a wide aggregation were built one after another on the thread that finished the run, and that is 7 to 17 ms of a 65 ms query at eight workers. Decoding is per group and the groups are settled once the order is, so slices of the order are independent work and the pool the scan just finished with is sitting idle. One hand per worker that ran, not one per core: a query asked for on one thread is answered on one thread, tail included, or the per core numbers in budgets.toml stop meaning what they say. Under four thousand groups a hand it stays on this thread, since a latch and a lock per hand cost more than the decoding they would split. --- crates/zu-exec/src/group.rs | 96 ++++++++++++++++++++++++++++++++----- crates/zu-exec/src/sink.rs | 5 +- 2 files changed, 89 insertions(+), 12 deletions(-) diff --git a/crates/zu-exec/src/group.rs b/crates/zu-exec/src/group.rs index 48fbc860..4c891464 100644 --- a/crates/zu-exec/src/group.rs +++ b/crates/zu-exec/src/group.rs @@ -46,6 +46,7 @@ //! them and decodes each group once, straight into its finished row. use std::cmp::Ordering; +use std::sync::Mutex; use zu_query::exec::Value; use zu_query::snapshot::TemporalLane; @@ -325,6 +326,13 @@ pub(crate) struct GroupTable { /// covers the many-groups-of-one queries that never grow past a handful. const INIT_SLOTS: usize = 64; +/// Groups a hand needs before the finished rows are worth building on +/// more than one thread. Four thousand of them is a few hundred +/// microseconds of decoding, which is well clear of what a latch and a +/// lock per hand cost, and it keeps every query that groups into tens +/// or hundreds on the one thread it was already answered on. +const SPLIT_ROWS: usize = 4096; + const IDX_MASK: u64 = (1 << 48) - 1; const SEED: u64 = 0x9E37_79B9_7F4A_7C15; @@ -828,11 +836,7 @@ impl GroupTable { /// thrown away a moment later. The order is settled over an index /// vector against the stored words, so nothing is decoded until the /// row is built and each row is built once, in its final place. - pub(crate) fn rows(mut self, item_agg: &[bool]) -> Vec> { - // Out of the table first, while it is still whole, so that the - // sort below can borrow the rest of it. - let states = std::mem::take(&mut self.accs); - let mut done: Vec = states.into_iter().map(Acc::finalize).collect(); + pub(crate) fn rows(mut self, item_agg: &[bool], hands: usize) -> Vec> { let mut order: Vec = if self.counting { std::mem::take(&mut self.order) } else { @@ -846,9 +850,62 @@ impl GroupTable { } else { order.sort_unstable_by(|&a, &b| self.cmp_keys(a as usize, b as usize)); } + // Decoding is per group and the groups are settled by now, so + // the slices of the order are independent and the pool the scan + // has just finished with is sitting there. `hands` is the run's + // worker count and not the machine's, because a query asked to + // run on one worker is answered on one worker, tail included. + // Under the split it is not worth a latch and two locks. + let hands = hands.min(order.len() / SPLIT_ROWS); + if hands < 2 { + return self.build(&order, item_agg); + } + let slices: Vec<&[u32]> = order.chunks(order.len().div_ceil(hands)).collect(); + let slots: Vec>>>> = + slices[1..].iter().map(|_| Mutex::new(None)).collect(); + let mine = { + let me = &self; + let jobs: Vec> = slices[1..] + .iter() + .zip(&slots) + .map(|(slice, slot)| { + Box::new(move || { + *slot.lock().unwrap() = Some(me.build(slice, item_agg)); + }) as Box + }) + .collect(); + let pending = crate::pool::submit(jobs); + let mine = self.build(slices[0], item_agg); + pending.wait(); + mine + }; + let mut out = mine; + out.reserve(order.len() - out.len()); + for (slot, slice) in slots.into_iter().zip(&slices[1..]) { + // A slot still empty after the latch means that hand + // panicked, and the rows it was building are gone. Building + // them here keeps the answer whole and costs the panic + // handler's path only. + match slot.into_inner().unwrap() { + Some(rows) => out.extend(rows), + None => out.extend(self.build(slice, item_agg)), + } + } + out + } + + /// The rows for one slice of the settled order. + /// + /// Takes the table by reference and reads each state out of it + /// rather than moving the states away first, so that slices of the + /// order can be built beside each other. An accumulator is a word or + /// two and is `Copy`, so reading one is cheaper than the bookkeeping + /// it would take to hand each thread the states it needs out of a + /// buffer indexed by group rather than by place in the order. + fn build(&self, order: &[u32], item_agg: &[bool]) -> Vec> { let mut out = Vec::with_capacity(order.len()); let mut keys = Vec::new(); - for &g in &order { + for &g in order { if self.counting { keys.push(self.counted_key(g)); } else { @@ -861,7 +918,7 @@ impl GroupTable { let v = if self.counting { Value::Int(self.slots[g as usize][1] as i64) } else { - std::mem::replace(&mut done[agg_at], Value::Null) + self.accs[agg_at].finalize() }; agg_at += 1; v @@ -1277,7 +1334,7 @@ mod tests { // count(*) first and the key second, which is what // `RETURN count(*), n.k` asks for, so a row that just // concatenated the two halves would come out backwards. - let rows = t.rows(&[true, false]); + let rows = t.rows(&[true, false], 1); assert_eq!(rows.len(), 4); assert_eq!( rows, @@ -1296,7 +1353,7 @@ mod tests { // order nor the answer. let mut t = GroupTable::counting(vec![PartKind::Int], 1); t.count_ints(&[7, 1, 4, 7, 7]); - let rows = t.rows(&[false, true]); + let rows = t.rows(&[false, true], 1); assert_eq!( rows, vec![ @@ -1356,7 +1413,24 @@ mod tests { vals }) .collect(); - assert_eq!(t.rows(&[false, false, false, true]), want); + assert_eq!(t.rows(&[false, false, false, true], 1), want); + } + + /// The split is a split of the work and not of the answer, so a + /// table wide enough to go to the pool has to put out exactly what + /// it puts out on one hand, in the same order. + #[test] + fn many_hands_answer_what_one_hand_answers() { + let vals: Vec = (0..200_000).map(|i| (i * 7919) % 20_000).collect(); + let mut t = GroupTable::new(vec![PartKind::Int], 1); + count_ints(&mut t, &vals); + let mut same = GroupTable::new(vec![PartKind::Int], 1); + same.merge_from(&t).unwrap(); + assert!( + t.groups() > SPLIT_ROWS * 2, + "wide enough that the split is taken" + ); + assert_eq!(t.rows(&[false, true], 8), same.rows(&[false, true], 1)); } /// A date is the one lane stored narrower than the word it rides in, @@ -1367,7 +1441,7 @@ mod tests { let lane = TemporalLane::Date; let mut t = GroupTable::new(vec![PartKind::Temporal(lane)], 1); count_ints(&mut t, &[10, -5, 0, -400]); - let rows = t.rows(&[false, true]); + let rows = t.rows(&[false, true], 1); let days: Vec = rows .iter() .map(|r| match r[0] { diff --git a/crates/zu-exec/src/sink.rs b/crates/zu-exec/src/sink.rs index 34fdbda7..8b76132f 100644 --- a/crates/zu-exec/src/sink.rs +++ b/crates/zu-exec/src/sink.rs @@ -1076,6 +1076,9 @@ pub(crate) fn finish_agg( partials: Vec, keys_empty: bool, ) -> Result { + // One hand per worker that ran, so a query asked for on one thread + // is finished on one thread too. + let hands = partials.len(); // A bare aggregate has one group whichever way the input went, so // it never needs the table: fold the per-worker state vectors and // emit the row even when no worker saw a row at all. @@ -1101,7 +1104,7 @@ pub(crate) fn finish_agg( "string property is not UTF-8".to_string(), )); } - let rows = merged.map(|t| t.rows(item_agg)).unwrap_or_default(); + let rows = merged.map(|t| t.rows(item_agg, hands)).unwrap_or_default(); Ok(QueryResult::new(columns, apply_post(post, rows))) } From 684a7398590e17efc910b6ddb62b083939d3a9ab Mon Sep 17 00:00:00 2001 From: Tam Nguyen Duc <1218621+tamnd@users.noreply.github.com> Date: Tue, 25 Aug 2026 15:07:10 +0700 Subject: [PATCH 3/4] docs: regenerate the api model Bounded list entities from #770, which merged without them. --- docs/api/model.json | 66 ++++++++++++++++++++++++++++++++------------- 1 file changed, 48 insertions(+), 18 deletions(-) diff --git a/docs/api/model.json b/docs/api/model.json index 87d1637e..5ab636af 100644 --- a/docs/api/model.json +++ b/docs/api/model.json @@ -10034,6 +10034,28 @@ "of": "zu::zu1::props::ListElement", "signature": "Word(u64)" }, + { + "id": "zu::zu1::props::ListRows", + "kind": "enum", + "name": "ListRows", + "doc": "How a row of a list column says how many elements it holds.\n\nThe two are not tellable apart from the bytes, which is why this is\ncarried rather than worked out: a bounded `LIST(768)` row of\n767 counted elements is the same length as one of 768 uncounted\nones. The column entry says, and a reader takes the answer." + }, + { + "id": "zu::zu1::props::ListRows::Counted", + "kind": "variant", + "name": "Counted", + "of": "zu::zu1::props::ListRows", + "signature": "Counted(usize)", + "doc": "Every row says its own count, in this many bytes." + }, + { + "id": "zu::zu1::props::ListRows::Fixed", + "kind": "variant", + "name": "Fixed", + "of": "zu::zu1::props::ListRows", + "signature": "Fixed(usize)", + "doc": "The directory says the count and no row carries one, which is a\ncolumn every row of which was written at its declared bound." + }, { "id": "zu::zu1::props::NodesWithProps", "kind": "type-alias", @@ -10124,6 +10146,14 @@ "kind": "struct", "name": "PropColumn" }, + { + "id": "zu::zu1::props::PropColumn::count_width", + "kind": "field", + "name": "count_width", + "of": "zu::zu1::props::PropColumn", + "signature": "u8", + "doc": "How many bytes each row of this list column spends saying its own\nelement count, where it says one at all. Four is what every\nversion before 12 wrote and what an unbounded list still writes.\n\nIt is carried rather than worked out from the declared bound\nbecause a column is read at the width its rows were written at,\nand those are two different questions once the rule has changed\nonce. A version 11 column of `LIST(4)` rows holds four byte\ncounts, and a fold that leaves those rows alone has to leave the\nfour with them, whatever a column written today would have used." + }, { "id": "zu::zu1::props::PropColumn::fixed_len", "kind": "field", @@ -10132,14 +10162,6 @@ "signature": "Option", "doc": "The element count every row of this list column holds, where the\nrows do not carry one themselves. `None` is every other column,\nand a list column whose rows say their own count.\n\nThis is the encoding half of schema/06 ยง2 written down: the type\nabove is the declaration, `LIST(768)`, and this says\nthat this column was written whole at exactly the bound, so the\nblock may drop what the declaration already gives. It cannot be\nworked out from the bytes, which is why it is a field: a bounded\n`LIST(768)` column of 767 counted elements a row is the\nsame 3072 bytes a row as one of 768 uncounted ones." }, - { - "id": "zu::zu1::props::PropColumn::fixed_list", - "kind": "method", - "name": "fixed_list", - "of": "zu::zu1::props::PropColumn", - "signature": "fn fixed_list(&self) -> Option", - "doc": "The count to read a row of this column at, which is what\n[`list_elements`] wants and `None` for every row that says its\nown." - }, { "id": "zu::zu1::props::PropColumn::is_lane", "kind": "method", @@ -10148,6 +10170,14 @@ "signature": "fn is_lane(&self) -> bool", "doc": "Whether this column stores in the fixed width lane, which is one\n64 bit word per row through the integer cascade, as against the\nblob segments a string or a byte string uses.\n\nEvery fixed stride type of eight bytes or less rides the lane:\nthe cascade encodes words and does not care what the bits mean,\nand the column's type is what says whether a word is a count, a\ntruth value, a float, a day or a nanosecond." }, + { + "id": "zu::zu1::props::PropColumn::list_rows", + "kind": "method", + "name": "list_rows", + "of": "zu::zu1::props::PropColumn", + "signature": "fn list_rows(&self) -> ListRows", + "doc": "How to read a row of this list column, which is what\n[`list_elements`] wants beside the bytes: the count the directory\nholds, or the width the row says its own in." + }, { "id": "zu::zu1::props::PropColumn::meta", "kind": "field", @@ -10428,14 +10458,6 @@ "of": "zu::zu1::props::PropsReader", "signature": "fn columns(&self) -> &[PropColumn]" }, - { - "id": "zu::zu1::props::PropsReader::fixed_list", - "kind": "method", - "name": "fixed_list", - "of": "zu::zu1::props::PropsReader", - "signature": "fn fixed_list(&self, col: usize) -> Option", - "doc": "The element count to read a row of `col` at, which is what\n[`list_elements`] wants beside the bytes." - }, { "id": "zu::zu1::props::PropsReader::gather_int", "kind": "method", @@ -10484,6 +10506,14 @@ "signature": "fn label_word(&mut self, db: &mut Zu1File, row: u64) -> Result>", "doc": "The label bitset of row `row`, `None` when the table stores\nnone, which is a table whose rows carry its name and nothing\nelse. The bits are dictionary positions, so a caller tests a\npattern's mask against the word with one AND." }, + { + "id": "zu::zu1::props::PropsReader::list_rows", + "kind": "method", + "name": "list_rows", + "of": "zu::zu1::props::PropsReader", + "signature": "fn list_rows(&self, col: usize) -> ListRows", + "doc": "How to read a row of list column `col`, which is what\n[`list_elements`] wants beside the bytes." + }, { "id": "zu::zu1::props::PropsReader::meta", "kind": "method", @@ -10665,8 +10695,8 @@ "id": "zu::zu1::props::list_elements", "kind": "function", "name": "list_elements", - "signature": "fn list_elements(elem: &zu_common::LogicalType, bytes: &[u8], fixed: Option) -> zu_common::Result>", - "doc": "Reads back a row written by `encode_list_row`.\n\n`fixed` is the element count a column whose rows do not carry one\nholds, which is [`PropColumn::fixed_list`], and `None` is every row\nthat says its own count. Which of the two a row is cannot be settled\nfrom the bytes the way the element width can: a bounded\n`LIST(768)` row of 767 counted elements is the same length as\none of 768 uncounted ones. So the directory says, and this takes the\nanswer rather than guessing at it.\n\nThe elements borrow the buffer they came out of, so a read of a list\ncolumn is the blob read and a walk over it, with no allocation per\nelement." + "signature": "fn list_elements(elem: &zu_common::LogicalType, bytes: &[u8], rows: ListRows) -> zu_common::Result>", + "doc": "Reads back a row written by `encode_list_row`.\n\n`rows` is [`PropColumn::list_rows`], which says whether the count is\nin the directory or at the head of the row and how wide the head is.\nNeither can be settled from the bytes the way the element width can:\na bounded `LIST(768)` row of 767 counted elements is the same\nlength as one of 768 uncounted ones, and the head that says 767 is\none, two or four bytes according to what the column was written at.\nSo the directory says, and this takes the answer rather than guessing\nat it.\n\nThe elements borrow the buffer they came out of, so a read of a list\ncolumn is the blob read and a walk over it, with no allocation per\nelement." }, { "id": "zu::zu1::props::load_props", From bb5f05e5f8dd5392cf04cc49e5ff2563002a6de0 Mon Sep 17 00:00:00 2001 From: Tam Nguyen Duc <1218621+tamnd@users.noreply.github.com> Date: Tue, 25 Aug 2026 15:29:48 +0700 Subject: [PATCH 4/4] bench: write down what the parallel finish is worth --- bench/budgets.toml | 18 ++++++++++++++++++ 1 file changed, 18 insertions(+) diff --git a/bench/budgets.toml b/bench/budgets.toml index c41f4ac4..cdc83299 100644 --- a/bench/budgets.toml +++ b/bench/budgets.toml @@ -717,6 +717,24 @@ exec_scale_8x = 2.0 # which is the check that this was the tail and not the scan: twelve # groups have no tail worth the name. # +# What was left of that tail went to the pool next. Folding the eight +# workers' tables into one is seven full width merges one after +# another, because the morsel scheduler hands every worker keys from +# all over the scan, so each partial is about as wide as the answer. +# That fold goes pairwise up a tree now, which is the same seven +# merges in three rounds. The row build went to slices of the settled +# order, one hand per worker that ran rather than one per core, since +# a query asked for on one thread has to be answered on one thread or +# this number stops measuring a core. On gamingpc, quiet, medians of +# six paired alternating runs of two separately built binaries, the +# hundred thousand group query went 40.8 ms to 33.9 at eight workers, +# 243 to 296 M rows/s, and the branch was faster in all six pairs. At +# one worker it read 99.5 against 99.4, which is the check that the +# hand count follows the run and not the machine. The twelve group +# and thousand group string shapes did not move either, the first +# because it has no tail and the second because a thousand groups is +# under the split. +# # The bench crosschecks the group count and the row total against the # generator on every run. Raise this floor, never lower it. exec_group_mrows_s_core = 8.0