Skip to main content

mz_timely_util/
operator.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//! Common operator transformations on timely streams and differential collections.
17
18use std::hash::{BuildHasher, Hash, Hasher};
19
20use columnation::Columnation;
21use differential_dataflow::consolidation::ConsolidatingContainerBuilder;
22use differential_dataflow::difference::{Multiply, Semigroup};
23use differential_dataflow::lattice::Lattice;
24use differential_dataflow::trace::Batcher;
25use differential_dataflow::{AsCollection, Collection, Hashable, VecCollection};
26use timely::container::{DrainContainer, PushInto};
27use timely::dataflow::channels::pact::{Exchange, ParallelizationContract, Pipeline};
28use timely::dataflow::operators::Capability;
29use timely::dataflow::operators::generic::builder_rc::{
30    OperatorBuilder as OperatorBuilderRc, OperatorBuilder,
31};
32use timely::dataflow::operators::generic::operator::{self, Operator};
33use timely::dataflow::operators::generic::{
34    InputHandleCore, OperatorInfo, OutputBuilder, OutputBuilderSession,
35};
36use timely::dataflow::{Scope, Stream, StreamVec};
37use timely::progress::operate::FrontierInterest;
38use timely::progress::{Antichain, Timestamp};
39use timely::{Container, ContainerBuilder, PartialOrder};
40
41use crate::columnation::{ColumnationChunker, ColumnationStack};
42
43/// Extension methods for timely [`Stream`]s.
44pub trait StreamExt<'scope, T, C1>
45where
46    T: Timestamp,
47    C1: Container + DrainContainer + Clone + 'static,
48{
49    /// Like `timely::dataflow::operators::generic::operator::Operator::unary`,
50    /// but the logic function can handle failures.
51    ///
52    /// Creates a new dataflow operator that partitions its input stream by a
53    /// parallelization strategy `pact` and repeatedly invokes `logic`, the
54    /// function returned by the function passed as `constructor`. The `logic`
55    /// function can read to the input stream and write to either of two output
56    /// streams, where the first output stream represents successful
57    /// computations and the second output stream represents failed
58    /// computations.
59    fn unary_fallible<DCB, ECB, B, P>(
60        self,
61        pact: P,
62        name: &str,
63        constructor: B,
64    ) -> (
65        Stream<'scope, T, DCB::Container>,
66        Stream<'scope, T, ECB::Container>,
67    )
68    where
69        DCB: ContainerBuilder,
70        ECB: ContainerBuilder,
71        B: FnOnce(
72            Capability<T>,
73            OperatorInfo,
74        ) -> Box<
75            dyn FnMut(
76                    &mut InputHandleCore<T, C1, P::Puller>,
77                    &mut OutputBuilderSession<'_, T, DCB>,
78                    &mut OutputBuilderSession<'_, T, ECB>,
79                ) + 'static,
80        >,
81        P: ParallelizationContract<T, C1>;
82
83    /// Like [`timely::dataflow::operators::vec::Map::flat_map`], but `logic`
84    /// is allowed to fail. The first returned stream will contain the
85    /// successful applications of `logic`, while the second returned stream
86    /// will contain the failed applications.
87    fn flat_map_fallible<DCB, ECB, D2, E, I, L>(
88        self,
89        name: &str,
90        logic: L,
91    ) -> (
92        Stream<'scope, T, DCB::Container>,
93        Stream<'scope, T, ECB::Container>,
94    )
95    where
96        DCB: ContainerBuilder + PushInto<D2>,
97        ECB: ContainerBuilder + PushInto<E>,
98        I: IntoIterator<Item = Result<D2, E>>,
99        L: for<'a> FnMut(C1::Item<'a>) -> I + 'static;
100
101    /// Routes each record to one of two output streams by a per-record predicate.
102    ///
103    /// Records for which `predicate` returns `true` go to the first output and all
104    /// others to the second. Both outputs are sent under the input capability. That
105    /// capability is only a lower bound on the times of the records it carries, so a
106    /// split that depends on a record's time must inspect the record, and cannot be
107    /// decided once per container.
108    fn partition_by<CB, L>(
109        self,
110        name: &str,
111        predicate: L,
112    ) -> (
113        Stream<'scope, T, CB::Container>,
114        Stream<'scope, T, CB::Container>,
115    )
116    where
117        CB: ContainerBuilder + for<'a> PushInto<C1::Item<'a>>,
118        L: for<'a> FnMut(&C1::Item<'a>) -> bool + 'static;
119
120    /// Block progress of the frontier at `expiration` time
121    fn expire_stream_at(self, name: &str, expiration: T) -> Stream<'scope, T, C1>;
122}
123
124/// Extension methods for differential [`Collection`]s.
125pub trait CollectionExt<'scope, T, D1, R>: Sized
126where
127    T: Timestamp,
128    R: Semigroup,
129{
130    /// Creates a new empty collection in `scope`.
131    fn empty(scope: Scope<'scope, T>) -> VecCollection<'scope, T, D1, R>;
132
133    /// Like [`Collection::map`], but `logic` is allowed to fail. The first
134    /// returned collection will contain successful applications of `logic`,
135    /// while the second returned collection will contain the failed
136    /// applications.
137    ///
138    /// Callers need to specify the following type parameters:
139    /// * `DCB`: The container builder for the `Ok` output.
140    /// * `ECB`: The container builder for the `Err` output.
141    fn map_fallible<DCB, ECB, D2, E, L>(
142        self,
143        name: &str,
144        mut logic: L,
145    ) -> (
146        VecCollection<'scope, T, D2, R>,
147        VecCollection<'scope, T, E, R>,
148    )
149    where
150        DCB: ContainerBuilder<Container = Vec<(D2, T, R)>> + PushInto<(D2, T, R)>,
151        ECB: ContainerBuilder<Container = Vec<(E, T, R)>> + PushInto<(E, T, R)>,
152        D2: Clone + 'static,
153        E: Clone + 'static,
154        L: FnMut(D1) -> Result<D2, E> + 'static,
155    {
156        self.flat_map_fallible::<DCB, ECB, _, _, _, _>(name, move |record| Some(logic(record)))
157    }
158
159    /// Like [`Collection::flat_map`], but `logic` is allowed to fail. The first
160    /// returned collection will contain the successful applications of `logic`,
161    /// while the second returned collection will contain the failed
162    /// applications.
163    fn flat_map_fallible<DCB, ECB, D2, E, I, L>(
164        self,
165        name: &str,
166        logic: L,
167    ) -> (
168        Collection<'scope, T, DCB::Container>,
169        Collection<'scope, T, ECB::Container>,
170    )
171    where
172        DCB: ContainerBuilder + PushInto<(D2, T, R)>,
173        ECB: ContainerBuilder + PushInto<(E, T, R)>,
174        D2: Clone + 'static,
175        E: Clone + 'static,
176        I: IntoIterator<Item = Result<D2, E>>,
177        L: FnMut(D1) -> I + 'static;
178
179    /// Block progress of the frontier at `expiration` time.
180    fn expire_collection_at(self, name: &str, expiration: T) -> VecCollection<'scope, T, D1, R>;
181
182    /// Replaces each record with another, with a new difference type.
183    ///
184    /// This method is most commonly used to take records containing aggregatable data (e.g. numbers to be summed)
185    /// and move the data into the difference component. This will allow differential dataflow to update in-place.
186    fn explode_one<D2, R2, L>(
187        self,
188        logic: L,
189    ) -> VecCollection<'scope, T, D2, <R2 as Multiply<R>>::Output>
190    where
191        D2: differential_dataflow::Data,
192        R2: Semigroup + Multiply<R>,
193        <R2 as Multiply<R>>::Output: Clone + 'static + Semigroup,
194        L: FnMut(D1) -> (D2, R2) + 'static,
195        T: Lattice;
196
197    /// Partitions the input into a monotonic collection and
198    /// non-monotone exceptions, with respect to differences.
199    ///
200    /// The exceptions are transformed by `into_err`.
201    fn ensure_monotonic<E, IE>(
202        self,
203        into_err: IE,
204    ) -> (
205        VecCollection<'scope, T, D1, R>,
206        VecCollection<'scope, T, E, R>,
207    )
208    where
209        E: Clone + 'static,
210        IE: Fn(D1, R) -> (E, R) + 'static,
211        R: num_traits::sign::Signed;
212
213    /// Consolidates the collection if `must_consolidate` is `true` and leaves it
214    /// untouched otherwise.
215    fn consolidate_named_if<Ba>(self, must_consolidate: bool, name: &str) -> Self
216    where
217        D1: differential_dataflow::ExchangeData + Hash + Columnation,
218        R: Semigroup + differential_dataflow::ExchangeData + Columnation,
219        T: Lattice + Columnation,
220        Ba: Batcher<Time = T, Output = ColumnationStack<((D1, ()), T, R)>> + 'static;
221
222    /// Consolidates the collection.
223    fn consolidate_named<Ba>(self, name: &str) -> Self
224    where
225        D1: differential_dataflow::ExchangeData + Hash + Columnation,
226        R: Semigroup + differential_dataflow::ExchangeData + Columnation,
227        T: Lattice + Columnation,
228        Ba: Batcher<Time = T, Output = ColumnationStack<((D1, ()), T, R)>> + 'static;
229}
230
231impl<'scope, T, C1> StreamExt<'scope, T, C1> for Stream<'scope, T, C1>
232where
233    T: Timestamp,
234    C1: Container + DrainContainer + Clone + 'static,
235{
236    fn unary_fallible<DCB, ECB, B, P>(
237        self,
238        pact: P,
239        name: &str,
240        constructor: B,
241    ) -> (
242        Stream<'scope, T, DCB::Container>,
243        Stream<'scope, T, ECB::Container>,
244    )
245    where
246        DCB: ContainerBuilder,
247        ECB: ContainerBuilder,
248        B: FnOnce(
249            Capability<T>,
250            OperatorInfo,
251        ) -> Box<
252            dyn FnMut(
253                    &mut InputHandleCore<T, C1, P::Puller>,
254                    &mut OutputBuilderSession<'_, T, DCB>,
255                    &mut OutputBuilderSession<'_, T, ECB>,
256                ) + 'static,
257        >,
258        P: ParallelizationContract<T, C1>,
259    {
260        let mut builder = OperatorBuilderRc::new(name.into(), self.scope());
261
262        let operator_info = builder.operator_info();
263
264        let mut input = builder.new_input(self.clone(), pact);
265        builder.set_notify_for(0, FrontierInterest::Never);
266        let (ok_output, ok_stream) = builder.new_output();
267        let mut ok_output = OutputBuilder::from(ok_output);
268        let (err_output, err_stream) = builder.new_output();
269        let mut err_output = OutputBuilder::from(err_output);
270
271        builder.build(move |mut capabilities| {
272            // `capabilities` should be a single-element vector.
273            let capability = capabilities.pop().unwrap();
274            let mut logic = constructor(capability, operator_info);
275            move |_frontiers| {
276                let mut ok_output_handle = ok_output.activate();
277                let mut err_output_handle = err_output.activate();
278                logic(&mut input, &mut ok_output_handle, &mut err_output_handle);
279            }
280        });
281
282        (ok_stream, err_stream)
283    }
284
285    // XXX(guswynn): file an minimization bug report for the logic flat_map
286    // false positive here
287    // TODO(guswynn): remove this after https://github.com/rust-lang/rust-clippy/issues/8098 is
288    // resolved. The `logic` `FnMut` needs to be borrowed in the `flat_map` call, not moved in
289    // so the simple `|d1| logic(d1)` closure is load-bearing
290    #[allow(clippy::redundant_closure)]
291    fn flat_map_fallible<DCB, ECB, D2, E, I, L>(
292        self,
293        name: &str,
294        mut logic: L,
295    ) -> (
296        Stream<'scope, T, DCB::Container>,
297        Stream<'scope, T, ECB::Container>,
298    )
299    where
300        DCB: ContainerBuilder + PushInto<D2>,
301        ECB: ContainerBuilder + PushInto<E>,
302        I: IntoIterator<Item = Result<D2, E>>,
303        L: for<'a> FnMut(C1::Item<'a>) -> I + 'static,
304    {
305        self.unary_fallible::<DCB, ECB, _, _>(Pipeline, name, move |_, _| {
306            Box::new(move |input, ok_output, err_output| {
307                input.for_each_time(|time, data| {
308                    let mut ok_session = ok_output.session_with_builder(&time);
309                    let mut err_session = err_output.session_with_builder(&time);
310                    for r in data
311                        .flat_map(DrainContainer::drain)
312                        .flat_map(|d1| logic(d1))
313                    {
314                        match r {
315                            Ok(d2) => ok_session.give(d2),
316                            Err(e) => err_session.give(e),
317                        }
318                    }
319                })
320            })
321        })
322    }
323
324    fn partition_by<CB, L>(
325        self,
326        name: &str,
327        mut predicate: L,
328    ) -> (
329        Stream<'scope, T, CB::Container>,
330        Stream<'scope, T, CB::Container>,
331    )
332    where
333        CB: ContainerBuilder + for<'a> PushInto<C1::Item<'a>>,
334        L: for<'a> FnMut(&C1::Item<'a>) -> bool + 'static,
335    {
336        self.unary_fallible::<CB, CB, _, _>(Pipeline, name, move |_, _| {
337            Box::new(move |input, matching_output, rest_output| {
338                input.for_each_time(|time, data| {
339                    let mut matching = matching_output.session_with_builder(&time);
340                    let mut rest = rest_output.session_with_builder(&time);
341                    for item in data.flat_map(DrainContainer::drain) {
342                        if predicate(&item) {
343                            matching.give(item);
344                        } else {
345                            rest.give(item);
346                        }
347                    }
348                })
349            })
350        })
351    }
352
353    fn expire_stream_at(self, name: &str, expiration: T) -> Stream<'scope, T, C1> {
354        let name = format!("expire_stream_at({name})");
355        self.unary_frontier(Pipeline, &name.clone(), move |cap, _| {
356            // Retain a capability for the expiration time, which we'll only drop if the token
357            // is dropped. Else, block progress at the expiration time to prevent downstream
358            // operators from making any statement about expiration time or any following time.
359            let cap = Some(cap.delayed(&expiration));
360            let mut warned = false;
361            move |(input, frontier), output| {
362                let _ = &cap;
363                let frontier = frontier.frontier();
364                if !frontier.less_than(&expiration) && !warned {
365                    // Here, we print a warning, not an error. The state is only a liveness
366                    // concern, but not relevant for correctness. Additionally, a race between
367                    // shutting down the dataflow and dropping the token can cause the dataflow
368                    // to shut down before we drop the token.  This can happen when dropping
369                    // the last remaining capability on a different worker.  We do not want to
370                    // log an error every time this happens.
371
372                    tracing::warn!(
373                        name = name,
374                        frontier = ?frontier,
375                        expiration = ?expiration,
376                        "frontier not less than expiration"
377                    );
378                    warned = true;
379                }
380                input.for_each(|time, data| {
381                    let mut session = output.session(&time);
382                    session.give_container(data);
383                });
384            }
385        })
386    }
387}
388
389impl<'scope, T, D1, R> CollectionExt<'scope, T, D1, R> for VecCollection<'scope, T, D1, R>
390where
391    T: Timestamp + Clone + 'static,
392    D1: Clone + 'static,
393    R: Semigroup + 'static,
394{
395    fn empty(scope: Scope<'scope, T>) -> VecCollection<'scope, T, D1, R> {
396        operator::empty(scope).as_collection()
397    }
398
399    fn flat_map_fallible<DCB, ECB, D2, E, I, L>(
400        self,
401        name: &str,
402        mut logic: L,
403    ) -> (
404        Collection<'scope, T, DCB::Container>,
405        Collection<'scope, T, ECB::Container>,
406    )
407    where
408        DCB: ContainerBuilder + PushInto<(D2, T, R)>,
409        ECB: ContainerBuilder + PushInto<(E, T, R)>,
410        D2: Clone + 'static,
411        E: Clone + 'static,
412        I: IntoIterator<Item = Result<D2, E>>,
413        L: FnMut(D1) -> I + 'static,
414    {
415        let (ok_stream, err_stream) =
416            self.inner
417                .flat_map_fallible::<DCB, ECB, _, _, _, _>(name, move |(d1, t, r)| {
418                    logic(d1).into_iter().map(move |res| match res {
419                        Ok(d2) => Ok((d2, t.clone(), r.clone())),
420                        Err(e) => Err((e, t.clone(), r.clone())),
421                    })
422                });
423        (ok_stream.as_collection(), err_stream.as_collection())
424    }
425
426    fn expire_collection_at(self, name: &str, expiration: T) -> VecCollection<'scope, T, D1, R> {
427        self.inner
428            .expire_stream_at(name, expiration)
429            .as_collection()
430    }
431
432    fn explode_one<D2, R2, L>(
433        self,
434        mut logic: L,
435    ) -> VecCollection<'scope, T, D2, <R2 as Multiply<R>>::Output>
436    where
437        D2: differential_dataflow::Data,
438        R2: Semigroup + Multiply<R>,
439        <R2 as Multiply<R>>::Output: Clone + 'static + Semigroup,
440        L: FnMut(D1) -> (D2, R2) + 'static,
441        T: Lattice,
442    {
443        self.inner
444            .clone()
445            .unary::<ConsolidatingContainerBuilder<_>, _, _, _>(
446                Pipeline,
447                "ExplodeOne",
448                move |_, _| {
449                    move |input, output| {
450                        input.for_each(|time, data| {
451                            output
452                                .session_with_builder(&time)
453                                .give_iterator(data.drain(..).map(|(x, t, d)| {
454                                    let (x, d2) = logic(x);
455                                    (x, t, d2.multiply(&d))
456                                }));
457                        });
458                    }
459                },
460            )
461            .as_collection()
462    }
463
464    fn ensure_monotonic<E, IE>(
465        self,
466        into_err: IE,
467    ) -> (
468        VecCollection<'scope, T, D1, R>,
469        VecCollection<'scope, T, E, R>,
470    )
471    where
472        E: Clone + 'static,
473        IE: Fn(D1, R) -> (E, R) + 'static,
474        R: num_traits::sign::Signed,
475    {
476        let (oks, errs) = self
477            .inner
478            .unary_fallible(Pipeline, "EnsureMonotonic", move |_, _| {
479                Box::new(move |input, ok_output, err_output| {
480                    input.for_each(|time, data| {
481                        let mut ok_session = ok_output.session(&time);
482                        let mut err_session = err_output.session(&time);
483                        for (x, t, d) in data.drain(..) {
484                            if d.is_positive() {
485                                ok_session.give((x, t, d))
486                            } else {
487                                let (e, d2) = into_err(x, d);
488                                err_session.give((e, t, d2))
489                            }
490                        }
491                    })
492                })
493            });
494        (oks.as_collection(), errs.as_collection())
495    }
496
497    fn consolidate_named_if<Ba>(self, must_consolidate: bool, name: &str) -> Self
498    where
499        D1: differential_dataflow::ExchangeData + Hash + Columnation,
500        R: Semigroup + differential_dataflow::ExchangeData + Columnation,
501        T: Lattice + Ord + Columnation,
502        Ba: Batcher<Time = T, Output = ColumnationStack<((D1, ()), T, R)>> + 'static,
503    {
504        if must_consolidate {
505            // We employ AHash below instead of the default hasher in DD to obtain
506            // a better distribution of data to workers. AHash claims empirically
507            // both speed and high quality, according to
508            // https://github.com/tkaitchuck/aHash/blob/master/compare/readme.md.
509            // TODO(vmarcos): Consider here if it is worth it to spend the time to
510            // implement twisted tabulation hashing as proposed in Mihai Patrascu,
511            // Mikkel Thorup: Twisted Tabulation Hashing. SODA 2013: 209-228, available
512            // at https://epubs.siam.org/doi/epdf/10.1137/1.9781611973105.16. The latter
513            // would provide good bounds for balls-into-bins problems when the number of
514            // bins is small (as is our case), so we'd have a theoretical guarantee.
515            // The seeds are fixed for determinism across builds; see
516            // [`crate::hash::fixed_state`].
517            let random_state = crate::hash::fixed_state();
518            let exchange = Exchange::new(move |update: &((D1, _), T, R)| {
519                let data = &(update.0).0;
520                let mut h = random_state.build_hasher();
521                data.hash(&mut h);
522                h.finish()
523            });
524            consolidate_pact::<ColumnationChunker<((D1, ()), T, R)>, Ba, _, _>(
525                self.map(|k| (k, ())).inner,
526                exchange,
527                name,
528            )
529            .unary(Pipeline, "unpack consolidated", |_, _| {
530                |input, output| {
531                    input.for_each(|time, data| {
532                        let mut session = output.session(&time);
533                        for ((k, ()), t, d) in data.iter().flatten().flat_map(|chunk| chunk.iter())
534                        {
535                            session.give((k.clone(), t.clone(), d.clone()))
536                        }
537                    })
538                }
539            })
540            .as_collection()
541        } else {
542            self
543        }
544    }
545
546    fn consolidate_named<Ba>(self, name: &str) -> Self
547    where
548        D1: differential_dataflow::ExchangeData + Hash + Columnation,
549        R: Semigroup + differential_dataflow::ExchangeData + Columnation,
550        T: Lattice + Ord + Columnation,
551        Ba: Batcher<Time = T, Output = ColumnationStack<((D1, ()), T, R)>> + 'static,
552    {
553        let exchange = Exchange::new(move |update: &((D1, ()), T, R)| (update.0).0.hashed());
554
555        consolidate_pact::<ColumnationChunker<((D1, ()), T, R)>, Ba, _, _>(
556            self.map(|k| (k, ())).inner,
557            exchange,
558            name,
559        )
560        .unary(Pipeline, &format!("Unpack {name}"), |_, _| {
561            |input, output| {
562                input.for_each(|time, data| {
563                    let mut session = output.session(&time);
564                    for ((k, ()), t, d) in data.iter().flatten().flat_map(|chunk| chunk.iter()) {
565                        session.give((k.clone(), t.clone(), d.clone()))
566                    }
567                })
568            }
569        })
570        .as_collection()
571    }
572}
573
574/// Aggregates the weights of equal records into at most one record.
575///
576/// Produces a stream of chains of records, partitioned according to `pact`. The
577/// data is sorted according to `Ba`. For each timestamp, it produces at most one chain.
578///
579/// The data are accumulated in place, each held back until their timestamp has completed.
580pub fn consolidate_pact<'scope, Chu, Ba, C, P>(
581    stream: Stream<'scope, Ba::Time, C>,
582    pact: P,
583    name: &str,
584) -> StreamVec<'scope, Ba::Time, Vec<Ba::Output>>
585where
586    Ba: Batcher + 'static,
587    Chu: ContainerBuilder<Container = Ba::Output> + for<'a> PushInto<&'a mut C> + 'static,
588    C: Container + Clone + 'static,
589    Ba::Output: Clone,
590    P: ParallelizationContract<Ba::Time, C>,
591{
592    let logger = stream
593        .scope()
594        .worker()
595        .logger_for("differential/arrange")
596        .map(Into::into);
597    stream.unary_frontier(pact, name, |_cap, info| {
598        // Acquire a logger for arrange events.
599
600        let mut batcher = Ba::new(logger, info.global_id);
601        // The chunker consolidates raw input containers into the chunks the
602        // batcher consumes.
603        let mut chunker = Chu::default();
604        // Capabilities for the lower envelope of updates in `batcher`.
605        let mut capabilities = Antichain::<Capability<Ba::Time>>::new();
606        let mut prev_frontier = Antichain::from_elem(Ba::Time::minimum());
607
608        move |(input, frontier), output| {
609            input.for_each(|cap, data| {
610                capabilities.insert(cap.retain(0));
611                chunker.push_into(data);
612                while let Some(chunk) = chunker.extract() {
613                    batcher.push_into(std::mem::take(chunk));
614                }
615            });
616
617            if prev_frontier.borrow() != frontier.frontier() {
618                // Flush any data the chunker is still accumulating into the
619                // batcher before we seal.
620                while let Some(chunk) = chunker.finish() {
621                    batcher.push_into(std::mem::take(chunk));
622                }
623
624                if capabilities
625                    .elements()
626                    .iter()
627                    .any(|c| !frontier.less_equal(c.time()))
628                {
629                    let mut upper = Antichain::new(); // re-used allocation for sealing batches.
630
631                    // For each capability not in advance of the input frontier ...
632                    for (index, capability) in capabilities.elements().iter().enumerate() {
633                        if !frontier.less_equal(capability.time()) {
634                            // Assemble the upper bound on times we can commit with this capabilities.
635                            // We must respect the input frontier, and *subsequent* capabilities, as
636                            // we are pretending to retire the capability changes one by one.
637                            upper.clear();
638                            for time in frontier.frontier().iter() {
639                                upper.insert(time.clone());
640                            }
641                            for other_capability in &capabilities.elements()[(index + 1)..] {
642                                upper.insert(other_capability.time().clone());
643                            }
644
645                            // send the batch to downstream consumers, empty or not.
646                            let mut session = output.session(&capabilities.elements()[index]);
647                            // Extract updates not in advance of `upper`.
648                            let (chain, _description) = batcher.seal(upper.clone());
649                            session.give(chain);
650                        }
651                    }
652
653                    // Having extracted and sent batches between each capability and the input frontier,
654                    // we should downgrade all capabilities to match the batcher's lower update frontier.
655                    // This may involve discarding capabilities, which is fine as any new updates arrive
656                    // in messages with new capabilities.
657
658                    let mut new_capabilities = Antichain::new();
659                    for time in batcher.frontier().iter() {
660                        if let Some(capability) = capabilities
661                            .elements()
662                            .iter()
663                            .find(|c| c.time().less_equal(time))
664                        {
665                            new_capabilities.insert(capability.delayed(time));
666                        } else {
667                            panic!("failed to find capability");
668                        }
669                    }
670
671                    capabilities = new_capabilities;
672                }
673
674                prev_frontier.clear();
675                prev_frontier.extend(frontier.frontier().iter().cloned());
676            }
677        }
678    })
679}
680
681/// Merge the contents of multiple streams and combine the containers using a container builder.
682pub trait ConcatenateFlatten<'scope, T: Timestamp, C: Container + DrainContainer> {
683    /// Merge the contents of multiple streams and use the provided container builder to form
684    /// output containers.
685    ///
686    /// # Examples
687    /// ```
688    /// use timely::container::CapacityContainerBuilder;
689    /// use timely::dataflow::operators::{ToStream, Inspect};
690    /// use mz_timely_util::operator::ConcatenateFlatten;
691    ///
692    /// timely::example(|scope| {
693    ///
694    ///     let streams: Vec<timely::dataflow::StreamVec<_, i32>> =
695    ///         vec![(0..10).to_stream(scope),
696    ///              (0..10).to_stream(scope),
697    ///              (0..10).to_stream(scope)];
698    ///
699    ///     scope.concatenate_flatten::<_, CapacityContainerBuilder<Vec<i32>>>(streams)
700    ///          .inspect(|x| println!("seen: {:?}", x));
701    /// });
702    /// ```
703    fn concatenate_flatten<I, CB>(&self, sources: I) -> Stream<'scope, T, CB::Container>
704    where
705        I: IntoIterator<Item = Stream<'scope, T, C>>,
706        CB: ContainerBuilder + for<'a> PushInto<C::Item<'a>>;
707}
708
709impl<'scope, T, C> ConcatenateFlatten<'scope, T, C> for Stream<'scope, T, C>
710where
711    T: Timestamp,
712    C: Container + DrainContainer + Clone + 'static,
713{
714    fn concatenate_flatten<I, CB>(&self, sources: I) -> Stream<'scope, T, CB::Container>
715    where
716        I: IntoIterator<Item = Stream<'scope, T, C>>,
717        CB: ContainerBuilder + for<'a> PushInto<C::Item<'a>>,
718    {
719        self.scope()
720            .concatenate_flatten::<_, CB>(Some(Clone::clone(self)).into_iter().chain(sources))
721    }
722}
723
724impl<'scope, T, C> ConcatenateFlatten<'scope, T, C> for Scope<'scope, T>
725where
726    T: Timestamp,
727    C: Container + DrainContainer,
728{
729    fn concatenate_flatten<I, CB>(&self, sources: I) -> Stream<'scope, T, CB::Container>
730    where
731        I: IntoIterator<Item = Stream<'scope, T, C>>,
732        CB: ContainerBuilder + for<'a> PushInto<C::Item<'a>>,
733    {
734        let mut builder = OperatorBuilder::new("ConcatenateFlatten".to_string(), self.clone());
735
736        // create new input handles for each input stream.
737        let mut handles = sources
738            .into_iter()
739            .map(|s| builder.new_input(s, Pipeline))
740            .collect::<Vec<_>>();
741        for i in 0..handles.len() {
742            builder.set_notify_for(i, FrontierInterest::Never);
743        }
744
745        // create one output handle for the concatenated results.
746        let (output, result) = builder.new_output::<CB::Container>();
747        let mut output = OutputBuilder::<_, CB>::from(output);
748
749        builder.build(move |_capability| {
750            move |_frontier| {
751                let mut output = output.activate();
752                for handle in handles.iter_mut() {
753                    handle.for_each_time(|time, data| {
754                        output
755                            .session_with_builder(&time)
756                            .give_iterator(data.flat_map(DrainContainer::drain));
757                    })
758                }
759            }
760        });
761
762        result
763    }
764}
765
766/// A trait for containers that can be cleared.
767pub trait ClearContainer {
768    /// Clear the contents of the container.
769    fn clear(&mut self);
770}
771
772impl<T> ClearContainer for Vec<T> {
773    fn clear(&mut self) {
774        Vec::clear(self)
775    }
776}
777
778#[cfg(test)]
779mod tests {
780    use timely::container::CapacityContainerBuilder;
781    use timely::dataflow::operators::Capture;
782    use timely::dataflow::operators::capture::Extract;
783    use timely::dataflow::operators::core::to_stream::ToStreamBuilder;
784    use timely::dataflow::operators::vec::ToStream;
785
786    use crate::columnar::Column;
787    use crate::columnar::builder::ColumnBuilder;
788
789    use super::*;
790
791    /// Updates whose own times straddle `SPLIT`, all carried under the single input
792    /// capability at time zero. A split decided per container would route them together.
793    const UPDATES: [(u64, u64, i64); 4] = [(0, 0, 1), (1, 3, 1), (2, 5, -1), (3, 7, 1)];
794    const SPLIT: u64 = 5;
795    const MATCHING: [(u64, u64, i64); 2] = [(2, 5, -1), (3, 7, 1)];
796    const REST: [(u64, u64, i64); 2] = [(0, 0, 1), (1, 3, 1)];
797
798    #[mz_ore::test]
799    fn partition_by_routes_vec_records_by_their_own_time() {
800        let (matching, rest) = timely::execute_directly(|worker| {
801            worker.dataflow::<u64, _, _>(|scope| {
802                let (matching, rest) = UPDATES
803                    .to_vec()
804                    .to_stream(scope)
805                    .partition_by::<CapacityContainerBuilder<Vec<_>>, _>("Test", |(_, time, _)| {
806                        *time >= SPLIT
807                    });
808                (matching.capture(), rest.capture())
809            })
810        });
811        let flatten = |captured: Vec<(u64, Vec<(u64, u64, i64)>)>| {
812            captured
813                .into_iter()
814                .flat_map(|(capability, updates)| {
815                    assert_eq!(capability, 0, "outputs keep the input capability");
816                    updates
817                })
818                .collect::<Vec<_>>()
819        };
820        assert_eq!(flatten(matching.extract()), MATCHING);
821        assert_eq!(flatten(rest.extract()), REST);
822    }
823
824    #[mz_ore::test]
825    fn partition_by_routes_columnar_records_by_their_own_time() {
826        let (matching, rest) = timely::execute_directly(|worker| {
827            worker.dataflow::<u64, _, _>(|scope| {
828                let (matching, rest) = UPDATES
829                    .to_vec()
830                    .to_stream_with_builder::<_, ColumnBuilder<(u64, u64, i64)>>(scope)
831                    .partition_by::<ColumnBuilder<(u64, u64, i64)>, _>("Test", |(_, time, _)| {
832                        **time >= SPLIT
833                    });
834                (matching.capture(), rest.capture())
835            })
836        });
837        let flatten = |captured: Vec<(u64, Column<(u64, u64, i64)>)>| {
838            captured
839                .into_iter()
840                .flat_map(|(capability, mut column)| {
841                    assert_eq!(capability, 0, "outputs keep the input capability");
842                    column
843                        .drain()
844                        .map(|(data, time, diff)| (*data, *time, *diff))
845                        .collect::<Vec<_>>()
846                })
847                .collect::<Vec<_>>()
848        };
849        assert_eq!(flatten(matching.extract()), MATCHING);
850        assert_eq!(flatten(rest.extract()), REST);
851    }
852}