Skip to main content

mz_timely_util/columnar/
body.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//! [`ColumnBody`]: a columnar chunk body at rest.
17//!
18//! [`Column`] is the container on dataflow edges: it is built by pushing, crosses an exchange as
19//! channel bytes, and is read on arrival. The chunk machinery behind an edge holds bodies instead:
20//! sorted, consolidated runs that a merge batcher chains, that a spill-backed chunk keeps on the
21//! heap or in the pool, and that a batch builder seals. A body is typed while it is written and
22//! serialized when a store handed it back, and nothing else. It never holds channel bytes, and
23//! where it lives is the business of whatever holds it.
24
25use std::io::Write;
26
27use columnar::bytes::indexed;
28use columnar::common::IterOwn;
29use columnar::{Borrow, BorrowedOf, Clear, Columnar, Container as _, FromBytes, Index, Len, Ref};
30use timely::Accountable;
31use timely::container::{DrainContainer, PushInto};
32
33use crate::columnar::{Column, at_serialized_capacity};
34
35/// A sorted, consolidated run of columnar records, typed while written and serialized once read
36/// back from a store.
37pub enum ColumnBody<C: Columnar> {
38    /// The typed containers, which is what a body is while it is written to.
39    Typed(C::Container),
40    /// The serialized form, as a store hands a body back: `u64`-aligned words holding the flat
41    /// [`columnar::bytes::indexed`] encoding of the typed containers. It is the only serialized
42    /// form a body has.
43    Words(Vec<u64>),
44}
45
46impl<C: Columnar> Default for ColumnBody<C> {
47    fn default() -> Self {
48        Self::Typed(Default::default())
49    }
50}
51
52impl<C: Columnar> Clone for ColumnBody<C>
53where
54    C::Container: Clone,
55{
56    fn clone(&self) -> Self {
57        match self {
58            ColumnBody::Typed(typed) => ColumnBody::Typed(typed.clone()),
59            ColumnBody::Words(words) => ColumnBody::Words(words.clone()),
60        }
61    }
62}
63
64impl<C: Columnar> ColumnBody<C> {
65    /// Borrows the body as a columnar view.
66    ///
67    /// The serialized form rebuilds its view from the encoded header on every call, so a caller
68    /// that reads more than one record hoists the view out of its loop.
69    // NOTE: `inline(always)` for the reason measured at `Column::borrow`.
70    #[inline(always)]
71    pub fn borrow(&self) -> BorrowedOf<'_, C> {
72        match self {
73            ColumnBody::Typed(typed) => typed.borrow(),
74            ColumnBody::Words(words) => borrow_words::<C>(words),
75        }
76    }
77
78    /// The number of records.
79    #[inline]
80    pub fn len(&self) -> usize {
81        match self {
82            ColumnBody::Typed(typed) => typed.len(),
83            ColumnBody::Words(words) => borrow_words::<C>(words).len(),
84        }
85    }
86
87    /// True when the body holds no records.
88    #[inline]
89    pub fn is_empty(&self) -> bool {
90        self.len() == 0
91    }
92
93    /// The size of the serialized form in bytes, whether or not the body is serialized.
94    pub fn length_in_bytes(&self) -> usize {
95        match self {
96            ColumnBody::Typed(typed) => indexed::length_in_bytes(&typed.borrow()),
97            ColumnBody::Words(words) => words.len() * std::mem::size_of::<u64>(),
98        }
99    }
100
101    /// Writes the serialized form, exactly [`ColumnBody::length_in_bytes`] bytes of it.
102    ///
103    /// The serialized form is written as it is, so a body that went through a store and back
104    /// round-trips byte-identically.
105    pub fn write_into<W: Write>(&self, writer: &mut W) -> std::io::Result<()> {
106        match self {
107            ColumnBody::Typed(typed) => indexed::write(writer, &typed.borrow()),
108            ColumnBody::Words(words) => writer.write_all(bytemuck::cast_slice(words)),
109        }
110    }
111
112    /// The typed containers, for writing to.
113    ///
114    /// A serialized body is copied into typed containers first, since the serialized form cannot
115    /// take pushes. The copy is a bulk per-leaf extension.
116    pub fn typed_mut(&mut self) -> &mut C::Container {
117        if let ColumnBody::Words(words) = &*self {
118            let typed = copy_typed::<C>(borrow_words::<C>(words));
119            *self = ColumnBody::Typed(typed);
120        }
121        let ColumnBody::Typed(typed) = self else {
122            unreachable!("a serialized body was materialized above");
123        };
124        typed
125    }
126
127    /// The typed containers, taking a typed body as it is and copying a serialized one.
128    pub fn into_typed(self) -> C::Container {
129        match self {
130            ColumnBody::Typed(typed) => typed,
131            ColumnBody::Words(words) => copy_typed::<C>(borrow_words::<C>(&words)),
132        }
133    }
134
135    /// An owned copy in the same form: serialized words are cloned, typed containers are copied
136    /// by bulk per-leaf extension.
137    pub fn duplicate(&self) -> Self {
138        match self {
139            ColumnBody::Typed(typed) => ColumnBody::Typed(copy_typed::<C>(typed.borrow())),
140            ColumnBody::Words(words) => ColumnBody::Words(words.clone()),
141        }
142    }
143
144    /// Empties the body, keeping a typed body's allocations for refilling.
145    ///
146    /// A serialized body owns no typed allocation, so it becomes an empty typed body.
147    #[inline]
148    pub fn clear(&mut self) {
149        match self {
150            ColumnBody::Typed(typed) => typed.clear(),
151            ColumnBody::Words(_) => *self = Default::default(),
152        }
153    }
154
155    /// True once the body is at the ship size the builder and merger cut chunks at.
156    ///
157    /// A serialized body is complete, so it is always at capacity.
158    #[inline]
159    pub fn at_capacity(&self) -> bool {
160        match self {
161            ColumnBody::Typed(typed) => at_serialized_capacity(&typed.borrow()),
162            ColumnBody::Words(_) => true,
163        }
164    }
165}
166
167/// Reconstructs the borrowed columnar view from serialized words, the same zero-copy decode
168/// [`Column::borrow`] performs on its `Align` variant.
169pub fn borrow_words<C: Columnar>(words: &[u64]) -> BorrowedOf<'_, C> {
170    <BorrowedOf<'_, C>>::from_bytes(&mut indexed::decode(words))
171}
172
173/// Copies a view into fresh typed containers by bulk per-leaf extension.
174fn copy_typed<C: Columnar>(view: BorrowedOf<'_, C>) -> C::Container {
175    let mut fresh = C::Container::default();
176    fresh.extend_from_self(view, 0..view.len());
177    fresh
178}
179
180impl<C: Columnar> Accountable for ColumnBody<C> {
181    #[inline]
182    fn record_count(&self) -> i64 {
183        i64::try_from(self.len()).expect("record count fits i64")
184    }
185}
186
187impl<C: Columnar> DrainContainer for ColumnBody<C> {
188    type Item<'a> = Ref<'a, C>;
189    type DrainIter<'a> = IterOwn<BorrowedOf<'a, C>>;
190    #[inline]
191    fn drain(&mut self) -> Self::DrainIter<'_> {
192        self.borrow().into_index_iter()
193    }
194}
195
196impl<C: Columnar, T> PushInto<T> for ColumnBody<C>
197where
198    C::Container: columnar::Push<T>,
199{
200    #[inline]
201    fn push_into(&mut self, item: T) {
202        use columnar::Push;
203        self.typed_mut().push(item);
204    }
205}
206
207/// A body leaves the edge container behind: typed data moves, and channel bytes are relocated
208/// into owned words, since a body never holds a channel's allocation.
209impl<C: Columnar> From<Column<C>> for ColumnBody<C> {
210    fn from(column: Column<C>) -> Self {
211        match column {
212            Column::Typed(typed) => ColumnBody::Typed(typed),
213            Column::Bytes(bytes) => {
214                assert_eq!(bytes.len() % 8, 0);
215                ColumnBody::Words(bytemuck::allocation::pod_collect_to_vec(&bytes))
216            }
217            Column::Align(words) => ColumnBody::Words(words),
218        }
219    }
220}
221
222/// A body goes back onto an edge as it is: both of its forms are forms of the edge container.
223impl<C: Columnar> From<ColumnBody<C>> for Column<C> {
224    fn from(body: ColumnBody<C>) -> Self {
225        match body {
226            ColumnBody::Typed(typed) => Column::Typed(typed),
227            ColumnBody::Words(words) => Column::Align(words),
228        }
229    }
230}
231
232#[cfg(test)]
233mod tests {
234    use super::*;
235
236    fn body(values: &[i32]) -> ColumnBody<i32> {
237        let mut body: ColumnBody<i32> = Default::default();
238        for value in values {
239            body.push_into(*value);
240        }
241        body
242    }
243
244    fn collect(body: &ColumnBody<i32>) -> Vec<i32> {
245        body.borrow().into_index_iter().copied().collect()
246    }
247
248    fn serialized(body: &ColumnBody<i32>) -> ColumnBody<i32> {
249        let mut bytes = Vec::new();
250        body.write_into(&mut bytes).expect("vec writes");
251        assert_eq!(bytes.len(), body.length_in_bytes());
252        ColumnBody::Words(bytemuck::allocation::pod_collect_to_vec(&bytes))
253    }
254
255    #[mz_ore::test]
256    fn serialized_body_reads_and_measures_like_typed() {
257        let typed = body(&[1, 2, 3]);
258        let words = serialized(&typed);
259        assert_eq!(collect(&words), vec![1, 2, 3]);
260        assert_eq!(words.len(), 3);
261        assert_eq!(words.length_in_bytes(), typed.length_in_bytes());
262        assert!(words.at_capacity(), "a serialized body is complete");
263    }
264
265    #[mz_ore::test]
266    fn serialized_body_round_trips_byte_identically() {
267        let words = serialized(&body(&[4, 5, 6]));
268        let again = serialized(&words);
269        let (ColumnBody::Words(a), ColumnBody::Words(b)) = (&words, &again) else {
270            panic!("serialized bodies are words");
271        };
272        assert_eq!(a, b);
273    }
274
275    #[mz_ore::test]
276    fn writing_a_serialized_body_materializes_it() {
277        let mut words = serialized(&body(&[7, 8]));
278        words.push_into(9);
279        assert!(matches!(words, ColumnBody::Typed(_)));
280        assert_eq!(collect(&words), vec![7, 8, 9]);
281    }
282
283    #[mz_ore::test]
284    fn duplicating_keeps_the_form() {
285        let typed = body(&[1, 2]);
286        assert!(matches!(typed.duplicate(), ColumnBody::Typed(_)));
287        let words = serialized(&typed);
288        let copy = words.duplicate();
289        assert!(matches!(copy, ColumnBody::Words(_)));
290        assert_eq!(collect(&copy), vec![1, 2]);
291    }
292
293    #[mz_ore::test]
294    fn clearing_keeps_a_typed_body_typed() {
295        let mut typed = body(&[1]);
296        typed.clear();
297        assert!(matches!(typed, ColumnBody::Typed(_)));
298        assert!(typed.is_empty());
299        let mut words = serialized(&body(&[1]));
300        words.clear();
301        assert!(matches!(words, ColumnBody::Typed(_)));
302        assert!(words.is_empty());
303    }
304
305    #[mz_ore::test]
306    fn edge_conversions_move_typed_and_relocate_bytes() {
307        let column: Column<i32> = Column::from(body(&[1, 2]));
308        assert!(matches!(column, Column::Typed(_)));
309        let column: Column<i32> = Column::from(serialized(&body(&[1, 2])));
310        assert!(matches!(column, Column::Align(_)));
311        let back = ColumnBody::from(column);
312        assert!(matches!(back, ColumnBody::Words(_)));
313        assert_eq!(collect(&back), vec![1, 2]);
314    }
315}