Skip to main content

mz_compute/extensions/
arrange.rs

1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Use of this software is governed by the Business Source License
4// included in the LICENSE file.
5//
6// As of the Change Date specified in that file, in accordance with
7// the Business Source License, use of this software will be governed
8// by the Apache License, Version 2.0.
9
10use std::collections::BTreeMap;
11use std::rc::Rc;
12use std::sync::{Arc, Weak};
13
14use differential_dataflow::difference::Semigroup;
15use differential_dataflow::lattice::Lattice;
16use differential_dataflow::operators::arrange::arrangement::arrange_core;
17use differential_dataflow::operators::arrange::{Arranged, TraceAgent};
18use differential_dataflow::trace::implementations::BatchContainer;
19use differential_dataflow::trace::implementations::spine_fueled::Spine;
20use differential_dataflow::trace::{Batch, Batcher, Builder, Trace, TraceReader};
21use differential_dataflow::{Collection, Data, ExchangeData, Hashable, VecCollection};
22use mz_compute_types::dyncfgs::{ENABLE_COLUMN_PAGED_BATCHER, ENABLE_COLUMNAR_MERGE_BATCHER};
23use mz_dyncfg::ConfigSet;
24use mz_row_spine::ArcBatch;
25use mz_timely_util::containers::HeapSize;
26use timely::Container;
27use timely::container::{ContainerBuilder, PushInto};
28use timely::dataflow::Stream;
29use timely::dataflow::channels::pact::{Exchange, ParallelizationContract, Pipeline};
30use timely::dataflow::operators::Operator;
31use timely::progress::Timestamp;
32
33use crate::logging::compute::{
34    ArrangementHeapAllocations, ArrangementHeapCapacity, ArrangementHeapSize,
35    ArrangementHeapSizeOperator, ComputeEvent, ComputeEventBuilder,
36};
37use crate::typedefs::{
38    KeyAgent, KeyValAgent, MzArrangeData, MzData, MzTimestamp, RowAgent, RowRowAgent, RowValAgent,
39};
40
41/// Which merge batcher an arrange site should instantiate.
42///
43/// The three parameters `mz_arrange_core` takes are one unit, not three
44/// knobs: the chunker's container and the builder's input are both pinned to
45/// `Batcher::Output`, so the chunker and builder follow from the batcher and
46/// a call site has to spell out a whole arm per variant.
47pub enum ArrangementBatcher {
48    /// `Chunker<ColumnationStack<_>>` + `Col2ValBatcher` + `RowRowBuilder`.
49    /// Chains are columnation stacks.
50    Columnation,
51    /// `ColumnChunker` + `Col2ValColBatcher` + `RowRowColPagedBuilder`.
52    /// Chains are resident `Column`s.
53    Columnar,
54    /// `ChunkChunker` + `ChunkBatcher` + `UnchunkBuilder<RowRowColPagedBuilder>`.
55    /// Chains are chunks whose bodies the process buffer pool spills while
56    /// the process chunk spill gate is set, and that stay resident
57    /// otherwise.
58    Chunked,
59}
60
61impl ArrangementBatcher {
62    /// Resolve the batcher from the replica's config set.
63    ///
64    /// `ENABLE_COLUMN_PAGED_BATCHER` wins over
65    /// `ENABLE_COLUMNAR_MERGE_BATCHER`, because it asks for the same columnar
66    /// chains plus the option to spill them. Call this once per arrange site
67    /// at operator construction time, so a dataflow keeps one batcher for its
68    /// whole life even if the flags flip underneath it.
69    pub fn from_config(config: &ConfigSet) -> Self {
70        if ENABLE_COLUMN_PAGED_BATCHER.get(config) {
71            Self::Chunked
72        } else if ENABLE_COLUMNAR_MERGE_BATCHER.get(config) {
73            Self::Columnar
74        } else {
75            Self::Columnation
76        }
77    }
78}
79
80/// Extension trait to arrange data.
81pub trait MzArrange<'scope>: MzArrangeCore<'scope> {
82    /// Arranges a stream of `(Key, Val)` updates by `Key` into a trace of type `Tr`.
83    ///
84    /// This operator arranges a stream of values into a shared trace, whose contents it maintains.
85    /// This trace is current for all times marked completed in the output stream, and probing this stream
86    /// is the correct way to determine that times in the shared trace are committed.
87    fn mz_arrange<Chu, Ba, Bu, Tr>(self, name: &str) -> Arranged<'scope, TraceAgent<Tr>>
88    where
89        Ba: Batcher<Time = Self::Timestamp> + 'static,
90        Chu: ContainerBuilder<Container = Ba::Output>
91            + for<'a> PushInto<&'a mut Self::Input>
92            + 'static,
93        Bu: Builder<Time = Self::Timestamp, Input = Ba::Output, Output = Tr::Batch>,
94        Tr: Trace + TraceReader<Time = Self::Timestamp> + 'static,
95        Tr::Batch: Batch,
96        Arranged<'scope, TraceAgent<Tr>>: ArrangementSize;
97}
98
99/// Extension trait to arrange data.
100pub trait MzArrangeCore<'scope> {
101    /// The current scope.
102    type Timestamp: Timestamp + Lattice;
103    /// The data input container type.
104    type Input: Container + Clone + 'static;
105
106    /// Arranges a stream of `(Key, Val)` updates by `Key` into a trace of type `Tr`. Partitions
107    /// the data according to `pact`.
108    ///
109    /// This operator arranges a stream of values into a shared trace, whose contents it maintains.
110    /// This trace is current for all times marked completed in the output stream, and probing this stream
111    /// is the correct way to determine that times in the shared trace are committed.
112    fn mz_arrange_core<P, Chu, Ba, Bu, Tr>(
113        self,
114        pact: P,
115        name: &str,
116    ) -> Arranged<'scope, TraceAgent<Tr>>
117    where
118        P: ParallelizationContract<Self::Timestamp, Self::Input>,
119        Ba: Batcher<Time = Self::Timestamp> + 'static,
120        Chu: ContainerBuilder<Container = Ba::Output>
121            + for<'a> PushInto<&'a mut Self::Input>
122            + 'static,
123        Bu: Builder<Time = Self::Timestamp, Input = Ba::Output, Output = Tr::Batch>,
124        Tr: Trace + TraceReader<Time = Self::Timestamp> + 'static,
125        Tr::Batch: Batch,
126        Arranged<'scope, TraceAgent<Tr>>: ArrangementSize;
127}
128
129impl<'scope, T, C> MzArrangeCore<'scope> for Stream<'scope, T, C>
130where
131    T: Timestamp + Lattice,
132    C: Container + Clone + 'static,
133{
134    type Timestamp = T;
135    type Input = C;
136
137    fn mz_arrange_core<P, Chu, Ba, Bu, Tr>(
138        self,
139        pact: P,
140        name: &str,
141    ) -> Arranged<'scope, TraceAgent<Tr>>
142    where
143        P: ParallelizationContract<T, Self::Input>,
144        Ba: Batcher<Time = T> + 'static,
145        Chu: ContainerBuilder<Container = Ba::Output>
146            + for<'a> PushInto<&'a mut Self::Input>
147            + 'static,
148        Bu: Builder<Time = T, Input = Ba::Output, Output = Tr::Batch>,
149        Tr: Trace + TraceReader<Time = T> + 'static,
150        Tr::Batch: Batch,
151        Arranged<'scope, TraceAgent<Tr>>: ArrangementSize,
152    {
153        // Allow access to `arrange_named` because we're within Mz's wrapper.
154        #[allow(clippy::disallowed_methods)]
155        arrange_core::<_, _, Chu, Ba, Bu, _>(self, pact, name).log_arrangement_size()
156    }
157}
158
159impl<'scope, T, K, V, R> MzArrange<'scope> for VecCollection<'scope, T, (K, V), R>
160where
161    T: Timestamp + Lattice,
162    K: ExchangeData + Hashable,
163    V: ExchangeData,
164    R: ExchangeData,
165{
166    fn mz_arrange<Chu, Ba, Bu, Tr>(self, name: &str) -> Arranged<'scope, TraceAgent<Tr>>
167    where
168        Ba: Batcher<Time = T> + 'static,
169        Chu: ContainerBuilder<Container = Ba::Output>
170            + for<'a> PushInto<&'a mut Self::Input>
171            + 'static,
172        Bu: Builder<Time = T, Input = Ba::Output, Output = Tr::Batch>,
173        Tr: Trace + TraceReader<Time = T> + 'static,
174        Tr::Batch: Batch,
175        Arranged<'scope, TraceAgent<Tr>>: ArrangementSize,
176    {
177        let exchange = Exchange::new(move |update: &((K, V), T, R)| (update.0).0.hashed().into());
178        self.mz_arrange_core::<_, Chu, Ba, Bu, _>(exchange, name)
179    }
180}
181
182impl<'scope, T, C> MzArrangeCore<'scope> for Collection<'scope, T, C>
183where
184    T: Timestamp + Lattice,
185    C: Container + Clone + 'static,
186{
187    type Timestamp = T;
188    type Input = C;
189
190    fn mz_arrange_core<P, Chu, Ba, Bu, Tr>(
191        self,
192        pact: P,
193        name: &str,
194    ) -> Arranged<'scope, TraceAgent<Tr>>
195    where
196        P: ParallelizationContract<T, Self::Input>,
197        Ba: Batcher<Time = T> + 'static,
198        Chu: ContainerBuilder<Container = Ba::Output>
199            + for<'a> PushInto<&'a mut Self::Input>
200            + 'static,
201        Bu: Builder<Time = T, Input = Ba::Output, Output = Tr::Batch>,
202        Tr: Trace + TraceReader<Time = T> + 'static,
203        Tr::Batch: Batch,
204        Arranged<'scope, TraceAgent<Tr>>: ArrangementSize,
205    {
206        self.inner.mz_arrange_core::<_, Chu, Ba, Bu, _>(pact, name)
207    }
208}
209
210/// A specialized collection where data only has a key, but no associated value.
211///
212/// Created by calling `collection.into()`.
213pub struct KeyCollection<'scope, T: Timestamp, K: 'static, R: 'static = usize>(
214    VecCollection<'scope, T, K, R>,
215);
216
217impl<'scope, T: Timestamp, K, R: Semigroup> From<VecCollection<'scope, T, K, R>>
218    for KeyCollection<'scope, T, K, R>
219{
220    fn from(value: VecCollection<'scope, T, K, R>) -> Self {
221        KeyCollection(value)
222    }
223}
224
225impl<'scope, T, K, R> MzArrange<'scope> for KeyCollection<'scope, T, K, R>
226where
227    T: Timestamp + Lattice,
228    K: ExchangeData + Hashable,
229    R: ExchangeData,
230{
231    fn mz_arrange<Chu, Ba, Bu, Tr>(self, name: &str) -> Arranged<'scope, TraceAgent<Tr>>
232    where
233        Ba: Batcher<Time = T> + 'static,
234        Chu: ContainerBuilder<Container = Ba::Output>
235            + for<'a> PushInto<&'a mut Self::Input>
236            + 'static,
237        Bu: Builder<Time = T, Input = Ba::Output, Output = Tr::Batch>,
238        Tr: Trace + TraceReader<Time = T> + 'static,
239        Tr::Batch: Batch,
240        Arranged<'scope, TraceAgent<Tr>>: ArrangementSize,
241    {
242        self.0.map(|d| (d, ())).mz_arrange::<Chu, Ba, Bu, _>(name)
243    }
244}
245
246impl<'scope, T, K, R> MzArrangeCore<'scope> for KeyCollection<'scope, T, K, R>
247where
248    T: Timestamp + Lattice,
249    K: Clone + 'static,
250    R: Clone + 'static,
251{
252    type Timestamp = T;
253    type Input = Vec<((K, ()), T, R)>;
254
255    fn mz_arrange_core<P, Chu, Ba, Bu, Tr>(
256        self,
257        pact: P,
258        name: &str,
259    ) -> Arranged<'scope, TraceAgent<Tr>>
260    where
261        P: ParallelizationContract<T, Self::Input>,
262        Ba: Batcher<Time = T> + 'static,
263        Chu: ContainerBuilder<Container = Ba::Output>
264            + for<'a> PushInto<&'a mut Self::Input>
265            + 'static,
266        Bu: Builder<Time = T, Input = Ba::Output, Output = Tr::Batch>,
267        Tr: Trace + TraceReader<Time = T> + 'static,
268        Tr::Batch: Batch,
269        Arranged<'scope, TraceAgent<Tr>>: ArrangementSize,
270    {
271        self.0
272            .map(|d| (d, ()))
273            .mz_arrange_core::<_, Chu, Ba, Bu, _>(pact, name)
274    }
275}
276
277/// A type that can log its heap size.
278pub trait ArrangementSize {
279    /// Install a logger to track the heap size of the target.
280    fn log_arrangement_size(self) -> Self;
281}
282
283/// Helper for [`ArrangementSize`] to install a common operator holding on to a trace.
284///
285/// * `arranged`: The arrangement to inspect.
286/// * `logic`: Closure that calculates the heap size/capacity/allocations for a batch. The return
287///    value are size and capacity in bytes, and number of allocations, all in absolute values.
288///
289/// Batch-size logging identifies each batch by the address of its backing allocation and holds a
290/// weak reference to it, so it needs the `Arc` underlying the spine's [`ArcBatch<B>`] batches;
291/// `batch.0` reaches straight through the newtype to it.
292fn log_arrangement_size_inner<'scope, B, L>(
293    arranged: Arranged<'scope, TraceAgent<Spine<ArcBatch<B>>>>,
294    mut logic: L,
295) -> Arranged<'scope, TraceAgent<Spine<ArcBatch<B>>>>
296where
297    B: Batch + 'static,
298    L: FnMut(&B) -> (usize, usize, usize) + 'static,
299{
300    let scope = arranged.stream.scope();
301    let Some(logger) = scope
302        .worker()
303        .logger_for::<ComputeEventBuilder>("materialize/compute")
304    else {
305        return arranged;
306    };
307    let operator_id = arranged.trace.operator().global_id;
308    let trace = Rc::downgrade(&arranged.trace.trace_box_unstable());
309
310    let (mut old_size, mut old_capacity, mut old_allocations) = (0isize, 0isize, 0isize);
311
312    let stream = arranged
313        .stream
314        .unary(Pipeline, "ArrangementSize", |_cap, info| {
315            let address = info.address;
316            logger.log(&ComputeEvent::ArrangementHeapSizeOperator(
317                ArrangementHeapSizeOperator {
318                    operator_id,
319                    address: address.to_vec(),
320                },
321            ));
322
323            // Weak references to batches, so we can observe batches outside the trace.
324            // Batches are immutable once sealed, so we compute their size exactly
325            // once (when first observed) and cache it alongside the weak reference.
326            // Subsequent activations only sum the cached values for live batches,
327            // avoiding a repeated walk of every batch's backing regions.
328            let mut batches: BTreeMap<*const B, (Weak<B>, (usize, usize, usize))> = BTreeMap::new();
329
330            move |input, output| {
331                input.for_each(|time, data| {
332                    for batch in data.iter() {
333                        batches
334                            .entry(Arc::as_ptr(&batch.0))
335                            .or_insert_with(|| (Arc::downgrade(&batch.0), logic(&batch.0)));
336                    }
337                    output.session(&time).give_container(data);
338                });
339                let Some(trace) = trace.upgrade() else {
340                    // Invariant: `batches` holds no entries once the trace is gone. Each entry's
341                    // `Weak` keeps its batch's `ArcInner` allocation reserved, and the `retain` below
342                    // that would drop it is unreachable on this path, so the entries have to go
343                    // here. The upgrade cannot start succeeding again, hence clearing on every
344                    // activation that takes this path also covers batches that arrive on the
345                    // input afterwards.
346                    batches.clear();
347                    return;
348                };
349
350                trace.borrow().trace().map_batches(|batch| {
351                    batches
352                        .entry(Arc::as_ptr(&batch.0))
353                        .or_insert_with(|| (Arc::downgrade(&batch.0), logic(&batch.0)));
354                });
355
356                let (mut size, mut capacity, mut allocations) = (0, 0, 0);
357                batches.retain(|_, (weak, cached)| {
358                    if weak.strong_count() > 0 {
359                        let (sz, c, a) = *cached;
360                        (size += sz, capacity += c, allocations += a);
361                        true
362                    } else {
363                        false
364                    }
365                });
366
367                let size = size.try_into().expect("must fit");
368                if size != old_size {
369                    logger.log(&ComputeEvent::ArrangementHeapSize(ArrangementHeapSize {
370                        operator_id,
371                        delta_size: size - old_size,
372                    }));
373                }
374
375                let capacity = capacity.try_into().expect("must fit");
376                if capacity != old_capacity {
377                    logger.log(&ComputeEvent::ArrangementHeapCapacity(
378                        ArrangementHeapCapacity {
379                            operator_id,
380                            delta_capacity: capacity - old_capacity,
381                        },
382                    ));
383                }
384
385                let allocations = allocations.try_into().expect("must fit");
386                if allocations != old_allocations {
387                    logger.log(&ComputeEvent::ArrangementHeapAllocations(
388                        ArrangementHeapAllocations {
389                            operator_id,
390                            delta_allocations: allocations - old_allocations,
391                        },
392                    ));
393                }
394
395                old_size = size;
396                old_capacity = capacity;
397                old_allocations = allocations;
398            }
399        });
400    Arranged {
401        trace: arranged.trace,
402        stream,
403    }
404}
405
406impl<'scope, T, K, V, R> ArrangementSize for Arranged<'scope, KeyValAgent<K, V, T, R>>
407where
408    T: MzTimestamp,
409    K: Data + MzData,
410    V: Data + MzData,
411    R: Semigroup + Ord + MzData + 'static,
412{
413    fn log_arrangement_size(self) -> Self {
414        log_arrangement_size_inner(self, |batch| {
415            let (mut size, mut capacity, mut allocations) = (0, 0, 0);
416            let mut callback = |siz, cap| {
417                size += siz;
418                capacity += cap;
419                allocations += usize::from(cap > 0);
420            };
421            batch.storage.keys.heap_size(&mut callback);
422            batch.storage.vals.offs.heap_size(&mut callback);
423            batch.storage.vals.vals.heap_size(&mut callback);
424            batch.storage.upds.offs.heap_size(&mut callback);
425            batch.storage.upds.times.heap_size(&mut callback);
426            batch.storage.upds.diffs.heap_size(&mut callback);
427            (size, capacity, allocations)
428        })
429    }
430}
431
432impl<'scope, T, K, R> ArrangementSize for Arranged<'scope, KeyAgent<K, T, R>>
433where
434    T: MzTimestamp,
435    K: Data + MzArrangeData,
436    R: Semigroup + Ord + MzData + 'static,
437{
438    fn log_arrangement_size(self) -> Self {
439        log_arrangement_size_inner(self, |batch| {
440            let (mut size, mut capacity, mut allocations) = (0, 0, 0);
441            let mut callback = |siz, cap| {
442                size += siz;
443                capacity += cap;
444                allocations += usize::from(cap > 0);
445            };
446            batch.storage.keys.heap_size(&mut callback);
447            batch.storage.upds.offs.heap_size(&mut callback);
448            batch.storage.upds.times.heap_size(&mut callback);
449            batch.storage.upds.diffs.heap_size(&mut callback);
450            (size, capacity, allocations)
451        })
452    }
453}
454
455impl<'scope, T, V, R> ArrangementSize for Arranged<'scope, RowValAgent<V, T, R>>
456where
457    T: MzTimestamp,
458    V: Data + MzArrangeData,
459    R: Semigroup + Ord + MzArrangeData + 'static,
460{
461    fn log_arrangement_size(self) -> Self {
462        log_arrangement_size_inner(self, |batch| {
463            let (mut size, mut capacity, mut allocations) = (0, 0, 0);
464            let mut callback = |siz, cap| {
465                size += siz;
466                capacity += cap;
467                allocations += usize::from(cap > 0);
468            };
469            batch.storage.keys.heap_size(&mut callback);
470            batch.storage.vals.offs.heap_size(&mut callback);
471            batch.storage.vals.vals.heap_size(&mut callback);
472            batch.storage.upds.offs.heap_size(&mut callback);
473            batch.storage.upds.times.heap_size(&mut callback);
474            batch.storage.upds.diffs.heap_size(&mut callback);
475            (size, capacity, allocations)
476        })
477    }
478}
479
480impl<'scope, T, R> ArrangementSize for Arranged<'scope, RowRowAgent<T, R>>
481where
482    T: MzTimestamp,
483    R: Semigroup + Ord + MzArrangeData + 'static,
484{
485    fn log_arrangement_size(self) -> Self {
486        log_arrangement_size_inner(self, |batch| {
487            let (mut size, mut capacity, mut allocations) = (0, 0, 0);
488            let mut callback = |siz, cap| {
489                size += siz;
490                capacity += cap;
491                allocations += usize::from(cap > 0);
492            };
493            batch.storage.keys.heap_size(&mut callback);
494            batch.storage.vals.offs.heap_size(&mut callback);
495            batch.storage.vals.vals.heap_size(&mut callback);
496            batch.storage.upds.offs.heap_size(&mut callback);
497            batch.storage.upds.times.heap_size(&mut callback);
498            batch.storage.upds.diffs.heap_size(&mut callback);
499            (size, capacity, allocations)
500        })
501    }
502}
503
504impl<'scope, T, DC> ArrangementSize for Arranged<'scope, RowAgent<T, DC::Owned, DC>>
505where
506    T: MzTimestamp,
507    DC: BatchContainer<Owned: Semigroup + 'static> + HeapSize,
508{
509    fn log_arrangement_size(self) -> Self {
510        log_arrangement_size_inner(self, |batch| {
511            let (mut size, mut capacity, mut allocations) = (0, 0, 0);
512            let mut callback = |siz, cap| {
513                size += siz;
514                capacity += cap;
515                allocations += usize::from(cap > 0);
516            };
517            batch.storage.keys.heap_size(&mut callback);
518            batch.storage.upds.offs.heap_size(&mut callback);
519            batch.storage.upds.times.heap_size(&mut callback);
520            batch.storage.upds.diffs.heap_size(&mut callback);
521            (size, capacity, allocations)
522        })
523    }
524}