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
74 changes: 67 additions & 7 deletions interactive/src/corgi/chunk.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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,
}
}

Expand Down Expand Up @@ -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();
}
Expand Down Expand Up @@ -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::<PointStamp<u64>, 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;
Expand Down
126 changes: 68 additions & 58 deletions interactive/src/corgi/join.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -134,7 +134,7 @@ impl<T: ColTime> ProxyJoinBackend<T, CBatch<T>, CBatch<T>> for CorgiJoinBackend<
) {
let chunks0 = side_chunks(&instance.batches0);
let chunks1 = side_chunks(&instance.batches1);
let keys0: Vec<Option<&CValue>> = chunks0.iter().map(|c| Some(c.keys())).collect();
let keys0: Vec<Option<&CValue>> = chunks0.iter().map(|c| Some(recover_key_ref(c.keys()))).collect();
let vals0: Vec<Option<&CValue>> = chunks0.iter().map(|c| Some(c.vals())).collect();
let vals1: Vec<Option<&CValue>> = chunks1.iter().map(|c| Some(c.vals())).collect();

Expand Down Expand Up @@ -186,13 +186,13 @@ impl<T: ColTime> ProxyJoinBackend<T, CBatch<T>, CBatch<T>> 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) };
Expand Down Expand Up @@ -253,10 +253,10 @@ fn side_chunks<T: ColTime>(batches: &[CBatch<T>]) -> Vec<&CorgiChunk<T, Diff>> {
/// 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<Vec<&CValue>> {
fn walk<'a>(col: &'a CValue, out: &mut Vec<&'a CValue>) -> bool {
fn leaf_lanes(col: &CValue) -> Option<Vec<&[u64]>> {
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,
Expand All @@ -266,47 +266,33 @@ fn leaf_lanes(col: &CValue) -> Option<Vec<&CValue>> {
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<T: ColTime>(chunks: &[&CorgiChunk<T, Diff>]) -> 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<Vec<u64>> {
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<T, Diff>,
cid: usize,
s: usize,
e: usize,
vals: Option<(&'a [Vec<u64>], 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,
}
}
Expand Down Expand Up @@ -464,16 +450,16 @@ fn block_ends<T: ColTime>(chunks: &[&CorgiChunk<T, Diff>], horizon: Option<u64>)
}

/// 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<T, Diff>,
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<Vec<Vec<u64>>>,
/// The chunk's full borrowed value lanes; `None` when vals are structured.
vals: Option<Vec<&'a [u64]>>,
/// Cursor within `keys`.
cur: usize,
}
Expand All @@ -483,7 +469,7 @@ const PULL: usize = 1 << 14;

impl<'a, T: ColTime> LeafView<'a, T> {
fn new(chunk: &'a CorgiChunk<T, Diff>, cid: usize, start: usize, end: usize, leaf_vals: bool) -> Self {
let vals = (leaf_vals && end > start).then(|| pull_lanes(chunk.vals(), &(start..end).collect::<Vec<_>>()));
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.
Expand All @@ -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<T, Diff>,
cid: usize,
lo: Vec<usize>,
hi: Vec<usize>,
vals: Option<Vec<Vec<u64>>>,
off: Vec<usize>,
vals: Option<Vec<&'a [u64]>>,
}

impl<'a, T: ColTime> Probe<'a, T> {
Expand All @@ -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<usize> = 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<RunRef<'_, T>> {
Expand All @@ -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(),
})
}
}
Expand Down Expand Up @@ -616,9 +590,7 @@ fn stage_collision<T: ColTime>(
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;
}
Expand Down Expand Up @@ -929,6 +901,44 @@ mod tests {
}
}

#[test]
fn borrowed_comparisons_match_structural_rows_after_skips() {
let make = |keys: Vec<u64>, 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;
Expand Down
Loading