diff --git a/interactive/src/corgi/chunk.rs b/interactive/src/corgi/chunk.rs index feb30441b..f008e79c8 100644 --- a/interactive/src/corgi/chunk.rs +++ b/interactive/src/corgi/chunk.rs @@ -237,9 +237,14 @@ where } } + // The survey can be large. Its last use precedes materializing the + // output columns, so do not retain it across that allocation peak. + drop(runs); if times.len() * 2 < n1 + n2 { times.shrink_to_fit(); } let srcs = [Some(&kv1), Some(&kv2)]; Self::emit(&srcs, &tags, &offs, times, diffs, out); + drop(tags); + drop(offs); // Push back the survivor's unconsumed suffix (all `>` the horizon), ahead of its deque. if p1 < n1 { @@ -578,12 +583,18 @@ pub fn key_is_hashed(keys: &CValue) -> bool { corgi::arrange::leaf_slice(keys).is_none() } -/// Undo [`present_key`]: the key as the rest of the system knows it. A corgi clone is an `Arc` -/// bump, so dropping the hash lane costs nothing. +/// Undo [`present_key`]: the key as the rest of the system knows it. +/// Cloning the column shares its primitive buffers. pub fn recover_key(keys: &CValue) -> CValue { + recover_key_ref(keys).clone() +} + +/// Borrow the declared key without its arrangement-only identifier lane. +/// Use before gathering a projection input to avoid copying a discarded hash column. +pub fn recover_key_ref(keys: &CValue) -> &CValue { match keys { - CValue::Prod(cols) if corgi::arrange::leaf_slice(keys).is_none() => cols[1].clone(), - _ => keys.clone(), + CValue::Prod(cols) if corgi::arrange::leaf_slice(keys).is_none() => &cols[1], + _ => keys, } } @@ -635,9 +646,16 @@ where } self.k_blocks.push(std::mem::replace(&mut c.keys, CValue::Unit(0))); self.v_blocks.push(std::mem::replace(&mut c.vals, CValue::Unit(0))); - self.times.push_range(&c.times, 0, c.times.len()); - c.times.clear(); - self.diffs.append(&mut c.diffs); + if self.times.is_empty() { + // The first block can donate its owned lanes. Return our empty + // buffers to the sender; later blocks append to the donated ones. + std::mem::swap(&mut self.times, &mut c.times); + std::mem::swap(&mut self.diffs, &mut c.diffs); + } else { + self.times.push_range(&c.times, 0, c.times.len()); + c.times.clear(); + self.diffs.append(&mut c.diffs); + } if self.times.len() >= INGEST { self.flush(); } @@ -726,6 +744,48 @@ mod test { } } + #[test] + fn chunker_owned_lanes_match_scalar_across_flushes() { + use differential_dataflow::dynamic::pointstamp::PointStamp; + use timely::container::{ContainerBuilder, PushInto}; + use crate::corgi::container::CorgiContainer; + let mut builder = CorgiChunker::, i64>::default(); + let mut source = CorgiContainer::default(); + let mut seed = 731; + for blocks in [1, 3, 1, 5] { + let mut expected = BTreeMap::new(); + // Empty containers do not consume the first-block opportunity. + builder.push_into(&mut source); + for block in 0..blocks { + let mut keys = Vec::new(); + for row in 0..64 { + let key = xorshift(&mut seed) % 8; + let time = PointStamp::new((0..(block + row) % 4) + .map(|_| xorshift(&mut seed) % 3).collect()); + let diff = (xorshift(&mut seed) % 5) as i64 - 2; + keys.push(key); + source.times.push(&time); + source.diffs.push(diff); + *expected.entry((key, time)).or_insert(0i64) += diff; + } + source.keys = CValue::u64(keys); + source.vals = CValue::Unit(64); + let donated = source.diffs.as_ptr(); + builder.push_into(&mut source); + if block == 0 { assert_eq!(builder.diffs.as_ptr(), donated); } + assert!(source.times.is_empty() && source.diffs.is_empty()); + } + expected.retain(|_, diff| *diff != 0); + let result = builder.finish().unwrap(); + let keys = corgi::arrange::leaf_slice(result.keys()).unwrap(); + let actual: BTreeMap<_, _> = (0..result.len_()).map(|i| + ((keys[i], result.times().get(i)), result.diffs()[i])).collect(); + assert_eq!(actual.len(), result.len_(), "duplicate output triples"); + assert_eq!(actual, expected); + assert!(builder.finish().is_none()); + } + } + #[test] fn cancelled_merge_does_not_retain_input_sized_diff_storage() { let mut retained_capacity = None; diff --git a/interactive/src/corgi/join.rs b/interactive/src/corgi/join.rs index fdaa222a7..35ca49f41 100644 --- a/interactive/src/corgi/join.rs +++ b/interactive/src/corgi/join.rs @@ -37,11 +37,11 @@ use differential_dataflow::operators::int_proxy::{KeyPosition, JoinInstance, Pro use differential_dataflow::operators::int_proxy::join::JoinMatches; use differential_dataflow::trace::chunk::{Chunk, ChunkBatch}; -use corgi::arrange::{compare_at, gather, gather_lanes, leaf_slice}; +use corgi::arrange::{compare_at, gather_lanes, leaf_slice}; use crate::corgi::search::matching_ranges; -use corgi::{shape_of_value, Shape, Value as CValue}; +use corgi::{shape_of_value, Value as CValue}; -use crate::corgi::chunk::{key_is_hashed, key_lane, recover_key, CorgiChunk}; +use crate::corgi::chunk::{key_is_hashed, key_lane, recover_key_ref, CorgiChunk}; use crate::corgi::col_times::ColTime; use crate::corgi::container::CorgiContainer; use crate::corgi::logic::compile_join_projection; @@ -134,7 +134,7 @@ impl ProxyJoinBackend, CBatch> for CorgiJoinBackend< ) { let chunks0 = side_chunks(&instance.batches0); let chunks1 = side_chunks(&instance.batches1); - let keys0: Vec> = chunks0.iter().map(|c| Some(c.keys())).collect(); + let keys0: Vec> = chunks0.iter().map(|c| Some(recover_key_ref(c.keys()))).collect(); let vals0: Vec> = chunks0.iter().map(|c| Some(c.vals())).collect(); let vals1: Vec> = chunks1.iter().map(|c| Some(c.vals())).collect(); @@ -186,13 +186,13 @@ impl ProxyJoinBackend, CBatch> for CorgiJoinBackend< start = end; continue; } - // The join's projection is written against the key the program declared, so drop - // the arrangement's leading identifier lane before evaluating it. The output goes to - // an arrange, which re-derives the identifier for the new key. + // Gather only the declared key borrowed above: the arrangement's leading + // identifier is not part of the projection input. A downstream arrange derives + // the identifier for the projected key. let ids = &matches.ids[start..end]; let kc = if primitive_keys { from_ids(chunks0[0].keys(), ids.iter().map(|x| x.0).collect()) - } else { recover_key(&gather_lanes(&keys0, &tag0, &off0)) }; + } else { gather_lanes(&keys0, &tag0, &off0) }; let v0 = if primitive0 { from_ids(chunks0[0].vals(), ids.iter().map(|x| x.1.0).collect()) } else { gather_lanes(&vals0, &tag0, &off0) }; @@ -253,10 +253,10 @@ fn side_chunks(batches: &[CBatch]) -> Vec<&CorgiChunk> { /// product of 64-bit leaves — the shape DDIR tuples transcode to — so row order is the /// lexicographic order of the lane tuples. `Sum`/`List`/narrow leaves give `None` /// (structural compares). -fn leaf_lanes(col: &CValue) -> Option> { - fn walk<'a>(col: &'a CValue, out: &mut Vec<&'a CValue>) -> bool { +fn leaf_lanes(col: &CValue) -> Option> { + fn walk<'a>(col: &'a CValue, out: &mut Vec<&'a [u64]>) -> bool { match col { - CValue::Prim(_) => { out.push(col); matches!(shape_of_value(col), Shape::Prim(64)) } + CValue::Prim(_) => match leaf_slice(col) { Some(lane) => { out.push(lane); true }, None => false }, CValue::Prod(fields) => fields.iter().all(|f| walk(f, out)), CValue::Unit(_) => true, _ => false, @@ -266,47 +266,33 @@ fn leaf_lanes(col: &CValue) -> Option> { if walk(col, &mut out) { Some(out) } else { None } } -/// Whether to copy comparison lanes. Primitive tokens read their source lane directly. +/// Whether to borrow comparison lanes. Primitive tokens read their source lane directly. fn leaf_valued(chunks: &[&CorgiChunk]) -> bool { !primitive_values(chunks[0]) && chunks.iter().filter(|c| c.len() > 0).all(|c| leaf_lanes(c.vals()).is_some()) } -/// Pull rows `idx` of the column's leaf lanes as `u64` buffers. -fn pull_lanes(col: &CValue, idx: &[usize]) -> Vec> { - leaf_lanes(col).expect("pull_lanes: leaf-laned column") - .into_iter() - .map(|lane| gather(lane, idx).into_u64("corgi join lane pull").unwrap()) - .collect() -} - /// One key's records in one chunk: absolute rows `[s, e)`. When the vals are leaf-laned, -/// `vals = (lanes, pos)` gives row `r`'s tuple as `lanes[.][pos + r - s]`; `None` falls -/// back to structural compares against the chunk itself. +/// `vals` borrows the chunk's full leaf lanes, indexed by absolute row; `None` +/// falls back to structural compares against the chunk itself. struct RunRef<'a, T: ColTime> { chunk: &'a CorgiChunk, cid: usize, s: usize, e: usize, - vals: Option<(&'a [Vec], usize)>, + vals: Option<&'a [&'a [u64]]>, } impl<'a, T: ColTime> RunRef<'a, T> { fn val_less(&self, row: usize, other: &Self, orow: usize) -> bool { match (self.vals, other.vals) { - (Some((a, ap)), Some((b, bp))) => { - let (i, j) = (ap + row - self.s, bp + orow - other.s); - a.iter().zip(b).map(|(la, lb)| (la[i], lb[j])).find(|(x, y)| x != y).is_some_and(|(x, y)| x < y) - } + (Some(a), Some(b)) => a.iter().zip(b).map(|(la, lb)| (la[row], lb[orow])).find(|(x, y)| x != y).is_some_and(|(x, y)| x < y), _ => compare_at(self.chunk.vals(), row, other.chunk.vals(), orow) == Ordering::Less, } } fn val_eq(&self, row: usize, other: &Self, orow: usize) -> bool { match (self.vals, other.vals) { - (Some((a, ap)), Some((b, bp))) => { - let (i, j) = (ap + row - self.s, bp + orow - other.s); - a.iter().zip(b).all(|(la, lb)| la[i] == lb[j]) - } + (Some(a), Some(b)) => a.iter().zip(b).all(|(la, lb)| la[row] == lb[orow]), _ => compare_at(self.chunk.vals(), row, other.chunk.vals(), orow) == Ordering::Equal, } } @@ -464,16 +450,16 @@ fn block_ends(chunks: &[&CorgiChunk], horizon: Option) } /// View over one leaf-keyed chunk's rows for THIS block, `[base, end)`. The identifiers are -/// borrowed from the chunk's own lane — the block reads them, it does not copy them — and the -/// vals are gathered once, over exactly those rows. +/// borrowed from the chunk's own lane. Value comparison lanes also borrow the chunk, +/// so selecting a block copies neither identifiers nor values. struct LeafView<'a, T: ColTime> { chunk: &'a CorgiChunk, cid: usize, /// Absolute row of `keys[0]`. base: usize, keys: &'a [u64], - /// Leaf-laned vals over the same rows; `None` when vals are structured. - vals: Option>>, + /// The chunk's full borrowed value lanes; `None` when vals are structured. + vals: Option>, /// Cursor within `keys`. cur: usize, } @@ -483,7 +469,7 @@ const PULL: usize = 1 << 14; impl<'a, T: ColTime> LeafView<'a, T> { fn new(chunk: &'a CorgiChunk, cid: usize, start: usize, end: usize, leaf_vals: bool) -> Self { - let vals = (leaf_vals && end > start).then(|| pull_lanes(chunk.vals(), &(start..end).collect::>())); + let vals = (leaf_vals && end > start).then(|| leaf_lanes(chunk.vals()).unwrap()); LeafView { chunk, cid, base: start, keys: &ident(chunk)[start..end], vals, cur: 0 } } /// The key under the cursor, if any remains in this block. @@ -505,21 +491,20 @@ impl<'a, T: ColTime> LeafView<'a, T> { cid: self.cid, s, e, - vals: self.vals.as_ref().map(|lanes| (&lanes[..], s - self.base)), + vals: self.vals.as_deref(), } } } /// The batched probe of one PROBEE-side chunk: per driver key, its equal-range in the -/// chunk, with the matched rows' vals gathered once as a `u64` buffer -/// when leaf-shaped (`off` gives each key's slice within it). +/// chunk, borrowing its full value lanes when leaf-shaped. Matched rows index the +/// original lanes directly, without a gather or an index buffer. struct Probe<'a, T: ColTime> { chunk: &'a CorgiChunk, cid: usize, lo: Vec, hi: Vec, - vals: Option>>, - off: Vec, + vals: Option>, } impl<'a, T: ColTime> Probe<'a, T> { @@ -532,19 +517,8 @@ impl<'a, T: ColTime> Probe<'a, T> { lo[j] = range.start; hi[j] = range.end; } - let mut off = Vec::with_capacity(lo.len() + 1); - let mut idx: Vec = Vec::new(); - off.push(0); - for i in 0..lo.len() { - idx.extend(lo[i]..hi[i]); - off.push(idx.len()); - } - let vals = if leaf_vals && !idx.is_empty() { - Some(pull_lanes(chunk.vals(), &idx)) - } else { - None - }; - Probe { chunk, cid, lo, hi, vals, off } + let vals = leaf_vals.then(|| leaf_lanes(chunk.vals()).unwrap()); + Probe { chunk, cid, lo, hi, vals } } /// The run of driver key `j` in this chunk, if any. fn run_ref(&self, j: usize) -> Option> { @@ -555,7 +529,7 @@ impl<'a, T: ColTime> Probe<'a, T> { cid: self.cid, s, e, - vals: self.vals.as_ref().map(|lanes| (&lanes[..], self.off[j])), + vals: self.vals.as_deref(), }) } } @@ -616,9 +590,7 @@ fn stage_collision( compare_at(run.chunk.keys(), candidate, reference, row) != Ordering::Equal }) .unwrap_or(run.e - start); - let vals = run - .vals - .map(|(lanes, offset)| (lanes, offset + start - run.s)); + let vals = run.vals; equal_runs.push(RunRef { chunk: run.chunk, cid: run.cid, s: start, e: end, vals }); positions[index] = end; } @@ -929,6 +901,44 @@ mod tests { } } + #[test] + fn borrowed_comparisons_match_structural_rows_after_skips() { + let make = |keys: Vec, salt: u64| { + let n = keys.len(); + CorgiChunk::from_columns( + CValue::u64(keys), + CValue::Prod(vec![ + CValue::u64((0..n).map(|i| (i as u64 ^ salt) % 3).collect()), + CValue::Prod(vec![CValue::Unit(n), CValue::u64((0..n).map(|i| u64::MAX - i as u64).collect())]), + ]), + (0..n).map(|_| 0u64).collect(), vec![1; n], + ) + }; + let left = make(vec![0, 0, 1, 2, 2, 2, 4, 4, 5], 1); + let right = make(vec![0, 1, 1, 4, 5], 2); + let mut view = LeafView::new(&left, 0, 2, 8, true); + assert_eq!(view.take_run(1), Some((2, 3))); + assert_eq!(view.take_run(2), Some((3, 6))); + assert_eq!(view.take_run(4), Some((6, 8))); + assert_eq!(view.cur_key(), None); + let probe = Probe::new(&right, 1, &[0, 1, 3, 5], true); + assert!(probe.run_ref(2).is_none()); + for (s, e) in [(2, 3), (3, 6), (6, 8)] { + let a = view.run_ref(s, e); + for j in [0, 1, 3] { + let b = probe.run_ref(j).unwrap(); + for i in a.s..a.e { + for k in b.s..b.e { + let expected = compare_at(a.chunk.vals(), i, b.chunk.vals(), k); + assert_eq!(a.val_eq(i, &b, k), expected == Ordering::Equal); + assert_eq!(a.val_less(i, &b, k), expected == Ordering::Less); + assert_eq!(b.val_less(k, &a, i), expected == Ordering::Greater); + } + } + } + } + } + #[test] fn primitive_tokens_consolidate_after_product_time_reordering() { use timely::order::Product;