1use std::collections::VecDeque;
31
32use columnar::{Clear, Columnar, Index, Len};
33use differential_dataflow::difference::Semigroup;
34use differential_dataflow::logging::{BatcherEvent, Logger};
35use differential_dataflow::trace::{Batcher, Description};
36use timely::Accountable;
37use timely::PartialOrder;
38use timely::container::{ContainerBuilder, PushInto, SizableContainer};
39use timely::dataflow::channels::ContainerBytes;
40use timely::progress::Timestamp;
41use timely::progress::frontier::{Antichain, AntichainRef};
42
43use crate::column_pager::{self, ColumnPager, PagedColumn};
44use crate::columnar::Column;
45use crate::columnar::batcher::ColumnChunker;
46use crate::columnar::body::ColumnBody;
47
48const STASH_CAP: usize = 2;
63
64const MAX_RECYCLE_BYTES: usize = 1 << 22;
71
72fn recycle_capped<C: Columnar>(chunk: Column<C>, stash: &mut Vec<Column<C>>) {
77 if stash.len() < STASH_CAP && chunk.length_in_bytes() <= MAX_RECYCLE_BYTES {
78 recycle_chunk(chunk, stash);
79 }
80}
81
82#[inline]
88fn empty_chunk<C: Columnar>(stash: &mut Vec<Column<C>>) -> Column<C> {
89 stash.pop().unwrap_or_default()
90}
91
92#[inline]
96fn recycle_chunk<C: Columnar>(mut chunk: Column<C>, stash: &mut Vec<Column<C>>) {
97 if let Column::Typed(c) = &mut chunk {
98 c.clear();
99 stash.push(chunk);
100 }
101}
102
103pub struct PagedChunker<U: Columnar> {
106 inner: ColumnChunker<U>,
107 staged: Column<U>,
108}
109
110impl<U: Columnar> Default for PagedChunker<U> {
111 fn default() -> Self {
112 Self {
113 inner: Default::default(),
114 staged: Default::default(),
115 }
116 }
117}
118
119impl<'a, D, T, R> PushInto<&'a mut Column<(D, T, R)>> for PagedChunker<(D, T, R)>
120where
121 D: Columnar,
122 T: Columnar,
123 R: Columnar,
124 ColumnChunker<(D, T, R)>: PushInto<&'a mut Column<(D, T, R)>>,
125{
126 fn push_into(&mut self, item: &'a mut Column<(D, T, R)>) {
127 self.inner.push_into(item);
128 }
129}
130
131impl<U: Columnar + 'static> ContainerBuilder for PagedChunker<U>
132where
133 U::Container: Clone + 'static,
134 ColumnChunker<U>: ContainerBuilder<Container = ColumnBody<U>>,
135{
136 type Container = Column<U>;
137
138 fn extract(&mut self) -> Option<&mut Self::Container> {
139 let body = self.inner.extract()?;
140 self.staged = Column::from(std::mem::take(body));
141 Some(&mut self.staged)
142 }
143
144 fn finish(&mut self) -> Option<&mut Self::Container> {
145 let body = self.inner.finish()?;
146 self.staged = Column::from(std::mem::take(body));
147 Some(&mut self.staged)
148 }
149}
150
151pub struct ColumnMergeBatcher<D, T, R>
164where
165 D: Columnar,
166 T: Columnar,
167 R: Columnar,
168{
169 chains: Vec<VecDeque<PagedColumn<(D, T, R)>>>,
170 lower: Antichain<T>,
171 frontier: Antichain<T>,
172 stash: Vec<Column<(D, T, R)>>,
180 pager_override: Option<ColumnPager>,
185 logger: Option<Logger>,
186 operator_id: usize,
187}
188
189impl<D, T, R> ColumnMergeBatcher<D, T, R>
190where
191 D: Columnar,
192 T: Columnar,
193 R: Columnar,
194{
195 pub fn set_pager(&mut self, pager: ColumnPager) {
199 self.pager_override = Some(pager);
200 }
201
202 fn pager(&self) -> ColumnPager {
206 self.pager_override
207 .clone()
208 .unwrap_or_else(column_pager::global_pager)
209 }
210
211 fn chain_push(&mut self, chain: VecDeque<PagedColumn<(D, T, R)>>) {
214 self.emit_account(&chain, 1);
215 self.chains.push(chain);
216 }
217
218 fn chain_pop(&mut self) -> Option<VecDeque<PagedColumn<(D, T, R)>>> {
229 let chain = self.chains.pop()?;
230 self.emit_account(&chain, -1);
231 Some(chain)
232 }
233
234 fn emit_account(&self, chain: &VecDeque<PagedColumn<(D, T, R)>>, diff: isize) {
237 let Some(logger) = &self.logger else {
238 return;
239 };
240 let (mut records, mut size, mut capacity, mut allocations) =
241 (0isize, 0isize, 0isize, 0isize);
242 for entry in chain {
243 let (r, s, c, a) = account_chunk(entry);
244 records = records.saturating_add_unsigned(r);
245 size = size.saturating_add_unsigned(s);
246 capacity = capacity.saturating_add_unsigned(c);
247 allocations = allocations.saturating_add_unsigned(a);
248 }
249 logger.log(BatcherEvent {
250 operator: self.operator_id,
251 records_diff: records.saturating_mul(diff),
252 size_diff: size.saturating_mul(diff),
253 capacity_diff: capacity.saturating_mul(diff),
254 allocations_diff: allocations.saturating_mul(diff),
255 });
256 }
257}
258
259impl<D, T, R> Drop for ColumnMergeBatcher<D, T, R>
260where
261 D: Columnar,
262 T: Columnar,
263 R: Columnar,
264{
265 fn drop(&mut self) {
266 while self.chain_pop().is_some() {}
269 }
270}
271
272fn account_chunk<C: Columnar>(entry: &PagedColumn<C>) -> (usize, usize, usize, usize) {
281 match entry {
282 PagedColumn::Resident(col, _) => {
283 let records = usize::try_from(col.record_count()).expect("non-negative");
284 let bytes = col.length_in_bytes();
285 (records, bytes, bytes, 1)
286 }
287 PagedColumn::Paged { .. } | PagedColumn::Compressed { .. } => (0, 0, 0, 0),
288 }
289}
290
291impl<D, T, R> Batcher for ColumnMergeBatcher<D, T, R>
292where
293 D: Columnar,
294 for<'a> columnar::Ref<'a, D>: Copy + Ord,
295 T: Columnar + Default + Timestamp + PartialOrder,
296 for<'a> columnar::Ref<'a, T>: Copy + Ord,
297 R: Columnar + Default + Semigroup + for<'a> Semigroup<columnar::Ref<'a, R>>,
298 for<'a> columnar::Ref<'a, R>: Ord,
299{
300 type Output = Column<(D, T, R)>;
301 type Time = T;
302
303 fn new(logger: Option<Logger>, operator_id: usize) -> Self {
304 Self {
305 chains: Vec::new(),
306 lower: Antichain::from_elem(T::minimum()),
307 frontier: Antichain::new(),
308 stash: Vec::new(),
309 pager_override: None,
310 logger,
311 operator_id,
312 }
313 }
314
315 fn seal(
316 &mut self,
317 upper: Antichain<Self::Time>,
318 ) -> (Vec<Self::Output>, Description<Self::Time>) {
319 let pager = self.pager();
320 while self.chains.len() > 1 {
322 let a = self.chain_pop().unwrap();
323 let b = self.chain_pop().unwrap();
324 let merged = self.merge_by(a, b);
325 self.chain_push(merged);
326 }
327 let merged = self.chain_pop().unwrap_or_default();
328
329 let mut readied: Vec<Column<(D, T, R)>> = Vec::new();
333 let mut kept_chain: VecDeque<PagedColumn<(D, T, R)>> = VecDeque::new();
334 self.frontier.clear();
335 {
336 let pager = &pager;
337 let frontier = &mut self.frontier;
338 let stash = &mut self.stash;
339 extract_chain(
340 FetchIter::new(merged, pager),
341 upper.borrow(),
342 frontier,
343 |paged| readied.push(pager.take(paged)),
344 |paged| kept_chain.push_back(paged),
345 stash,
346 );
347 }
348
349 if !kept_chain.is_empty() {
350 self.chain_push(kept_chain);
351 }
352
353 let description = Description::new(
354 self.lower.clone(),
355 upper.clone(),
356 Antichain::from_elem(T::minimum()),
357 );
358 self.lower = upper;
359
360 self.stash.clear();
365
366 (readied, description)
367 }
368
369 fn frontier(&mut self) -> AntichainRef<'_, Self::Time> {
370 self.frontier.borrow()
371 }
372}
373
374impl<D, T, R> PushInto<Column<(D, T, R)>> for ColumnMergeBatcher<D, T, R>
375where
376 D: Columnar,
377 for<'a> columnar::Ref<'a, D>: Copy + Ord,
378 T: Columnar + Default + Clone + PartialOrder,
379 for<'a> columnar::Ref<'a, T>: Copy + Ord,
380 R: Columnar + Default + Semigroup + for<'a> Semigroup<columnar::Ref<'a, R>>,
381{
382 fn push_into(&mut self, mut chunk: Column<(D, T, R)>) {
385 let pager = self.pager();
386 let paged = pager.page(&mut chunk);
387 self.insert_chain(VecDeque::from([paged]));
388 }
389}
390
391impl<D, T, R> ColumnMergeBatcher<D, T, R>
392where
393 D: Columnar,
394 for<'a> columnar::Ref<'a, D>: Copy + Ord,
395 T: Columnar + Default + Clone + PartialOrder,
396 for<'a> columnar::Ref<'a, T>: Copy + Ord,
397 R: Columnar + Default + Semigroup + for<'a> Semigroup<columnar::Ref<'a, R>>,
398{
399 fn insert_chain(&mut self, chain: VecDeque<PagedColumn<(D, T, R)>>) {
402 if chain.is_empty() {
403 return;
404 }
405 self.chain_push(chain);
406 while self.chains.len() > 1
407 && self.chains[self.chains.len() - 1].len()
408 >= self.chains[self.chains.len() - 2].len() / 2
409 {
410 let a = self.chain_pop().unwrap();
411 let b = self.chain_pop().unwrap();
412 let merged = self.merge_by(a, b);
413 self.chain_push(merged);
414 }
415 }
416
417 fn merge_by(
421 &mut self,
422 a: VecDeque<PagedColumn<(D, T, R)>>,
423 b: VecDeque<PagedColumn<(D, T, R)>>,
424 ) -> VecDeque<PagedColumn<(D, T, R)>> {
425 let mut output: VecDeque<PagedColumn<(D, T, R)>> = VecDeque::new();
426 let pager = self.pager();
427 let pager = &pager;
428 let stash = &mut self.stash;
429 merge_chains(
430 FetchIter::new(a, pager),
431 FetchIter::new(b, pager),
432 |paged| output.push_back(paged),
433 stash,
434 );
435 output
436 }
437}
438
439pub struct FetchIter<'a, D, T, R>
445where
446 (D, T, R): Columnar,
447{
448 queue: VecDeque<PagedColumn<(D, T, R)>>,
449 pager: &'a ColumnPager,
450}
451
452impl<'a, D, T, R> FetchIter<'a, D, T, R>
453where
454 (D, T, R): Columnar,
455{
456 pub fn new(queue: VecDeque<PagedColumn<(D, T, R)>>, pager: &'a ColumnPager) -> Self {
458 Self { queue, pager }
459 }
460
461 pub fn pager(&self) -> &'a ColumnPager {
466 self.pager
467 }
468
469 pub fn into_paged(self) -> std::collections::vec_deque::IntoIter<PagedColumn<(D, T, R)>> {
474 self.queue.into_iter()
475 }
476}
477
478impl<D, T, R> Iterator for FetchIter<'_, D, T, R>
479where
480 (D, T, R): Columnar,
481{
482 type Item = Column<(D, T, R)>;
483
484 fn next(&mut self) -> Option<Self::Item> {
485 self.queue.pop_front().map(|p| self.pager.take(p))
486 }
487}
488
489pub fn merge_chains<D, T, R, Sink>(
504 list1: FetchIter<'_, D, T, R>,
505 list2: FetchIter<'_, D, T, R>,
506 mut sink: Sink,
507 stash: &mut Vec<Column<(D, T, R)>>,
508) where
509 D: Columnar,
510 for<'a> columnar::Ref<'a, D>: Copy + Ord,
511 T: Columnar + Default + Clone + PartialOrder,
512 for<'a> columnar::Ref<'a, T>: Copy + Ord,
513 R: Columnar + Default + Semigroup + for<'a> Semigroup<columnar::Ref<'a, R>>,
514 Sink: FnMut(PagedColumn<(D, T, R)>),
515{
516 let pager = list1.pager();
517 let mut list1 = list1;
518 let mut list2 = list2;
519
520 let mut heads = [
521 list1.next().unwrap_or_default(),
522 list2.next().unwrap_or_default(),
523 ];
524 let mut positions = [0usize, 0usize];
525 let mut result: Column<(D, T, R)> = empty_chunk(stash);
526
527 loop {
528 let upper_l = heads[0].borrow().len();
529 let upper_r = heads[1].borrow().len();
530 if positions[0] >= upper_l || positions[1] >= upper_r {
531 break;
532 }
533
534 let lhs_passthrough = positions[0] == 0 && upper_l > 0 && {
536 let lhs = heads[0].borrow();
537 let rhs = heads[1].borrow();
538 let last_l = (lhs.0.get(upper_l - 1), lhs.1.get(upper_l - 1));
539 let cur_r = (rhs.0.get(positions[1]), rhs.1.get(positions[1]));
540 last_l < cur_r
541 };
542 if lhs_passthrough {
543 if !result.is_empty() {
544 sink(pager.page(&mut result));
545 if let Some(reuse) = stash.pop() {
546 result = reuse;
547 }
548 }
549 let mut head = std::mem::replace(&mut heads[0], list1.next().unwrap_or_default());
550 sink(pager.page(&mut head));
551 positions[0] = 0;
552 continue;
553 }
554
555 let rhs_passthrough = positions[1] == 0 && upper_r > 0 && {
556 let lhs = heads[0].borrow();
557 let rhs = heads[1].borrow();
558 let last_r = (rhs.0.get(upper_r - 1), rhs.1.get(upper_r - 1));
559 let cur_l = (lhs.0.get(positions[0]), lhs.1.get(positions[0]));
560 last_r < cur_l
561 };
562 if rhs_passthrough {
563 if !result.is_empty() {
564 sink(pager.page(&mut result));
565 if let Some(reuse) = stash.pop() {
566 result = reuse;
567 }
568 }
569 let mut head = std::mem::replace(&mut heads[1], list2.next().unwrap_or_default());
570 sink(pager.page(&mut head));
571 positions[1] = 0;
572 continue;
573 }
574
575 let yielded = result.merge_from(&mut heads, &mut positions);
576
577 if positions[0] >= heads[0].borrow().len() {
578 let old = std::mem::replace(&mut heads[0], list1.next().unwrap_or_default());
579 recycle_capped(old, stash);
580 positions[0] = 0;
581 }
582 if positions[1] >= heads[1].borrow().len() {
583 let old = std::mem::replace(&mut heads[1], list2.next().unwrap_or_default());
584 recycle_capped(old, stash);
585 positions[1] = 0;
586 }
587 if yielded || result.at_capacity() {
588 sink(pager.page(&mut result));
589 if let Some(reuse) = stash.pop() {
595 result = reuse;
596 }
597 }
598 }
599
600 drain_side(
604 &mut heads[0],
605 &mut positions[0],
606 list1,
607 &mut result,
608 &mut sink,
609 pager,
610 stash,
611 );
612 drain_side(
613 &mut heads[1],
614 &mut positions[1],
615 list2,
616 &mut result,
617 &mut sink,
618 pager,
619 stash,
620 );
621
622 if !result.is_empty() {
623 sink(pager.page(&mut result));
624 } else {
625 recycle_capped(result, stash);
629 }
630 let [h0, h1] = heads;
634 recycle_capped(h0, stash);
635 recycle_capped(h1, stash);
636}
637
638fn drain_side<D, T, R, Sink>(
642 head: &mut Column<(D, T, R)>,
643 pos: &mut usize,
644 rest: FetchIter<'_, D, T, R>,
645 result: &mut Column<(D, T, R)>,
646 sink: &mut Sink,
647 pager: &ColumnPager,
648 stash: &mut Vec<Column<(D, T, R)>>,
649) where
650 D: Columnar,
651 for<'a> columnar::Ref<'a, D>: Copy + Ord,
652 T: Columnar + Default + Clone + PartialOrder,
653 for<'a> columnar::Ref<'a, T>: Copy + Ord,
654 R: Columnar + Default + Semigroup + for<'a> Semigroup<columnar::Ref<'a, R>>,
655 Sink: FnMut(PagedColumn<(D, T, R)>),
656{
657 if *pos < head.borrow().len() {
658 let _ = result.merge_from(std::slice::from_mut(head), std::slice::from_mut(pos));
660 }
661 if !result.is_empty() {
662 sink(pager.page(result));
663 if let Some(reuse) = stash.pop() {
664 *result = reuse;
665 }
666 }
667 for paged in rest.into_paged() {
668 sink(paged);
669 }
670}
671
672pub fn extract_chain<D, T, R, SinkShip, SinkKeep>(
681 merged: FetchIter<'_, D, T, R>,
682 upper: AntichainRef<T>,
683 frontier: &mut Antichain<T>,
684 mut ship: SinkShip,
685 mut keep: SinkKeep,
686 stash: &mut Vec<Column<(D, T, R)>>,
687) where
688 D: Columnar,
689 for<'a> columnar::Ref<'a, D>: Copy + Ord,
690 T: Columnar + Default + Clone + PartialOrder,
691 for<'a> columnar::Ref<'a, T>: Copy + Ord,
692 R: Columnar + Default + Semigroup + for<'a> Semigroup<columnar::Ref<'a, R>>,
693 SinkShip: FnMut(PagedColumn<(D, T, R)>),
694 SinkKeep: FnMut(PagedColumn<(D, T, R)>),
695{
696 let pager = merged.pager();
697 let mut keep_buf: Column<(D, T, R)> = empty_chunk(stash);
698 let mut ship_buf: Column<(D, T, R)> = empty_chunk(stash);
699
700 for mut buffer in merged {
701 let mut position = 0;
702 let len = buffer.borrow().len();
703 while position < len {
704 buffer.extract(&mut position, upper, frontier, &mut keep_buf, &mut ship_buf);
705 if keep_buf.at_capacity() {
706 keep(pager.page(&mut keep_buf));
707 if let Some(reuse) = stash.pop() {
708 keep_buf = reuse;
709 }
710 }
711 if ship_buf.at_capacity() {
712 ship(pager.page(&mut ship_buf));
713 if let Some(reuse) = stash.pop() {
714 ship_buf = reuse;
715 }
716 }
717 }
718 recycle_capped(buffer, stash);
720 }
721 if !keep_buf.is_empty() {
722 keep(pager.page(&mut keep_buf));
723 } else {
724 recycle_capped(keep_buf, stash);
725 }
726 if !ship_buf.is_empty() {
727 ship(pager.page(&mut ship_buf));
728 } else {
729 recycle_capped(ship_buf, stash);
730 }
731}
732
733#[cfg(test)]
734#[allow(clippy::clone_on_ref_ptr)]
735mod tests {
736 use std::sync::Arc;
737
738 use columnar::Index;
739
740 use super::*;
741 use crate::column_pager::{PageDecision, PageEvent, PageHint, PagingPolicy};
742
743 type KvUpdate = ((u64, u64), u64, i64);
744
745 fn col(rows: &[KvUpdate]) -> Column<KvUpdate> {
746 let mut c: Column<KvUpdate> = Default::default();
747 for &t in rows {
748 c.push_into(t);
749 }
750 c
751 }
752
753 fn collect_pc(chunks: &[PagedColumn<KvUpdate>], pager: &ColumnPager) -> Vec<KvUpdate> {
754 chunks
757 .iter()
758 .flat_map(|p| {
759 let view: Column<KvUpdate> = match p {
760 PagedColumn::Resident(c, _) => clone_column(c),
761 _ => pager.take(clone_paged(p)),
762 };
763 collect_column(&view).into_iter()
764 })
765 .collect()
766 }
767
768 fn collect_column(c: &Column<KvUpdate>) -> Vec<KvUpdate> {
769 c.borrow()
770 .into_index_iter()
771 .map(|((k, v), t, r)| {
772 (
773 (u64::into_owned(k), u64::into_owned(v)),
774 u64::into_owned(t),
775 i64::into_owned(r),
776 )
777 })
778 .collect()
779 }
780
781 fn clone_column(c: &Column<KvUpdate>) -> Column<KvUpdate> {
782 c.clone()
786 }
787
788 fn clone_paged(p: &PagedColumn<KvUpdate>) -> PagedColumn<KvUpdate> {
792 match p {
793 PagedColumn::Resident(c, _) => {
794 let mut c = c.clone();
796 ColumnPager::disabled().page(&mut c)
797 }
798 _ => panic!("clone_paged only supports Resident"),
801 }
802 }
803
804 struct ForcePagePolicy {
807 out: std::sync::atomic::AtomicUsize,
808 r#in: std::sync::atomic::AtomicUsize,
809 }
810 impl ForcePagePolicy {
811 fn new() -> Arc<Self> {
812 Arc::new(Self {
813 out: std::sync::atomic::AtomicUsize::new(0),
814 r#in: std::sync::atomic::AtomicUsize::new(0),
815 })
816 }
817 }
818 impl PagingPolicy for ForcePagePolicy {
819 fn decide(&self, _hint: PageHint) -> PageDecision {
820 PageDecision::Page {
821 backend: mz_ore::pager::Backend::Swap,
822 codec: None,
823 }
824 }
825 fn record(&self, event: PageEvent) {
826 use std::sync::atomic::Ordering;
827 match event {
828 PageEvent::PagedOut { .. } => {
829 self.out.fetch_add(1, Ordering::Relaxed);
830 }
831 PageEvent::PagedIn { .. } => {
832 self.r#in.fetch_add(1, Ordering::Relaxed);
833 }
834 _ => {}
835 }
836 }
837 }
838
839 fn to_chain(
841 cols: Vec<Column<KvUpdate>>,
842 pager: &ColumnPager,
843 ) -> VecDeque<PagedColumn<KvUpdate>> {
844 cols.into_iter().map(|mut c| pager.page(&mut c)).collect()
845 }
846
847 fn drive_merge(chain1: Vec<Column<KvUpdate>>, chain2: Vec<Column<KvUpdate>>) -> Vec<KvUpdate> {
849 let pager = ColumnPager::disabled();
850 let q1 = to_chain(chain1, &pager);
851 let q2 = to_chain(chain2, &pager);
852 let mut output: Vec<PagedColumn<KvUpdate>> = Vec::new();
853 let mut stash: Vec<Column<KvUpdate>> = Vec::new();
854 merge_chains(
855 FetchIter::new(q1, &pager),
856 FetchIter::new(q2, &pager),
857 |paged| output.push(paged),
858 &mut stash,
859 );
860 collect_pc(&output, &pager)
861 }
862
863 #[mz_ore::test]
867 fn merge_chains_disjoint_ranges() {
868 let out = drive_merge(
869 vec![
870 col(&[((0, 0), 0, 1), ((1, 0), 0, 1)]),
871 col(&[((2, 0), 0, 1), ((3, 0), 0, 1)]),
872 ],
873 vec![
874 col(&[((10, 0), 0, 1), ((11, 0), 0, 1)]),
875 col(&[((12, 0), 0, 1), ((13, 0), 0, 1)]),
876 ],
877 );
878 let expected: Vec<_> = (0..4u64)
879 .map(|d| ((d, 0u64), 0u64, 1i64))
880 .chain((10..14u64).map(|d| ((d, 0u64), 0u64, 1i64)))
881 .collect();
882 assert_eq!(out, expected);
883 }
884
885 #[mz_ore::test]
886 fn merge_chains_interleaved() {
887 let out = drive_merge(
888 vec![
889 col(&[((0, 0), 0, 1), ((2, 0), 0, 1)]),
890 col(&[((4, 0), 0, 1), ((6, 0), 0, 1)]),
891 ],
892 vec![
893 col(&[((1, 0), 0, 1), ((3, 0), 0, 1)]),
894 col(&[((5, 0), 0, 1), ((7, 0), 0, 1)]),
895 ],
896 );
897 let expected: Vec<_> = (0..8u64).map(|d| ((d, 0u64), 0u64, 1i64)).collect();
898 assert_eq!(out, expected);
899 }
900
901 #[mz_ore::test]
905 fn merge_chains_equal_boundary() {
906 let out = drive_merge(
907 vec![col(&[((0, 0), 0, 1), ((5, 0), 0, 1)])],
908 vec![col(&[((5, 0), 0, 1), ((10, 0), 0, 1)])],
909 );
910 assert_eq!(out, vec![((0, 0), 0, 1), ((5, 0), 0, 2), ((10, 0), 0, 1)]);
911 }
912
913 #[mz_ore::test]
921 fn merge_chains_ships_fitting_align() {
922 let pager = ColumnPager::disabled();
923 let q1 = to_chain(vec![col(&[((0, 0), 0, 1), ((2, 0), 0, 1)])], &pager);
927 let q2 = to_chain(vec![col(&[((1, 0), 0, 1), ((3, 0), 0, 1)])], &pager);
928
929 let mut output: Vec<PagedColumn<KvUpdate>> = Vec::new();
930 let mut stash: Vec<Column<KvUpdate>> = Vec::new();
931 merge_chains(
932 FetchIter::new(q1, &pager),
933 FetchIter::new(q2, &pager),
934 |paged| output.push(paged),
935 &mut stash,
936 );
937
938 assert!(!output.is_empty(), "merge produced no chunks");
939 for entry in &output {
940 match entry {
941 PagedColumn::Resident(col, _) => assert!(
942 matches!(col, Column::Align(_)),
943 "shipped chunk parked as non-Align resident: {:?}",
944 std::mem::discriminant(col),
945 ),
946 other => panic!(
947 "disabled pager should ship Resident, got a paged variant: {:?}",
948 std::mem::discriminant(other)
949 ),
950 }
951 }
952
953 assert_eq!(
955 collect_pc(&output, &pager),
956 vec![
957 ((0, 0), 0, 1),
958 ((1, 0), 0, 1),
959 ((2, 0), 0, 1),
960 ((3, 0), 0, 1)
961 ],
962 );
963 }
964
965 #[mz_ore::test]
968 fn merge_chains_force_paged_round_trip() {
969 let policy = ForcePagePolicy::new();
970 let pager = ColumnPager::new(policy.clone());
971 let q1 = to_chain(vec![col(&[((0, 0), 0, 1), ((2, 0), 0, 1)])], &pager);
972 let q2 = to_chain(vec![col(&[((1, 0), 0, 1), ((3, 0), 0, 1)])], &pager);
973
974 assert!(matches!(q1.front().unwrap(), PagedColumn::Paged { .. }));
976 assert!(matches!(q2.front().unwrap(), PagedColumn::Paged { .. }));
977
978 let mut output: Vec<PagedColumn<KvUpdate>> = Vec::new();
979 let mut stash: Vec<Column<KvUpdate>> = Vec::new();
980 merge_chains(
981 FetchIter::new(q1, &pager),
982 FetchIter::new(q2, &pager),
983 |paged| output.push(paged),
984 &mut stash,
985 );
986
987 for p in &output {
989 assert!(matches!(p, PagedColumn::Paged { .. }));
990 }
991
992 let mut collected = Vec::new();
994 for p in output {
995 let c = pager.take(p);
996 collected.extend(collect_column(&c));
997 }
998 let expected: Vec<_> = (0..4u64).map(|d| ((d, 0u64), 0u64, 1i64)).collect();
999 assert_eq!(collected, expected);
1000 }
1001
1002 #[mz_ore::test]
1003 fn extract_chain_partitions_by_frontier() {
1004 let pager = ColumnPager::disabled();
1005 let data = vec![
1006 ((0, 0), 0u64, 1i64),
1007 ((1, 0), 1, 1),
1008 ((2, 0), 2, 1),
1009 ((3, 0), 3, 1),
1010 ];
1011 let chain = to_chain(vec![col(&data)], &pager);
1012 let upper = Antichain::from_elem(2u64);
1013 let mut frontier: Antichain<u64> = Antichain::new();
1014 let mut ship: Vec<PagedColumn<KvUpdate>> = Vec::new();
1015 let mut keep: Vec<PagedColumn<KvUpdate>> = Vec::new();
1016 let mut stash: Vec<Column<KvUpdate>> = Vec::new();
1017
1018 extract_chain(
1019 FetchIter::new(chain, &pager),
1020 upper.borrow(),
1021 &mut frontier,
1022 |p| ship.push(p),
1023 |p| keep.push(p),
1024 &mut stash,
1025 );
1026
1027 let shipped = collect_pc(&ship, &pager);
1028 let kept = collect_pc(&keep, &pager);
1029 for (_, t, _) in &shipped {
1030 assert!(*t < 2, "shipped time {t} should be < upper");
1031 }
1032 for (_, t, _) in &kept {
1033 assert!(*t >= 2, "kept time {t} should be >= upper");
1034 }
1035 assert_eq!(shipped.len() + kept.len(), data.len());
1036 }
1037
1038 #[mz_ore::test]
1039 fn batcher_seal_round_trip() {
1040 let mut b: ColumnMergeBatcher<(u64, u64), u64, i64> =
1041 differential_dataflow::trace::Batcher::new(None, 0);
1042 let input1 = col(&[((1, 1), 0, 1), ((2, 0), 0, 1), ((3, 0), 0, 1)]);
1046 let input2 = col(&[((2, 0), 0, 2), ((4, 0), 0, 1)]);
1047 b.push_into(input1);
1048 b.push_into(input2);
1049
1050 let upper = Antichain::from_elem(u64::MAX);
1052 let (chain, _description) = differential_dataflow::trace::Batcher::seal(&mut b, upper);
1053 let out: Vec<KvUpdate> = chain.iter().flat_map(collect_column).collect();
1054
1055 let mut expected = vec![
1057 ((1u64, 1u64), 0u64, 1i64),
1058 ((2, 0), 0, 3),
1059 ((3, 0), 0, 1),
1060 ((4, 0), 0, 1),
1061 ];
1062 expected.sort();
1063 let mut out_sorted = out.clone();
1064 out_sorted.sort();
1065 assert_eq!(out_sorted, expected);
1066 }
1067
1068 #[mz_ore::test]
1069 fn account_chunk_resident_vs_paged() {
1070 let policy = ForcePagePolicy::new();
1071 let pager_paged = ColumnPager::new(policy.clone());
1072 let pager_res = ColumnPager::disabled();
1073
1074 let mut c1 = col(&[((1, 1), 0, 1), ((2, 0), 0, 1), ((3, 0), 0, 1)]);
1075 let resident = pager_res.page(&mut c1);
1076 let (records, size, capacity, allocations) = account_chunk(&resident);
1077 assert_eq!(records, 3);
1078 assert!(size > 0);
1079 assert_eq!(size, capacity);
1080 assert_eq!(allocations, 1);
1081
1082 let mut c2 = col(&[((1, 1), 0, 1), ((2, 0), 0, 1)]);
1083 let paged = pager_paged.page(&mut c2);
1084 assert!(matches!(paged, PagedColumn::Paged { .. }));
1085 assert_eq!(account_chunk(&paged), (0, 0, 0, 0));
1087 }
1088
1089 #[mz_ore::test]
1090 fn batcher_seal_keeps_kept_chain_paged() {
1091 let policy = ForcePagePolicy::new();
1094 let pager = ColumnPager::new(policy.clone());
1095
1096 let mut b: ColumnMergeBatcher<(u64, u64), u64, i64> =
1097 differential_dataflow::trace::Batcher::new(None, 0);
1098 b.set_pager(pager);
1099
1100 let n: u64 = 200;
1103 for i in 0..n {
1104 let input = col(&[((i, 0), i % 10, 1)]);
1105 b.push_into(input);
1106 }
1107 let upper = Antichain::from_elem(5u64);
1108 let _ = differential_dataflow::trace::Batcher::seal(&mut b, upper);
1109
1110 let kept_records: usize = b
1112 .chains
1113 .iter()
1114 .flat_map(|c| c.iter())
1115 .map(|p| match p {
1116 PagedColumn::Paged { meta, .. } => {
1117 let _ = meta;
1120 1
1121 }
1122 PagedColumn::Compressed { meta, .. } => {
1123 let _ = meta;
1124 1
1125 }
1126 PagedColumn::Resident(_, _) => {
1127 panic!("kept chain entry was Resident under ForcePagePolicy");
1128 }
1129 })
1130 .sum();
1131 assert!(kept_records > 0, "expected at least one kept paged entry");
1133 assert!(policy.out.load(std::sync::atomic::Ordering::Relaxed) > 0);
1134 let _ = n;
1135 }
1136}