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::consolidate::ConsolidatingColumnBuilder;
37use mz_timely_util::columnar::{
38 Col2ValBatcher, Col2ValColBatcher, Col2ValPagedBatcher, columnar_exchange,
39};
40use mz_timely_util::columnation::ColumnationChunker;
41use timely::ContainerBuilder;
42use timely::container::CapacityContainerBuilder;
43use timely::dataflow::channels::pact::{ExchangeCore, Pipeline};
44use timely::dataflow::operators::Capability;
45use timely::dataflow::operators::generic::builder_rc::OperatorBuilder;
46use timely::dataflow::operators::generic::{OutputBuilder, OutputBuilderSession};
47use timely::dataflow::{Scope, Stream};
48use timely::progress::operate::FrontierInterest;
49use timely::progress::{Antichain, Timestamp};
50
51use crate::compute_state::ComputeState;
52use crate::extensions::arrange::{ArrangementBatcher, KeyCollection, MzArrange, MzArrangeCore};
53use crate::extensions::reduce::MzReduce;
54use crate::render::columnar::{CollectionEdge, flat_map_datums};
55use crate::render::errors::{DataflowErrorSer, ErrorLogger};
56use crate::render::{LinearJoinSpec, MaybeBucketByTime, RenderTimestamp};
57use crate::typedefs::{
58 ErrAgent, ErrBatcher, ErrBuilder, ErrEnter, ErrSpine, RowRowAgent, RowRowEnter, RowRowSpine,
59};
60use mz_row_spine::{RowRowBuilder, RowRowColPagedBuilder};
61
62pub struct Context<'scope, T: RenderTimestamp> {
70 pub(crate) scope: Scope<'scope, T>,
74 pub debug_name: String,
76 pub dataflow_id: usize,
78 pub export_ids: Vec<GlobalId>,
80 pub as_of_frontier: Antichain<mz_repr::Timestamp>,
85 pub until: Antichain<mz_repr::Timestamp>,
88 pub bindings: BTreeMap<Id, CollectionBundle<'scope, T>>,
90 pub(super) compute_logger: Option<crate::logging::compute::Logger>,
92 pub(super) linear_join_spec: LinearJoinSpec,
94 pub dataflow_expiration: Antichain<mz_repr::Timestamp>,
97 pub config_set: Rc<ConfigSet>,
99}
100
101impl<'scope, T: RenderTimestamp> Context<'scope, T> {
102 pub fn for_dataflow_in<Plan>(
104 dataflow: &DataflowDescription<Plan, CollectionMetadata>,
105 scope: Scope<'scope, T>,
106 compute_state: &ComputeState,
107 until: Antichain<mz_repr::Timestamp>,
108 dataflow_expiration: Antichain<mz_repr::Timestamp>,
109 ) -> Self {
110 use mz_ore::collections::CollectionExt as IteratorExt;
111 let dataflow_id = *scope.addr().into_first();
112 let as_of_frontier = dataflow
113 .as_of
114 .clone()
115 .unwrap_or_else(|| Antichain::from_elem(Timestamp::minimum()));
116
117 let export_ids = dataflow.export_ids().collect();
118
119 let compute_logger = if dataflow.is_transient() {
123 None
124 } else {
125 compute_state.compute_logger.clone()
126 };
127
128 Self {
129 scope,
130 debug_name: dataflow.debug_name.clone(),
131 dataflow_id,
132 export_ids,
133 as_of_frontier,
134 until,
135 bindings: BTreeMap::new(),
136 compute_logger,
137 linear_join_spec: compute_state.linear_join_spec,
138 dataflow_expiration,
139 config_set: Rc::clone(&compute_state.worker_config),
140 }
141 }
142}
143
144impl<'scope, T: RenderTimestamp> Context<'scope, T> {
145 pub fn insert_id(
150 &mut self,
151 id: Id,
152 collection: CollectionBundle<'scope, T>,
153 ) -> Option<CollectionBundle<'scope, T>> {
154 self.bindings.insert(id, collection)
155 }
156 pub fn remove_id(&mut self, id: Id) -> Option<CollectionBundle<'scope, T>> {
160 self.bindings.remove(&id)
161 }
162 pub fn update_id(&mut self, id: Id, collection: CollectionBundle<'scope, T>) {
164 if !self.bindings.contains_key(&id) {
165 self.bindings.insert(id, collection);
166 } else {
167 let binding = self
168 .bindings
169 .get_mut(&id)
170 .expect("Binding verified to exist");
171 if collection.collection.is_some() {
172 binding.collection = collection.collection;
173 }
174 for (key, flavor) in collection.arranged.into_iter() {
175 binding.arranged.insert(key, flavor);
176 }
177 }
178 }
179 pub fn lookup_id(&self, id: Id) -> Option<CollectionBundle<'scope, T>> {
181 self.bindings.get(&id).cloned()
182 }
183
184 pub(super) fn error_logger(&self) -> ErrorLogger {
185 ErrorLogger::new(self.debug_name.clone())
186 }
187}
188
189impl<'scope, T: RenderTimestamp> Context<'scope, T> {
190 pub fn enter_region<'a>(
192 &self,
193 region: Scope<'a, T>,
194 bindings: Option<&std::collections::BTreeSet<Id>>,
195 ) -> Context<'a, T> {
196 let bindings = self
197 .bindings
198 .iter()
199 .filter(|(key, _)| bindings.as_ref().map(|b| b.contains(key)).unwrap_or(true))
200 .map(|(key, bundle)| (*key, bundle.enter_region(region)))
201 .collect();
202
203 Context {
204 scope: region,
205 debug_name: self.debug_name.clone(),
206 dataflow_id: self.dataflow_id.clone(),
207 export_ids: self.export_ids.clone(),
208 as_of_frontier: self.as_of_frontier.clone(),
209 until: self.until.clone(),
210 compute_logger: self.compute_logger.clone(),
211 linear_join_spec: self.linear_join_spec.clone(),
212 bindings,
213 dataflow_expiration: self.dataflow_expiration.clone(),
214 config_set: Rc::clone(&self.config_set),
215 }
216 }
217}
218
219#[derive(Clone)]
221pub enum ArrangementFlavor<'scope, T: RenderTimestamp> {
222 Local(
224 Arranged<'scope, RowRowAgent<T, Diff>>,
225 Arranged<'scope, ErrAgent<T, Diff>>,
226 ),
227 Trace(
232 GlobalId,
233 Arranged<'scope, RowRowEnter<mz_repr::Timestamp, Diff, T>>,
234 Arranged<'scope, ErrEnter<mz_repr::Timestamp, T>>,
235 ),
236}
237
238impl<'scope, T: RenderTimestamp> ArrangementFlavor<'scope, T> {
239 pub fn flat_map<DCB, L>(
272 &self,
273 key: Option<&Row>,
274 max_demand: usize,
275 logic: L,
276 ) -> (
277 Stream<'scope, T, DCB::Container>,
278 VecCollection<'scope, T, DataflowErrorSer, Diff>,
279 )
280 where
281 DCB: ContainerBuilder,
282 L: for<'a, 'b> FnMut(
283 &'a mut DatumVecBorrow<'b>,
284 T,
285 Diff,
286 &mut Session<T, DCB>,
287 &mut Session<T, ECB<T>>,
288 ) -> usize
289 + 'static,
290 {
291 match &self {
294 ArrangementFlavor::Local(oks, errs) => {
295 let (oks, mfp_errs) = CollectionBundle::<T>::flat_map_core_fallible::<_, DCB, _>(
296 oks.clone(),
297 key,
298 max_demand,
299 logic,
300 REFUEL,
301 );
302 let errs = errs.clone().as_collection(|k, &()| k.clone());
303 let errs = errs.concat(mfp_errs.as_collection());
304 (oks, errs)
305 }
306 ArrangementFlavor::Trace(_, oks, errs) => {
307 let (oks, mfp_errs) = CollectionBundle::<T>::flat_map_core_fallible::<_, DCB, _>(
308 oks.clone(),
309 key,
310 max_demand,
311 logic,
312 REFUEL,
313 );
314 let errs = errs.clone().as_collection(|k, &()| k.clone());
315 let errs = errs.concat(mfp_errs.as_collection());
316 (oks, errs)
317 }
318 }
319 }
320
321 pub fn flat_map_ok<DCB, L>(
326 &self,
327 key: Option<&Row>,
328 max_demand: usize,
329 logic: L,
330 ) -> (
331 Stream<'scope, T, DCB::Container>,
332 VecCollection<'scope, T, DataflowErrorSer, Diff>,
333 )
334 where
335 DCB: ContainerBuilder,
338 L: for<'a, 'b> FnMut(&'a mut DatumVecBorrow<'b>, T, Diff, &mut Session<T, DCB>) -> usize
339 + 'static,
340 {
341 match &self {
342 ArrangementFlavor::Local(oks, errs) => {
343 let oks = CollectionBundle::<T>::flat_map_core_ok::<_, DCB, _>(
344 oks.clone(),
345 key,
346 max_demand,
347 logic,
348 REFUEL,
349 );
350 let errs = errs.clone().as_collection(|k, &()| k.clone());
351 (oks, errs)
352 }
353 ArrangementFlavor::Trace(_, oks, errs) => {
354 let oks = CollectionBundle::<T>::flat_map_core_ok::<_, DCB, _>(
355 oks.clone(),
356 key,
357 max_demand,
358 logic,
359 REFUEL,
360 );
361 let errs = errs.clone().as_collection(|k, &()| k.clone());
362 (oks, errs)
363 }
364 }
365 }
366}
367impl<'scope, T: RenderTimestamp> ArrangementFlavor<'scope, T> {
368 pub fn scope(&self) -> Scope<'scope, T> {
370 match self {
371 ArrangementFlavor::Local(oks, _errs) => oks.stream.scope(),
372 ArrangementFlavor::Trace(_gid, oks, _errs) => oks.stream.scope(),
373 }
374 }
375
376 pub fn enter_region<'a>(&self, region: Scope<'a, T>) -> ArrangementFlavor<'a, T> {
378 match self {
379 ArrangementFlavor::Local(oks, errs) => ArrangementFlavor::Local(
380 oks.clone().enter_region(region),
381 errs.clone().enter_region(region),
382 ),
383 ArrangementFlavor::Trace(gid, oks, errs) => ArrangementFlavor::Trace(
384 *gid,
385 oks.clone().enter_region(region),
386 errs.clone().enter_region(region),
387 ),
388 }
389 }
390}
391impl<'scope, T: RenderTimestamp> ArrangementFlavor<'scope, T> {
392 pub fn leave_region<'outer>(&self, outer: Scope<'outer, T>) -> ArrangementFlavor<'outer, T> {
394 match self {
395 ArrangementFlavor::Local(oks, errs) => ArrangementFlavor::Local(
396 oks.clone().leave_region(outer),
397 errs.clone().leave_region(outer),
398 ),
399 ArrangementFlavor::Trace(gid, oks, errs) => ArrangementFlavor::Trace(
400 *gid,
401 oks.clone().leave_region(outer),
402 errs.clone().leave_region(outer),
403 ),
404 }
405 }
406}
407pub(crate) fn distinct_arranged_errs<'a, T: RenderTimestamp>(
419 errs: Arranged<'a, ErrAgent<T, Diff>>,
420 name: &str,
421) -> Arranged<'a, ErrAgent<T, Diff>> {
422 errs.mz_reduce_abelian::<_, ErrBuilder<_, _>, ErrSpine<_, _>, _>(
423 name,
424 |_err, _input, output| output.push(((), Diff::ONE)),
425 )
426}
427
428pub(crate) fn distinct_errs_collection<'a, T: RenderTimestamp>(
433 errs: VecCollection<'a, T, DataflowErrorSer, Diff>,
434) -> VecCollection<'a, T, DataflowErrorSer, Diff> {
435 let errs: KeyCollection<_, _, _> = errs.into();
436 let errs = errs
437 .mz_arrange::<ColumnationChunker<_>, ErrBatcher<_, _>, ErrBuilder<_, _>, ErrSpine<_, _>>(
438 "Arrange errors",
439 );
440 distinct_arranged_errs(errs, "Distinct errors").as_collection(|err, _| err.clone())
441}
442
443#[derive(Clone)]
448pub struct CollectionBundle<'scope, T: RenderTimestamp> {
449 pub collection: Option<(
450 CollectionEdge<'scope, T>,
451 VecCollection<'scope, T, DataflowErrorSer, Diff>,
452 )>,
453 pub arranged: BTreeMap<Vec<LirScalarExpr>, ArrangementFlavor<'scope, T>>,
454}
455
456impl<'scope, T: RenderTimestamp> CollectionBundle<'scope, T> {
457 pub fn from_edge(
459 oks: CollectionEdge<'scope, T>,
460 errs: VecCollection<'scope, T, DataflowErrorSer, Diff>,
461 ) -> Self {
462 Self {
463 collection: Some((oks, errs)),
464 arranged: BTreeMap::default(),
465 }
466 }
467
468 pub fn from_expressions(
470 exprs: Vec<LirScalarExpr>,
471 arrangements: ArrangementFlavor<'scope, T>,
472 ) -> Self {
473 let mut arranged = BTreeMap::new();
474 arranged.insert(exprs, arrangements);
475 Self {
476 collection: None,
477 arranged,
478 }
479 }
480
481 pub fn from_columns<I: IntoIterator<Item = usize>>(
483 columns: I,
484 arrangements: ArrangementFlavor<'scope, T>,
485 ) -> Self {
486 let mut keys = Vec::new();
487 for column in columns {
488 keys.push(LirScalarExpr::column(column));
489 }
490 Self::from_expressions(keys, arrangements)
491 }
492
493 pub fn scope(&self) -> Scope<'scope, T> {
495 if let Some((oks, _errs)) = &self.collection {
496 oks.scope()
497 } else {
498 self.arranged
499 .values()
500 .next()
501 .expect("Must contain a valid collection")
502 .scope()
503 }
504 }
505
506 pub fn distinct_errs(mut self) -> Self {
532 if let Some((oks, errs)) = self.collection.take() {
533 self.collection = Some((oks, distinct_errs_collection(errs)));
534 }
535 for (key, flavor) in std::mem::take(&mut self.arranged) {
536 let flavor = match flavor {
537 ArrangementFlavor::Local(oks, errs) => {
538 let name = format!("Distinct errors[{key:?}]");
541 ArrangementFlavor::Local(oks, distinct_arranged_errs(errs, &name))
542 }
543 flavor @ ArrangementFlavor::Trace(..) => flavor,
544 };
545 self.arranged.insert(key, flavor);
546 }
547 self
548 }
549
550 pub fn enter_region<'inner>(&self, region: Scope<'inner, T>) -> CollectionBundle<'inner, T> {
552 CollectionBundle {
553 collection: self.collection.as_ref().map(|(oks, errs)| {
554 (
555 oks.clone().enter_region(region),
556 errs.clone().enter_region(region),
557 )
558 }),
559 arranged: self
560 .arranged
561 .iter()
562 .map(|(key, bundle)| (key.clone(), bundle.enter_region(region)))
563 .collect(),
564 }
565 }
566}
567
568impl<'scope, T: RenderTimestamp> CollectionBundle<'scope, T> {
569 pub fn leave_region<'outer>(&self, outer: Scope<'outer, T>) -> CollectionBundle<'outer, T> {
571 CollectionBundle {
572 collection: self.collection.as_ref().map(|(oks, errs)| {
573 (
574 oks.clone().leave_region(outer),
575 errs.clone().leave_region(outer),
576 )
577 }),
578 arranged: self
579 .arranged
580 .iter()
581 .map(|(key, bundle)| (key.clone(), bundle.leave_region(outer)))
582 .collect(),
583 }
584 }
585}
586
587impl<'scope, T: RenderTimestamp> CollectionBundle<'scope, T> {
588 pub fn as_specific_collection(
603 &self,
604 key: Option<&[LirScalarExpr]>,
605 ) -> (
606 CollectionEdge<'scope, T>,
607 VecCollection<'scope, T, DataflowErrorSer, Diff>,
608 ) {
609 match key {
615 None => self
616 .collection
617 .clone()
618 .expect("The unarranged collection doesn't exist."),
619 Some(key) => {
620 let arranged = self.arranged.get(key).unwrap_or_else(|| {
621 panic!("The collection arranged by {:?} doesn't exist.", key)
622 });
623 let (ok, err) =
627 arranged.flat_map_ok::<ColumnBuilder<(Row, T, Diff)>, _>(None, usize::MAX, {
628 let mut row_buf = Row::default();
631 move |borrow, t, r, ok_session| {
632 row_buf.packer().extend(borrow.iter());
633 ok_session.give((&row_buf, &t, &r));
634 1
635 }
636 });
637 (ok.as_collection(), err)
638 }
639 }
640 }
641
642 pub fn flat_map<DCB, L>(
658 &self,
659 key_val: Option<(Vec<LirScalarExpr>, Option<Row>)>,
660 max_demand: usize,
661 logic: L,
662 ) -> (
663 Stream<'scope, T, DCB::Container>,
664 VecCollection<'scope, T, DataflowErrorSer, Diff>,
665 )
666 where
667 DCB: ContainerBuilder,
668 L: for<'a> FnMut(
669 &'a mut DatumVecBorrow<'_>,
670 T,
671 Diff,
672 &mut Session<T, DCB>,
673 &mut Session<T, ECB<T>>,
674 ) -> usize
675 + 'static,
676 {
677 if let Some((key, val)) = key_val {
681 self.arrangement(&key)
682 .expect("Should have ensured during planning that this arrangement exists.")
683 .flat_map::<DCB, _>(val.as_ref(), max_demand, logic)
684 } else {
685 let (oks, errs) = self
686 .collection
687 .clone()
688 .expect("Invariant violated: CollectionBundle contains no collection.");
689 let (ok_stream, err_stream) = flat_map_datums::<_, DCB, _>(oks, 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 CollectionEdge<'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: CollectionEdge<'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 CollectionEdge<'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 = OutputBuilder::<
1195 _,
1196 CapacityContainerBuilder<Column<(Row, T, Diff)>>,
1197 >::from(passthrough_output);
1198 let mut input = builder.new_input(oks.inner, Pipeline);
1199 builder.set_notify_for(0, FrontierInterest::Never);
1200 builder.build(move |_capabilities| {
1201 let mut key_buf = Row::default();
1202 let mut val_buf = Row::default();
1203 let mut datums = DatumVec::new();
1204 move |_frontiers| {
1205 let mut temp_storage = RowArena::new();
1209 let mut ok_output = ok_output.activate();
1210 let mut err_output = err_output.activate();
1211 let mut passthrough_output = passthrough_output.activate();
1212 input.for_each(|time, data| {
1213 let mut ok_session = ok_output.session_with_builder(&time);
1214 let mut err_session = err_output.session(&time);
1215 for (row, t, d) in data.borrow().into_index_iter() {
1218 temp_storage.clear();
1219 let datums = datums.borrow_with(row);
1220 let key_iter = key.iter().map(|k| k.eval(&datums, &temp_storage));
1221 match key_buf.packer().try_extend(key_iter) {
1222 Ok(()) => {
1223 let val_datum_iter = thinning.iter().map(|c| datums[*c]);
1224 val_buf.packer().extend(val_datum_iter);
1225 ok_session.give(((&*key_buf, &*val_buf), t, d));
1226 }
1227 Err(e) => {
1228 err_session.give((
1229 e.into(),
1230 Columnar::into_owned(t),
1231 Columnar::into_owned(d),
1232 ));
1233 }
1234 }
1235 }
1236 passthrough_output.session(&time).give_container(data);
1237 });
1238 }
1239 });
1240 (ok_stream, err_stream, passthrough_stream.as_collection())
1241 };
1242
1243 let exchange =
1244 ExchangeCore::<ColumnBuilder<_>, _>::new_core(columnar_exchange::<Row, Row, T, Diff>);
1245 let oks = match batcher {
1246 ArrangementBatcher::ColumnarPaged => ok_stream.mz_arrange_core::<
1247 _,
1248 batcher::ColumnChunker<_>,
1249 Col2ValPagedBatcher<_, _, _, _>,
1250 RowRowColPagedBuilder<_, _>,
1251 RowRowSpine<_, _>,
1252 >(exchange, name),
1253 ArrangementBatcher::Columnar => ok_stream.mz_arrange_core::<
1254 _,
1255 batcher::ColumnChunker<_>,
1256 Col2ValColBatcher<_, _, _, _>,
1257 RowRowColPagedBuilder<_, _>,
1258 RowRowSpine<_, _>,
1259 >(exchange, name),
1260 ArrangementBatcher::Columnation => ok_stream.mz_arrange_core::<
1261 _,
1262 batcher::Chunker<_>,
1263 Col2ValBatcher<_, _, _, _>,
1264 RowRowBuilder<_, _>,
1265 RowRowSpine<_, _>,
1266 >(exchange, name),
1267 };
1268 (oks, err_stream.as_collection(), passthrough)
1269 }
1270}
1271
1272pub(crate) type Session<'a, 'b, T, CB> =
1276 timely::dataflow::operators::generic::Session<'a, 'b, T, CB, Capability<T>>;
1277
1278pub(crate) type ECB<T> = ConsolidatingContainerBuilder<Vec<(DataflowErrorSer, T, Diff)>>;
1283
1284const REFUEL: usize = 1_000_000;
1288
1289struct PendingWork<C>
1290where
1291 C: Cursor,
1292{
1293 ok_capability: Capability<C::Time>,
1295 err_capability: Capability<C::Time>,
1297 cursor: C,
1298 batch: C::Storage,
1299}
1300
1301impl<C> PendingWork<C>
1302where
1303 C: Cursor<KeyContainer: BatchContainer<Owned: PartialEq + Sized>>,
1304{
1305 fn new(
1308 ok_capability: Capability<C::Time>,
1309 err_capability: Capability<C::Time>,
1310 cursor: C,
1311 batch: C::Storage,
1312 ) -> Self {
1313 Self {
1314 ok_capability,
1315 err_capability,
1316 cursor,
1317 batch,
1318 }
1319 }
1320 fn do_work<DCB, L>(
1323 &mut self,
1324 key: Option<&C::Key<'_>>,
1325 logic: &mut L,
1326 fuel: &mut usize,
1327 ok_output: &mut OutputBuilderSession<'_, C::Time, DCB>,
1328 err_output: &mut OutputBuilderSession<'_, C::Time, ECB<C::Time>>,
1329 ) where
1330 DCB: ContainerBuilder,
1331 L: FnMut(
1332 C::Key<'_>,
1333 C::Val<'_>,
1334 C::Time,
1335 C::Diff,
1336 &mut Session<C::Time, DCB>,
1337 &mut Session<C::Time, ECB<C::Time>>,
1338 ) -> usize,
1339 {
1340 let mut ok_session = ok_output.session_with_builder(&self.ok_capability);
1341 let mut err_session = err_output.session_with_builder(&self.err_capability);
1342 walk_cursor(&mut self.cursor, &self.batch, key, fuel, |k, v, t, d| {
1343 logic(k, v, t, d, &mut ok_session, &mut err_session)
1344 });
1345 }
1346}
1347
1348struct PendingWorkOk<C>
1351where
1352 C: Cursor,
1353{
1354 capability: Capability<C::Time>,
1355 cursor: C,
1356 batch: C::Storage,
1357}
1358
1359impl<C> PendingWorkOk<C>
1360where
1361 C: Cursor<KeyContainer: BatchContainer<Owned: PartialEq + Sized>>,
1362{
1363 fn new(capability: Capability<C::Time>, cursor: C, batch: C::Storage) -> Self {
1364 Self {
1365 capability,
1366 cursor,
1367 batch,
1368 }
1369 }
1370
1371 fn do_work<DCB, L>(
1374 &mut self,
1375 key: Option<&C::Key<'_>>,
1376 logic: &mut L,
1377 fuel: &mut usize,
1378 ok_output: &mut OutputBuilderSession<'_, C::Time, DCB>,
1379 ) where
1380 DCB: ContainerBuilder,
1381 L: FnMut(C::Key<'_>, C::Val<'_>, C::Time, C::Diff, &mut Session<C::Time, DCB>) -> usize,
1382 {
1383 let mut ok_session = ok_output.session_with_builder(&self.capability);
1384 walk_cursor(&mut self.cursor, &self.batch, key, fuel, |k, v, t, d| {
1385 logic(k, v, t, d, &mut ok_session)
1386 });
1387 }
1388}
1389
1390fn walk_cursor<C, F>(
1400 cursor: &mut C,
1401 batch: &C::Storage,
1402 key: Option<&C::Key<'_>>,
1403 fuel: &mut usize,
1404 mut emit: F,
1405) where
1406 C: Cursor<KeyContainer: BatchContainer<Owned: PartialEq + Sized>>,
1407 F: FnMut(C::Key<'_>, C::Val<'_>, C::Time, C::Diff) -> usize,
1408{
1409 use differential_dataflow::consolidation::consolidate;
1410
1411 let mut work: usize = 0;
1412 let mut buffer = Vec::new();
1413 if let Some(key) = key {
1414 let key = C::KeyContainer::reborrow(*key);
1415 if cursor.get_key(batch).map(|k| k == key) != Some(true) {
1416 cursor.seek_key(batch, key);
1417 }
1418 if cursor.get_key(batch).map(|k| k == key) == Some(true) {
1419 let key = cursor.key(batch);
1420 while let Some(val) = cursor.get_val(batch) {
1421 cursor.map_times(batch, |time, diff| {
1422 buffer.push((C::owned_time(time), C::owned_diff(diff)));
1423 });
1424 consolidate(&mut buffer);
1425 for (time, diff) in buffer.drain(..) {
1426 work += emit(key, val, time, diff);
1427 }
1428 cursor.step_val(batch);
1429 if work >= *fuel {
1430 *fuel = 0;
1431 return;
1432 }
1433 }
1434 }
1435 } else {
1436 while let Some(key) = cursor.get_key(batch) {
1437 while let Some(val) = cursor.get_val(batch) {
1438 cursor.map_times(batch, |time, diff| {
1439 buffer.push((C::owned_time(time), C::owned_diff(diff)));
1440 });
1441 consolidate(&mut buffer);
1442 for (time, diff) in buffer.drain(..) {
1443 work += emit(key, val, time, diff);
1444 }
1445 cursor.step_val(batch);
1446 if work >= *fuel {
1447 *fuel = 0;
1448 return;
1449 }
1450 }
1451 cursor.step_key(batch);
1452 }
1453 }
1454 *fuel -= work;
1455}
1456
1457#[cfg(test)]
1458mod tests {
1459 use differential_dataflow::input::Input;
1460 use mz_expr::{EvalError, MapFilterProject};
1461 use mz_repr::{Datum, ReprScalarType, Timestamp};
1462 use timely::dataflow::operators::Capture;
1463 use timely::dataflow::operators::capture::{Event, Extract};
1464
1465 use super::*;
1466 use crate::render::columnar::{columnar_to_vec, vec_to_columnar};
1467
1468 type OkUpdate = ((Row, Row), Timestamp, Diff);
1469 type ErrUpdate = (DataflowErrorSer, Timestamp, Diff);
1470 type Captured<D> = std::sync::mpsc::Receiver<Event<Timestamp, Vec<D>>>;
1471
1472 fn extract_ok(captured: Captured<OkUpdate>) -> Vec<OkUpdate> {
1473 let mut updates: Vec<_> = captured
1474 .extract()
1475 .into_iter()
1476 .flat_map(|(_, data)| data)
1477 .collect();
1478 updates.sort();
1479 updates
1480 }
1481
1482 fn extract_err(captured: Captured<ErrUpdate>) -> Vec<(String, Timestamp, Diff)> {
1484 let mut updates: Vec<_> = captured
1485 .extract()
1486 .into_iter()
1487 .flat_map(|(_, data)| data)
1488 .map(|(e, t, d)| (format!("{e:?}"), t, d))
1489 .collect();
1490 updates.sort();
1491 updates
1492 }
1493
1494 fn arrange_columnar(
1496 rows: Vec<(Row, u64)>,
1497 key: Vec<LirScalarExpr>,
1498 ) -> (Vec<OkUpdate>, Vec<(String, Timestamp, Diff)>) {
1499 let thinning = vec![0, 1];
1500 let (ok, err) = timely::execute_directly(move |worker| {
1501 worker.dataflow::<Timestamp, _, _>(|scope| {
1502 let (mut input, collection) = scope.new_collection();
1503 let (arranged, errs, _passthrough) =
1504 CollectionBundle::<Timestamp>::arrange_collection(
1505 &"col".to_string(),
1506 vec_to_columnar(collection),
1507 key,
1508 thinning,
1509 ArrangementBatcher::Columnation,
1510 );
1511 let ok = arranged
1512 .as_collection(|k, v| (k.to_row(), v.to_row()))
1513 .inner
1514 .capture();
1515 let err = errs.inner.capture();
1516
1517 let max_time = rows.iter().map(|(_, t)| *t).max().unwrap_or(0);
1518 for (row, time) in rows {
1519 input.update_at(row, Timestamp::from(time), Diff::ONE);
1520 }
1521 input.advance_to(Timestamp::from(max_time + 1));
1522 input.flush();
1523 (ok, err)
1524 })
1525 });
1526 (extract_ok(ok), extract_err(err))
1527 }
1528
1529 fn test_rows() -> Vec<(Row, u64)> {
1532 vec![
1533 (Row::pack_slice(&[Datum::Int32(1), Datum::String("a")]), 0),
1534 (Row::pack_slice(&[Datum::Int32(2), Datum::String("b")]), 1),
1535 (Row::pack_slice(&[Datum::Int32(1), Datum::String("a")]), 1),
1536 (Row::pack_slice(&[Datum::Int32(3), Datum::Null]), 2),
1537 ]
1538 }
1539
1540 #[mz_ore::test]
1543 fn arrange_collection_keys_correctly() {
1544 let rows = test_rows();
1545 let mut expected: Vec<OkUpdate> = rows
1546 .iter()
1547 .map(|(row, t)| {
1548 let key = Row::pack_slice(&[row.iter().next().unwrap()]);
1549 ((key, row.clone()), Timestamp::from(*t), Diff::ONE)
1550 })
1551 .collect();
1552 expected.sort();
1553
1554 let (ok, err) = arrange_columnar(rows, vec![LirScalarExpr::column(0)]);
1555 assert_eq!(ok, expected);
1556 assert!(err.is_empty());
1557 }
1558
1559 #[mz_ore::test]
1561 fn arrange_collection_error_path() {
1562 let key = vec![LirScalarExpr::literal(
1563 Err(EvalError::DivisionByZero),
1564 ReprScalarType::Int32,
1565 )];
1566 let (ok, err) = arrange_columnar(test_rows(), key);
1567
1568 assert!(ok.is_empty());
1569 assert!(!err.is_empty());
1570 }
1571
1572 fn extract_row_updates(
1573 captured: Captured<(Row, Timestamp, Diff)>,
1574 ) -> Vec<(Row, Timestamp, Diff)> {
1575 let mut updates: Vec<_> = captured
1576 .extract()
1577 .into_iter()
1578 .flat_map(|(_, data)| data)
1579 .collect();
1580 updates.sort();
1581 updates
1582 }
1583
1584 #[mz_ore::test]
1587 fn get_arrange_by_produces_projected_rows() {
1588 let rows = vec![
1589 (Row::pack_slice(&[Datum::Int64(1), Datum::Int64(10)]), 0u64),
1590 (Row::pack_slice(&[Datum::Int64(2), Datum::Int64(20)]), 1),
1591 (Row::pack_slice(&[Datum::Int64(1), Datum::Int64(10)]), 1),
1592 ];
1593 let mfp = MapFilterProject::<LirScalarExpr>::new(2)
1596 .project(vec![0])
1597 .into_plan()
1598 .expect("mfp");
1599 let mut expected: Vec<(Row, Timestamp, Diff)> = rows
1600 .iter()
1601 .map(|(r, t)| {
1602 let col0 = r.iter().next().unwrap();
1603 (Row::pack_slice(&[col0]), Timestamp::from(*t), Diff::ONE)
1604 })
1605 .collect();
1606 expected.sort();
1607
1608 let produced = timely::execute_directly(move |worker| {
1609 worker.dataflow::<Timestamp, _, _>(|scope| {
1610 let (mut input, collection) = scope.new_collection();
1611 let (_err_input, errs) = scope.new_collection::<DataflowErrorSer, Diff>();
1612 let bundle = CollectionBundle::from_edge(vec_to_columnar(collection), errs);
1613 let (edge, _errs) = bundle.as_collection_core(mfp, None, Antichain::new());
1614 let produced = columnar_to_vec(edge.clone()).inner.capture();
1615 let (_arranged, _arrange_errs, _passthrough) =
1616 CollectionBundle::<Timestamp>::arrange_collection(
1617 &"arrange".to_string(),
1618 edge,
1619 vec![LirScalarExpr::column(0)],
1620 vec![0],
1621 ArrangementBatcher::Columnation,
1622 );
1623
1624 let max_time = rows.iter().map(|(_, t)| *t).max().unwrap();
1625 for (row, time) in rows {
1626 input.update_at(row, Timestamp::from(time), Diff::ONE);
1627 }
1628 input.advance_to(Timestamp::from(max_time + 1));
1629 input.flush();
1630 produced
1631 })
1632 });
1633
1634 assert_eq!(extract_row_updates(produced), expected);
1635 }
1636
1637 #[mz_ore::test]
1640 fn as_collection_core_identity_passes_edge_through() {
1641 let expected = vec![(
1642 Row::pack_slice(&[Datum::Int64(1)]),
1643 Timestamp::from(0u64),
1644 Diff::ONE,
1645 )];
1646 let captured = timely::execute_directly(move |worker| {
1647 worker.dataflow::<Timestamp, _, _>(|scope| {
1648 let (mut input, collection) = scope.new_collection::<Row, Diff>();
1649 let (_err_input, errs) = scope.new_collection::<DataflowErrorSer, Diff>();
1650 let bundle = CollectionBundle::from_edge(vec_to_columnar(collection), errs);
1651 let identity = MapFilterProject::<LirScalarExpr>::new(1)
1652 .into_plan()
1653 .expect("identity mfp");
1654 let (out, _errs) = bundle.as_collection_core(identity, None, Antichain::new());
1655 let captured = columnar_to_vec(out).inner.capture();
1656 input.update_at(
1657 Row::pack_slice(&[Datum::Int64(1)]),
1658 Timestamp::from(0u64),
1659 Diff::ONE,
1660 );
1661 input.advance_to(Timestamp::from(1u64));
1662 input.flush();
1663 captured
1664 })
1665 });
1666 assert_eq!(extract_row_updates(captured), expected);
1667 }
1668
1669 #[mz_ore::test]
1672 fn as_collection_core_consolidates_within_batch() {
1673 let rows = vec![
1677 Row::pack_slice(&[Datum::Int64(1), Datum::Int64(10)]),
1678 Row::pack_slice(&[Datum::Int64(1), Datum::Int64(20)]),
1679 Row::pack_slice(&[Datum::Int64(2), Datum::Int64(30)]),
1680 ];
1681 let mfp = MapFilterProject::<LirScalarExpr>::new(2)
1682 .project(vec![0])
1683 .into_plan()
1684 .expect("mfp");
1685 let expected = vec![
1686 (
1687 Row::pack_slice(&[Datum::Int64(1)]),
1688 Timestamp::from(0u64),
1689 Diff::from(2),
1690 ),
1691 (
1692 Row::pack_slice(&[Datum::Int64(2)]),
1693 Timestamp::from(0u64),
1694 Diff::ONE,
1695 ),
1696 ];
1697
1698 let captured = timely::execute_directly(move |worker| {
1699 worker.dataflow::<Timestamp, _, _>(|scope| {
1700 let (mut input, collection) = scope.new_collection();
1701 let (_err_input, errs) = scope.new_collection::<DataflowErrorSer, Diff>();
1702 let bundle = CollectionBundle::from_edge(vec_to_columnar(collection), errs);
1703 let (edge, _errs) = bundle.as_collection_core(mfp, None, Antichain::new());
1704 let captured = columnar_to_vec(edge).inner.capture();
1705 for row in rows {
1708 input.update_at(row, Timestamp::from(0u64), Diff::ONE);
1709 }
1710 input.advance_to(Timestamp::from(1u64));
1711 input.flush();
1712 captured
1713 })
1714 });
1715
1716 assert_eq!(extract_row_updates(captured), expected);
1717 }
1718
1719 #[mz_ore::test]
1722 fn as_specific_collection_materializes_columnar() {
1723 let rows = test_rows();
1724 let key = vec![LirScalarExpr::column(0)];
1725 let mut expected: Vec<(Row, Timestamp, Diff)> = rows
1726 .iter()
1727 .map(|(r, t)| (r.clone(), Timestamp::from(*t), Diff::ONE))
1728 .collect();
1729 expected.sort();
1730
1731 let captured = timely::execute_directly(move |worker| {
1732 worker.dataflow::<Timestamp, _, _>(|scope| {
1733 let (mut input, collection) = scope.new_collection();
1734 let (arranged, arr_errs, _passthrough) =
1735 CollectionBundle::<Timestamp>::arrange_collection(
1736 &"agg".to_string(),
1737 vec_to_columnar(collection),
1738 key.clone(),
1739 vec![1],
1740 ArrangementBatcher::Columnation,
1741 );
1742 let err_arranged = {
1743 let kc: KeyCollection<_, _, _> = arr_errs.into();
1744 kc.mz_arrange::<
1745 ColumnationChunker<_>,
1746 ErrBatcher<_, _>,
1747 ErrBuilder<_, _>,
1748 ErrSpine<_, _>,
1749 >("agg-errs")
1750 };
1751 let bundle = CollectionBundle::from_columns(
1753 0..1,
1754 ArrangementFlavor::Local(arranged, err_arranged),
1755 );
1756 let (edge, _errs) = bundle.as_specific_collection(Some(&key));
1757 let captured = columnar_to_vec(edge).inner.capture();
1758
1759 let max_time = rows.iter().map(|(_, t)| *t).max().unwrap();
1760 for (row, time) in rows {
1761 input.update_at(row, Timestamp::from(time), Diff::ONE);
1762 }
1763 input.advance_to(Timestamp::from(max_time + 1));
1764 input.flush();
1765 captured
1766 })
1767 });
1768
1769 assert_eq!(extract_row_updates(captured), expected);
1770 }
1771}