From e89a427ac38e280c4a70786b674bd560f0579b53 Mon Sep 17 00:00:00 2001 From: Frank McSherry Date: Sun, 27 Sep 2026 02:54:26 -0400 Subject: [PATCH 1/4] Borrow join comparison lanes instead of gathering copies The join's comparison views gathered temporary key columns before comparing them. Borrow the chunk's u64 lanes with absolute row indices instead; the owning chunks outlive the views. The structural fallback and the output gathers are unchanged. Adds an oracle against structural row comparison for sparse and skipped ranges. Co-Authored-By: Astra (OpenAI Codex) Co-Authored-By: Claude Opus 5.5 (1M context) --- interactive/src/corgi/join.rs | 114 ++++++++++++++++++---------------- 1 file changed, 62 insertions(+), 52 deletions(-) diff --git a/interactive/src/corgi/join.rs b/interactive/src/corgi/join.rs index fdaa222a7..07d841c36 100644 --- a/interactive/src/corgi/join.rs +++ b/interactive/src/corgi/join.rs @@ -37,9 +37,9 @@ 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::col_times::ColTime; @@ -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; From 0ddcafb81c8f1a49db1bc974d1bb640c19c8eb6c Mon Sep 17 00:00:00 2001 From: Frank McSherry Date: Sun, 27 Sep 2026 06:54:54 -0400 Subject: [PATCH 2/4] Join output gathers skip the hash column they would discard A compound key's arrangement carries an identifier (hash) column that join output projections drop. Borrow the declared key before gathering output instead of gathering the presented key and discarding its hash. Presented keys stay for matching and collision checks. Co-Authored-By: Astra (OpenAI Codex) Co-Authored-By: Claude Opus 5.5 (1M context) --- interactive/src/corgi/chunk.rs | 14 ++++++++++---- interactive/src/corgi/join.rs | 12 ++++++------ 2 files changed, 16 insertions(+), 10 deletions(-) diff --git a/interactive/src/corgi/chunk.rs b/interactive/src/corgi/chunk.rs index feb30441b..57764b893 100644 --- a/interactive/src/corgi/chunk.rs +++ b/interactive/src/corgi/chunk.rs @@ -578,12 +578,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, } } diff --git a/interactive/src/corgi/join.rs b/interactive/src/corgi/join.rs index 07d841c36..35ca49f41 100644 --- a/interactive/src/corgi/join.rs +++ b/interactive/src/corgi/join.rs @@ -41,7 +41,7 @@ use corgi::arrange::{compare_at, gather_lanes, leaf_slice}; use crate::corgi::search::matching_ranges; 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) }; From 795a11cb837d9c68656ba0b94b95850d14fff8d7 Mon Sep 17 00:00:00 2001 From: Frank McSherry Date: Sat, 26 Sep 2026 21:05:52 -0400 Subject: [PATCH 3/4] Move the first arrangement input time/diff buffers into the chunker The chunker already took ownership of payload columns, but copied times and diffs even when starting an empty batch. Swap the first input's owned buffers into the empty accumulator and hand its empty buffers back to the sender; later blocks append as before. Adds an oracle over single and multiple block flushes, mixed PointStamp widths, and cancellation. Co-Authored-By: Astra (OpenAI Codex) Co-Authored-By: Claude Opus 5.5 (1M context) --- interactive/src/corgi/chunk.rs | 55 ++++++++++++++++++++++++++++++++-- 1 file changed, 52 insertions(+), 3 deletions(-) diff --git a/interactive/src/corgi/chunk.rs b/interactive/src/corgi/chunk.rs index 57764b893..1c177056a 100644 --- a/interactive/src/corgi/chunk.rs +++ b/interactive/src/corgi/chunk.rs @@ -641,9 +641,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(); } @@ -732,6 +739,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; From 0ba24980fbee92dd70e195707a50fbeae49d0415 Mon Sep 17 00:00:00 2001 From: Frank McSherry Date: Sat, 26 Sep 2026 22:04:19 -0400 Subject: [PATCH 4/4] Release merge planning buffers before materializing output columns The survey result is not read once output begins, and the source/offset indices are not read once the surviving suffix is gathered. Drop them at their last use rather than at the end of the merge, so they no longer overlap the output's allocations. Output and API unchanged. Co-Authored-By: Astra (OpenAI Codex) Co-Authored-By: Claude Opus 5.5 (1M context) --- interactive/src/corgi/chunk.rs | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/interactive/src/corgi/chunk.rs b/interactive/src/corgi/chunk.rs index 1c177056a..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 {