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}