diff --git a/interactive/src/corgi/join.rs b/interactive/src/corgi/join.rs index cbb6e33b2..fdaa222a7 100644 --- a/interactive/src/corgi/join.rs +++ b/interactive/src/corgi/join.rs @@ -355,9 +355,9 @@ impl SideScratch { return; } if let [run] = runs { - // Single-chunk run: values grouped, times ascending within a value; advanced times - // stay ascending (join is monotone), so consolidation is adjacent. Coordinates of - // successive value runs ascend, so entries are born sorted. + // Single-chunk run: values are grouped and coordinates of successive value + // runs ascend. Joining product times need not preserve their total order, + // so adjacent consolidation suffices only if advanced times stay ordered. let coord_hi = (run.cid as u64) << COORD_BITS; let mut v_start = run.s; for row in run.s..run.e { @@ -366,6 +366,9 @@ impl SideScratch { } self.push(coord_hi | v_start as u64, run.chunk.times().get(row).join(lower), run.chunk.diffs()[row]); } + if self.entries.windows(2).any(|w| w[0].0 == w[1].0 && w[0].1 > w[1].1) { + differential_dataflow::consolidation::consolidate_updates(&mut self.entries); + } } else { // Cross-chunk merge by value content: heads are (run index, row); each step takes // the least value among heads, drains every chunk's sub-run of it into `tds`, @@ -942,6 +945,23 @@ mod tests { assert_eq!(scratch.entries, vec![(11, Product::new(1, 0), 1)]); } + #[test] + fn compound_tokens_consolidate_after_product_time_reordering() { + use timely::order::Product; + let chunk = CorgiChunk::from_columns( + CValue::u64(vec![7; 3]), + CValue::Prod(vec![CValue::u64(vec![11; 3]), CValue::u64(vec![12; 3])]), + [Product::new(0, 2), Product::new(1, 0), Product::new(1, 2)].into_iter().collect(), + vec![1, 1, -1], + ); + let mut scratch = SideScratch::new(); + scratch.stage_runs( + &[RunRef { chunk: &chunk, cid: 0, s: 0, e: 3, vals: None }], + &Product::new(1u64, 0u64), + ); + assert_eq!(scratch.entries, vec![(0, Product::new(1, 0), 1)]); + } + #[test] fn collision_after_completed_block_filters_real_keys_without_restarting() { let collision = PULL as u64 + 1;