1use std::collections::BTreeMap;
14use std::rc::Rc;
15
16use columnar::{Columnar, Index};
17use differential_dataflow::consolidation::ConsolidatingContainerBuilder;
18use differential_dataflow::operators::arrange::Arranged;
19use differential_dataflow::trace::cursor::{BatchCursor, BatchKey, BatchVal};
20use differential_dataflow::trace::implementations::BatchContainer;
21use differential_dataflow::trace::{Cursor, Navigable, TraceReader};
22use differential_dataflow::{AsCollection, VecCollection};
23use mz_compute_types::dataflows::DataflowDescription;
24use mz_compute_types::dyncfgs::{ENABLE_COMPUTE_TEMPORAL_BUCKETING, TEMPORAL_BUCKETING_SUMMARY};
25use mz_compute_types::plan::scalar::{LirScalarExpr, mfp_mir_to_lir_plan, mfp_plan_lir_to_mir};
26use mz_compute_types::plan::{ArrangementStrategy, AvailableCollections};
27use mz_dyncfg::ConfigSet;
28use mz_expr::{Eval, Id, MfpPlan};
29use mz_ore::soft_assert_or_log;
30use mz_repr::fixed_length::ExtendDatums;
31use mz_repr::{DatumVec, DatumVecBorrow, Diff, GlobalId, Row, RowArena, SharedRow, StableRow};
32use mz_storage_types::controller::CollectionMetadata;
33use mz_timely_util::columnar::Column;
34use mz_timely_util::columnar::batcher;
35use mz_timely_util::columnar::builder::ColumnBuilder;
36use mz_timely_util::columnar::chunk::{AccountedChunkBatcher, ChunkChunker, UnchunkBuilder};
37use mz_timely_util::columnar::consolidate::ConsolidatingColumnBuilder;
38use mz_timely_util::columnar::{Col2ValBatcher, Col2ValColBatcher, columnar_exchange};
39use mz_timely_util::columnation::ColumnationChunker;
40use timely::ContainerBuilder;
41use timely::container::NoopBuilder;
42use timely::dataflow::channels::pact::{ExchangeCore, Pipeline};
43use timely::dataflow::operators::Capability;
44use timely::dataflow::operators::generic::builder_rc::OperatorBuilder;
45use timely::dataflow::operators::generic::{OutputBuilder, OutputBuilderSession};
46use timely::dataflow::{Scope, Stream};
47use timely::progress::operate::FrontierInterest;
48use timely::progress::{Antichain, Timestamp};
49
50use crate::compute_state::ComputeState;
51use crate::extensions::arrange::{ArrangementBatcher, KeyCollection, MzArrange, MzArrangeCore};
52use crate::extensions::reduce::MzReduce;
53use crate::render::columnar::{ColCollection, flat_map_datums};
54use crate::render::errors::{DataflowErrorSer, ErrorLogger};
55use crate::render::{LinearJoinSpec, MaybeBucketByTime, RenderTimestamp};
56use crate::typedefs::{
57 ErrAgent, ErrBatcher, ErrBuilder, ErrEnter, ErrSpine, RowRowAgent, RowRowEnter, RowRowSpine,
58};
59use mz_row_spine::{RowRowBuilder, RowRowColPagedBuilder};
60
61pub struct Context<'scope, T: RenderTimestamp> {
69 pub(crate) scope: Scope<'scope, T>,
73 pub debug_name: String,
75 pub dataflow_id: usize,
77 pub export_ids: Vec<GlobalId>,
79 pub as_of_frontier: Antichain<mz_repr::Timestamp>,
84 pub until: Antichain<mz_repr::Timestamp>,
87 pub bindings: BTreeMap<Id, CollectionBundle<'scope, T>>,
89 pub(super) compute_logger: Option<crate::logging::compute::Logger>,
91 pub(super) linear_join_spec: LinearJoinSpec,
93 pub dataflow_expiration: Antichain<mz_repr::Timestamp>,
96 pub config_set: Rc<ConfigSet>,
98}
99
100impl<'scope, T: RenderTimestamp> Context<'scope, T> {
101 pub fn for_dataflow_in<Plan>(
103 dataflow: &DataflowDescription<Plan, CollectionMetadata>,
104 scope: Scope<'scope, T>,
105 compute_state: &ComputeState,
106 until: Antichain<mz_repr::Timestamp>,
107 dataflow_expiration: Antichain<mz_repr::Timestamp>,
108 ) -> Self {
109 use mz_ore::collections::CollectionExt as IteratorExt;
110 let dataflow_id = *scope.addr().into_first();
111 let as_of_frontier = dataflow
112 .as_of
113 .clone()
114 .unwrap_or_else(|| Antichain::from_elem(Timestamp::minimum()));
115
116 let export_ids = dataflow.export_ids().collect();
117
118 let compute_logger = if dataflow.is_transient() {
122 None
123 } else {
124 compute_state.compute_logger.clone()
125 };
126
127 Self {
128 scope,
129 debug_name: dataflow.debug_name.clone(),
130 dataflow_id,
131 export_ids,
132 as_of_frontier,
133 until,
134 bindings: BTreeMap::new(),
135 compute_logger,
136 linear_join_spec: compute_state.linear_join_spec,
137 dataflow_expiration,
138 config_set: Rc::clone(&compute_state.worker_config),
139 }
140 }
141}
142
143impl<'scope, T: RenderTimestamp> Context<'scope, T> {
144 pub fn insert_id(
149 &mut self,
150 id: Id,
151 collection: CollectionBundle<'scope, T>,
152 ) -> Option<CollectionBundle<'scope, T>> {
153 self.bindings.insert(id, collection)
154 }
155 pub fn remove_id(&mut self, id: Id) -> Option<CollectionBundle<'scope, T>> {
159 self.bindings.remove(&id)
160 }
161 pub fn update_id(&mut self, id: Id, collection: CollectionBundle<'scope, T>) {
163 if !self.bindings.contains_key(&id) {
164 self.bindings.insert(id, collection);
165 } else {
166 let binding = self
167 .bindings
168 .get_mut(&id)
169 .expect("Binding verified to exist");
170 if collection.collection.is_some() {
171 binding.collection = collection.collection;
172 }
173 for (key, flavor) in collection.arranged.into_iter() {
174 binding.arranged.insert(key, flavor);
175 }
176 }
177 }
178 pub fn lookup_id(&self, id: Id) -> Option<CollectionBundle<'scope, T>> {
180 self.bindings.get(&id).cloned()
181 }
182
183 pub(super) fn error_logger(&self) -> ErrorLogger {
184 ErrorLogger::new(self.debug_name.clone())
185 }
186}
187
188impl<'scope, T: RenderTimestamp> Context<'scope, T> {
189 pub fn enter_region<'a>(
191 &self,
192 region: Scope<'a, T>,
193 bindings: Option<&std::collections::BTreeSet<Id>>,
194 ) -> Context<'a, T> {
195 let bindings = self
196 .bindings
197 .iter()
198 .filter(|(key, _)| bindings.as_ref().map(|b| b.contains(key)).unwrap_or(true))
199 .map(|(key, bundle)| (*key, bundle.enter_region(region)))
200 .collect();
201
202 Context {
203 scope: region,
204 debug_name: self.debug_name.clone(),
205 dataflow_id: self.dataflow_id.clone(),
206 export_ids: self.export_ids.clone(),
207 as_of_frontier: self.as_of_frontier.clone(),
208 until: self.until.clone(),
209 compute_logger: self.compute_logger.clone(),
210 linear_join_spec: self.linear_join_spec.clone(),
211 bindings,
212 dataflow_expiration: self.dataflow_expiration.clone(),
213 config_set: Rc::clone(&self.config_set),
214 }
215 }
216}
217
218#[derive(Clone)]
220pub enum ArrangementFlavor<'scope, T: RenderTimestamp> {
221 Local(
223 Arranged<'scope, RowRowAgent<T, Diff>>,
224 Arranged<'scope, ErrAgent<T, Diff>>,
225 ),
226 Trace(
231 GlobalId,
232 Arranged<'scope, RowRowEnter<mz_repr::Timestamp, Diff, T>>,
233 Arranged<'scope, ErrEnter<mz_repr::Timestamp, T>>,
234 ),
235}
236
237impl<'scope, T: RenderTimestamp> ArrangementFlavor<'scope, T> {
238 pub fn flat_map<DCB, L>(
271 &self,
272 key: Option<&Row>,
273 max_demand: usize,
274 logic: L,
275 ) -> (
276 Stream<'scope, T, DCB::Container>,
277 VecCollection<'scope, T, DataflowErrorSer, Diff>,
278 )
279 where
280 DCB: ContainerBuilder,
281 L: for<'a, 'b> FnMut(
282 &'a mut DatumVecBorrow<'b>,
283 T,
284 Diff,
285 &mut Session<T, DCB>,
286 &mut Session<T, ECB<T>>,
287 ) -> usize
288 + 'static,
289 {
290 match &self {
293 ArrangementFlavor::Local(oks, errs) => {
294 let (oks, mfp_errs) = CollectionBundle::<T>::flat_map_core_fallible::<_, DCB, _>(
295 oks.clone(),
296 key,
297 max_demand,
298 logic,
299 REFUEL,
300 );
301 let errs = errs.clone().as_collection(|k, &()| k.clone());
302 let errs = errs.concat(mfp_errs.as_collection());
303 (oks, errs)
304 }
305 ArrangementFlavor::Trace(_, oks, errs) => {
306 let (oks, mfp_errs) = CollectionBundle::<T>::flat_map_core_fallible::<_, DCB, _>(
307 oks.clone(),
308 key,
309 max_demand,
310 logic,
311 REFUEL,
312 );
313 let errs = errs.clone().as_collection(|k, &()| k.clone());
314 let errs = errs.concat(mfp_errs.as_collection());
315 (oks, errs)
316 }
317 }
318 }
319
320 pub fn flat_map_ok<DCB, L>(
325 &self,
326 key: Option<&Row>,
327 max_demand: usize,
328 logic: L,
329 ) -> (
330 Stream<'scope, T, DCB::Container>,
331 VecCollection<'scope, T, DataflowErrorSer, Diff>,
332 )
333 where
334 DCB: ContainerBuilder,
337 L: for<'a, 'b> FnMut(&'a mut DatumVecBorrow<'b>, T, Diff, &mut Session<T, DCB>) -> usize
338 + 'static,
339 {
340 match &self {
341 ArrangementFlavor::Local(oks, errs) => {
342 let oks = CollectionBundle::<T>::flat_map_core_ok::<_, DCB, _>(
343 oks.clone(),
344 key,
345 max_demand,
346 logic,
347 REFUEL,
348 );
349 let errs = errs.clone().as_collection(|k, &()| k.clone());
350 (oks, errs)
351 }
352 ArrangementFlavor::Trace(_, oks, errs) => {
353 let oks = CollectionBundle::<T>::flat_map_core_ok::<_, DCB, _>(
354 oks.clone(),
355 key,
356 max_demand,
357 logic,
358 REFUEL,
359 );
360 let errs = errs.clone().as_collection(|k, &()| k.clone());
361 (oks, errs)
362 }
363 }
364 }
365}
366impl<'scope, T: RenderTimestamp> ArrangementFlavor<'scope, T> {
367 pub fn scope(&self) -> Scope<'scope, T> {
369 match self {
370 ArrangementFlavor::Local(oks, _errs) => oks.stream.scope(),
371 ArrangementFlavor::Trace(_gid, oks, _errs) => oks.stream.scope(),
372 }
373 }
374
375 pub fn enter_region<'a>(&self, region: Scope<'a, T>) -> ArrangementFlavor<'a, T> {
377 match self {
378 ArrangementFlavor::Local(oks, errs) => ArrangementFlavor::Local(
379 oks.clone().enter_region(region),
380 errs.clone().enter_region(region),
381 ),
382 ArrangementFlavor::Trace(gid, oks, errs) => ArrangementFlavor::Trace(
383 *gid,
384 oks.clone().enter_region(region),
385 errs.clone().enter_region(region),
386 ),
387 }
388 }
389}
390impl<'scope, T: RenderTimestamp> ArrangementFlavor<'scope, T> {
391 pub fn leave_region<'outer>(&self, outer: Scope<'outer, T>) -> ArrangementFlavor<'outer, T> {
393 match self {
394 ArrangementFlavor::Local(oks, errs) => ArrangementFlavor::Local(
395 oks.clone().leave_region(outer),
396 errs.clone().leave_region(outer),
397 ),
398 ArrangementFlavor::Trace(gid, oks, errs) => ArrangementFlavor::Trace(
399 *gid,
400 oks.clone().leave_region(outer),
401 errs.clone().leave_region(outer),
402 ),
403 }
404 }
405}
406pub(crate) fn distinct_arranged_errs<'a, T: RenderTimestamp>(
418 errs: Arranged<'a, ErrAgent<T, Diff>>,
419 name: &str,
420) -> Arranged<'a, ErrAgent<T, Diff>> {
421 errs.mz_reduce_abelian::<_, ErrBuilder<_, _>, ErrSpine<_, _>, _>(
422 name,
423 |_err, _input, output| output.push(((), Diff::ONE)),
424 )
425}
426
427pub(crate) fn distinct_errs_collection<'a, T: RenderTimestamp>(
432 errs: VecCollection<'a, T, DataflowErrorSer, Diff>,
433) -> VecCollection<'a, T, DataflowErrorSer, Diff> {
434 let errs: KeyCollection<_, _, _> = errs.into();
435 let errs = errs
436 .mz_arrange::<ColumnationChunker<_>, ErrBatcher<_, _>, ErrBuilder<_, _>, ErrSpine<_, _>>(
437 "Arrange errors",
438 );
439 distinct_arranged_errs(errs, "Distinct errors").as_collection(|err, _| err.clone())
440}
441
442#[derive(Clone)]
447pub struct CollectionBundle<'scope, T: RenderTimestamp> {
448 pub collection: Option<(
449 ColCollection<'scope, T>,
450 VecCollection<'scope, T, DataflowErrorSer, Diff>,
451 )>,
452 pub arranged: BTreeMap<Vec<LirScalarExpr>, ArrangementFlavor<'scope, T>>,
453}
454
455impl<'scope, T: RenderTimestamp> CollectionBundle<'scope, T> {
456 pub fn from_edge(
458 oks: ColCollection<'scope, T>,
459 errs: VecCollection<'scope, T, DataflowErrorSer, Diff>,
460 ) -> Self {
461 Self {
462 collection: Some((oks, errs)),
463 arranged: BTreeMap::default(),
464 }
465 }
466
467 pub fn from_expressions(
469 exprs: Vec<LirScalarExpr>,
470 arrangements: ArrangementFlavor<'scope, T>,
471 ) -> Self {
472 let mut arranged = BTreeMap::new();
473 arranged.insert(exprs, arrangements);
474 Self {
475 collection: None,
476 arranged,
477 }
478 }
479
480 pub fn from_columns<I: IntoIterator<Item = usize>>(
482 columns: I,
483 arrangements: ArrangementFlavor<'scope, T>,
484 ) -> Self {
485 let mut keys = Vec::new();
486 for column in columns {
487 keys.push(LirScalarExpr::column(column));
488 }
489 Self::from_expressions(keys, arrangements)
490 }
491
492 pub fn scope(&self) -> Scope<'scope, T> {
494 if let Some((oks, _errs)) = &self.collection {
495 oks.scope()
496 } else {
497 self.arranged
498 .values()
499 .next()
500 .expect("Must contain a valid collection")
501 .scope()
502 }
503 }
504
505 pub fn distinct_errs(mut self) -> Self {
531 if let Some((oks, errs)) = self.collection.take() {
532 self.collection = Some((oks, distinct_errs_collection(errs)));
533 }
534 for (key, flavor) in std::mem::take(&mut self.arranged) {
535 let flavor = match flavor {
536 ArrangementFlavor::Local(oks, errs) => {
537 let name = format!("Distinct errors[{key:?}]");
540 ArrangementFlavor::Local(oks, distinct_arranged_errs(errs, &name))
541 }
542 flavor @ ArrangementFlavor::Trace(..) => flavor,
543 };
544 self.arranged.insert(key, flavor);
545 }
546 self
547 }
548
549 pub fn enter_region<'inner>(&self, region: Scope<'inner, T>) -> CollectionBundle<'inner, T> {
551 CollectionBundle {
552 collection: self.collection.as_ref().map(|(oks, errs)| {
553 (
554 oks.clone().enter_region(region),
555 errs.clone().enter_region(region),
556 )
557 }),
558 arranged: self
559 .arranged
560 .iter()
561 .map(|(key, bundle)| (key.clone(), bundle.enter_region(region)))
562 .collect(),
563 }
564 }
565}
566
567impl<'scope, T: RenderTimestamp> CollectionBundle<'scope, T> {
568 pub fn leave_region<'outer>(&self, outer: Scope<'outer, T>) -> CollectionBundle<'outer, T> {
570 CollectionBundle {
571 collection: self.collection.as_ref().map(|(oks, errs)| {
572 (
573 oks.clone().leave_region(outer),
574 errs.clone().leave_region(outer),
575 )
576 }),
577 arranged: self
578 .arranged
579 .iter()
580 .map(|(key, bundle)| (key.clone(), bundle.leave_region(outer)))
581 .collect(),
582 }
583 }
584}
585
586impl<'scope, T: RenderTimestamp> CollectionBundle<'scope, T> {
587 pub fn as_specific_collection(
602 &self,
603 key: Option<&[LirScalarExpr]>,
604 ) -> (
605 ColCollection<'scope, T>,
606 VecCollection<'scope, T, DataflowErrorSer, Diff>,
607 ) {
608 match key {
614 None => self
615 .collection
616 .clone()
617 .expect("The unarranged collection doesn't exist."),
618 Some(key) => {
619 let arranged = self.arranged.get(key).unwrap_or_else(|| {
620 panic!("The collection arranged by {:?} doesn't exist.", key)
621 });
622 let (ok, err) =
626 arranged.flat_map_ok::<ColumnBuilder<(Row, T, Diff)>, _>(None, usize::MAX, {
627 let mut row_buf = Row::default();
630 move |borrow, t, r, ok_session| {
631 row_buf.packer().extend(borrow.iter());
632 ok_session.give((&row_buf, &t, &r));
633 1
634 }
635 });
636 (ok.as_collection(), err)
637 }
638 }
639 }
640
641 pub fn flat_map<DCB, L>(
657 &self,
658 key_val: Option<(Vec<LirScalarExpr>, Option<Row>)>,
659 max_demand: usize,
660 logic: L,
661 ) -> (
662 Stream<'scope, T, DCB::Container>,
663 VecCollection<'scope, T, DataflowErrorSer, Diff>,
664 )
665 where
666 DCB: ContainerBuilder,
667 L: for<'a> FnMut(
668 &'a mut DatumVecBorrow<'_>,
669 T,
670 Diff,
671 &mut Session<T, DCB>,
672 &mut Session<T, ECB<T>>,
673 ) -> usize
674 + 'static,
675 {
676 if let Some((key, val)) = key_val {
680 self.arrangement(&key)
681 .expect("Should have ensured during planning that this arrangement exists.")
682 .flat_map::<DCB, _>(val.as_ref(), max_demand, logic)
683 } else {
684 let (oks, errs) = self
685 .collection
686 .clone()
687 .expect("Invariant violated: CollectionBundle contains no collection.");
688 let (ok_stream, err_stream) =
689 flat_map_datums::<_, DCB, _>(oks, "CollectionFlatMap", max_demand, logic);
690 let errs = errs.concat(err_stream.as_collection());
691 (ok_stream, errs)
692 }
693 }
694
695 fn flat_map_core_fallible<Tr, DCB, L>(
706 trace: Arranged<'scope, Tr>,
707 key: Option<&<<BatchCursor<Tr> as Cursor>::KeyContainer as BatchContainer>::Owned>,
708 max_demand: usize,
709 mut logic: L,
710 refuel: usize,
711 ) -> (
712 Stream<'scope, T, DCB::Container>,
713 Stream<'scope, T, Vec<(DataflowErrorSer, T, Diff)>>,
714 )
715 where
716 Tr: TraceReader<Batch: Navigable, Time = T> + Clone + 'static,
717 for<'a> BatchCursor<Tr>:
718 Cursor<Key<'a>: ExtendDatums, Val<'a>: ExtendDatums, Time = T, Diff = mz_repr::Diff>,
719 <<BatchCursor<Tr> as Cursor>::KeyContainer as BatchContainer>::Owned: PartialEq,
720 DCB: ContainerBuilder,
721 L: for<'a, 'b> FnMut(
725 &'a mut DatumVecBorrow<'b>,
726 T,
727 mz_repr::Diff,
728 &mut Session<T, DCB>,
729 &mut Session<T, ECB<T>>,
730 ) -> usize
731 + 'static,
732 {
733 let scope = trace.stream.scope();
734
735 let mut key_con = <BatchCursor<Tr> as Cursor>::KeyContainer::with_capacity(1);
736 if let Some(key) = &key {
737 key_con.push_own(key);
738 }
739 let mode = if key.is_some() { "index" } else { "scan" };
740 let name = format!("ArrangementFlatMap({})", mode);
741
742 let mut builder = OperatorBuilder::new(name, scope.clone());
743 let (ok_output, ok_stream) = builder.new_output();
744 let mut ok_output = OutputBuilder::<_, DCB>::from(ok_output);
745 let (err_output, err_stream) = builder.new_output();
746 let mut err_output = OutputBuilder::<_, ECB<T>>::from(err_output);
747 let mut input = builder.new_input(trace.stream.clone(), Pipeline);
748 let operator_info = builder.operator_info();
749
750 builder.build(move |_capabilities| {
751 let activator = scope.activator_for(operator_info.address);
753 let mut todo = std::collections::VecDeque::new();
755 move |_frontiers| {
756 let key = key_con.get(0);
757 let mut ok_output = ok_output.activate();
758 let mut err_output = err_output.activate();
759
760 input.for_each(|time, data| {
762 let ok_cap = time.retain(0);
765 let err_cap = time.retain(1);
766 for batch in data.iter() {
767 todo.push_back(PendingWork::new(
768 ok_cap.clone(),
769 err_cap.clone(),
770 batch.cursor(),
771 batch.clone(),
772 ));
773 }
774 });
775
776 let mut temp_storage = RowArena::new();
781 let mut datums = DatumVec::new();
782 let mut decode_logic =
783 |k: BatchKey<'_, Tr>,
784 v: BatchVal<'_, Tr>,
785 t: T,
786 d: mz_repr::Diff,
787 ok_session: &mut Session<T, DCB>,
788 err_session: &mut Session<T, ECB<T>>| {
789 temp_storage.clear();
790 let mut datums_borrow = datums.borrow();
791 k.extend_datums(&temp_storage, &mut datums_borrow, Some(max_demand));
792 let remaining = max_demand.saturating_sub(datums_borrow.len());
793 v.extend_datums(&temp_storage, &mut datums_borrow, Some(remaining));
794 logic(&mut datums_borrow, t, d, ok_session, err_session)
795 };
796
797 let mut fuel = refuel;
799 while !todo.is_empty() && fuel > 0 {
800 todo.front_mut().unwrap().do_work(
801 key.as_ref(),
802 &mut decode_logic,
803 &mut fuel,
804 &mut ok_output,
805 &mut err_output,
806 );
807 if fuel > 0 {
808 todo.pop_front();
809 }
810 }
811 if !todo.is_empty() {
813 activator.activate();
814 }
815 }
816 });
817
818 (ok_stream, err_stream)
819 }
820
821 fn flat_map_core_ok<Tr, DCB, L>(
827 trace: Arranged<'scope, Tr>,
828 key: Option<&<<BatchCursor<Tr> as Cursor>::KeyContainer as BatchContainer>::Owned>,
829 max_demand: usize,
830 mut logic: L,
831 refuel: usize,
832 ) -> Stream<'scope, T, DCB::Container>
833 where
834 Tr: TraceReader<Batch: Navigable, Time = T> + Clone + 'static,
835 for<'a> BatchCursor<Tr>:
836 Cursor<Key<'a>: ExtendDatums, Val<'a>: ExtendDatums, Time = T, Diff = mz_repr::Diff>,
837 <<BatchCursor<Tr> as Cursor>::KeyContainer as BatchContainer>::Owned: PartialEq,
838 DCB: ContainerBuilder,
839 L: for<'a, 'b> FnMut(
840 &'a mut DatumVecBorrow<'b>,
841 T,
842 mz_repr::Diff,
843 &mut Session<T, DCB>,
844 ) -> usize
845 + 'static,
846 {
847 let scope = trace.stream.scope();
848
849 let mut key_con = <BatchCursor<Tr> as Cursor>::KeyContainer::with_capacity(1);
850 if let Some(key) = &key {
851 key_con.push_own(key);
852 }
853 let mode = if key.is_some() { "index" } else { "scan" };
854 let name = format!("ArrangementFlatMapOk({})", mode);
855
856 let mut builder = OperatorBuilder::new(name, scope.clone());
857 let (ok_output, ok_stream) = builder.new_output();
858 let mut ok_output = OutputBuilder::<_, DCB>::from(ok_output);
859 let mut input = builder.new_input(trace.stream.clone(), Pipeline);
860 let operator_info = builder.operator_info();
861
862 builder.build(move |_capabilities| {
863 let activator = scope.activator_for(operator_info.address);
864 let mut todo = std::collections::VecDeque::new();
865 move |_frontiers| {
866 let key = key_con.get(0);
867 let mut ok_output = ok_output.activate();
868
869 input.for_each(|time, data| {
870 let cap = time.retain(0);
871 for batch in data.iter() {
872 todo.push_back(PendingWorkOk::new(
873 cap.clone(),
874 batch.cursor(),
875 batch.clone(),
876 ));
877 }
878 });
879
880 let mut temp_storage = RowArena::new();
882 let mut datums = DatumVec::new();
883 let mut decode_logic =
884 |k: BatchKey<'_, Tr>,
885 v: BatchVal<'_, Tr>,
886 t: T,
887 d: mz_repr::Diff,
888 ok_session: &mut Session<T, DCB>| {
889 temp_storage.clear();
890 let mut datums_borrow = datums.borrow();
891 k.extend_datums(&temp_storage, &mut datums_borrow, Some(max_demand));
892 let remaining = max_demand.saturating_sub(datums_borrow.len());
893 v.extend_datums(&temp_storage, &mut datums_borrow, Some(remaining));
894 logic(&mut datums_borrow, t, d, ok_session)
895 };
896
897 let mut fuel = refuel;
898 while !todo.is_empty() && fuel > 0 {
899 todo.front_mut().unwrap().do_work(
900 key.as_ref(),
901 &mut decode_logic,
902 &mut fuel,
903 &mut ok_output,
904 );
905 if fuel > 0 {
906 todo.pop_front();
907 }
908 }
909 if !todo.is_empty() {
910 activator.activate();
911 }
912 }
913 });
914
915 ok_stream
916 }
917
918 pub fn arrangement(&self, key: &[LirScalarExpr]) -> Option<ArrangementFlavor<'scope, T>> {
923 self.arranged.get(key).map(|x| x.clone())
924 }
925}
926
927impl<'scope, T: RenderTimestamp> CollectionBundle<'scope, T> {
928 pub fn as_collection_core(
937 &self,
938 mfp_plan: MfpPlan<LirScalarExpr>,
939 key_val: Option<(Vec<LirScalarExpr>, Option<StableRow>)>,
940 until: Antichain<mz_repr::Timestamp>,
941 ) -> (
942 ColCollection<'scope, T>,
943 VecCollection<'scope, T, DataflowErrorSer, Diff>,
944 ) {
945 let key_val = key_val.map(|(key, val)| (key, val.map(|val| val.0)));
948 let has_key_val = if let Some((_key, Some(_val))) = &key_val {
954 true
955 } else {
956 false
957 };
958
959 if mfp_plan.is_identity() && !has_key_val {
960 let key = key_val.map(|(k, _v)| k);
961 return match key {
962 None => self
965 .collection
966 .clone()
967 .expect("The unarranged collection doesn't exist."),
968 Some(key) => self.as_specific_collection(Some(&key)),
969 };
970 }
971
972 let (mfp_plan, max_demand) = {
977 let mut mir_mfp = mfp_plan_lir_to_mir(mfp_plan).into_map_filter_project();
978 let max_demand = mir_mfp.demand().last().map(|x| *x + 1).unwrap_or(0);
979 mir_mfp.permute_fn(|c| c, max_demand);
980 mir_mfp.optimize();
981 let plan = mfp_mir_to_lir_plan(mir_mfp);
982 (plan, max_demand)
983 };
984
985 let mut datum_vec = DatumVec::new();
986 let until = std::rc::Rc::new(until);
988
989 let (stream, errors) = self.flat_map::<ConsolidatingColumnBuilder<Row, T, Diff>, _>(
993 key_val,
994 max_demand,
995 move |row_datums, time, diff, ok_session, err_session| {
996 let mut row_builder = SharedRow::get();
997 let until = std::rc::Rc::clone(&until);
998 let temp_storage = RowArena::new();
999 let row_iter = row_datums.iter();
1000 let mut datums_local = datum_vec.borrow();
1001 datums_local.extend(row_iter);
1002 let event_time = time.event_time();
1003 let mut work: usize = 0;
1004 for result in mfp_plan.evaluate(
1005 &mut datums_local,
1006 &temp_storage,
1007 event_time,
1008 diff.clone(),
1009 move |time| !until.less_equal(time),
1010 &mut row_builder,
1011 ) {
1012 work += 1;
1013 match result {
1014 Ok((row, event_time, diff)) => {
1015 let mut time: T = time.clone();
1017 *time.event_time_mut() = event_time;
1018 ok_session.give((row, time, diff));
1019 }
1020 Err((e, event_time, diff)) => {
1021 let mut time: T = time.clone();
1023 *time.event_time_mut() = event_time;
1024 err_session.give((e, time, diff));
1025 }
1026 }
1027 }
1028 work
1029 },
1030 );
1031
1032 (stream.as_collection(), errors)
1033 }
1034 pub fn ensure_collections(
1035 mut self,
1036 collections: AvailableCollections,
1037 input_key: Option<Vec<LirScalarExpr>>,
1038 input_mfp: MfpPlan<LirScalarExpr>,
1039 as_of: Antichain<mz_repr::Timestamp>,
1040 until: Antichain<mz_repr::Timestamp>,
1041 config_set: &ConfigSet,
1042 strategy: ArrangementStrategy,
1043 ) -> Self
1044 where
1045 T: MaybeBucketByTime,
1046 {
1047 if collections == Default::default() {
1048 return self;
1049 }
1050 for (key, _, _) in collections.arranged.iter() {
1059 soft_assert_or_log!(
1060 !self.arranged.contains_key(key),
1061 "LIR ArrangeBy tried to create an existing arrangement"
1062 );
1063 }
1064
1065 let mut bucketed = false;
1068
1069 let will_create_arrangement = collections
1073 .arranged
1074 .iter()
1075 .any(|(key, _, _)| !self.arranged.contains_key(key));
1076
1077 let form_raw_collection = collections.raw || will_create_arrangement;
1079 if form_raw_collection && self.collection.is_none() {
1080 let (oks, errs) =
1081 self.as_collection_core(input_mfp, input_key.map(|k| (k, None)), until);
1082 let effective_strategy = if will_create_arrangement {
1086 strategy
1087 } else {
1088 ArrangementStrategy::Direct
1089 };
1090 let oks = if matches!(effective_strategy, ArrangementStrategy::TemporalBucketing)
1091 && ENABLE_COMPUTE_TEMPORAL_BUCKETING.get(config_set)
1092 {
1093 let summary: mz_repr::Timestamp = TEMPORAL_BUCKETING_SUMMARY
1094 .get(config_set)
1095 .try_into()
1096 .expect("must fit");
1097 bucketed = true;
1098 T::maybe_apply_temporal_bucketing(oks.inner, as_of.clone(), summary)
1100 } else {
1101 oks
1102 };
1103 self.collection = Some((oks, errs));
1104 }
1105 for (key, _, thinning) in collections.arranged {
1106 if !self.arranged.contains_key(&key) {
1107 let name = format!("ArrangeBy[{:?}]", key);
1109
1110 let (oks, errs) = self
1111 .collection
1112 .take()
1113 .expect("Collection constructed above");
1114 let effective_strategy = if bucketed {
1119 ArrangementStrategy::Direct
1120 } else {
1121 strategy
1122 };
1123 let oks = if matches!(effective_strategy, ArrangementStrategy::TemporalBucketing)
1124 && ENABLE_COMPUTE_TEMPORAL_BUCKETING.get(config_set)
1125 {
1126 let summary: mz_repr::Timestamp = TEMPORAL_BUCKETING_SUMMARY
1127 .get(config_set)
1128 .try_into()
1129 .expect("must fit");
1130 bucketed = true;
1131 T::maybe_apply_temporal_bucketing(oks.inner, as_of.clone(), summary)
1133 } else {
1134 oks
1135 };
1136 let batcher = ArrangementBatcher::from_config(config_set);
1137 let (oks, errs_keyed, passthrough) =
1138 Self::arrange_collection(&name, oks, key.clone(), thinning.clone(), batcher);
1139 let errs_concat: KeyCollection<_, _, _> = errs.clone().concat(errs_keyed).into();
1140 self.collection = Some((passthrough, errs));
1141 let errs =
1142 errs_concat.mz_arrange::<
1143 ColumnationChunker<_>,
1144 ErrBatcher<_, _>,
1145 ErrBuilder<_, _>,
1146 ErrSpine<_, _>,
1147 >(
1148 &format!("{}-errors", name),
1149 );
1150 self.arranged
1151 .insert(key, ArrangementFlavor::Local(oks, errs));
1152 }
1153 }
1154 self
1155 }
1156
1157 fn arrange_collection(
1168 name: &String,
1169 oks: ColCollection<'scope, T>,
1170 key: Vec<LirScalarExpr>,
1171 thinning: Vec<usize>,
1172 batcher: ArrangementBatcher,
1173 ) -> (
1174 Arranged<'scope, RowRowAgent<T, Diff>>,
1175 VecCollection<'scope, T, DataflowErrorSer, Diff>,
1176 ColCollection<'scope, T>,
1177 ) {
1178 let (ok_stream, err_stream, passthrough) = {
1184 let mut builder =
1185 OperatorBuilder::new("FormArrangementKey".to_string(), oks.inner.scope());
1186 let (ok_output, ok_stream) = builder.new_output();
1187 let mut ok_output =
1188 OutputBuilder::<_, ColumnBuilder<((Row, Row), T, Diff)>>::from(ok_output);
1189 let (err_output, err_stream) = builder.new_output();
1190 let mut err_output = OutputBuilder::from(err_output);
1191 let (passthrough_output, passthrough_stream) = builder.new_output();
1192 let mut passthrough_output =
1193 OutputBuilder::<_, NoopBuilder<Column<(Row, T, Diff)>>>::from(passthrough_output);
1194 let mut input = builder.new_input(oks.inner, Pipeline);
1195 builder.set_notify_for(0, FrontierInterest::Never);
1196 builder.build(move |_capabilities| {
1197 let mut key_buf = Row::default();
1198 let mut val_buf = Row::default();
1199 let mut datums = DatumVec::new();
1200 move |_frontiers| {
1201 let mut temp_storage = RowArena::new();
1205 let mut ok_output = ok_output.activate();
1206 let mut err_output = err_output.activate();
1207 let mut passthrough_output = passthrough_output.activate();
1208 input.for_each(|time, data| {
1209 let mut ok_session = ok_output.session_with_builder(&time);
1210 let mut err_session = err_output.session(&time);
1211 for (row, t, d) in data.borrow().into_index_iter() {
1214 temp_storage.clear();
1215 let datums = datums.borrow_with(row);
1216 let key_iter = key.iter().map(|k| k.eval(&datums, &temp_storage));
1217 match key_buf.packer().try_extend(key_iter) {
1218 Ok(()) => {
1219 let val_datum_iter = thinning.iter().map(|c| datums[*c]);
1220 val_buf.packer().extend(val_datum_iter);
1221 ok_session.give(((&*key_buf, &*val_buf), t, d));
1222 }
1223 Err(e) => {
1224 err_session.give((
1225 e.into(),
1226 Columnar::into_owned(t),
1227 Columnar::into_owned(d),
1228 ));
1229 }
1230 }
1231 }
1232 passthrough_output
1233 .session_with_builder(&time)
1234 .give_container(data);
1235 });
1236 }
1237 });
1238 (ok_stream, err_stream, passthrough_stream.as_collection())
1239 };
1240
1241 let exchange =
1242 ExchangeCore::<ColumnBuilder<_>, _>::new_core(columnar_exchange::<Row, Row, T, Diff>);
1243 let oks = match batcher {
1244 ArrangementBatcher::Chunked => ok_stream.mz_arrange_core::<
1245 _,
1246 ChunkChunker<(Row, Row), T, Diff>,
1247 AccountedChunkBatcher<(Row, Row), T, Diff>,
1248 UnchunkBuilder<RowRowColPagedBuilder<T, Diff>, (Row, Row), T, Diff>,
1249 RowRowSpine<_, _>,
1250 >(exchange, name),
1251 ArrangementBatcher::Columnar => ok_stream.mz_arrange_core::<
1252 _,
1253 batcher::ColumnChunker<_>,
1254 Col2ValColBatcher<_, _, _, _>,
1255 RowRowColPagedBuilder<_, _>,
1256 RowRowSpine<_, _>,
1257 >(exchange, name),
1258 ArrangementBatcher::Columnation => ok_stream.mz_arrange_core::<
1259 _,
1260 batcher::Chunker<_>,
1261 Col2ValBatcher<_, _, _, _>,
1262 RowRowBuilder<_, _>,
1263 RowRowSpine<_, _>,
1264 >(exchange, name),
1265 };
1266 (oks, err_stream.as_collection(), passthrough)
1267 }
1268}
1269
1270pub(crate) type Session<'a, 'b, T, CB> =
1274 timely::dataflow::operators::generic::Session<'a, 'b, T, CB, Capability<T>>;
1275
1276pub(crate) type ECB<T> = ConsolidatingContainerBuilder<Vec<(DataflowErrorSer, T, Diff)>>;
1281
1282const REFUEL: usize = 1_000_000;
1286
1287struct PendingWork<C>
1288where
1289 C: Cursor,
1290{
1291 ok_capability: Capability<C::Time>,
1293 err_capability: Capability<C::Time>,
1295 cursor: C,
1296 batch: C::Storage,
1297}
1298
1299impl<C> PendingWork<C>
1300where
1301 C: Cursor<KeyContainer: BatchContainer<Owned: PartialEq + Sized>>,
1302{
1303 fn new(
1306 ok_capability: Capability<C::Time>,
1307 err_capability: Capability<C::Time>,
1308 cursor: C,
1309 batch: C::Storage,
1310 ) -> Self {
1311 Self {
1312 ok_capability,
1313 err_capability,
1314 cursor,
1315 batch,
1316 }
1317 }
1318 fn do_work<DCB, L>(
1321 &mut self,
1322 key: Option<&C::Key<'_>>,
1323 logic: &mut L,
1324 fuel: &mut usize,
1325 ok_output: &mut OutputBuilderSession<'_, C::Time, DCB>,
1326 err_output: &mut OutputBuilderSession<'_, C::Time, ECB<C::Time>>,
1327 ) where
1328 DCB: ContainerBuilder,
1329 L: FnMut(
1330 C::Key<'_>,
1331 C::Val<'_>,
1332 C::Time,
1333 C::Diff,
1334 &mut Session<C::Time, DCB>,
1335 &mut Session<C::Time, ECB<C::Time>>,
1336 ) -> usize,
1337 {
1338 let mut ok_session = ok_output.session_with_builder(&self.ok_capability);
1339 let mut err_session = err_output.session_with_builder(&self.err_capability);
1340 walk_cursor(&mut self.cursor, &self.batch, key, fuel, |k, v, t, d| {
1341 logic(k, v, t, d, &mut ok_session, &mut err_session)
1342 });
1343 }
1344}
1345
1346struct PendingWorkOk<C>
1349where
1350 C: Cursor,
1351{
1352 capability: Capability<C::Time>,
1353 cursor: C,
1354 batch: C::Storage,
1355}
1356
1357impl<C> PendingWorkOk<C>
1358where
1359 C: Cursor<KeyContainer: BatchContainer<Owned: PartialEq + Sized>>,
1360{
1361 fn new(capability: Capability<C::Time>, cursor: C, batch: C::Storage) -> Self {
1362 Self {
1363 capability,
1364 cursor,
1365 batch,
1366 }
1367 }
1368
1369 fn do_work<DCB, L>(
1372 &mut self,
1373 key: Option<&C::Key<'_>>,
1374 logic: &mut L,
1375 fuel: &mut usize,
1376 ok_output: &mut OutputBuilderSession<'_, C::Time, DCB>,
1377 ) where
1378 DCB: ContainerBuilder,
1379 L: FnMut(C::Key<'_>, C::Val<'_>, C::Time, C::Diff, &mut Session<C::Time, DCB>) -> usize,
1380 {
1381 let mut ok_session = ok_output.session_with_builder(&self.capability);
1382 walk_cursor(&mut self.cursor, &self.batch, key, fuel, |k, v, t, d| {
1383 logic(k, v, t, d, &mut ok_session)
1384 });
1385 }
1386}
1387
1388fn walk_cursor<C, F>(
1398 cursor: &mut C,
1399 batch: &C::Storage,
1400 key: Option<&C::Key<'_>>,
1401 fuel: &mut usize,
1402 mut emit: F,
1403) where
1404 C: Cursor<KeyContainer: BatchContainer<Owned: PartialEq + Sized>>,
1405 F: FnMut(C::Key<'_>, C::Val<'_>, C::Time, C::Diff) -> usize,
1406{
1407 use differential_dataflow::consolidation::consolidate;
1408
1409 let mut work: usize = 0;
1410 let mut buffer = Vec::new();
1411 if let Some(key) = key {
1412 let key = C::KeyContainer::reborrow(*key);
1413 if cursor.get_key(batch).map(|k| k == key) != Some(true) {
1414 cursor.seek_key(batch, key);
1415 }
1416 if cursor.get_key(batch).map(|k| k == key) == Some(true) {
1417 let key = cursor.key(batch);
1418 while let Some(val) = cursor.get_val(batch) {
1419 cursor.map_times(batch, |time, diff| {
1420 buffer.push((C::owned_time(time), C::owned_diff(diff)));
1421 });
1422 consolidate(&mut buffer);
1423 for (time, diff) in buffer.drain(..) {
1424 work += emit(key, val, time, diff);
1425 }
1426 cursor.step_val(batch);
1427 if work >= *fuel {
1428 *fuel = 0;
1429 return;
1430 }
1431 }
1432 }
1433 } else {
1434 while let Some(key) = cursor.get_key(batch) {
1435 while let Some(val) = cursor.get_val(batch) {
1436 cursor.map_times(batch, |time, diff| {
1437 buffer.push((C::owned_time(time), C::owned_diff(diff)));
1438 });
1439 consolidate(&mut buffer);
1440 for (time, diff) in buffer.drain(..) {
1441 work += emit(key, val, time, diff);
1442 }
1443 cursor.step_val(batch);
1444 if work >= *fuel {
1445 *fuel = 0;
1446 return;
1447 }
1448 }
1449 cursor.step_key(batch);
1450 }
1451 }
1452 *fuel -= work;
1453}
1454
1455#[cfg(test)]
1456mod tests {
1457 use differential_dataflow::input::Input;
1458 use mz_expr::{EvalError, MapFilterProject};
1459 use mz_repr::{Datum, ReprScalarType, Timestamp};
1460 use timely::dataflow::operators::Capture;
1461 use timely::dataflow::operators::capture::{Event, Extract};
1462
1463 use super::*;
1464 use crate::render::columnar::{columnar_to_vec, vec_to_columnar};
1465
1466 type OkUpdate = ((Row, Row), Timestamp, Diff);
1467 type ErrUpdate = (DataflowErrorSer, Timestamp, Diff);
1468 type Captured<D> = std::sync::mpsc::Receiver<Event<Timestamp, Vec<D>>>;
1469
1470 fn extract_ok(captured: Captured<OkUpdate>) -> Vec<OkUpdate> {
1471 let mut updates: Vec<_> = captured
1472 .extract()
1473 .into_iter()
1474 .flat_map(|(_, data)| data)
1475 .collect();
1476 updates.sort();
1477 updates
1478 }
1479
1480 fn extract_err(captured: Captured<ErrUpdate>) -> Vec<(String, Timestamp, Diff)> {
1482 let mut updates: Vec<_> = captured
1483 .extract()
1484 .into_iter()
1485 .flat_map(|(_, data)| data)
1486 .map(|(e, t, d)| (format!("{e:?}"), t, d))
1487 .collect();
1488 updates.sort();
1489 updates
1490 }
1491
1492 fn arrange_columnar(
1494 rows: Vec<(Row, u64)>,
1495 key: Vec<LirScalarExpr>,
1496 ) -> (Vec<OkUpdate>, Vec<(String, Timestamp, Diff)>) {
1497 let thinning = vec![0, 1];
1498 let (ok, err) = timely::execute_directly(move |worker| {
1499 worker.dataflow::<Timestamp, _, _>(|scope| {
1500 let (mut input, collection) = scope.new_collection();
1501 let (arranged, errs, _passthrough) =
1502 CollectionBundle::<Timestamp>::arrange_collection(
1503 &"col".to_string(),
1504 vec_to_columnar(collection),
1505 key,
1506 thinning,
1507 ArrangementBatcher::Columnation,
1508 );
1509 let ok = arranged
1510 .as_collection(|k, v| (k.to_row(), v.to_row()))
1511 .inner
1512 .capture();
1513 let err = errs.inner.capture();
1514
1515 let max_time = rows.iter().map(|(_, t)| *t).max().unwrap_or(0);
1516 for (row, time) in rows {
1517 input.update_at(row, Timestamp::from(time), Diff::ONE);
1518 }
1519 input.advance_to(Timestamp::from(max_time + 1));
1520 input.flush();
1521 (ok, err)
1522 })
1523 });
1524 (extract_ok(ok), extract_err(err))
1525 }
1526
1527 fn test_rows() -> Vec<(Row, u64)> {
1530 vec![
1531 (Row::pack_slice(&[Datum::Int32(1), Datum::String("a")]), 0),
1532 (Row::pack_slice(&[Datum::Int32(2), Datum::String("b")]), 1),
1533 (Row::pack_slice(&[Datum::Int32(1), Datum::String("a")]), 1),
1534 (Row::pack_slice(&[Datum::Int32(3), Datum::Null]), 2),
1535 ]
1536 }
1537
1538 #[mz_ore::test]
1541 fn arrange_collection_keys_correctly() {
1542 let rows = test_rows();
1543 let mut expected: Vec<OkUpdate> = rows
1544 .iter()
1545 .map(|(row, t)| {
1546 let key = Row::pack_slice(&[row.iter().next().unwrap()]);
1547 ((key, row.clone()), Timestamp::from(*t), Diff::ONE)
1548 })
1549 .collect();
1550 expected.sort();
1551
1552 let (ok, err) = arrange_columnar(rows, vec![LirScalarExpr::column(0)]);
1553 assert_eq!(ok, expected);
1554 assert!(err.is_empty());
1555 }
1556
1557 #[mz_ore::test]
1559 fn arrange_collection_error_path() {
1560 let key = vec![LirScalarExpr::literal(
1561 Err(EvalError::DivisionByZero),
1562 ReprScalarType::Int32,
1563 )];
1564 let (ok, err) = arrange_columnar(test_rows(), key);
1565
1566 assert!(ok.is_empty());
1567 assert!(!err.is_empty());
1568 }
1569
1570 #[mz_ore::test]
1574 fn arrange_collection_passthrough_forwards_input() {
1575 let rows = test_rows();
1576 let mut expected: Vec<(Row, Timestamp, Diff)> = rows
1577 .iter()
1578 .map(|(row, t)| (row.clone(), Timestamp::from(*t), Diff::ONE))
1579 .collect();
1580 expected.sort();
1581
1582 let key = vec![LirScalarExpr::literal(
1583 Err(EvalError::DivisionByZero),
1584 ReprScalarType::Int32,
1585 )];
1586 let captured = timely::execute_directly(move |worker| {
1587 worker.dataflow::<Timestamp, _, _>(|scope| {
1588 let (mut input, collection) = scope.new_collection();
1589 let (_arranged, _errs, passthrough) =
1590 CollectionBundle::<Timestamp>::arrange_collection(
1591 &"col".to_string(),
1592 vec_to_columnar(collection),
1593 key,
1594 vec![0, 1],
1595 ArrangementBatcher::Columnation,
1596 );
1597 let captured = columnar_to_vec(passthrough).inner.capture();
1598
1599 let max_time = rows.iter().map(|(_, t)| *t).max().unwrap_or(0);
1600 for (row, time) in rows {
1601 input.update_at(row, Timestamp::from(time), Diff::ONE);
1602 }
1603 input.advance_to(Timestamp::from(max_time + 1));
1604 input.flush();
1605 captured
1606 })
1607 });
1608
1609 assert_eq!(extract_row_updates(captured), expected);
1610 }
1611
1612 fn extract_row_updates(
1613 captured: Captured<(Row, Timestamp, Diff)>,
1614 ) -> Vec<(Row, Timestamp, Diff)> {
1615 let mut updates: Vec<_> = captured
1616 .extract()
1617 .into_iter()
1618 .flat_map(|(_, data)| data)
1619 .collect();
1620 updates.sort();
1621 updates
1622 }
1623
1624 #[mz_ore::test]
1627 fn get_arrange_by_produces_projected_rows() {
1628 let rows = vec![
1629 (Row::pack_slice(&[Datum::Int64(1), Datum::Int64(10)]), 0u64),
1630 (Row::pack_slice(&[Datum::Int64(2), Datum::Int64(20)]), 1),
1631 (Row::pack_slice(&[Datum::Int64(1), Datum::Int64(10)]), 1),
1632 ];
1633 let mfp = MapFilterProject::<LirScalarExpr>::new(2)
1636 .project(vec![0])
1637 .into_plan()
1638 .expect("mfp");
1639 let mut expected: Vec<(Row, Timestamp, Diff)> = rows
1640 .iter()
1641 .map(|(r, t)| {
1642 let col0 = r.iter().next().unwrap();
1643 (Row::pack_slice(&[col0]), Timestamp::from(*t), Diff::ONE)
1644 })
1645 .collect();
1646 expected.sort();
1647
1648 let produced = timely::execute_directly(move |worker| {
1649 worker.dataflow::<Timestamp, _, _>(|scope| {
1650 let (mut input, collection) = scope.new_collection();
1651 let (_err_input, errs) = scope.new_collection::<DataflowErrorSer, Diff>();
1652 let bundle = CollectionBundle::from_edge(vec_to_columnar(collection), errs);
1653 let (edge, _errs) = bundle.as_collection_core(mfp, None, Antichain::new());
1654 let produced = columnar_to_vec(edge.clone()).inner.capture();
1655 let (_arranged, _arrange_errs, _passthrough) =
1656 CollectionBundle::<Timestamp>::arrange_collection(
1657 &"arrange".to_string(),
1658 edge,
1659 vec![LirScalarExpr::column(0)],
1660 vec![0],
1661 ArrangementBatcher::Columnation,
1662 );
1663
1664 let max_time = rows.iter().map(|(_, t)| *t).max().unwrap();
1665 for (row, time) in rows {
1666 input.update_at(row, Timestamp::from(time), Diff::ONE);
1667 }
1668 input.advance_to(Timestamp::from(max_time + 1));
1669 input.flush();
1670 produced
1671 })
1672 });
1673
1674 assert_eq!(extract_row_updates(produced), expected);
1675 }
1676
1677 #[mz_ore::test]
1680 fn as_collection_core_identity_passes_edge_through() {
1681 let expected = vec![(
1682 Row::pack_slice(&[Datum::Int64(1)]),
1683 Timestamp::from(0u64),
1684 Diff::ONE,
1685 )];
1686 let captured = timely::execute_directly(move |worker| {
1687 worker.dataflow::<Timestamp, _, _>(|scope| {
1688 let (mut input, collection) = scope.new_collection::<Row, Diff>();
1689 let (_err_input, errs) = scope.new_collection::<DataflowErrorSer, Diff>();
1690 let bundle = CollectionBundle::from_edge(vec_to_columnar(collection), errs);
1691 let identity = MapFilterProject::<LirScalarExpr>::new(1)
1692 .into_plan()
1693 .expect("identity mfp");
1694 let (out, _errs) = bundle.as_collection_core(identity, None, Antichain::new());
1695 let captured = columnar_to_vec(out).inner.capture();
1696 input.update_at(
1697 Row::pack_slice(&[Datum::Int64(1)]),
1698 Timestamp::from(0u64),
1699 Diff::ONE,
1700 );
1701 input.advance_to(Timestamp::from(1u64));
1702 input.flush();
1703 captured
1704 })
1705 });
1706 assert_eq!(extract_row_updates(captured), expected);
1707 }
1708
1709 #[mz_ore::test]
1712 fn as_collection_core_consolidates_within_batch() {
1713 let rows = vec![
1717 Row::pack_slice(&[Datum::Int64(1), Datum::Int64(10)]),
1718 Row::pack_slice(&[Datum::Int64(1), Datum::Int64(20)]),
1719 Row::pack_slice(&[Datum::Int64(2), Datum::Int64(30)]),
1720 ];
1721 let mfp = MapFilterProject::<LirScalarExpr>::new(2)
1722 .project(vec![0])
1723 .into_plan()
1724 .expect("mfp");
1725 let expected = vec![
1726 (
1727 Row::pack_slice(&[Datum::Int64(1)]),
1728 Timestamp::from(0u64),
1729 Diff::from(2),
1730 ),
1731 (
1732 Row::pack_slice(&[Datum::Int64(2)]),
1733 Timestamp::from(0u64),
1734 Diff::ONE,
1735 ),
1736 ];
1737
1738 let captured = timely::execute_directly(move |worker| {
1739 worker.dataflow::<Timestamp, _, _>(|scope| {
1740 let (mut input, collection) = scope.new_collection();
1741 let (_err_input, errs) = scope.new_collection::<DataflowErrorSer, Diff>();
1742 let bundle = CollectionBundle::from_edge(vec_to_columnar(collection), errs);
1743 let (edge, _errs) = bundle.as_collection_core(mfp, None, Antichain::new());
1744 let captured = columnar_to_vec(edge).inner.capture();
1745 for row in rows {
1748 input.update_at(row, Timestamp::from(0u64), Diff::ONE);
1749 }
1750 input.advance_to(Timestamp::from(1u64));
1751 input.flush();
1752 captured
1753 })
1754 });
1755
1756 assert_eq!(extract_row_updates(captured), expected);
1757 }
1758
1759 #[mz_ore::test]
1762 fn as_specific_collection_materializes_columnar() {
1763 let rows = test_rows();
1764 let key = vec![LirScalarExpr::column(0)];
1765 let mut expected: Vec<(Row, Timestamp, Diff)> = rows
1766 .iter()
1767 .map(|(r, t)| (r.clone(), Timestamp::from(*t), Diff::ONE))
1768 .collect();
1769 expected.sort();
1770
1771 let captured = timely::execute_directly(move |worker| {
1772 worker.dataflow::<Timestamp, _, _>(|scope| {
1773 let (mut input, collection) = scope.new_collection();
1774 let (arranged, arr_errs, _passthrough) =
1775 CollectionBundle::<Timestamp>::arrange_collection(
1776 &"agg".to_string(),
1777 vec_to_columnar(collection),
1778 key.clone(),
1779 vec![1],
1780 ArrangementBatcher::Columnation,
1781 );
1782 let err_arranged = {
1783 let kc: KeyCollection<_, _, _> = arr_errs.into();
1784 kc.mz_arrange::<
1785 ColumnationChunker<_>,
1786 ErrBatcher<_, _>,
1787 ErrBuilder<_, _>,
1788 ErrSpine<_, _>,
1789 >("agg-errs")
1790 };
1791 let bundle = CollectionBundle::from_columns(
1793 0..1,
1794 ArrangementFlavor::Local(arranged, err_arranged),
1795 );
1796 let (edge, _errs) = bundle.as_specific_collection(Some(&key));
1797 let captured = columnar_to_vec(edge).inner.capture();
1798
1799 let max_time = rows.iter().map(|(_, t)| *t).max().unwrap();
1800 for (row, time) in rows {
1801 input.update_at(row, Timestamp::from(time), Diff::ONE);
1802 }
1803 input.advance_to(Timestamp::from(max_time + 1));
1804 input.flush();
1805 captured
1806 })
1807 });
1808
1809 assert_eq!(extract_row_updates(captured), expected);
1810 }
1811}