Skip to main content

mz_timely_util/columnar/
merge_batcher.rs

1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License in the LICENSE file at the
6// root of this repository, or online at
7//
8//     http://www.apache.org/licenses/LICENSE-2.0
9//
10// Unless required by applicable law or agreed to in writing, software
11// distributed under the License is distributed on an "AS IS" BASIS,
12// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13// See the License for the specific language governing permissions and
14// limitations under the License.
15
16//! Merge-batcher for [`Column`] chunks with per-chunk paging.
17//!
18//! Forks the [`differential_dataflow`] merge-batcher framework so chains can
19//! hold [`PagedColumn`] entries — letting the [`ColumnPager`] page chunks
20//! out as they're produced and fetch them back lazily during merge / extract.
21//!
22//! Reuses the resident building blocks from [`super::batcher`]: the inherent
23//! `Column::merge_from` / `Column::extract` methods (per-chunk merge / split).
24//! Input consolidation happens upstream: the chunker ([`PagedChunker`]) is
25//! supplied to the arrange operator separately, so this batcher receives
26//! already-consolidated [`Column`] chunks via [`PushInto`].
27//!
28//! [`differential_dataflow`]: differential_dataflow::trace::implementations::merge_batcher
29
30use 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
48/// Max recycled empty chunks held in the per-batcher stash. Deliberately
49/// tight: the stash is a hot-buffer cache for the result/keep/ship churn,
50/// not a hoard. Stash entries are cleared `Column::Typed` allocations that
51/// retain capacity but are *not* tracked by [`ColumnPager`]'s
52/// `ResidentTicket` accounting, so each one is a chunk's worth of resident
53/// bytes the pager's budget doesn't see. There's one stash per arrange
54/// batcher per worker, so this multiplies fast.
55///
56/// 2 covers steady-state reuse for both code paths: `merge_chains` ships
57/// `result` and immediately pulls a refill; `extract_chain` ships `keep` /
58/// `ship` and pulls a refill for whichever was at capacity. Heads that
59/// drain mid-loop arrive resident from `FetchIter`, so the whole-chunk
60/// passthrough fast path keeps most of them off the merge inner loop
61/// entirely — only a small minority ever flow back through the stash.
62const STASH_CAP: usize = 2;
63
64/// Don't park a buffer larger than this in the free-list. A transiently
65/// oversize merge buffer (post-explosion, past the natural ship threshold)
66/// held resident would compete with the pager's budget; drop it and let a
67/// fresh default regrow. 2 × the natural ship word count (≈ 4 MiB
68/// serialized) keeps normal ship-sized chunks while excluding pathological
69/// ones.
70const MAX_RECYCLE_BYTES: usize = 1 << 22;
71
72/// Recycle `chunk` only if the stash isn't already at [`STASH_CAP`] and the
73/// chunk isn't oversize per [`MAX_RECYCLE_BYTES`]. `length_in_bytes` is
74/// measured before clear, so it reflects the data the chunk was carrying
75/// (a proxy for the capacity we'd park).
76fn 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/// Pop a chunk from `stash` or allocate a fresh one. Stashed chunks are
83/// already cleared via `recycle_chunk`, so they're ready for push.
84///
85/// The [`Column`] counterpart of the merger's body helpers, which this batcher
86/// needs until the column pager is retired.
87#[inline]
88fn empty_chunk<C: Columnar>(stash: &mut Vec<Column<C>>) -> Column<C> {
89    stash.pop().unwrap_or_default()
90}
91
92/// Reset `chunk` to an empty `Typed` and push it to `stash` for reuse. Only
93/// typed chunks carry an allocation worth keeping, so a serialized chunk is
94/// dropped instead.
95#[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
103/// [`ColumnChunker`] for this batcher, whose chains are still [`Column`]s:
104/// each body the chunker produces goes back onto the edge container, a move.
105pub 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
151/// Drives the merge-batcher over [`Column`] chunks routed through a
152/// [`ColumnPager`].
153///
154/// Chains hold [`PagedColumn`] entries rather than resident [`Column`]s, so
155/// each insert / merge / extract step can hand its output to the pager and
156/// store whatever the policy returns (resident, paged, or compressed). Reads
157/// during merge materialize lazily via [`FetchIter`].
158///
159/// Resolves its pager lazily per call via [`column_pager::global_pager`], so
160/// late-arriving dyncfg updates (e.g. `enable_column_paged_batcher` flipping
161/// on after the batcher was constructed) take effect without rebuilding the
162/// operator. Tests may override that lookup via [`Self::set_pager`].
163pub 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    /// Recycled empty `Column::Typed` chunks. Drained heads and shipped result
173    /// buffers feed in here; subsequent merge / extract calls pop from here
174    /// instead of starting from a zero-capacity `Column::default()`. Mirrors
175    /// the stash carried by the upstream `differential_dataflow` merge-batcher
176    /// framework, which this type forks. Without it, each shipped chunk
177    /// triggers a fresh per-leaf grow cycle and per-merge-round allocation
178    /// dominates the inner loop.
179    stash: Vec<Column<(D, T, R)>>,
180    /// Optional override. `None` means "read [`column_pager::global_pager`]
181    /// fresh on every use" — the production path, so worker_config dyncfg
182    /// changes that re-install the process-global pager take effect on the
183    /// very next chunk this batcher processes.
184    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    /// Pin the pager this batcher uses, overriding the thread-local lookup.
196    /// Mainly for tests; production should leave the override unset so
197    /// dyncfg-driven re-installs take effect immediately.
198    pub fn set_pager(&mut self, pager: ColumnPager) {
199        self.pager_override = Some(pager);
200    }
201
202    /// Current pager — override if set, else the process-global pager
203    /// installed by `apply_worker_config`. `ColumnPager` is cheaply
204    /// cloneable (Arc inside).
205    fn pager(&self) -> ColumnPager {
206        self.pager_override
207            .clone()
208            .unwrap_or_else(column_pager::global_pager)
209    }
210
211    /// Push a chain into `self.chains`, emitting a positive `BatcherEvent`
212    /// covering its resident entries.
213    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    /// Pop a chain from `self.chains`, emitting a negative `BatcherEvent`
219    /// retracting its resident entries.
220    ///
221    /// Invariant for the retract to reconcile against the matching
222    /// `chain_push`: chain entries are never mutated in place between push
223    /// and pop. The only allowed mutation is a full pop / push pair (see
224    /// `insert_chain` and `merge_by`), so each entry's accounting category
225    /// — `Resident` vs `Paged` vs `Compressed` — is the same at both ends.
226    /// If a future change ever pages an entry out in place after push, this
227    /// path silently double-counts.
228    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    /// Emit a single `BatcherEvent` summing resident accounting across
235    /// `chain` with the given sign. No-op when no logger is attached.
236    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        // Retract accounting for any chains still resident at drop time so
267        // the BatcherEvent counters end at zero per-operator.
268        while self.chain_pop().is_some() {}
269    }
270}
271
272/// Resident-only accounting. Returns `(records, size_bytes, capacity_bytes,
273/// allocations)` for a single chain entry; paged-out entries contribute 0
274/// across the board.
275///
276/// `BatcherEvent` feeds the `mz_arrangement_batcher_*_raw` introspection
277/// tables, which downstream surface as memory-resource dashboards. Bytes
278/// living on swap or in a pager file aren't part of RSS and shouldn't be
279/// reported there.
280fn 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        // Merge all remaining chains into one.
321        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        // Extract `merged` into `readied` (ship side, materialized for the
330        // builder) and `kept_chain` (keep side, stays paged for the next
331        // round).
332        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        // Drop the recycle stash now that this round's hot work is done:
361        // the next merge re-pays one chunk's worth of leaf grow tax, and in
362        // exchange the leaf bytes are not held resident across what may be
363        // a quiet stretch.
364        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    /// Accept an already-consolidated chunk from the upstream chunker, route
383    /// it through the pager, and insert it as a singleton chain.
384    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    /// Insert `chain` and rebalance: while the youngest chain is at least
400    /// half the size of its predecessor, merge them.
401    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    /// Merge two sorted chains. Outputs are routed through `self.pager.page`
418    /// per chunk produced, so the result chain holds `PagedColumn`s and the
419    /// caller never sees a fully materialized merge result.
420    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
439/// Streaming materializer over a chain of [`PagedColumn`] entries.
440///
441/// `next` consumes one entry and calls [`ColumnPager::take`] to produce a
442/// resident [`Column`]. Bounds materialized chunks to whatever the consumer
443/// holds (typically one head per chain in [`merge_chains`]).
444pub 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    /// Wraps `queue` for streaming materialization through `pager`.
457    pub fn new(queue: VecDeque<PagedColumn<(D, T, R)>>, pager: &'a ColumnPager) -> Self {
458        Self { queue, pager }
459    }
460
461    /// Borrow the pager backing this iter so drivers can route output chunks
462    /// back through `page()` without threading a separate `&pager`. The
463    /// returned reference is tied to the outer `'a`, not to `&self`, so it
464    /// stays valid across subsequent `next()` calls.
465    pub fn pager(&self) -> &'a ColumnPager {
466        self.pager
467    }
468
469    /// Drain remaining queued entries as `PagedColumn`s without materializing.
470    /// Used by `merge_chains`'s drain-tail phase: once the other side is
471    /// exhausted, the remaining entries on this side can pass straight to the
472    /// output sink.
473    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
489/// Two-way merge driver. Reuses today's per-chunk gallop / ship-threshold
490/// logic from `Column::merge_from`, but pulls heads from [`FetchIter`] and
491/// emits finished output chunks through `sink` after routing them through
492/// the pager exposed by [`FetchIter::pager`].
493///
494/// `stash` is a pool of empty `Column::Typed` chunks. Drained heads and
495/// shipped result buffers get recycled into it; the next result chunk is
496/// pulled from it instead of starting from a zero-capacity default. This
497/// matches the recycling discipline the upstream `differential_dataflow`
498/// merge-batcher carries via `Merger::merge`'s `stash` parameter.
499///
500/// Whole-chunk passthrough mirrors the fast path in `super::batcher`'s
501/// `Merger::merge`: a head that sorts entirely before the other side's
502/// current record ships wholesale.
503pub 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        // Whole-chunk passthrough. Two probes on already-resident heads.
535        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            // `pager.page` either took `result`'s allocation (Skip path leaves
590            // a zero-cap default) or kept the Typed buffer (Paged / Compressed
591            // paths clear in place). Pull a fresh chunk from the stash so the
592            // next `merge_from` starts with retained capacity; if the stash is
593            // empty, fall back to whatever `result` already is.
594            if let Some(reuse) = stash.pop() {
595                result = reuse;
596            }
597        }
598    }
599
600    // Drain remaining: copy partial head through `merge_from`'s 1-input
601    // dispatch, then hand the rest of the chain's `PagedColumn`s straight to
602    // the sink without materializing.
603    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        // Empty `result` may still carry a useful Typed allocation; recycle
626        // so subsequent calls (next `merge_by`, the seal `extract_chain`)
627        // can pick it up.
628        recycle_capped(result, stash);
629    }
630    // Recycle the now-exhausted (or default) head slots too — for `Resident`
631    // heads that finished naturally, this preserves their Typed allocation
632    // for the next call.
633    let [h0, h1] = heads;
634    recycle_capped(h0, stash);
635    recycle_capped(h1, stash);
636}
637
638/// Helper for `merge_chains`'s drain phase: copy a partially-consumed head
639/// into `result` (via 1-input `merge_from`), ship `result` if non-empty, then
640/// pass the remaining queued `PagedColumn`s straight through.
641fn 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        // 1-input dispatch — bulk copy that runs to completion.
659        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
672/// Streaming extract: walks `merged` chunk-by-chunk via `Column::extract`,
673/// routing each filled keep/ship chunk through its sink after pageing.
674/// Mirrors the per-chunk ship-threshold yield already inside
675/// `Column::extract`.
676///
677/// `stash` carries recycled `Column::Typed` buffers in and out so the
678/// per-chunk extract loop doesn't restart from zero capacity each time
679/// `keep_buf` / `ship_buf` ships and the source `buffer` is dropped.
680pub 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        // Buffer fully consumed; recycle whatever Typed allocation it had.
719        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        // `collect_pc` peeks via materialization on a side path so the test's
755        // assertions don't consume the chain.
756        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        // `Column` is `Clone` when `C::Container: Clone`, which is true for
783        // tuple-of-primitive containers. Used so test helpers can peek at a
784        // chain without consuming it.
785        c.clone()
786    }
787
788    /// Helper that bypasses `pager.take` for non-`Resident` variants by
789    /// taking and re-pageing. Only used in test inspection paths where the
790    /// extra round-trip is acceptable.
791    fn clone_paged(p: &PagedColumn<KvUpdate>) -> PagedColumn<KvUpdate> {
792        match p {
793            PagedColumn::Resident(c, _) => {
794                // Wrap via a disabled pager so the ticket is fresh.
795                let mut c = c.clone();
796                ColumnPager::disabled().page(&mut c)
797            }
798            // For paged/compressed variants we can't clone without
799            // re-reading; the tests below only inspect Resident chains.
800            _ => panic!("clone_paged only supports Resident"),
801        }
802    }
803
804    /// Always-page policy: bypasses any resident shortcut so we can assert
805    /// the chains remain in `Paged` form regardless of memory pressure.
806    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    /// Wrap a Vec<Column> as a paged chain for `FetchIter`.
840    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    /// Drive `merge_chains` with a disabled pager and return owned tuples.
848    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    /// Disjoint chains: same data as the legacy passthrough test. Without
864    /// passthrough, the merger runs per-record but should still produce the
865    /// fully ordered output.
866    #[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    /// Equal-key consolidation across chunk boundaries: chain1's last record
902    /// shares `(d, t)` with chain2's first; sum of diffs should land on a
903    /// single output record.
904    #[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    /// Regression: under the disabled (always-resident) pager, shipped chunks
914    /// must be serialized into a fitting `Column::Align`, never parked as
915    /// `Column::Typed`. A `Typed` result carries `Column::merge_from`'s
916    /// worst-case `reserve_for` capacity; leaving it in the chain across merge
917    /// rounds was the dominant source of merge-batcher resident memory. Only
918    /// the live accumulator (`result`) and not-yet-shipped heads may be
919    /// `Typed` — every entry that reaches the sink should be `Align`.
920    #[mz_ore::test]
921    fn merge_chains_ships_fitting_align() {
922        let pager = ColumnPager::disabled();
923        // Interleaved keys force the per-record merge path: records flow
924        // through the `result` accumulator and ship as a merged chunk rather
925        // than passing a head through wholesale.
926        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        // Sanity: data round-trips through the fitting Align buffers.
954        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    /// Same merge, force-paged: chains stay in `Paged` form throughout, and
966    /// the consolidated result still matches.
967    #[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        // Confirm the chains started paged-out (not Resident).
975        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        // Output entries should also have been routed through the pager.
988        for p in &output {
989            assert!(matches!(p, PagedColumn::Paged { .. }));
990        }
991
992        // Materialize the output and check correctness.
993        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        // Two pushes; second has an equal-key collision with the first.
1043        // Inputs arrive pre-consolidated chunk-by-chunk, as from the upstream
1044        // chunker.
1045        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        // Seal everything (upper = ∞-ish, here just past any time we used).
1051        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        // (2, 0)@0 was pushed with +1 then +2; sums to +3 after consolidation.
1056        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        // Paged variants contribute zero to memory accounting.
1086        assert_eq!(account_chunk(&paged), (0, 0, 0, 0));
1087    }
1088
1089    #[mz_ore::test]
1090    fn batcher_seal_keeps_kept_chain_paged() {
1091        // Force-page policy; verify that after seal, the kept chain in
1092        // self.chains contains only Paged entries (no Resident).
1093        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        // Push records straddling an upper of 5 — half should be kept, half
1101        // shipped. Use enough records to fill at least one chunk.
1102        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        // Anything kept (times >= 5) should be sitting in b.chains as paged.
1111        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                    // Records aren't directly available here; sanity-check
1118                    // that no Resident snuck in.
1119                    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        // We expect *some* kept entries (times in [5..10) loop slot).
1132        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}