Skip to main content

mz_timely_util/columnar/
builder_input.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//! `BuilderInput` impls for [`ColumnBody`] and [`Column`] so DD `Builder`s can
17//! drain a batcher's chunks without an extra container conversion.
18//!
19//! Mirrors the impl on [`ColumnationStack`](crate::columnation::ColumnationStack)
20//! at `columnation.rs`, but the `Item<'a>` here is a columnar `Ref` tuple
21//! rather than a borrowed owned tuple, so:
22//!
23//! - `Key<'a>` / `Val<'a>` are `Ref<'a, K>` / `Ref<'a, V>` — no `Owned`
24//!   round-trip on the read side.
25//! - `Time` / `Diff` materialize as owned on `into_parts` (the trait
26//!   contract requires owned for these).
27//!
28//! Distinct-counts (`key_val_upd_counts`) tally per chunk and sum, accepting
29//! at most `chain.len()` over-counts at chunk boundaries. The downstream
30//! consumer uses these as capacity hints, so a small over-estimate is
31//! cheaper than the alternative (snapshotting `K::Owned` / `V::Owned`
32//! across chunk boundaries).
33//!
34//! The body impl serves the chunk batchers, whose chains are bodies. The
35//! column impl serves the column pager's batcher, whose chains are still
36//! edge containers.
37
38use columnar::{BorrowedOf, Columnar, Index, Len};
39use differential_dataflow::difference::Semigroup;
40use differential_dataflow::lattice::Lattice;
41use differential_dataflow::trace::implementations::{BatchContainer, BuilderInput};
42use timely::progress::Timestamp;
43
44use crate::columnar::Column;
45use crate::columnar::body::ColumnBody;
46
47/// Distinct key, distinct `(key, val)`, and update counts over a chain of
48/// views, deduplicated within each view and summed.
49fn key_val_upd_counts<'a, K, V, T, R>(
50    chain: impl Iterator<Item = BorrowedOf<'a, ((K, V), T, R)>>,
51) -> (usize, usize, usize)
52where
53    K: Columnar,
54    V: Columnar,
55    T: Columnar,
56    R: Columnar,
57    for<'b> columnar::Ref<'b, K>: Copy + Ord,
58    for<'b> columnar::Ref<'b, V>: Copy + Ord,
59{
60    // Per-chunk dedup, summed. Skips cross-chunk equality checks; the
61    // counts may over-count by up to `chain.len()` (one boundary per
62    // chunk). Capacity-hint consumers tolerate over-estimates.
63    let mut keys = 0;
64    let mut vals = 0;
65    let mut upds = 0;
66    for view in chain {
67        let len = view.len();
68        if len == 0 {
69            continue;
70        }
71        let mut prev: Option<(columnar::Ref<'_, K>, columnar::Ref<'_, V>)> = None;
72        for i in 0..len {
73            let ((k, v), _, _) = view.get(i);
74            match prev {
75                None => {
76                    keys += 1;
77                    vals += 1;
78                }
79                Some((pk, pv)) => {
80                    if pk != k {
81                        keys += 1;
82                        vals += 1;
83                    } else if pv != v {
84                        vals += 1;
85                    }
86                }
87            }
88            upds += 1;
89            prev = Some((k, v));
90        }
91    }
92    (keys, vals, upds)
93}
94
95macro_rules! impl_builder_input {
96    ($container:ident) => {
97        impl<KBC, VBC, K, V, T, R> BuilderInput<KBC, VBC> for $container<((K, V), T, R)>
98        where
99            K: Columnar,
100            V: Columnar,
101            T: Columnar + Timestamp + Lattice + Clone,
102            R: Columnar + Ord + Semigroup + Clone,
103            for<'a> columnar::Ref<'a, K>: Copy + Ord,
104            for<'a> columnar::Ref<'a, V>: Copy + Ord,
105            KBC: BatchContainer,
106            VBC: BatchContainer,
107            for<'a, 'b> KBC::ReadItem<'a>: PartialEq<columnar::Ref<'b, K>>,
108            for<'a, 'b> VBC::ReadItem<'a>: PartialEq<columnar::Ref<'b, V>>,
109        {
110            type Key<'a> = columnar::Ref<'a, K>;
111            type Val<'a> = columnar::Ref<'a, V>;
112            type Time = T;
113            type Diff = R;
114
115            fn into_parts<'a>(
116                item: Self::Item<'a>,
117            ) -> (Self::Key<'a>, Self::Val<'a>, Self::Time, Self::Diff) {
118                let ((key, val), time, diff) = item;
119                (key, val, T::into_owned(time), R::into_owned(diff))
120            }
121
122            fn key_eq(this: &columnar::Ref<'_, K>, other: KBC::ReadItem<'_>) -> bool {
123                KBC::reborrow(other) == *this
124            }
125
126            fn val_eq(this: &columnar::Ref<'_, V>, other: VBC::ReadItem<'_>) -> bool {
127                VBC::reborrow(other) == *this
128            }
129
130            fn key_val_upd_counts(chain: &[Self]) -> (usize, usize, usize) {
131                key_val_upd_counts::<K, V, T, R>(chain.iter().map(|chunk| chunk.borrow()))
132            }
133        }
134    };
135}
136
137impl_builder_input!(ColumnBody);
138impl_builder_input!(Column);