diff --git a/bench/budgets.toml b/bench/budgets.toml index 926b23f7..c41f4ac4 100644 --- a/bench/budgets.toml +++ b/bench/budgets.toml @@ -699,9 +699,26 @@ exec_scale_8x = 2.0 # wide enough to miss. These boxes do not have the memory bandwidth # to feed one core a scan this size, and no amount of work on the # sink moves that. So the floor sits under server3's worst run and -# the number that answers perf/05 is the bare metal one. The bench -# crosschecks the group count and the row total against the generator -# on every run. Raise this floor, never lower it. +# the number that answers perf/05 is the bare metal one. +# +# At d745f07 the sink stopped building the groups before it knew where +# any of them went. It used to drain the folded table into a pair of +# vectors per group, sort those by comparing the values they decoded +# to, and then walk the result a third time to put the keys and the +# aggregates back in the order the RETURN clause named them, which for +# a hundred thousand groups is three hundred thousand small +# allocations made and thrown away in the tail of the query. The table +# orders an index vector against its own packed key words instead, and +# decodes each group once, straight into its finished row. The hundred +# thousand group query went 136.7 ms to 111.9 at one worker and 80.6 +# to 64.9 at eight, medians of six paired alternating runs of two +# separately built binaries on the local M series, and the branch was +# faster in every pair. The twelve group shape beside it did not move, +# which is the check that this was the tail and not the scan: twelve +# groups have no tail worth the name. +# +# 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 # The same sink at one worker on a string key, M rows/s: the same ten diff --git a/crates/zu-exec/src/group.rs b/crates/zu-exec/src/group.rs index 51126a6c..48fbc860 100644 --- a/crates/zu-exec/src/group.rs +++ b/crates/zu-exec/src/group.rs @@ -36,9 +36,16 @@ //! once instead of one, and the same trick works on the accumulator //! updates, which the caller runs as its own loop over group indices. //! -//! Ordering is not the table's job. Groups come out in insertion order -//! and the sink sorts them at the end, which is where the old engine's -//! ascending-key output is reproduced. +//! Ordering is not the probe's job, but it is the table's. Groups are +//! built in insertion order and put out by key ascending, which is the +//! old engine's output, and [`GroupTable::rows`] is where that happens. +//! It is here rather than in the sink because the sink only had the +//! decoded `Value`s to sort, and sorting those meant building every +//! group's key and states before knowing where any of them went. The +//! table has the packed words, so it orders an index vector against +//! them and decodes each group once, straight into its finished row. + +use std::cmp::Ordering; use zu_query::exec::Value; use zu_query::snapshot::TemporalLane; @@ -781,55 +788,180 @@ impl GroupTable { } /// Every group as its key values and its states, insertion ordered. - pub(crate) fn drain(self) -> Vec<(Vec, Vec)> { + #[cfg(test)] + pub(crate) fn drain(mut self) -> Vec<(Vec, Vec)> { let mut out = Vec::with_capacity(self.groups); if self.counting { // Every counter over the same group counted the same rows, // so the one count in the slot answers all of them. for &at in &self.order { let slot = self.slots[at as usize]; - let key = match self.parts[0] { - PartKind::Int => Value::Int(slot[0] as i64), - PartKind::Temporal(lane) => Value::Temporal(lane.value(slot[0] as i64)), - PartKind::Node | PartKind::Str => { - unreachable!("counting mode is one fixed-width word") - } - }; - out.push((vec![key], vec![Acc::Count(slot[1] as i64); self.n_aggs])); + out.push(( + vec![self.counted_key(at)], + vec![Acc::Count(slot[1] as i64); self.n_aggs], + )); } return out; } - let mut accs = self.accs.into_iter(); - let mut scratch = [0u8; INLINE_MAX + 1]; + let mut accs = std::mem::take(&mut self.accs).into_iter(); + let mut vals = Vec::new(); for g in 0..self.groups { - let base = g * self.stride; - let mut vals = Vec::with_capacity(self.parts.len()); - let mut w = 0; - for &part in &self.parts { - vals.push(match part { - PartKind::Int => Value::Int(self.keys[base + w] as i64), - PartKind::Temporal(lane) => { - Value::Temporal(lane.value(self.keys[base + w] as i64)) - } - PartKind::Node => Value::Node { - table: self.keys[base + w] as u32, - offset: self.keys[base + w + 1], - }, - PartKind::Str => { - let (w0, w1) = (self.keys[base + w], self.keys[base + w + 1]); - // The caller checked these bytes for UTF-8 on - // the way in, so the lossy read never replaces - // anything. - let bytes = str_bytes(w0, w1, &self.heap, &mut scratch); - Value::Str(String::from_utf8_lossy(bytes).into_owned()) - } + self.key_into(g, &mut vals); + out.push(( + std::mem::take(&mut vals), + accs.by_ref().take(self.n_aggs).collect(), + )); + } + out + } + + /// Every group as one answer row, key columns ascending, with the + /// keys and the finalized aggregates put back in the order the + /// RETURN clause named them. `item_agg` is that order: true where + /// the item is an aggregate and false where it is a key. + /// + /// One pass rather than a drain and a sort over what the drain + /// built, because the drain built two vectors per group and the + /// sort then compared groups through the `Value`s they decoded to. + /// A hundred thousand groups was three hundred thousand small + /// allocations and a sort chasing a pointer per compare, all of it + /// 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(); + let mut order: Vec = if self.counting { + std::mem::take(&mut self.order) + } else { + (0..self.groups as u32).collect() + }; + if self.counting { + let part = self.parts[0]; + order.sort_unstable_by(|&a, &b| { + cmp_word(part, self.slots[a as usize][0], self.slots[b as usize][0]) + }); + } else { + order.sort_unstable_by(|&a, &b| self.cmp_keys(a as usize, b as usize)); + } + let mut out = Vec::with_capacity(order.len()); + let mut keys = Vec::new(); + for &g in &order { + if self.counting { + keys.push(self.counted_key(g)); + } else { + self.key_into(g as usize, &mut keys); + } + let (mut key_at, mut agg_at) = (0, g as usize * self.n_aggs); + let mut row = Vec::with_capacity(item_agg.len()); + for &is_agg in item_agg { + row.push(if is_agg { + 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) + }; + agg_at += 1; + v + } else { + key_at += 1; + std::mem::replace(&mut keys[key_at - 1], Value::Null) }); - w += part.words(); } - out.push((vals, accs.by_ref().take(self.n_aggs).collect())); + keys.clear(); + out.push(row); } out } + + /// The key of a group counting mode holds in slot `at`, which is + /// the whole group: the slot carries the key rather than an index + /// into the key buffer. + fn counted_key(&self, at: u32) -> Value { + let word = self.slots[at as usize][0] as i64; + match self.parts[0] { + PartKind::Int => Value::Int(word), + PartKind::Temporal(lane) => Value::Temporal(lane.value(word)), + PartKind::Node | PartKind::Str => { + unreachable!("counting mode is one fixed-width word") + } + } + } + + /// The stored key of group `g`, part by part, appended to `out`. + fn key_into(&self, g: usize, out: &mut Vec) { + let base = g * self.stride; + let mut scratch = [0u8; INLINE_MAX + 1]; + let mut w = 0; + for &part in &self.parts { + out.push(match part { + PartKind::Int => Value::Int(self.keys[base + w] as i64), + PartKind::Temporal(lane) => Value::Temporal(lane.value(self.keys[base + w] as i64)), + PartKind::Node => Value::Node { + table: self.keys[base + w] as u32, + offset: self.keys[base + w + 1], + }, + PartKind::Str => { + let (w0, w1) = (self.keys[base + w], self.keys[base + w + 1]); + // The caller checked these bytes for UTF-8 on the + // way in, so the lossy read never replaces anything. + let bytes = str_bytes(w0, w1, &self.heap, &mut scratch); + Value::Str(String::from_utf8_lossy(bytes).into_owned()) + } + }); + w += part.words(); + } + } + + /// Group order: the stored keys compared left to right, which is + /// the order the sink used to reach by decoding both groups and + /// comparing the `Value`s. Reading the words gives the same answer + /// for every kind a key part can be. An integer orders by its word. + /// A temporal lane orders by its word too, because every group in + /// one table shares the lane, so the kind that ranks ahead of the + /// number in the general compare is the same on both sides. A node + /// orders by table and then offset, which is how the two words sit. + /// A string orders by its bytes, which is what comparing the two + /// `String`s came to. + fn cmp_keys(&self, a: usize, b: usize) -> Ordering { + let (mut i, mut j) = (a * self.stride, b * self.stride); + let mut sa = [0u8; INLINE_MAX + 1]; + let mut sb = [0u8; INLINE_MAX + 1]; + for &part in &self.parts { + let ord = match part { + PartKind::Int | PartKind::Temporal(_) => cmp_word(part, self.keys[i], self.keys[j]), + PartKind::Node => (self.keys[i] as u32, self.keys[i + 1]) + .cmp(&(self.keys[j] as u32, self.keys[j + 1])), + PartKind::Str => str_bytes(self.keys[i], self.keys[i + 1], &self.heap, &mut sa) + .cmp(str_bytes( + self.keys[j], + self.keys[j + 1], + &self.heap, + &mut sb, + )), + }; + if ord != Ordering::Equal { + return ord; + } + i += part.words(); + j += part.words(); + } + Ordering::Equal + } +} + +/// Two words of a one word key part, compared as the values they decode +/// to. A date is the one lane narrower than the word it rides in, so it +/// compares through the same narrowing the decode does rather than over +/// the bits above it. +fn cmp_word(part: PartKind, a: u64, b: u64) -> Ordering { + if part == PartKind::Temporal(TemporalLane::Date) { + (a as i32).cmp(&(b as i32)) + } else { + (a as i64).cmp(&(b as i64)) + } } /// Hash of a string key part: eight bytes at a time through the same @@ -850,6 +982,8 @@ fn hash_bytes(b: &[u8]) -> u64 { #[cfg(test)] mod tests { + use zu_query::exec::OrdValue; + use super::*; /// Feeds ints one per row and counts them, the shape the sink drives. @@ -1121,4 +1255,126 @@ mod tests { assert_eq!(rows[0].0[0], Value::Int(0)); assert_eq!(count_of(&rows[0].1[0]), 2); } + + /// The order the drain used to be put in by the sink, so that the + /// word compare can be held to it. + fn by_value(rows: &mut [(Vec, Vec)]) { + rows.sort_by(|a, b| { + a.0.iter() + .zip(&b.0) + .map(|(x, y)| OrdValue(x.clone()).cmp(&OrdValue(y.clone()))) + .find(|o| *o != std::cmp::Ordering::Equal) + .unwrap_or(std::cmp::Ordering::Equal) + }); + } + + #[test] + fn rows_come_out_by_key_ascending_and_in_clause_order() { + let mut t = GroupTable::new(vec![PartKind::Int], 1); + // Negative keys among them, because the words are unsigned and + // the order is not. + count_ints(&mut t, &[5, -3, 5, 0, 12, -3, -3]); + // 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]); + assert_eq!(rows.len(), 4); + assert_eq!( + rows, + vec![ + vec![Value::Int(3), Value::Int(-3)], + vec![Value::Int(1), Value::Int(0)], + vec![Value::Int(2), Value::Int(5)], + vec![Value::Int(1), Value::Int(12)], + ] + ); + } + + #[test] + fn a_counting_table_orders_its_slots_too() { + // Insertion order here is 7, 1, 4, which is neither the slot + // 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]); + assert_eq!( + rows, + vec![ + vec![Value::Int(1), Value::Int(1)], + vec![Value::Int(4), Value::Int(1)], + vec![Value::Int(7), Value::Int(3)], + ] + ); + } + + /// The point of the whole change: ordering over the packed words has + /// to be the ordering over the values they decode to, column by + /// column, for a key that has one of every kind in it. + #[test] + fn the_word_order_is_the_value_order() { + let parts = vec![PartKind::Node, PartKind::Str, PartKind::Int]; + let specs = [AggSpec::CountStar]; + // Both string forms are here on purpose: "a longer string past + // the inline limit" lives in the heap and the rest live in their + // words, and they still have to sort against each other by their + // bytes rather than by where the bytes are. + let keys: [(u64, u64, &str, i64); 8] = [ + (2, 5, "x", 9), + (1, 5, "x", 9), + (1, 5, "x", -8), + (1, 4, "x", 9), + (1, 5, "a longer string past the inline limit", 9), + (1, 5, "y", 9), + (1, 5, "", 9), + (0, 0, "zz", 0), + ]; + let mut t = GroupTable::new(parts.clone(), 1); + let mut batch = KeyBatch::default(); + batch.reset(t.stride(), keys.len()); + for (row, &(table, offset, s, n)) in keys.iter().enumerate() { + let (words, stride) = batch.words_mut(); + words[row * stride] = table; + words[row * stride + 1] = offset; + words[row * stride + 4] = n as u64; + batch.set_str(row, 2, s.as_bytes()); + } + let mut gids = Vec::new(); + t.probe(&batch, &specs, &mut gids); + for &g in &gids { + t.accs_mut()[g as usize].add_star(1); + } + // The same table twice, since one of the two ways of reading it + // consumes it. + let mut same = GroupTable::new(parts, 1); + same.merge_from(&t).unwrap(); + let mut expect = same.drain(); + by_value(&mut expect); + let want: Vec> = expect + .into_iter() + .map(|(mut vals, accs)| { + vals.push(accs.into_iter().next().expect("one count").finalize()); + vals + }) + .collect(); + assert_eq!(t.rows(&[false, false, false, true]), want); + } + + /// A date is the one lane stored narrower than the word it rides in, + /// so a negative one is the case where comparing the raw word and + /// comparing the value part ways. + #[test] + fn dates_before_the_epoch_sort_before_the_ones_after() { + 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 days: Vec = rows + .iter() + .map(|r| match r[0] { + Value::Temporal(zu_common::Temporal::Date(d)) => i64::from(d), + _ => panic!("date key"), + }) + .collect(); + assert_eq!(days, [-400, -5, 0, 10]); + } } diff --git a/crates/zu-exec/src/sink.rs b/crates/zu-exec/src/sink.rs index 2c8cd35c..94086a9a 100644 --- a/crates/zu-exec/src/sink.rs +++ b/crates/zu-exec/src/sink.rs @@ -130,7 +130,7 @@ impl Acc { Ok(()) } - fn finalize(self) -> Value { + pub(crate) fn finalize(self) -> Value { match self { Acc::Count(n) => Value::Int(n), Acc::Sum(acc) => Value::Int(acc.unwrap_or(0)), @@ -375,17 +375,6 @@ fn val_cmp(x: &Value, y: &Value) -> Ordering { } } -/// Group order: the key columns compared left to right. -fn key_cmp(a: &[Value], b: &[Value]) -> Ordering { - for (x, y) in a.iter().zip(b) { - let ord = val_cmp(x, y); - if ord != Ordering::Equal { - return ord; - } - } - Ordering::Equal -} - /// How many of the sorted rows the steps above the sort can still use. /// A SKIP of n under a LIMIT of k needs n + k of them and nothing more, /// which is what turns an ORDER BY under a LIMIT into a selection. @@ -1075,9 +1064,9 @@ fn distinct(mut rows: Vec>) -> Vec> { } /// Merges keyed aggregation partials into the final result: fold the -/// group tables, produce the empty-input row for a bare aggregate, -/// order groups by key ascending like the old BTreeMap sink, and -/// interleave keys and aggregates back into clause order. +/// group tables, produce the empty-input row for a bare aggregate, and +/// hand the folded table the clause order so that it can put the rows +/// out by key ascending, the way the old BTreeMap sink did. pub(crate) fn finish_agg( columns: Vec, item_agg: &[bool], @@ -1118,22 +1107,7 @@ pub(crate) fn finish_agg( "string property is not UTF-8".to_string(), )); } - let mut groups = merged.map(GroupTable::drain).unwrap_or_default(); - groups.sort_by(|a, b| key_cmp(&a.0, &b.0)); - let mut rows = Vec::with_capacity(groups.len()); - for (keyvals, states) in groups { - let mut kit = keyvals.into_iter(); - let mut sit = states.into_iter(); - let mut row = Vec::with_capacity(item_agg.len()); - for &is_agg in item_agg { - row.push(if is_agg { - sit.next().expect("one state per aggregate item").finalize() - } else { - kit.next().expect("one value per key item") - }); - } - rows.push(row); - } + let rows = merged.map(|t| t.rows(item_agg)).unwrap_or_default(); Ok(QueryResult::new(columns, apply_post(post, rows))) } diff --git a/docs/api/model.json b/docs/api/model.json index 4898ff42..87d1637e 100644 --- a/docs/api/model.json +++ b/docs/api/model.json @@ -4650,6 +4650,14 @@ "signature": "Chain(std::sync::Arc)", "doc": "A PMR chain (docs/07 ยง5): the executor-internal form of a\nvariable-length path. [`settle`] turns it into the edge list\nbefore any value leaves the pipeline, so results never hold one." }, + { + "id": "zu::query::Value::Decimal", + "kind": "variant", + "name": "Decimal", + "of": "zu::query::Value", + "signature": "Decimal(zu_common::Decimal)", + "doc": "GV17, an exact decimal: an integer of units and the scale that\nsays how large a unit is.\n\nIt is an arm of its own rather than a float because a tenth is\nnot a binary fraction, so a price held as binary64 is not the\nprice. Two decimals written at different scales are one value\nand print differently, which is why the scale rides on the\nnumber: a decimal in a `RETURN` has no column to ask." + }, { "id": "zu::query::Value::Float", "kind": "variant", @@ -5036,6 +5044,14 @@ "signature": "DayTime", "doc": "Nanoseconds." }, + { + "id": "zu::query::column::ColumnType::Decimal", + "kind": "variant", + "name": "Decimal", + "of": "zu::query::column::ColumnType", + "signature": "Decimal { scale: u16 }", + "doc": "GV17, exact decimals at one scale. The scale is the column's\nand not each value's: a column mixing `1.5` and `1.25` settles\non hundredths, because the wider scale holds both exactly and\nthe narrower one would round a value away." + }, { "id": "zu::query::column::ColumnType::Float", "kind": "variant",