Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 18 additions & 0 deletions bench/budgets.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
96 changes: 85 additions & 11 deletions crates/zu-exec/src/group.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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<Vec<Value>> {
// 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<Value> = states.into_iter().map(Acc::finalize).collect();
pub(crate) fn rows(mut self, item_agg: &[bool], hands: usize) -> Vec<Vec<Value>> {
let mut order: Vec<u32> = if self.counting {
std::mem::take(&mut self.order)
} else {
Expand All @@ -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<Mutex<Option<Vec<Vec<Value>>>>> =
slices[1..].iter().map(|_| Mutex::new(None)).collect();
let mine = {
let me = &self;
let jobs: Vec<Box<dyn FnOnce() + Send + '_>> = slices[1..]
.iter()
.zip(&slots)
.map(|(slice, slot)| {
Box::new(move || {
*slot.lock().unwrap() = Some(me.build(slice, item_agg));
}) as Box<dyn FnOnce() + Send + '_>
})
.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<Vec<Value>> {
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 {
Expand All @@ -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
Expand Down Expand Up @@ -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,
Expand All @@ -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![
Expand Down Expand Up @@ -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<i64> = (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,
Expand All @@ -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<i64> = rows
.iter()
.map(|r| match r[0] {
Expand Down
88 changes: 72 additions & 16 deletions crates/zu-exec/src/sink.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Expand Down Expand Up @@ -1075,6 +1076,9 @@ pub(crate) fn finish_agg(
partials: Vec<SinkState>,
keys_empty: bool,
) -> Result<QueryResult> {
// 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.
Expand All @@ -1091,14 +1095,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<GroupTable> = 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.
Expand All @@ -1107,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)))
}

Expand All @@ -1118,15 +1115,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<SinkState>) -> Result<i64> {
let mut merged: Option<GroupTable> = 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<SinkState>) -> Result<Option<GroupTable>> {
let mut tables: Vec<GroupTable> = 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<Mutex<Option<Result<GroupTable>>>> =
pairs.iter().map(|_| Mutex::new(None)).collect();
let here = {
let jobs: Vec<Box<dyn FnOnce() + Send + '_>> = 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<dyn FnOnce() + Send + '_>
})
.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.
Expand Down
Loading
Loading