diff --git a/interactive/src/corgi/reduce.rs b/interactive/src/corgi/reduce.rs index ff539a2c4..618c4e1f8 100644 --- a/interactive/src/corgi/reduce.rs +++ b/interactive/src/corgi/reduce.rs @@ -5,7 +5,9 @@ //! //! * ids — `key_hash`/`value_id` are value-as-id for primitive columns (the value IS the id) and //! the canonical native `corgi::hash` for compound columns (columnar, content-addressed, so ids -//! coincide across the output→input boundary); DD never hashes. +//! coincide across the output→input boundary); DD never hashes. The exception is the input under +//! primitive keys, whose compound values get ordinals from a merge in value order +//! (`present_input_merged`). //! * the value callback — `reduce_many` runs ONE crossing per retire over every `(key, time)` //! bracket, building the output value COLUMNS directly (Count → a `u64` prim, Distinct → a //! `Unit`, Min → the chosen input rows, Collect → a `List`), never through DDIR rows. @@ -28,7 +30,7 @@ use std::collections::HashMap; use std::hash::{BuildHasherDefault, Hasher}; use std::rc::Rc; -use differential_dataflow::consolidation::consolidate_updates; +use differential_dataflow::consolidation::{consolidate, consolidate_updates}; use differential_dataflow::trace::Description; use differential_dataflow::trace::chunk::ChunkBatch; use differential_dataflow::operators::int_proxy::{KeyPosition, ProxyBridge}; @@ -39,7 +41,7 @@ use corgi::{ArithOp, Bounds, NumOp, OpLike, Value as CValue}; use crate::corgi::col_times::{ColTime, ColTimes}; use crate::corgi::search::MatchingRanges; -use crate::corgi::chunk::{columns_to_batch, key_ids, key_lane, CorgiChunk}; +use crate::corgi::chunk::{columns_to_batch, key_ids, key_is_hashed, key_lane, CorgiChunk}; use crate::ir::{Diff, Reducer}; type CBatch = Rc>>; @@ -137,6 +139,10 @@ struct IdPool { blocks: Vec, index: IdMap, len: usize, + /// The columns ordinal IDs name rows of, when this pool resolves ordinals. + columns: Option>, + /// Ordinal `i` is row `rows[i].1` of `columns[rows[i].0]`. + rows: Vec<(u32, u32)>, } impl IdPool { fn clear(&mut self) { @@ -144,6 +150,8 @@ impl IdPool { self.blocks.clear(); self.index.clear(); self.len = 0; + self.columns = None; + self.rows.clear(); } fn register(&mut self, col: CValue, ids: &[u64]) { if col.len() == 0 { return; } @@ -157,7 +165,18 @@ impl IdPool { self.blocks.push(col); } } + /// Mint the next ordinal ID, for row `row` of `columns[column]`. + fn mint(&mut self, column: usize, row: usize) -> u64 { + self.rows.push((column as u32, row as u32)); + (self.rows.len() - 1) as u64 + } fn gather(&self, ids: &[u64]) -> CValue { + if let Some(columns) = &self.columns { + let srcs: Vec> = columns.iter().map(Some).collect(); + let tags: Vec = ids.iter().map(|&id| self.rows[id as usize].0 as usize).collect(); + let offs: Vec = ids.iter().map(|&id| self.rows[id as usize].1 as usize).collect(); + return gather_lanes(&srcs, &tags, &offs); + } if let Some(depth) = self.depth { let mut col = CValue::u64(ids.to_vec()); for _ in 0..depth { col = CValue::Prod(vec![col]); } @@ -238,6 +257,29 @@ fn ids(col: &CValue) -> Vec { corgi::hash(col) } +/// Borrow a column's `u64` leaves in field order when it is a (possibly nested) product of them, +/// a `Unit` contributing none. Row order is then lexicographic over the leaves, which is corgi's +/// structural order for such a column. Returns `false` for any other shape. +fn leaf_fields<'a>(col: &'a CValue, into: &mut Vec<&'a [u64]>) -> bool { + match col { + CValue::Prod(fields) => fields.iter().all(|field| leaf_fields(field, into)), + CValue::Unit(_) => true, + other => match corgi::arrange::leaf_slice(other) { + Some(leaf) => { into.push(leaf); true } + None => false, + }, + } +} + +/// Each chunk's `(key index, rows)` matches of the ascending `keys`, in key order. +fn search(chunks: &[&CorgiChunk], keys: &[u64]) -> Vec)>> { + chunks.iter().map(|chunk| { + if chunk.diffs().is_empty() { return Vec::new(); } + let lane = corgi::arrange::leaf_slice(key_lane(chunk.keys())).expect("the identifier lane is a u64 leaf"); + MatchingRanges::new(keys, lane).collect() + }).collect() +} + /// Concatenate the records of the `changed` keys across a run of chunks into parallel /// `(keys_col, vals_col)` corgi columns plus per-record `(key_hash, time, diff)`. `changed` is the /// ASCENDING set of changed key ids; a row is kept iff its key id is in it. @@ -408,6 +450,81 @@ where } } + /// Present the merged input run as `present_input` does, but by merging each key's runs across + /// the chunks in value order: each distinct value is met once with all of its records, which + /// consolidate there, and gets its id there. A primitive value is its own id; any other gets + /// the next ordinal, which resolves to the chunk row where the merge met it. Ids ascend with + /// values, so the bridge is built in order, with no hash, map, or gather. + /// + /// Returns `false`, having done nothing, for keys with a hash lane (one id may cover several + /// keys) or for values that are not products of `u64` leaves. + fn present_input_merged( + &mut self, + chunks: &[&CorgiChunk], + keys: &[u64], + bridge: &mut ProxyBridge, + ) -> bool { + let Some(first) = chunks.iter().find(|chunk| !chunk.diffs().is_empty()) else { return true }; + if chunks.iter().any(|chunk| !chunk.diffs().is_empty() && key_is_hashed(chunk.keys())) { + return false; + } + let mut leaves: Vec> = Vec::with_capacity(chunks.len()); + for chunk in chunks.iter() { + let mut fields = Vec::new(); + if !chunk.diffs().is_empty() && !leaf_fields(chunk.vals(), &mut fields) { + return false; + } + leaves.push(fields); + } + let width = leaves.iter().map(Vec::len).max().unwrap_or(0); + let primitive = corgi::arrange::leaf_slice(first.vals()).is_some(); + // Primitive columns register only their shape. + self.keys.register(first.keys().clone(), &[]); + if primitive { + self.input.register(first.vals().clone(), &[]); + } else { + self.input.columns = Some(chunks.iter().map(|chunk| chunk.vals().clone()).collect()); + } + + // Row `a` of one chunk against row `b` of another, lexicographically over value leaves. + let order = |(ca, ra): (usize, usize), (cb, rb): (usize, usize)| { + (0..width).map(|f| leaves[ca][f][ra].cmp(&leaves[cb][f][rb])).find(|o| o.is_ne()).unwrap_or(std::cmp::Ordering::Equal) + }; + let matches = search(chunks, keys); + // At most one record per matched row, as in `merge_present`. + bridge.reserve(matches.iter().flatten().map(|(_, rows)| rows.len()).sum()); + let mut cursors = vec![0; chunks.len()]; + // A key's runs yet to merge, as `(chunk, next row, end)`, and one value's `(time, diff)`s. + let (mut heads, mut entries) = (Vec::new(), Vec::new()); + for (index, &key) in keys.iter().enumerate() { + for (chunk, list) in matches.iter().enumerate() { + if let Some((_, rows)) = list.get(cursors[chunk]).filter(|(matched, _)| *matched == index) { + heads.push((chunk, rows.start, rows.end)); + cursors[chunk] += 1; + } + } + // The least value at any head, and every head's run of it. + while let Some(at) = heads.iter().map(|&(chunk, row, _)| (chunk, row)).min_by(|a, b| order(*a, *b)) { + let mut h = 0; + while h < heads.len() { + let (chunk, row, end) = heads[h]; + if order((chunk, row), at).is_ne() { h += 1; continue; } + let stop = (row + 1..end).find(|&r| order((chunk, r), at).is_ne()).unwrap_or(end); + debug_assert!(stop == end || order((chunk, stop), at).is_gt(), "a key's rows ascend by value"); + let (times, diffs) = (chunks[chunk].times(), chunks[chunk].diffs()); + entries.extend((row..stop).map(|r| (times.get(r), diffs[r]))); + if stop < end { heads[h].1 = stop; h += 1; } else { heads.swap_remove(h); } + } + consolidate(&mut entries); + if !entries.is_empty() { + let id = if primitive { leaves[at.0][0][at.1] } else { self.input.mint(at.0, at.1) }; + bridge.extend(entries.drain(..).map(|(time, diff)| ((key, id), time, diff))); + } + } + } + true + } + /// The one value crossing for a retire: every `(key, time)` bracket at once. Builds the output /// value COLUMN directly per reducer, registers it (id → row) into the val pool, and returns the /// proxy `(value_id, diff)` deltas with per-bracket ends. `input[k] = (value_id, accumulated diff)`; the bracket `i` is `input[ends[i-1]..ends[i]]`, non-empty. @@ -645,12 +762,14 @@ where } // ONE merged input presentation: novel and prior together, netted by the consolidation — - // equal values share a content-hash id, so an exactly cancelling pair vanishes here, and + // equal values share an id, so an exactly cancelling pair vanishes here, and // its time survives in `window.seeds` above. The input pool resolves values // needed by Min and Collect. let mut in_chunks = chunks_of(instance.source_batches); in_chunks.extend(novel_chunks.iter().copied()); - self.present_input(&in_chunks, &keys, &mut window.input); + if !self.present_input_merged(&in_chunks, &keys, &mut window.input) { + self.present_input(&in_chunks, &keys, &mut window.input); + } // Output-history presentation, same keys (register keys + values for correction resolution). let (o_keys, o_vals, o_khs, mut o_times, o_diffs, o_run_ends) = collect_present(&chunks_of(instance.output_batches), &keys); @@ -846,6 +965,72 @@ mod tests { assert_eq!(bridge, vec![((2, vids[0]), 0, 5), ((3, vids[0]), 0, 2)]); } + /// The merged input presentation agrees with the content-hash one up to the naming of values, + /// in order and with one id per value, across value shapes, chunkings, and cancellations. + #[test] + fn merged_input_presentation_matches_hashed() { + use std::collections::BTreeMap; + use differential_dataflow::dynamic::pointstamp::PointStamp; + use crate::ir::Time; + let stamp = |outer, coords: &[u64]| Time::new(outer, PointStamp::new(coords.iter().copied().collect())); + let times = [stamp(0, &[]), stamp(1, &[2, 3]), stamp(2, &[1, 4]), stamp(3, &[2])]; + let mut state = 0x9E37_79B9_7F4A_7C15u64; + let mut next = |n: u64| { + state = state.wrapping_mul(6364136223846793005).wrapping_add(1442695040888963407); + (state >> 33) % n + }; + let rows: Vec<_> = (0..300).map(|_| (next(6), (next(3), next(3)), times[next(4) as usize].clone(), [1, -1, 2][next(3) as usize])).collect(); + let keys = [0, 2, 3, 5, 9]; + // Each shape's value column, and the leaves a row of it compares by. + type Shape = (fn(&[(u64, u64)]) -> CValue, fn((u64, u64)) -> Vec); + let shapes: [Shape; 3] = [ + (|v| CValue::u64(v.iter().map(|v| v.0).collect()), |v| vec![v.0]), + (|v| CValue::Prod(vec![CValue::u64(v.iter().map(|v| v.0).collect()), CValue::u64(v.iter().map(|v| v.1).collect())]), |v| vec![v.0, v.1]), + (|v| CValue::Unit(v.len()), |_| vec![]), + ]; + for (column, leaves) in shapes { + let mut expected = BTreeMap::new(); + for (key, val, time, diff) in rows.iter().filter(|r| keys.contains(&r.0)) { + *expected.entry(((*key, leaves(*val)), time.clone())).or_insert(0) += diff; + } + expected.retain(|_, diff| *diff != 0); + for size in [1, 7, rows.len()] { + let chunks: Vec<_> = rows.chunks(size).chain([&rows[..0]]).map(|rows| CorgiChunk::from_columns( + CValue::u64(rows.iter().map(|r| r.0).collect()), + column(&rows.iter().map(|r| r.1).collect::>()), + rows.iter().map(|r| r.2.clone()).collect(), + rows.iter().map(|r| r.3).collect(), + )).collect(); + let chunks: Vec<_> = chunks.iter().collect(); + let (mut merged, mut hashed) = (CorgiReduceBackend::