Skip to main content

mz_timely_util/
columnar.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//! Container for columnar data.
17
18#![deny(missing_docs)]
19
20pub mod batcher;
21pub mod body;
22pub mod builder;
23pub mod builder_input;
24pub mod chunk;
25pub mod consolidate;
26pub mod merge_batcher;
27pub mod unload;
28
29use std::hash::{BuildHasher, Hash, Hasher};
30use std::sync::LazyLock;
31
32use columnar::Borrow;
33use columnar::bytes::indexed;
34use columnar::common::IterOwn;
35use columnar::{Clear, FromBytes, Index, Len};
36use columnar::{Columnar, Ref};
37use differential_dataflow::Hashable;
38use differential_dataflow::collection::containers::{Enter, ResultsIn};
39use differential_dataflow::trace::implementations::merge_batcher::MergeBatcher;
40use itertools::Itertools;
41use timely::Accountable;
42use timely::bytes::arc::Bytes;
43use timely::container::{DrainContainer, PushInto, SizableContainer};
44use timely::dataflow::channels::ContainerBytes;
45use timely::progress::Timestamp;
46use timely::progress::timestamp::Refines;
47
48use crate::columnation::ColInternalMerger;
49
50/// A batcher for columnar storage.
51///
52/// The chunker is supplied to the arrange operator separately. Callers pass
53/// it explicitly: [`ColumnationChunker`](crate::columnation::ColumnationChunker)
54/// for `Vec<_>` input, or [`batcher::Chunker`] (over a `ColumnationStack<_>`) for
55/// [`Column`] input.
56pub type Col2ValBatcher<K, V, T, R> = MergeBatcher<ColInternalMerger<(K, V), T, R>>;
57/// A batcher for columnar storage with unit values.
58pub type Col2KeyBatcher<K, T, R> = Col2ValBatcher<K, (), T, R>;
59
60/// Pageable counterpart to [`Col2ValBatcher`]. Routes every chunk produced
61/// by chunking, merging, or extract through a [`crate::column_pager::ColumnPager`],
62/// so memory pressure can spill chains to a backing store without touching
63/// the merge / extract bodies.
64///
65/// Drop-in shape at the type level: both aliases take `(K, V, T, R)` and
66/// produce a `Batcher<Input = Column<((K, V), T, R)>, Output = Column<((K,
67/// V), T, R)>>`. Call sites can swap with `cargo fix`–style renaming once
68/// downstream `Trace`/`Builder` impls have been wired up. The pager itself
69/// defaults to [`crate::column_pager::ColumnPager::disabled`]; inject a
70/// real one via [`merge_batcher::ColumnMergeBatcher::set_pager`].
71pub type Col2ValPagedBatcher<K, V, T, R> = merge_batcher::ColumnMergeBatcher<(K, V), T, R>;
72
73/// Columnar-native counterpart to [`Col2ValBatcher`], holding
74/// [`body::ColumnBody`] chunks rather than columnation stacks and merging
75/// them through [`batcher::ColumnMerger`].
76///
77/// Pairs with [`batcher::ColumnChunker`] and any builder whose `Input` is
78/// `ColumnBody<((K, V), T, R)>`. Unlike [`Col2ValPagedBatcher`] the chains
79/// stay resident, so this arm carries no pager and no spill budget.
80pub type Col2ValColBatcher<K, V, T, R> = MergeBatcher<batcher::ColumnMerger<(K, V), T, R>>;
81
82/// A container based on a columnar store, encoded in aligned bytes: the container on dataflow
83/// edges.
84///
85/// The type can represent typed data, bytes from Timely, or an aligned allocation. The name
86/// is singular to express that the preferred format is [`Column::Align`]. The [`Column::Typed`]
87/// variant is used to construct the container, and it owns potentially multiple columns of data.
88///
89/// Data at rest behind an edge, in a merge batcher's chains, a spill-backed chunk, or a batch
90/// builder, is a [`body::ColumnBody`] instead, which has no channel-bytes form.
91pub enum Column<C: Columnar> {
92    /// The typed variant of the container.
93    Typed(C::Container),
94    /// The binary variant of the container.
95    Bytes(Bytes),
96    /// Relocated, aligned binary data, if `Bytes` doesn't work for some reason.
97    ///
98    /// Reasons could include misalignment, cloning of data, or wanting
99    /// to release the `Bytes` as a scarce resource.
100    ///
101    /// `Vec<u64>` guarantees `u64` alignment for the contained bytes.
102    Align(Vec<u64>),
103}
104
105impl<C: Columnar> Column<C> {
106    /// Empties the column, retaining the `Typed` variant's allocation so the
107    /// caller can refill it.
108    ///
109    /// [`columnar::Clear`] clears the typed container in place without
110    /// releasing its capacity. The serialized variants (`Bytes`/`Align`) own
111    /// no reusable typed buffer, so they are reset to an empty `Typed`.
112    #[inline]
113    pub fn clear(&mut self) {
114        match self {
115            Column::Typed(t) => t.clear(),
116            Column::Bytes(_) | Column::Align(_) => *self = Default::default(),
117        }
118    }
119
120    /// True when the column holds no records.
121    ///
122    /// The `Typed` variant answers from the container itself. The serialized
123    /// variants have to reconstruct their borrowed view to reach a length, so
124    /// there this costs as much as [`Column::borrow`].
125    #[inline]
126    pub fn is_empty(&self) -> bool {
127        match self {
128            Column::Typed(t) => t.is_empty(),
129            Column::Bytes(_) | Column::Align(_) => self.borrow().is_empty(),
130        }
131    }
132
133    /// Borrows the container as a reference.
134    // NOTE: `inline(always)` measured 7.1 ns to 5.8 ns per call for a per-element borrow of an
135    // `Align` column. The header decode stays per call, so loops still hoist the borrow.
136    #[inline(always)]
137    pub fn borrow(&self) -> <C::Container as Borrow>::Borrowed<'_> {
138        match self {
139            Column::Typed(t) => t.borrow(),
140            Column::Bytes(b) => <<C::Container as Borrow>::Borrowed<'_>>::from_bytes(
141                &mut indexed::decode(bytemuck::cast_slice(b)),
142            ),
143            Column::Align(a) => {
144                <<C::Container as Borrow>::Borrowed<'_>>::from_bytes(&mut indexed::decode(a))
145            }
146        }
147    }
148}
149
150impl<C: Columnar> Default for Column<C> {
151    fn default() -> Self {
152        Self::Typed(Default::default())
153    }
154}
155
156impl<C: Columnar> Clone for Column<C>
157where
158    C::Container: Clone,
159{
160    fn clone(&self) -> Self {
161        match self {
162            // Typed stays typed, although we would have the option to move to aligned data.
163            // If we did it might be confusing why we couldn't push into a cloned column.
164            Column::Typed(t) => Column::Typed(t.clone()),
165            Column::Bytes(b) => {
166                assert_eq!(b.len() % 8, 0);
167                Self::Align(bytemuck::allocation::pod_collect_to_vec(b))
168            }
169            Column::Align(a) => Column::Align(a.clone()),
170        }
171    }
172}
173
174impl<C: Columnar> Accountable for Column<C> {
175    #[inline]
176    fn record_count(&self) -> i64 {
177        self.borrow().len().try_into().expect("Must fit")
178    }
179}
180impl<C: Columnar> DrainContainer for Column<C> {
181    type Item<'a> = Ref<'a, C>;
182    type DrainIter<'a> = IterOwn<<C::Container as Borrow>::Borrowed<'a>>;
183    #[inline]
184    fn drain(&mut self) -> Self::DrainIter<'_> {
185        self.borrow().into_index_iter()
186    }
187}
188
189impl<C: Columnar, T> PushInto<T> for Column<C>
190where
191    C::Container: columnar::Push<T>,
192{
193    #[inline]
194    fn push_into(&mut self, item: T) {
195        use columnar::Push;
196        match self {
197            Column::Typed(t) => t.push(item),
198            Column::Align(_) | Column::Bytes(_) => {
199                // We really oughtn't be calling this in this case.
200                // We could convert to owned, but need more constraints on `C`.
201                unimplemented!("Pushing into Column::Bytes without first clearing");
202            }
203        }
204    }
205}
206
207/// Re-encodes a column's times into a refining timestamp, so a columnar collection can enter
208/// an iterative scope.
209impl<D, T1, T2, R> Enter<T1, T2> for Column<(D, T1, R)>
210where
211    D: Columnar,
212    T1: Columnar + Timestamp,
213    T2: Columnar + Refines<T1>,
214    R: Columnar,
215    (D, T1, R): Columnar<Container = (D::Container, T1::Container, R::Container)>,
216    (D, T2, R): Columnar<Container = (D::Container, T2::Container, R::Container)>,
217    for<'a> D::Container: columnar::Push<Ref<'a, D>>,
218    for<'a> T2::Container: columnar::Push<&'a T2>,
219    for<'a> R::Container: columnar::Push<Ref<'a, R>>,
220{
221    type InnerContainer = Column<(D, T2, R)>;
222
223    fn enter(self) -> Self::InnerContainer {
224        use columnar::Push;
225        match self {
226            // Only the times change, so the data and diff columns move across whole.
227            Column::Typed((data, times, diffs)) => {
228                let mut inner = T2::Container::default();
229                for time in times.borrow().into_index_iter() {
230                    inner.push(&T2::to_inner(T1::into_owned(time)));
231                }
232                Column::Typed((data, inner, diffs))
233            }
234            // A serialized column owns no typed sub-containers, so only the times are
235            // materialized. The data and diff columns go from their borrowed views
236            // straight into the output allocation, a copy per column rather than a
237            // decode and re-encode per record.
238            serialized => {
239                let (borrowed_data, borrowed_times, borrowed_diffs) = serialized.borrow();
240                let mut times = T2::Container::default();
241                for time in borrowed_times.into_index_iter() {
242                    times.push(&T2::to_inner(T1::into_owned(time)));
243                }
244                let view = (borrowed_data, times.borrow(), borrowed_diffs);
245                let words = indexed::length_in_words(&view);
246                let mut alloc: Vec<u64> = Vec::with_capacity(words);
247                indexed::encode(&mut alloc, &view);
248                Column::Align(alloc)
249            }
250        }
251    }
252}
253
254/// Advances a column's times by `step`, so a columnar collection can be a loop variable.
255///
256/// A time that the summary does not advance drops its record, so the data and diff columns
257/// cannot move across whole the way [`Enter`] moves them: a record's position in one column
258/// has to keep matching its position in the others. The times are read and advanced first,
259/// and the rest is rebuilt only when the step actually dropped something.
260impl<D, T, R> ResultsIn<T::Summary> for Column<(D, T, R)>
261where
262    D: Columnar,
263    T: Columnar + Timestamp,
264    R: Columnar,
265    (D, T, R): Columnar<Container = (D::Container, T::Container, R::Container)>,
266    for<'a> D::Container: columnar::Push<Ref<'a, D>>,
267    for<'a> T::Container: columnar::Push<&'a T>,
268    for<'a> R::Container: columnar::Push<Ref<'a, R>>,
269{
270    fn results_in(self, step: &T::Summary) -> Self {
271        use columnar::Push;
272        use timely::progress::PathSummary;
273
274        // Advance the times first. Whether the step dropped any record decides whether the
275        // data and diff columns still line up with the new times, which is what lets them
276        // be reused rather than rebuilt.
277        let (times, kept) = {
278            let (_, borrowed_times, _) = self.borrow();
279            let mut times = T::Container::default();
280            let mut time = T::minimum();
281            let kept: Vec<bool> = borrowed_times
282                .into_index_iter()
283                .map(|reference| {
284                    time.copy_from(reference);
285                    match step.results_in(&time) {
286                        Some(time) => {
287                            times.push(&time);
288                            true
289                        }
290                        None => false,
291                    }
292                })
293                .collect();
294            (times, kept)
295        };
296
297        if kept.iter().all(|kept| *kept) {
298            match self {
299                // Only the times changed, so the data and diff columns move across whole.
300                Column::Typed((data, _, diffs)) => Column::Typed((data, times, diffs)),
301                // A serialized column owns no typed sub-containers, so its data and diff
302                // columns go from their borrowed views straight into the output allocation.
303                serialized => {
304                    let (borrowed_data, _, borrowed_diffs) = serialized.borrow();
305                    let view = (borrowed_data, times.borrow(), borrowed_diffs);
306                    let words = indexed::length_in_words(&view);
307                    let mut alloc: Vec<u64> = Vec::with_capacity(words);
308                    indexed::encode(&mut alloc, &view);
309                    Column::Align(alloc)
310                }
311            }
312        } else {
313            let (borrowed_data, _, borrowed_diffs) = self.borrow();
314            let mut data = D::Container::default();
315            let mut diffs = R::Container::default();
316            let records = borrowed_data
317                .into_index_iter()
318                .zip_eq(borrowed_diffs.into_index_iter())
319                .zip_eq(kept.iter());
320            for ((datum, diff), _) in records.filter(|(_, kept)| **kept) {
321                data.push(datum);
322                diffs.push(diff);
323            }
324            Column::Typed((data, times, diffs))
325        }
326    }
327}
328
329/// Words per 2 MiB. `length_in_words` returns serialized size in `u64` units,
330/// so this is the page count we round up to. Picked to match
331/// [`builder::ColumnBuilder`]'s output granularity so chunks shipped from the
332/// merger and chunks shipped from the builder are sized comparably.
333const SHIP_WORDS: usize = 1 << 18;
334
335/// Returns true once the serialized size of `borrow` reaches 10% under
336/// `SHIP_WORDS`.
337///
338/// Monotone in size, deliberately not a window below the boundary. A single
339/// record wider than a window steps clear over it, and a ship signal that
340/// un-fires past the boundary lets a chunk grow until it exceeds the buffer
341/// pool's largest size class, past which a spilled body degrades to
342/// permanently resident. The same heuristic as [`builder::ColumnBuilder`]'s
343/// ship point, lifted out so the builder, the merger, and the
344/// `SizableContainer` impl agree on the signal.
345#[inline]
346pub(crate) fn at_serialized_capacity<'a, A>(borrow: &A) -> bool
347where
348    A: columnar::AsBytes<'a>,
349{
350    indexed::length_in_words(borrow) >= SHIP_WORDS - SHIP_WORDS / 10
351}
352
353impl<C: Columnar> SizableContainer for Column<C> {
354    fn at_capacity(&self) -> bool {
355        // Match `ColumnBuilder`'s ship heuristic: serialized size at the
356        // 2 MiB ship threshold. Aligns chunk-size choices across the two
357        // paths and keeps recipients dealing with a single granularity.
358        //
359        // Serialized chunks (`Bytes` / `Align`) have no typed builder to push
360        // into, so they're trivially "at capacity" — there's no further work
361        // they can absorb.
362        match self {
363            Column::Typed(c) => at_serialized_capacity(&c.borrow()),
364            Column::Bytes(_) | Column::Align(_) => true,
365        }
366    }
367
368    fn ensure_capacity(&mut self, _stash: &mut Option<Self>) {
369        // No pre-reservation: chunks are recycled by the merge framework, so
370        // leaf capacities settle to steady-state after the first round and
371        // there is nothing useful to reserve up front. The `SizableContainer`
372        // impl exists so `at_capacity` is callable on result chunks during
373        // `Merger::merge` orchestration; `ensure_capacity` is a required
374        // method on the trait but has no work to do here.
375    }
376}
377
378impl<C: Columnar> ContainerBytes for Column<C> {
379    #[inline]
380    fn from_bytes(bytes: Bytes) -> Self {
381        // Our expectation / hope is that `bytes` is `u64` aligned and sized.
382        // If the alignment is borked, we can relocate. If the size is borked,
383        // not sure what we do in that case. An incorrect size indicates a problem
384        // of `into_bytes`, or a failure of the communication layer, both of which
385        // are unrecoverable.
386        assert_eq!(bytes.len() % 8, 0);
387        if let Ok(_) = bytemuck::try_cast_slice::<_, u64>(&bytes) {
388            Self::Bytes(bytes)
389        } else {
390            // We failed to cast the slice, so we'll reallocate. `Vec<u64>`
391            // is u64-aligned by construction.
392            Self::Align(bytemuck::allocation::pod_collect_to_vec(&bytes[..]))
393        }
394    }
395
396    #[inline]
397    fn length_in_bytes(&self) -> usize {
398        match self {
399            Column::Typed(t) => indexed::length_in_bytes(&t.borrow()),
400            Column::Bytes(b) => b.len(),
401            Column::Align(a) => 8 * a.len(),
402        }
403    }
404
405    #[inline]
406    fn into_bytes<W: ::std::io::Write>(&self, writer: &mut W) {
407        match self {
408            Column::Typed(t) => indexed::write(writer, &t.borrow()).unwrap(),
409            Column::Bytes(b) => writer.write_all(b).unwrap(),
410            Column::Align(a) => writer.write_all(bytemuck::cast_slice(a)).unwrap(),
411        }
412    }
413}
414
415/// An exchange function for columnar tuples of the form `((K, V), T, D)`. Rust has a hard
416/// time to figure out the lifetimes of the elements when specified as a closure, so we rather
417/// specify it as a function.
418#[inline(always)]
419pub fn columnar_exchange<K, V, T, D>(((k, _), _, _): &Ref<'_, ((K, V), T, D)>) -> u64
420where
421    K: Columnar,
422    for<'a> Ref<'a, K>: Hash,
423    V: Columnar,
424    D: Columnar,
425    T: Columnar,
426{
427    k.hashed()
428}
429
430/// Routes a `(D, T, R)` column by the hash of its data column.
431///
432/// Counterpart to [`columnar_exchange`] for collections whose data is not a
433/// key/value pair. Consolidation sites want [`columnar_consolidate_exchange`]
434/// instead.
435pub fn columnar_exchange_data<D, T, R>((d, _, _): &Ref<'_, (D, T, R)>) -> u64
436where
437    D: Columnar,
438    for<'a> Ref<'a, D>: Hash,
439    T: Columnar,
440    R: Columnar,
441{
442    d.hashed()
443}
444
445/// Routes a `(D, T, R)` column for consolidation, by a fixed-seed AHash of its data column.
446///
447/// Worker assignment is `hash % workers`, so the low bits alone decide the split and the
448/// [`Hashable`] default (FNV) the other exchange functions use diffuses them poorly. The
449/// seed is fixed, so routing is identical across builds and replicas.
450///
451/// Spelled as a function rather than a closure over the hasher state, because the argument
452/// is higher-ranked in its lifetime and closure inference cannot express that.
453pub fn columnar_consolidate_exchange<D, T, R>((d, _, _): &Ref<'_, (D, T, R)>) -> u64
454where
455    D: Columnar,
456    for<'a> Ref<'a, D>: Hash,
457    T: Columnar,
458    R: Columnar,
459{
460    static STATE: LazyLock<ahash::RandomState> = LazyLock::new(crate::hash::fixed_state);
461    let mut hasher = STATE.build_hasher();
462    d.hash(&mut hasher);
463    hasher.finish()
464}
465
466#[cfg(test)]
467mod tests {
468    use timely::bytes::arc::BytesMut;
469    use timely::container::PushInto;
470    use timely::dataflow::channels::ContainerBytes;
471
472    use super::*;
473
474    /// Produce some bytes that are in columnar format.
475    fn raw_columnar_bytes() -> Vec<u8> {
476        let mut raw = Vec::new();
477        raw.extend(16_u64.to_le_bytes()); // offsets
478        raw.extend(28_u64.to_le_bytes()); // length
479        raw.extend(1_i32.to_le_bytes());
480        raw.extend(2_i32.to_le_bytes());
481        raw.extend(3_i32.to_le_bytes());
482        raw.extend([0, 0, 0, 0]); // padding
483        raw
484    }
485
486    #[mz_ore::test]
487    fn test_column_clone() {
488        let columns = Columnar::as_columns([1, 2, 3].iter());
489        let column_typed: Column<i32> = Column::Typed(columns);
490        let column_typed2 = column_typed.clone();
491
492        assert_eq!(
493            column_typed2.borrow().into_index_iter().collect::<Vec<_>>(),
494            vec![&1, &2, &3]
495        );
496
497        let bytes = BytesMut::from(raw_columnar_bytes()).freeze();
498        let column_bytes: Column<i32> = Column::Bytes(bytes);
499        let column_bytes2 = column_bytes.clone();
500
501        assert_eq!(
502            column_bytes2.borrow().into_index_iter().collect::<Vec<_>>(),
503            vec![&1, &2, &3]
504        );
505
506        let raw = raw_columnar_bytes();
507        let mut region: Vec<u64> = vec![0; raw.len() / 8];
508        let region_bytes = bytemuck::cast_slice_mut(&mut region[..]);
509        region_bytes[..raw.len()].copy_from_slice(&raw);
510        let column_align: Column<i32> = Column::Align(region);
511        let column_align2 = column_align.clone();
512
513        assert_eq!(
514            column_align2.borrow().into_index_iter().collect::<Vec<_>>(),
515            vec![&1, &2, &3]
516        );
517    }
518
519    /// Assert the desired contents of raw_columnar_bytes so that diagnosing test failures is
520    /// easier.
521    #[mz_ore::test]
522    fn test_column_known_bytes() {
523        let mut column: Column<i32> = Default::default();
524        column.push_into(1);
525        column.push_into(2);
526        column.push_into(3);
527        let mut data = Vec::new();
528        column.into_bytes(&mut std::io::Cursor::new(&mut data));
529        assert_eq!(data, raw_columnar_bytes());
530    }
531
532    #[mz_ore::test]
533    fn test_column_from_bytes() {
534        let raw = raw_columnar_bytes();
535
536        let buf = vec![0; raw.len() + 8];
537        let align = buf.as_ptr().align_offset(std::mem::size_of::<u64>());
538        let mut bytes_mut = BytesMut::from(buf);
539        let _ = bytes_mut.extract_to(align);
540        bytes_mut[..raw.len()].copy_from_slice(&raw);
541        let aligned_bytes = bytes_mut.extract_to(raw.len());
542
543        let column: Column<i32> = Column::from_bytes(aligned_bytes);
544        assert!(matches!(column, Column::Bytes(_)));
545        assert_eq!(
546            column.borrow().into_index_iter().collect::<Vec<_>>(),
547            vec![&1, &2, &3]
548        );
549
550        let buf = vec![0; raw.len() + 8];
551        let align = buf.as_ptr().align_offset(std::mem::size_of::<u64>());
552        let mut bytes_mut = BytesMut::from(buf);
553        let _ = bytes_mut.extract_to(align + 1);
554        bytes_mut[..raw.len()].copy_from_slice(&raw);
555        let unaligned_bytes = bytes_mut.extract_to(raw.len());
556
557        let column: Column<i32> = Column::from_bytes(unaligned_bytes);
558        assert!(matches!(column, Column::Align(_)));
559        assert_eq!(
560            column.borrow().into_index_iter().collect::<Vec<_>>(),
561            vec![&1, &2, &3]
562        );
563    }
564
565    /// The ship signal is monotone: once it fires it stays fired, even when
566    /// a single wide record steps far past the 2 MiB boundary in one push.
567    #[mz_ore::test]
568    #[cfg_attr(miri, ignore)] // too slow
569    fn ship_threshold_monotone() {
570        use columnar::Push;
571        let mut container = <Vec<u64> as Columnar>::Container::default();
572        // Wider than 10% of any boundary a 25 MiB run can reach.
573        let wide: Vec<u64> = vec![0u64; 50_000];
574        let mut fired = false;
575        for pushes in 1..=64 {
576            container.push(&wide);
577            let now = at_serialized_capacity(&container.borrow());
578            if fired {
579                assert!(now, "ship signal un-fired at {pushes} records");
580            }
581            fired = fired || now;
582        }
583        assert!(fired, "ship signal never fired");
584    }
585}