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);