Skip to main content

mz_compute/extensions/
temporal_bucket.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
10//! Utilities and stream extensions for temporal bucketing.
11
12use std::hash::Hash;
13
14use columnar::{Columnar, Index, Len, Push};
15use differential_dataflow::Hashable;
16use differential_dataflow::difference::Semigroup;
17use differential_dataflow::lattice::Lattice;
18use differential_dataflow::trace::Batcher;
19use mz_timely_util::columnar::Column;
20use mz_timely_util::columnar::batcher::ColumnChunker;
21use mz_timely_util::columnar::builder::ColumnBuilder;
22use mz_timely_util::columnar::chunk::{AccountedChunkBatcher, ColumnChunk};
23use mz_timely_util::columnar::columnar_exchange_data;
24use mz_timely_util::temporal::{Bucket, BucketChain, BucketRange, BucketTimestamp};
25use timely::Accountable;
26use timely::ExchangeData;
27use timely::container::{CapacityContainerBuilder, PushInto};
28use timely::dataflow::channels::pact::{Exchange, ExchangeCore};
29use timely::dataflow::operators::Operator;
30use timely::dataflow::{Stream, StreamVec};
31use timely::order::TotalOrder;
32use timely::progress::{Antichain, PathSummary, Timestamp};
33
34use crate::typedefs::MzData;
35
36/// Sort outstanding updates into a [`BucketChain`], and reveal data not in advance of the input
37/// frontier. Retains a capability at the last input frontier to retain the right to produce data
38/// at times between the last input frontier and the current input frontier.
39pub trait TemporalBucketing<'scope, T: Timestamp>: Sized {
40    /// Construct a new stream that stores updates into a [`BucketChain`] and reveals data
41    /// not in advance of the frontier. Data that is within `threshold` distance of the input
42    /// frontier or the `as_of` is passed through without being stored in the chain.
43    ///
44    /// The output container matches the input's, so a caller keeps whichever
45    /// representation it had.
46    fn bucket(self, as_of: Antichain<T>, threshold: T::Summary) -> Self;
47}
48
49/// Implementation for streams in scopes where timestamps define a total order.
50impl<'scope, T, D> TemporalBucketing<'scope, T> for Stream<'scope, T, Column<(D, T, mz_repr::Diff)>>
51where
52    T: Timestamp + Default + ExchangeData + MzData + BucketTimestamp + TotalOrder + Lattice,
53    for<'a> columnar::Ref<'a, T>: Copy + Ord,
54    D: ExchangeData + MzData + Ord + Clone + std::fmt::Debug + Hashable,
55    for<'a> columnar::Ref<'a, D>: Copy + Ord + Hash,
56    for<'a> columnar::Ref<'a, mz_repr::Diff>: Ord,
57    for<'a> <(D, T, mz_repr::Diff) as Columnar>::Container:
58        Push<columnar::Ref<'a, (D, T, mz_repr::Diff)>>,
59{
60    fn bucket(self, as_of: Antichain<T>, threshold: T::Summary) -> Self {
61        let scope = self.scope();
62        let logger = scope
63            .worker()
64            .logger_for("differential/arrange")
65            .map(Into::into);
66
67        type CB<D, T> = CapacityContainerBuilder<Column<(D, T, mz_repr::Diff)>>;
68
69        let pact = ExchangeCore::<ColumnBuilder<_>, _>::new_core(
70            columnar_exchange_data::<D, T, mz_repr::Diff>,
71        );
72        self.unary_frontier::<CB<D, T>, _, _, _>(pact, "Temporal delay", |cap, info| {
73            let mut chain = BucketChain::new(MergeBatcherWrapper::new(logger, info.global_id));
74            let activator = scope.activator_for(info.address);
75
76            // Cap tracking the lower bound of potentially outstanding data.
77            let mut cap = Some(cap);
78
79            // Holds one bucket's worth of updates on the way into the chain.
80            // Reused across activations for its allocation.
81            let mut buffer: Column<(D, T, mz_repr::Diff)> = Default::default();
82            // Reused input permutation, ordered by time.
83            let mut permutation: Vec<usize> = Vec::new();
84            // Reused so reading a record's time does not allocate: an iterative `T`
85            // owns a `PointStamp`'s allocation.
86            let mut time_buf = T::minimum();
87
88            move |(input, frontier), output| {
89                // The upper frontier is the join of the input frontier and the `as_of` frontier,
90                // with the `threshold` summary applied to it.
91                let mut upper = Antichain::new();
92                for time1 in &frontier.frontier() {
93                    for time2 in as_of.elements() {
94                        // TODO: Use `join_assign` if we ever use a timestamp with allocations.
95                        if let Some(time) = threshold.results_in(&time1.join(time2)) {
96                            upper.insert(time);
97                        }
98                    }
99                }
100
101                input.for_each_time(|time, data| {
102                    let mut session = output.session_with_builder(&time);
103                    for data in data {
104                        let borrowed = data.borrow();
105
106                        // Pass through data about to be revealed, and retain the
107                        // index of everything the chain has to hold. Only the
108                        // retained records need ordering, and in steady state the
109                        // pass-through share is the larger one.
110                        permutation.clear();
111                        for index in 0..borrowed.len() {
112                            let update = borrowed.get(index);
113                            time_buf.copy_from(update.1);
114                            if upper.less_equal(&time_buf) {
115                                permutation.push(index);
116                            } else {
117                                session.give(update);
118                            }
119                        }
120
121                        // Order the retained records by time so each bucket's
122                        // records land contiguously below. Sorting indices keeps
123                        // the records in place.
124                        permutation.sort_unstable_by_key(|index| borrowed.get(*index).1);
125
126                        // The range `buffer`'s contents belong to, `None` while empty.
127                        let mut buffered_range = None;
128                        for index in permutation.drain(..) {
129                            let update = borrowed.get(index);
130                            time_buf.copy_from(update.1);
131
132                            // Ship the buffer whenever the bucket changes, which
133                            // the time order makes a single transition per bucket.
134                            let contained = match &buffered_range {
135                                Some(range) => BucketRange::contains(range, &time_buf),
136                                None => false,
137                            };
138                            if !contained {
139                                if let Some(range) = buffered_range.take() {
140                                    let bucket = chain.find_mut(&range.start).expect("Must exist");
141                                    bucket.push_container(&mut buffer);
142                                }
143                                buffered_range =
144                                    Some(chain.range_of(&time_buf).expect("Must exist"));
145                            }
146                            buffer.push_into(update);
147                        }
148
149                        // Handle leftover data in the buffer.
150                        if let Some(range) = buffered_range.take() {
151                            let bucket = chain.find_mut(&range.start).expect("Must exist");
152                            bucket.push_container(&mut buffer);
153                        }
154                    }
155                });
156
157                // Check for data that is ready to be revealed.
158                let peeled = chain.peel(upper.borrow());
159                if let Some(cap) = cap.as_ref() {
160                    let mut session = output.session_with_builder(cap);
161                    // The chain hands back chunks whose bodies go back onto the
162                    // edge container as a move, so each one ships as a container.
163                    for chunk in peeled.into_iter().flat_map(|x| x.done()) {
164                        let mut column = Column::from(chunk.into_body());
165                        session.give_container(&mut column);
166                    }
167                } else {
168                    // If we don't have a cap, we should not have any data to reveal.
169                    assert!(
170                        peeled
171                            .into_iter()
172                            .flat_map(|x| x.done())
173                            .all(|chunk| chunk.record_count() == 0),
174                        "Unexpected data revealed without a cap."
175                    );
176                }
177
178                // Downgrade the cap to the current input frontier.
179                if frontier.is_empty() || upper.is_empty() {
180                    cap = None;
181                } else if let Some(cap) = cap.as_mut() {
182                    // TODO: This assumes that the time is total ordered.
183                    cap.downgrade(&upper[0]);
184                }
185
186                // Maintain the bucket chain by restoring it with fuel.
187                let mut fuel = 1_000_000;
188                chain.restore(&mut fuel);
189                if fuel <= 0 {
190                    // If we run out of fuel, we activate the operator to continue processing.
191                    activator.activate();
192                }
193            }
194        })
195    }
196}
197
198/// Implementation for `Vec` streams in scopes where timestamps define a total order.
199///
200/// A caller whose consumer wants owned records keeps a `Vec`-native operator, because
201/// staging the whole stream through a column would copy every pass-through record and
202/// allocate it again on the way out. Only records that enter the chain are encoded, which
203/// they were anyway: the chain's batcher is columnar. The reduce key-value path is the one
204/// such caller, and this implementation goes away once its consumer reads columns.
205impl<'scope, T, D> TemporalBucketing<'scope, T> for StreamVec<'scope, T, (D, T, mz_repr::Diff)>
206where
207    T: Timestamp + Default + ExchangeData + MzData + BucketTimestamp + TotalOrder + Lattice,
208    D: ExchangeData + MzData + Ord + Clone + std::fmt::Debug + Hashable,
209    for<'a> <(D, T, mz_repr::Diff) as Columnar>::Container: Push<&'a (D, T, mz_repr::Diff)>,
210{
211    fn bucket(self, as_of: Antichain<T>, threshold: T::Summary) -> Self {
212        let scope = self.scope();
213        let logger = scope
214            .worker()
215            .logger_for("differential/arrange")
216            .map(Into::into);
217
218        let pact = Exchange::new(|(d, _, _): &(D, T, mz_repr::Diff)| d.hashed().into());
219        self.unary_frontier::<CapacityContainerBuilder<Vec<(D, T, mz_repr::Diff)>>, _, _, _>(
220            pact,
221            "Temporal delay",
222            |cap, info| {
223                let mut chain = BucketChain::new(MergeBatcherWrapper::new(logger, info.global_id));
224                let activator = scope.activator_for(info.address);
225
226                // Cap tracking the lower bound of potentially outstanding data.
227                let mut cap = Some(cap);
228
229                // Staging column for the records of one bucket. The chain's batcher is
230                // columnar, so a stored record is encoded either way.
231                let mut buffer: Column<(D, T, mz_repr::Diff)> = Default::default();
232
233                move |(input, frontier), output| {
234                    // The upper frontier is the join of the input frontier and the `as_of`
235                    // frontier, with the `threshold` summary applied to it.
236                    let mut upper = Antichain::new();
237                    for time1 in &frontier.frontier() {
238                        for time2 in as_of.elements() {
239                            // TODO: Use `join_assign` if we ever use a timestamp with allocations.
240                            if let Some(time) = threshold.results_in(&time1.join(time2)) {
241                                upper.insert(time);
242                            }
243                        }
244                    }
245
246                    input.for_each_time(|time, data| {
247                        let mut session = output.session_with_builder(&time);
248                        for data in data {
249                            // Skip data that is about to be revealed.
250                            let pass_through =
251                                data.extract_if(.., |(_, t, _)| !upper.less_equal(t));
252                            session.give_iterator(pass_through);
253
254                            // Sort data by time, then drain it into a buffer that contains data
255                            // for a single bucket. We scan the data for ranges of time that fall
256                            // into the same bucket so we can push batches of data at once.
257                            data.sort_unstable_by(|(_, t, _), (_, t2, _)| t.cmp(t2));
258
259                            let mut drain = data.drain(..);
260                            if let Some(update) = drain.next() {
261                                let mut range = chain.range_of(&update.1).expect("Must exist");
262                                buffer.push_into(&update);
263                                for update in drain {
264                                    // If we have a range, check if the time is not within it.
265                                    if !range.contains(&update.1) {
266                                        // If the time is outside the range, push the current
267                                        // buffer to the chain and reset the range.
268                                        if !buffer.is_empty() {
269                                            let bucket =
270                                                chain.find_mut(&range.start).expect("Must exist");
271                                            bucket.push_container(&mut buffer);
272                                        }
273                                        range = chain.range_of(&update.1).expect("Must exist");
274                                    }
275                                    buffer.push_into(&update);
276                                }
277
278                                // Handle leftover data in the buffer.
279                                if !buffer.is_empty() {
280                                    let bucket = chain.find_mut(&range.start).expect("Must exist");
281                                    bucket.push_container(&mut buffer);
282                                }
283                            }
284                        }
285                    });
286
287                    // Check for data that is ready to be revealed.
288                    let peeled = chain.peel(upper.borrow());
289                    if let Some(cap) = cap.as_ref() {
290                        let mut session = output.session_with_builder(cap);
291                        for chunk in peeled.into_iter().flat_map(|x| x.done()) {
292                            let body = chunk.into_body();
293                            session.give_iterator(
294                                body.borrow()
295                                    .into_index_iter()
296                                    .map(<(D, T, mz_repr::Diff)>::into_owned),
297                            );
298                        }
299                    } else {
300                        // If we don't have a cap, we should not have any data to reveal.
301                        assert!(
302                            peeled
303                                .into_iter()
304                                .flat_map(|x| x.done())
305                                .all(|chunk| chunk.record_count() == 0),
306                            "Unexpected data revealed without a cap."
307                        );
308                    }
309
310                    // Downgrade the cap to the current input frontier.
311                    if frontier.is_empty() || upper.is_empty() {
312                        cap = None;
313                    } else if let Some(cap) = cap.as_mut() {
314                        // TODO: This assumes that the time is total ordered.
315                        cap.downgrade(&upper[0]);
316                    }
317
318                    // Maintain the bucket chain by restoring it with fuel.
319                    let mut fuel = 1_000_000;
320                    chain.restore(&mut fuel);
321                    if fuel <= 0 {
322                        // If we run out of fuel, we activate the operator to continue processing.
323                        activator.activate();
324                    }
325                }
326            },
327        )
328    }
329}
330
331/// A wrapper around [`AccountedChunkBatcher`] that implements the bucketing API.
332///
333/// This is the same merge batcher the arrange sites' chunked arm uses, so the
334/// bucket chain and those arrangements share one merge-batcher implementation.
335/// The choice is unconditional here: the bucket chain does not consult the
336/// arrange batcher selector.
337///
338/// The batcher consumes sorted, consolidated [`ColumnChunk`] input, so this
339/// wrapper carries a [`ColumnChunker`] that sorts and consolidates the input
340/// columns into the [`Column`]s it wraps as chunks.
341struct MergeBatcherWrapper<D, T, R>
342where
343    D: MzData,
344    T: MzData + Ord + Default + Timestamp + Lattice,
345    R: MzData + Semigroup + Default + for<'a> Semigroup<columnar::Ref<'a, R>>,
346    for<'a> columnar::Ref<'a, R>: Ord,
347{
348    logger: Option<differential_dataflow::logging::Logger>,
349    operator_id: usize,
350    chunker: ColumnChunker<(D, T, R)>,
351    inner: AccountedChunkBatcher<D, T, R>,
352}
353
354impl<D, T, R> MergeBatcherWrapper<D, T, R>
355where
356    D: MzData + Ord + Clone,
357    T: MzData + Ord + Clone + Default + Timestamp + Lattice,
358    R: MzData + Semigroup + Default + for<'a> Semigroup<columnar::Ref<'a, R>>,
359    for<'a> columnar::Ref<'a, R>: Ord,
360    for<'a> <D as Columnar>::Container: Push<columnar::Ref<'a, D>>,
361    for<'a> <T as Columnar>::Container: Push<columnar::Ref<'a, T>>,
362    for<'a> <R as Columnar>::Container: Push<&'a R>,
363    for<'a> <(D, T, R) as Columnar>::Container: Push<&'a (D, T, R)>,
364{
365    /// Construct a new `MergeBatcherWrapper` with the given logger and operator ID.
366    fn new(logger: Option<differential_dataflow::logging::Logger>, operator_id: usize) -> Self {
367        Self {
368            logger: logger.clone(),
369            operator_id,
370            chunker: ColumnChunker::default(),
371            inner: AccountedChunkBatcher::new(logger, operator_id),
372        }
373    }
374
375    /// Consolidate `buffer` through the chunker and feed any complete chunks to
376    /// the batcher. Leaves `buffer` empty, retaining its allocation.
377    fn push_container(&mut self, buffer: &mut Column<(D, T, R)>) {
378        use timely::container::{ContainerBuilder as _, PushInto as _};
379        if buffer.is_empty() {
380            return;
381        }
382        self.chunker.push_into(buffer);
383        buffer.clear();
384        while let Some(chunk) = self.chunker.extract() {
385            self.inner
386                .push_into(ColumnChunk::from_body(std::mem::take(chunk)));
387        }
388    }
389
390    /// Flush any partial chunk still held by the chunker into the batcher.
391    fn flush(&mut self) {
392        use timely::container::{ContainerBuilder as _, PushInto as _};
393        while let Some(chunk) = self.chunker.finish() {
394            self.inner
395                .push_into(ColumnChunk::from_body(std::mem::take(chunk)));
396        }
397    }
398
399    /// Reveal the contents of the merge batcher, returning a vector of chunks.
400    fn done(mut self) -> Vec<ColumnChunk<D, T, R>> {
401        self.flush();
402        let (chain, _description) = self.inner.seal(Antichain::new());
403        chain
404    }
405}
406
407impl<D, T, R> Bucket for MergeBatcherWrapper<D, T, R>
408where
409    D: MzData + Ord + Clone,
410    T: MzData + Ord + Clone + Default + Lattice + BucketTimestamp,
411    R: MzData + Semigroup + Default + for<'a> Semigroup<columnar::Ref<'a, R>>,
412    for<'a> columnar::Ref<'a, R>: Ord,
413    for<'a> <D as Columnar>::Container: Push<columnar::Ref<'a, D>>,
414    for<'a> <T as Columnar>::Container: Push<columnar::Ref<'a, T>>,
415    for<'a> <R as Columnar>::Container: Push<&'a R>,
416    for<'a> <(D, T, R) as Columnar>::Container: Push<&'a (D, T, R)>,
417{
418    type Timestamp = T;
419
420    fn split(mut self, timestamp: &Self::Timestamp, fuel: &mut i64) -> (Self, Self) {
421        use timely::container::PushInto as _;
422        self.flush();
423        let upper = Antichain::from_elem(timestamp.clone());
424        let mut lower = Self::new(self.logger.clone(), self.operator_id);
425        // Sealing at `timestamp` ships exactly the updates strictly less than it,
426        // as sorted, consolidated chunks; feed them to the lower batcher whole,
427        // so spilled bodies move without being loaded. The `fuel` charge covers
428        // the sealing work those records paid for.
429        let (chain, _description) = self.inner.seal(upper);
430        for chunk in chain {
431            *fuel = fuel.saturating_sub(chunk.record_count());
432            lower.inner.push_into(chunk);
433        }
434        (lower, self)
435    }
436}