1use std::any::Any;
104use std::cell::RefCell;
105use std::collections::{BTreeMap, BTreeSet};
106use std::convert::Infallible;
107use std::future::Future;
108use std::pin::Pin;
109use std::rc::{Rc, Weak};
110use std::sync::Arc;
111use std::task::Poll;
112
113use ::columnar::{Columnar as ColumnarData, Index as ColumnarIndex, Push as ColumnarPush};
114use differential_dataflow::dynamic::pointstamp::PointStamp;
115use differential_dataflow::lattice::Lattice;
116use differential_dataflow::operators::arrange::Arranged;
117use differential_dataflow::operators::arrange::ShutdownButton;
118use differential_dataflow::operators::iterate::Variable;
119use differential_dataflow::trace::cursor::{BatchCursor, BatchDiff, BatchKey, BatchVal};
120use differential_dataflow::trace::{BatchReader, Cursor, Navigable, TraceReader};
121use differential_dataflow::{AsCollection, Collection, Data, VecCollection};
122use futures::FutureExt;
123use futures::channel::oneshot;
124use itertools::Itertools;
125use mz_compute_types::dataflows::{DataflowDescription, IndexDesc};
126use mz_compute_types::dyncfgs::{
127 COMPUTE_APPLY_COLUMN_DEMANDS, COMPUTE_LOGICAL_BACKPRESSURE_INFLIGHT_SLACK,
128 COMPUTE_LOGICAL_BACKPRESSURE_MAX_RETAINED_CAPABILITIES, ENABLE_COMPUTE_LOGICAL_BACKPRESSURE,
129 ENABLE_COMPUTE_TEMPORAL_BUCKETING, ENABLE_ERROR_DISTINCT, SUBSCRIBE_SNAPSHOT_OPTIMIZATION,
130 TEMPORAL_BUCKETING_SUMMARY,
131};
132use mz_compute_types::plan::render_plan::{
133 self, BindStage, LetBind, LetFreePlan, RecBind, RenderPlan,
134};
135use mz_compute_types::plan::scalar::LirScalarExpr;
136use mz_compute_types::plan::{ArrangementStrategy, LirId};
137use mz_expr::{EvalError, Id, LocalId, permutation_for_arrangement};
138use mz_persist_client::operators::shard_source::{ErrorHandler, SnapshotMode};
139use mz_repr::explain::DummyHumanizer;
140use mz_repr::fixed_length::ExtendDatums;
141use mz_repr::{Datum, DatumVec, Diff, GlobalId, ReprRelationType, Row, RowArena, SharedRow};
142use mz_storage_operators::persist_source;
143use mz_storage_types::controller::CollectionMetadata;
144use mz_timely_util::columnar::Column;
145use mz_timely_util::columnar::builder::ColumnBuilder;
146use mz_timely_util::columnation::ColumnationChunker;
147use mz_timely_util::operator::StreamExt;
148use mz_timely_util::probe::{Handle as MzProbeHandle, ProbeNotify};
149use mz_timely_util::scope_label::ScopeExt;
150use timely::PartialOrder;
151use timely::container::CapacityContainerBuilder;
152use timely::dataflow::channels::pact::Pipeline;
153use timely::dataflow::operators::core::to_stream::ToStreamBuilder;
154use timely::dataflow::operators::vec::Filter;
155use timely::dataflow::operators::vec::ToStream;
156use timely::dataflow::operators::{Capability, Operator, Probe, probe};
157use timely::dataflow::{Scope, Stream, StreamVec};
158use timely::order::{Product, TotalOrder};
159use timely::progress::timestamp::Refines;
160use timely::progress::{Antichain, Timestamp};
161use timely::scheduling::ActivateOnDrop;
162use timely::worker::Worker as TimelyWorker;
163
164use crate::arrangement::manager::TraceBundle;
165use crate::compute_state::ComputeState;
166use crate::extensions::arrange::{KeyCollection, MzArrange};
167use crate::extensions::reduce::MzReduce;
168use crate::extensions::temporal_bucket::TemporalBucketing;
169use crate::logging::compute::{
170 ComputeEvent, DataflowGlobal, LirMapping, LirMetadata, LogDataflowErrors, OperatorHydration,
171};
172use crate::render::columnar::{
173 ColCollection, RecTimestamp, columnar_consolidate, columnar_leave_dynamic, columnar_negate,
174 concat_many, flat_map_datums,
175};
176use crate::render::context::{ArrangementFlavor, Context};
177use crate::render::errors::DataflowErrorSer;
178use crate::typedefs::{ErrBatcher, ErrBuilder, ErrSpine, MzTimestamp};
179use mz_row_spine::{DatumSeq, RowRowBatcher, RowRowBuilder};
180use mz_timely_util::columnar::consolidate::ConsolidatingColumnBuilder;
181
182pub(crate) mod columnar;
183pub mod context;
184pub(crate) mod errors;
185mod flat_map;
186mod join;
187mod reduce;
188pub mod sinks;
189mod threshold;
190mod top_k;
191
192pub use context::CollectionBundle;
193pub use join::LinearJoinSpec;
194
195struct PressOnDrop<T>(ShutdownButton<T>);
199
200impl<T> Drop for PressOnDrop<T> {
201 fn drop(&mut self) {
202 self.0.press();
203 }
204}
205
206pub fn build_compute_dataflow(
212 timely_worker: &mut TimelyWorker,
213 compute_state: &mut ComputeState,
214 dataflow: DataflowDescription<RenderPlan, CollectionMetadata>,
215 start_signal: StartSignal,
216 until: Antichain<mz_repr::Timestamp>,
217 dataflow_expiration: Antichain<mz_repr::Timestamp>,
218) {
219 let recursive = dataflow
221 .objects_to_build
222 .iter()
223 .any(|object| object.plan.is_recursive());
224
225 let indexes = dataflow
227 .index_exports
228 .iter()
229 .map(|(idx_id, (idx, _typ))| (*idx_id, dataflow.depends_on(idx.on_id), idx.as_lir()))
230 .collect::<Vec<_>>();
231
232 let sinks = dataflow
234 .sink_exports
235 .iter()
236 .map(|(sink_id, sink)| (*sink_id, dataflow.depends_on(sink.from), sink.clone()))
237 .collect::<Vec<_>>();
238
239 let worker_logging = timely_worker.logger_for("timely").map(Into::into);
240 let apply_demands = COMPUTE_APPLY_COLUMN_DEMANDS.get(&compute_state.worker_config);
241 let subscribe_snapshot_optimization =
242 SUBSCRIBE_SNAPSHOT_OPTIMIZATION.get(&compute_state.worker_config);
243
244 let name = format!("Dataflow: {}", dataflow.debug_name);
245 let input_name = format!("InputRegion: {}", dataflow.debug_name);
246 let build_name = format!("BuildRegion: {}", dataflow.debug_name);
247
248 timely_worker.dataflow_core(&name, worker_logging, Box::new(()), |_, scope| {
249 let scope = scope.with_label();
250
251 let mut imported_sources = Vec::new();
256 let mut tokens: BTreeMap<_, Rc<dyn Any>> = BTreeMap::new();
257 let output_probe = MzProbeHandle::default();
258
259 scope.clone().region_named(&input_name, |region| {
260 for (source_id, import) in dataflow.source_imports.iter() {
262 region.region_named(&format!("Source({:?})", source_id), |inner| {
263 let mut read_schema = None;
264 let mut mfp = import.desc.arguments.operators.clone().map(|mut ops| {
265 if apply_demands {
268 let demands = ops.demand();
269 let new_desc = import
270 .desc
271 .storage_metadata
272 .relation_desc
273 .apply_demand(&demands);
274 let new_arity = demands.len();
275 let remap: BTreeMap<_, _> = demands
276 .into_iter()
277 .enumerate()
278 .map(|(new, old)| (old, new))
279 .collect();
280 ops.permute_fn(|old_idx| remap[&old_idx], new_arity);
281 read_schema = Some(new_desc);
282 }
283
284 mz_expr::MfpPlan::create_from(ops)
285 .expect("Linear operators should always be valid")
286 });
287
288 let snapshot_mode = if import.with_snapshot || !subscribe_snapshot_optimization
289 {
290 SnapshotMode::Include
291 } else {
292 compute_state.metrics.inc_subscribe_snapshot_optimization();
293 SnapshotMode::Exclude
294 };
295 let suppress_early_progress_as_of = dataflow.as_of.clone();
296
297 let (mut ok_stream, err_stream, token) = persist_source::persist_source::<
300 DataflowErrorSer,
301 ConsolidatingColumnBuilder<Row, mz_repr::Timestamp, Diff>,
302 >(
303 inner,
304 *source_id,
305 Arc::clone(&compute_state.persist_clients),
306 &compute_state.txns_ctx,
307 import.desc.storage_metadata.clone(),
308 read_schema,
309 dataflow.as_of.clone(),
310 snapshot_mode,
311 until.clone(),
312 mfp.as_mut(),
313 compute_state.dataflow_max_inflight_bytes(),
314 start_signal.clone().into_send_future(),
315 ErrorHandler::Halt("compute_import"),
316 );
317
318 assert!(mfp.map(|x| x.is_identity()).unwrap_or(true));
321
322 if let Some(as_of) = suppress_early_progress_as_of {
326 ok_stream = suppress_early_progress(ok_stream, as_of);
327 }
328
329 if ENABLE_COMPUTE_LOGICAL_BACKPRESSURE.get(&compute_state.worker_config) {
330 let limit = COMPUTE_LOGICAL_BACKPRESSURE_MAX_RETAINED_CAPABILITIES
332 .get(&compute_state.worker_config);
333 let slack = COMPUTE_LOGICAL_BACKPRESSURE_INFLIGHT_SLACK
334 .get(&compute_state.worker_config)
335 .as_millis()
336 .try_into()
337 .expect("must fit");
338
339 let stream = ok_stream.limit_progress(
340 output_probe.clone(),
341 slack,
342 limit,
343 import.upper.clone(),
344 name.clone(),
345 );
346 ok_stream = stream;
347 }
348
349 let input_probe =
351 compute_state.input_probe_for(*source_id, dataflow.export_ids());
352 ok_stream = ok_stream.probe_with(&input_probe);
353
354 let (oks, errs) = (
355 ok_stream
356 .as_collection()
357 .leave_region(region)
358 .leave_region(scope),
359 err_stream
360 .as_collection()
361 .leave_region(region)
362 .leave_region(scope),
363 );
364
365 imported_sources.push((mz_expr::Id::Global(*source_id), (oks, errs)));
366
367 tokens.insert(*source_id, Rc::new(token));
369 });
370 }
371 });
372
373 if recursive {
376 scope.clone().iterative::<PointStamp<u64>, _, _>(|region| {
377 let mut context = Context::for_dataflow_in(
378 &dataflow,
379 region.clone(),
380 compute_state,
381 until,
382 dataflow_expiration,
383 );
384
385 for (id, (oks, errs)) in imported_sources.into_iter() {
386 let bundle = crate::render::CollectionBundle::from_edge(
387 oks.enter(region),
388 errs.enter(region),
389 );
390 context.insert_id(id, bundle);
392 }
393
394 for (idx_id, idx) in &dataflow.index_imports {
396 let input_probe = compute_state.input_probe_for(*idx_id, dataflow.export_ids());
397 let snapshot_mode = if idx.with_snapshot || !subscribe_snapshot_optimization {
398 SnapshotMode::Include
399 } else {
400 compute_state.metrics.inc_subscribe_snapshot_optimization();
401 SnapshotMode::Exclude
402 };
403 context.import_index(
404 scope,
405 compute_state,
406 &mut tokens,
407 input_probe,
408 *idx_id,
409 &idx.desc.as_lir(),
410 &idx.typ,
411 snapshot_mode,
412 start_signal.clone(),
413 );
414 }
415
416 for object in dataflow.objects_to_build {
418 let bundle = context.scope.clone().region_named(
419 &format!("BuildingObject({:?})", object.id),
420 |region| {
421 let depends = object.plan.depends();
422 let in_let = object.plan.is_recursive();
423 context
424 .enter_region(region, Some(&depends))
425 .render_recursive_plan(
426 object.id,
427 0,
428 object.plan,
429 BindingInfo::Body { in_let },
431 )
432 .leave_region(context.scope)
433 },
434 );
435 let global_id = object.id;
436
437 context.log_dataflow_global_id(
438 *bundle
439 .scope()
440 .addr()
441 .first()
442 .expect("Dataflow root id must exist"),
443 global_id,
444 );
445 context.insert_id(Id::Global(object.id), bundle);
446 }
447
448 for (idx_id, dependencies, idx) in indexes {
450 context.export_index_iterative(
451 scope,
452 compute_state,
453 &tokens,
454 dependencies,
455 idx_id,
456 &idx,
457 &output_probe,
458 );
459 }
460
461 for (sink_id, dependencies, sink) in sinks {
463 context.export_sink(
464 compute_state,
465 &tokens,
466 dependencies,
467 sink_id,
468 &sink,
469 start_signal.clone(),
470 &output_probe,
471 scope,
472 );
473 }
474 });
475 } else {
476 scope.clone().region_named(&build_name, |region| {
477 let mut context = Context::for_dataflow_in(
478 &dataflow,
479 region.clone(),
480 compute_state,
481 until,
482 dataflow_expiration,
483 );
484
485 for (id, (oks, errs)) in imported_sources.into_iter() {
486 let bundle = crate::render::CollectionBundle::from_edge(
487 oks.enter_region(region),
488 errs.enter_region(region),
489 );
490 context.insert_id(id, bundle);
492 }
493
494 for (idx_id, idx) in &dataflow.index_imports {
496 let input_probe = compute_state.input_probe_for(*idx_id, dataflow.export_ids());
497 let snapshot_mode = if idx.with_snapshot || !subscribe_snapshot_optimization {
498 SnapshotMode::Include
499 } else {
500 compute_state.metrics.inc_subscribe_snapshot_optimization();
501 SnapshotMode::Exclude
502 };
503 context.import_index(
504 scope,
505 compute_state,
506 &mut tokens,
507 input_probe,
508 *idx_id,
509 &idx.desc.as_lir(),
510 &idx.typ,
511 snapshot_mode,
512 start_signal.clone(),
513 );
514 }
515
516 for object in dataflow.objects_to_build {
518 let bundle = context.scope.clone().region_named(
519 &format!("BuildingObject({:?})", object.id),
520 |region| {
521 let depends = object.plan.depends();
522 context
523 .enter_region(region, Some(&depends))
524 .render_plan(object.id, object.plan)
525 .leave_region(context.scope)
526 },
527 );
528 let global_id = object.id;
529 context.log_dataflow_global_id(
530 *bundle
531 .scope()
532 .addr()
533 .first()
534 .expect("Dataflow root id must exist"),
535 global_id,
536 );
537 context.insert_id(Id::Global(object.id), bundle);
538 }
539
540 for (idx_id, dependencies, idx) in indexes {
542 context.export_index(
543 compute_state,
544 &tokens,
545 dependencies,
546 idx_id,
547 &idx,
548 &output_probe,
549 );
550 }
551
552 for (sink_id, dependencies, sink) in sinks {
554 context.export_sink(
555 compute_state,
556 &tokens,
557 dependencies,
558 sink_id,
559 &sink,
560 start_signal.clone(),
561 &output_probe,
562 scope,
563 );
564 }
565 });
566 }
567 });
568}
569
570impl<'g, T> Context<'g, T>
573where
574 T: Refines<mz_repr::Timestamp> + RenderTimestamp,
575{
576 fn import_filtered_index_collection<
580 'outer,
581 Tr: TraceReader<Time = mz_repr::Timestamp, Batch: Navigable> + Clone,
582 V: Data,
583 >(
584 &self,
585 arranged: Arranged<'outer, Tr>,
586 start_signal: StartSignal,
587 mut logic: impl FnMut(BatchKey<'_, Tr>, BatchVal<'_, Tr>) -> V + 'static,
588 ) -> VecCollection<'g, T, V, BatchDiff<Tr>>
589 where
590 mz_repr::Timestamp: TotalOrder,
593 BatchCursor<Tr>: Cursor<Time = mz_repr::Timestamp>,
594 {
595 let oks = arranged.stream.with_start_signal(start_signal).filter({
596 let as_of = self.as_of_frontier.clone();
597 move |b| !<Antichain<mz_repr::Timestamp> as PartialOrder>::less_equal(b.upper(), &as_of)
598 });
599 Arranged::<'outer, Tr>::flat_map_batches(oks, move |a, b| [logic(a, b)]).enter(self.scope)
600 }
601
602 fn import_filtered_index_edge<'outer, Tr>(
611 &self,
612 arranged: Arranged<'outer, Tr>,
613 start_signal: StartSignal,
614 mut logic: impl FnMut(BatchKey<'_, Tr>, BatchVal<'_, Tr>, &mut Row) + 'static,
615 ) -> ColCollection<'g, T>
616 where
617 Tr: TraceReader<Time = mz_repr::Timestamp, Batch: Navigable> + Clone,
618 mz_repr::Timestamp: TotalOrder,
619 BatchCursor<Tr>: Cursor<Time = mz_repr::Timestamp, Diff = Diff>,
620 {
621 let as_of = self.as_of_frontier.clone();
622 arranged
623 .stream
624 .with_start_signal(start_signal)
625 .unary::<ColumnBuilder<(Row, mz_repr::Timestamp, Diff)>, _, _, _>(
626 Pipeline,
627 "IndexToColumnar",
628 move |_cap, _info| {
629 let mut row_buf = Row::default();
630 move |input, output| {
631 input.for_each(|time, data| {
632 let mut session = output.session_with_builder(&time);
633 for batch in data.iter() {
634 if <Antichain<mz_repr::Timestamp> as PartialOrder>::less_equal(
635 batch.upper(),
636 &as_of,
637 ) {
638 continue;
639 }
640 let mut cursor = batch.cursor();
641 while let Some(key) = cursor.get_key(batch) {
642 while let Some(val) = cursor.get_val(batch) {
643 logic(key, val, &mut row_buf);
644 cursor.map_times(batch, |t, d| {
645 let t = <BatchCursor<Tr> as Cursor>::owned_time(t);
646 let d = <BatchCursor<Tr> as Cursor>::owned_diff(d);
647 session.give((&*row_buf, &t, &d));
648 });
649 cursor.step_val(batch);
650 }
651 cursor.step_key(batch);
652 }
653 }
654 });
655 }
656 },
657 )
658 .as_collection()
659 .enter(self.scope)
660 }
661
662 pub(crate) fn import_index<'outer>(
663 &mut self,
664 outer: Scope<'outer, mz_repr::Timestamp>,
665 compute_state: &mut ComputeState,
666 tokens: &mut BTreeMap<GlobalId, Rc<dyn std::any::Any>>,
667 input_probe: probe::Handle<mz_repr::Timestamp>,
668 idx_id: GlobalId,
669 idx: &IndexDesc<LirScalarExpr>,
670 typ: &ReprRelationType,
671 snapshot_mode: SnapshotMode,
672 start_signal: StartSignal,
673 ) {
674 if let Some(traces) = compute_state.traces.get_mut(&idx_id) {
675 assert!(
676 PartialOrder::less_equal(&traces.compaction_frontier(), &self.as_of_frontier),
677 "Index {idx_id} has been allowed to compact beyond the dataflow as_of"
678 );
679
680 let token = traces.to_drop().clone();
681
682 let (mut oks, ok_button) = traces.oks_mut().import_frontier_core(
683 outer,
684 &format!("Index({}, {:?})", idx.on_id, idx.key),
685 self.as_of_frontier.clone(),
686 self.until.clone(),
687 );
688
689 oks.stream = oks.stream.probe_with(&input_probe);
690
691 let (err_arranged, err_button) = traces.errs_mut().import_frontier_core(
692 outer,
693 &format!("ErrIndex({}, {:?})", idx.on_id, idx.key),
694 self.as_of_frontier.clone(),
695 self.until.clone(),
696 );
697
698 let bundle = match snapshot_mode {
699 SnapshotMode::Include => {
700 let ok_arranged = oks
701 .enter(self.scope)
702 .with_start_signal(start_signal.clone());
703 let err_arranged = err_arranged
704 .enter(self.scope)
705 .with_start_signal(start_signal);
706 CollectionBundle::from_expressions(
707 idx.key.clone(),
708 ArrangementFlavor::Trace(idx_id, ok_arranged, err_arranged),
709 )
710 }
711 SnapshotMode::Exclude => {
712 let oks = {
719 let mut datums = DatumVec::new();
720 let (permutation, _thinning) =
721 permutation_for_arrangement(&idx.key, typ.arity());
722 self.import_filtered_index_edge(
723 oks,
724 start_signal.clone(),
725 move |k: DatumSeq, v: DatumSeq, row: &mut Row| {
726 let temp_storage = RowArena::new();
727 let mut datums_borrow = datums.borrow();
728 k.extend_datums(&temp_storage, &mut datums_borrow, None);
729 v.extend_datums(&temp_storage, &mut datums_borrow, None);
730 row.packer()
731 .extend(permutation.iter().map(|i| datums_borrow[*i]));
732 },
733 )
734 };
735 let errs = self.import_filtered_index_collection(
736 err_arranged,
737 start_signal,
738 |e, _| e.clone(),
739 );
740 CollectionBundle::from_edge(oks, errs)
741 }
742 };
743 self.update_id(Id::Global(idx.on_id), bundle);
744 tokens.insert(
745 idx_id,
746 Rc::new((PressOnDrop(ok_button), PressOnDrop(err_button), token)),
747 );
748 } else {
749 panic!(
750 "import of index {} failed while building dataflow {}",
751 idx_id, self.dataflow_id
752 );
753 }
754 }
755}
756
757impl<'g> Context<'g, mz_repr::Timestamp> {
760 pub(crate) fn export_index(
761 &self,
762 compute_state: &mut ComputeState,
763 tokens: &BTreeMap<GlobalId, Rc<dyn std::any::Any>>,
764 dependency_ids: BTreeSet<GlobalId>,
765 idx_id: GlobalId,
766 idx: &IndexDesc<LirScalarExpr>,
767 output_probe: &MzProbeHandle<mz_repr::Timestamp>,
768 ) {
769 let mut needed_tokens = Vec::new();
771 for dep_id in dependency_ids {
772 if let Some(token) = tokens.get(&dep_id) {
773 needed_tokens.push(Rc::clone(token));
774 }
775 }
776 let bundle = self.lookup_id(Id::Global(idx_id)).unwrap_or_else(|| {
777 panic!(
778 "Arrangement alarmingly absent! id: {:?}",
779 Id::Global(idx_id)
780 )
781 });
782
783 let key = &idx.key;
784 match bundle.arrangement(key) {
785 Some(ArrangementFlavor::Local(mut oks, mut errs)) => {
786 if let Some(&expiration) = self.dataflow_expiration.as_option() {
800 oks.stream = oks.stream.expire_stream_at(
801 &format!("{}_export_index_oks", self.debug_name),
802 expiration,
803 );
804 errs.stream = errs.stream.expire_stream_at(
805 &format!("{}_export_index_errs", self.debug_name),
806 expiration,
807 );
808 }
809
810 oks.stream = oks.stream.probe_notify_with(vec![output_probe.clone()]);
811
812 if let Some(logger) = compute_state.compute_logger.clone() {
814 errs.stream = errs.stream.log_dataflow_errors(logger, idx_id);
815 }
816
817 compute_state.traces.set(
818 idx_id,
819 TraceBundle::new(oks.trace, errs.trace).with_drop(needed_tokens),
820 );
821 }
822 Some(ArrangementFlavor::Trace(gid, _, _)) => {
823 let trace = compute_state.traces.get(&gid).unwrap().clone();
826 compute_state.traces.set(idx_id, trace);
827 }
828 None => {
829 println!("collection available: {:?}", bundle.collection.is_none());
830 println!(
831 "keys available: {:?}",
832 bundle.arranged.keys().collect::<Vec<_>>()
833 );
834 panic!(
835 "Arrangement alarmingly absent! id: {:?}, keys: {:?}",
836 Id::Global(idx_id),
837 key
838 );
839 }
840 };
841 }
842}
843
844impl<'g, T> Context<'g, T>
847where
848 T: RenderTimestamp,
849{
850 pub(crate) fn export_index_iterative<'outer>(
851 &self,
852 outer: Scope<'outer, mz_repr::Timestamp>,
853 compute_state: &mut ComputeState,
854 tokens: &BTreeMap<GlobalId, Rc<dyn std::any::Any>>,
855 dependency_ids: BTreeSet<GlobalId>,
856 idx_id: GlobalId,
857 idx: &IndexDesc<LirScalarExpr>,
858 output_probe: &MzProbeHandle<mz_repr::Timestamp>,
859 ) {
860 let mut needed_tokens = Vec::new();
862 for dep_id in dependency_ids {
863 if let Some(token) = tokens.get(&dep_id) {
864 needed_tokens.push(Rc::clone(token));
865 }
866 }
867 let bundle = self.lookup_id(Id::Global(idx_id)).unwrap_or_else(|| {
868 panic!(
869 "Arrangement alarmingly absent! id: {:?}",
870 Id::Global(idx_id)
871 )
872 });
873
874 let key = &idx.key;
875 match bundle.arrangement(key) {
876 Some(ArrangementFlavor::Local(oks, errs)) => {
877 let mut oks = oks
881 .as_collection(|k, v| (k.to_row(), v.to_row()))
882 .leave(outer)
883 .mz_arrange::<
884 ColumnationChunker<_>,
885 RowRowBatcher<_, _>,
886 RowRowBuilder<_, _>,
887 _,
888 >(
889 "Arrange export iterative",
890 );
891
892 let mut errs = errs
893 .as_collection(|k, v| (k.clone(), v.clone()))
894 .leave(outer)
895 .mz_arrange::<ColumnationChunker<_>, ErrBatcher<_, _>, ErrBuilder<_, _>, _>(
896 "Arrange export iterative err",
897 );
898
899 if let Some(&expiration) = self.dataflow_expiration.as_option() {
902 oks.stream = oks.stream.expire_stream_at(
903 &format!("{}_export_index_iterative_oks", self.debug_name),
904 expiration,
905 );
906 errs.stream = errs.stream.expire_stream_at(
907 &format!("{}_export_index_iterative_err", self.debug_name),
908 expiration,
909 );
910 }
911
912 oks.stream = oks.stream.probe_notify_with(vec![output_probe.clone()]);
913
914 if let Some(logger) = compute_state.compute_logger.clone() {
916 errs.stream = errs.stream.log_dataflow_errors(logger, idx_id);
917 }
918
919 compute_state.traces.set(
920 idx_id,
921 TraceBundle::new(oks.trace, errs.trace).with_drop(needed_tokens),
922 );
923 }
924 Some(ArrangementFlavor::Trace(gid, _, _)) => {
925 let trace = compute_state.traces.get(&gid).unwrap().clone();
928 compute_state.traces.set(idx_id, trace);
929 }
930 None => {
931 println!("collection available: {:?}", bundle.collection.is_none());
932 println!(
933 "keys available: {:?}",
934 bundle.arranged.keys().collect::<Vec<_>>()
935 );
936 panic!(
937 "Arrangement alarmingly absent! id: {:?}, keys: {:?}",
938 Id::Global(idx_id),
939 key,
940 );
941 }
942 };
943 }
944}
945
946enum BindingInfo {
952 Body { in_let: bool },
953 Let { id: LocalId, last: bool },
954 LetRec { id: LocalId, last: bool },
955}
956
957impl<'scope> Context<'scope, Product<mz_repr::Timestamp, PointStamp<u64>>> {
958 fn render_recursive_plan(
971 &mut self,
972 object_id: GlobalId,
973 level: usize,
974 plan: RenderPlan,
975 binding: BindingInfo,
976 ) -> CollectionBundle<'scope, Product<mz_repr::Timestamp, PointStamp<u64>>> {
977 for BindStage { lets, recs } in plan.binds {
978 let mut let_iter = lets.into_iter().peekable();
980 while let Some(LetBind { id, value }) = let_iter.next() {
981 let bundle =
982 self.scope
983 .clone()
984 .region_named(&format!("Binding({:?})", id), |region| {
985 let depends = value.depends();
986 let last = let_iter.peek().is_none();
987 let binding = BindingInfo::Let { id, last };
988 self.enter_region(region, Some(&depends))
989 .render_letfree_plan(object_id, value, binding)
990 .leave_region(self.scope)
991 });
992 let bundle = self.distinct_binding_errs(bundle);
993 self.insert_id(Id::Local(id), bundle);
994 }
995
996 let rec_ids: Vec<_> = recs.iter().map(|r| r.id).collect();
997
998 let mut variable_read = BTreeSet::new();
1005 let mut read_so_far = BTreeSet::new();
1006 for rec in recs.iter() {
1007 read_so_far.extend(rec.value.depends());
1008 if read_so_far.contains(&Id::Local(rec.id)) {
1009 variable_read.insert(rec.id);
1010 }
1011 }
1012
1013 let mut variables = BTreeMap::new();
1017 for id in rec_ids.iter() {
1018 use differential_dataflow::dynamic::feedback_summary;
1019 let inner = feedback_summary::<u64>(level + 1, 1);
1020 let (oks_v, oks_collection) =
1021 Variable::new(self.scope, Product::new(Default::default(), inner.clone()));
1022 let (err_v, err_collection) =
1023 Variable::new(self.scope, Product::new(Default::default(), inner));
1024
1025 if variable_read.contains(id) {
1026 self.insert_id(
1027 Id::Local(*id),
1028 CollectionBundle::from_edge(oks_collection, err_collection),
1029 );
1030 }
1031 variables.insert(Id::Local(*id), (oks_v, err_v));
1032 }
1033 let mut bound_oks = BTreeMap::new();
1036 let mut rec_iter = recs.into_iter().peekable();
1037 while let Some(RecBind { id, value, limit }) = rec_iter.next() {
1038 let last = rec_iter.peek().is_none();
1039 let binding = BindingInfo::LetRec { id, last };
1040 let bundle = self.render_recursive_plan(object_id, level + 1, value, binding);
1041 let (oks, mut err) = bundle.collection.clone().unwrap();
1044 bound_oks.insert(id, oks.clone());
1045 let bundle = self.distinct_binding_errs(bundle);
1050 self.insert_id(Id::Local(id), bundle);
1051 let (oks_v, err_v) = variables.remove(&Id::Local(id)).unwrap();
1052
1053 let mut oks = columnar_consolidate(oks, "LetRecConsolidation");
1055
1056 if let Some(limit) = limit {
1057 let mut time = RecTimestamp::default();
1068 let (over_limit, in_limit) = oks
1069 .inner
1070 .partition_by::<ColumnBuilder<(Row, RecTimestamp, Diff)>, _>(
1071 "LetRecLimit",
1072 move |(_data, reference, _diff)| {
1073 time.copy_from(*reference);
1074 let iteration_index = time.inner.get(level).copied().unwrap_or(0);
1077 iteration_index + 1 >= limit.max_iters.into()
1079 },
1080 );
1081 oks = Collection::new(in_limit);
1082 if !limit.return_at_limit {
1083 let (_, over_limit) = flat_map_datums::<
1087 _,
1088 CapacityContainerBuilder<Vec<(Row, RecTimestamp, Diff)>>,
1089 _,
1090 >(
1091 Collection::new(over_limit),
1092 "LetRecLimitExceeded",
1093 0,
1094 move |_datums, time, diff, _ok_session, err_session| {
1095 err_session.give((
1096 DataflowErrorSer::from(EvalError::LetRecLimitExceeded(
1097 format!("{}", limit.max_iters.get()).into(),
1098 )),
1099 time,
1100 diff,
1101 ));
1102 1
1103 },
1104 );
1105 err = err.concat(over_limit.as_collection());
1106 }
1107 }
1108
1109 let err: KeyCollection<_, _, _> = err.into();
1115 let errs = err
1116 .mz_arrange::<
1117 ColumnationChunker<_>,
1118 ErrBatcher<_, _>,
1119 ErrBuilder<_, _>,
1120 ErrSpine<_, _>,
1121 >("Arrange recursive err")
1122 .mz_reduce_abelian::<_, ErrBuilder<_, _>, ErrSpine<_, _>, _>(
1123 "Distinct recursive err",
1124 move |_k, _s, t| t.push(((), Diff::ONE)),
1125 )
1126 .as_collection(|k, _| k.clone());
1127
1128 oks_v.set(oks);
1129 err_v.set(errs);
1130 }
1131 for id in rec_ids.into_iter() {
1133 let bundle = self.remove_id(Id::Local(id)).unwrap();
1134 let (_, err) = bundle.collection.unwrap();
1135 let oks = bound_oks
1136 .remove(&id)
1137 .expect("rec binding bound while rendering above");
1138 self.insert_id(
1139 Id::Local(id),
1140 CollectionBundle::from_edge(
1141 columnar_leave_dynamic(oks, level + 1),
1142 err.leave_dynamic(level + 1),
1143 ),
1144 );
1145 }
1146 }
1147
1148 self.render_letfree_plan(object_id, plan.body, binding)
1149 }
1150}
1151
1152impl<'scope, T: RenderTimestamp + MaybeBucketByTime> Context<'scope, T> {
1153 fn render_plan(
1164 &mut self,
1165 object_id: GlobalId,
1166 plan: RenderPlan,
1167 ) -> CollectionBundle<'scope, T> {
1168 let mut in_let = false;
1169 for BindStage { lets, recs } in plan.binds {
1170 assert!(recs.is_empty());
1171
1172 let mut let_iter = lets.into_iter().peekable();
1173 while let Some(LetBind { id, value }) = let_iter.next() {
1174 in_let = true;
1176 let bundle =
1177 self.scope
1178 .clone()
1179 .region_named(&format!("Binding({:?})", id), |region| {
1180 let depends = value.depends();
1181 let last = let_iter.peek().is_none();
1182 let binding = BindingInfo::Let { id, last };
1183 self.enter_region(region, Some(&depends))
1184 .render_letfree_plan(object_id, value, binding)
1185 .leave_region(self.scope)
1186 });
1187 let bundle = self.distinct_binding_errs(bundle);
1188 self.insert_id(Id::Local(id), bundle);
1189 }
1190 }
1191
1192 self.scope.clone().region_named("Main Body", |region| {
1193 let depends = plan.body.depends();
1194 self.enter_region(region, Some(&depends))
1195 .render_letfree_plan(object_id, plan.body, BindingInfo::Body { in_let })
1196 .leave_region(self.scope)
1197 })
1198 }
1199
1200 fn distinct_binding_errs(
1209 &self,
1210 bundle: CollectionBundle<'scope, T>,
1211 ) -> CollectionBundle<'scope, T> {
1212 if ENABLE_ERROR_DISTINCT.get(&self.config_set) {
1213 bundle.distinct_errs()
1214 } else {
1215 bundle
1216 }
1217 }
1218
1219 fn render_letfree_plan(
1221 &self,
1222 object_id: GlobalId,
1223 plan: LetFreePlan,
1224 binding: BindingInfo,
1225 ) -> CollectionBundle<'scope, T> {
1226 let (mut nodes, root_id, topological_order) = plan.destruct();
1227
1228 let mut collections = BTreeMap::new();
1230
1231 let should_compute_lir_metadata = self.compute_logger.is_some();
1237 let mut lir_mapping_metadata = if should_compute_lir_metadata {
1238 Some(Vec::with_capacity(nodes.len()))
1239 } else {
1240 None
1241 };
1242
1243 let mut topo_iter = topological_order.into_iter().peekable();
1244 while let Some(lir_id) = topo_iter.next() {
1245 let node = nodes.remove(&lir_id).unwrap();
1246
1247 let metadata = if should_compute_lir_metadata {
1251 let operator = node.expr.humanize(&DummyHumanizer);
1252
1253 let operator = if topo_iter.peek().is_none() {
1255 match &binding {
1256 BindingInfo::Body { in_let: true } => format!("Returning {operator}"),
1257 BindingInfo::Body { in_let: false } => operator,
1258 BindingInfo::Let { id, last: true } => {
1259 format!("With {id} = {operator}")
1260 }
1261 BindingInfo::Let { id, last: false } => {
1262 format!("{id} = {operator}")
1263 }
1264 BindingInfo::LetRec { id, last: true } => {
1265 format!("With Recursive {id} = {operator}")
1266 }
1267 BindingInfo::LetRec { id, last: false } => {
1268 format!("{id} = {operator}")
1269 }
1270 }
1271 } else {
1272 operator
1273 };
1274
1275 let operator_id_start = self.scope.worker().peek_identifier();
1276 Some((operator, operator_id_start))
1277 } else {
1278 None
1279 };
1280
1281 let mut bundle = self.render_plan_expr(node.expr, &collections);
1282
1283 if let Some((operator, operator_id_start)) = metadata {
1284 let operator_id_end = self.scope.worker().peek_identifier();
1285 let operator_span = (operator_id_start, operator_id_end);
1286
1287 if let Some(lir_mapping_metadata) = &mut lir_mapping_metadata {
1288 lir_mapping_metadata.push((
1289 lir_id,
1290 LirMetadata::new(operator, node.parent, node.nesting, operator_span),
1291 ))
1292 }
1293 }
1294
1295 self.log_operator_hydration(&mut bundle, lir_id);
1296
1297 collections.insert(lir_id, bundle);
1298 }
1299
1300 if let Some(lir_mapping_metadata) = lir_mapping_metadata {
1301 self.log_lir_mapping(object_id, lir_mapping_metadata);
1302 }
1303
1304 collections
1305 .remove(&root_id)
1306 .expect("LetFreePlan invariant (1)")
1307 }
1308
1309 fn render_plan_expr(
1316 &self,
1317 expr: render_plan::Expr,
1318 collections: &BTreeMap<LirId, CollectionBundle<'scope, T>>,
1319 ) -> CollectionBundle<'scope, T> {
1320 use render_plan::Expr::*;
1321
1322 let expect_input = |id| {
1323 collections
1324 .get(&id)
1325 .cloned()
1326 .unwrap_or_else(|| panic!("missing input collection: {id}"))
1327 };
1328
1329 match expr {
1330 Constant { rows } => {
1331 let (rows, errs) = match rows {
1333 Ok(rows) => (rows, Vec::new()),
1334 Err(e) => (Vec::new(), vec![e]),
1335 };
1336
1337 let as_of_frontier = self.as_of_frontier.clone();
1339 let until = self.until.clone();
1340 let ok_collection = rows
1344 .into_iter()
1345 .filter_map(move |(row, mut time, diff)| {
1346 time.advance_by(as_of_frontier.borrow());
1347 if !until.less_equal(&time) {
1348 Some((
1349 row.0,
1350 <T as Refines<mz_repr::Timestamp>>::to_inner(time),
1351 diff,
1352 ))
1353 } else {
1354 None
1355 }
1356 })
1357 .to_stream_with_builder::<_, ConsolidatingColumnBuilder<Row, T, Diff>>(
1358 self.scope,
1359 )
1360 .as_collection();
1361
1362 let mut error_time: mz_repr::Timestamp = Timestamp::minimum();
1363 error_time.advance_by(self.as_of_frontier.borrow());
1364 let err_collection = errs
1365 .into_iter()
1366 .map(move |e| {
1367 (
1368 DataflowErrorSer::from(e),
1369 <T as Refines<mz_repr::Timestamp>>::to_inner(error_time),
1370 Diff::ONE,
1371 )
1372 })
1373 .to_stream(self.scope)
1374 .as_collection();
1375
1376 CollectionBundle::from_edge(ok_collection, err_collection)
1377 }
1378 Get { id, keys, plan } => {
1379 let mut collection = self
1382 .lookup_id(id)
1383 .unwrap_or_else(|| panic!("Get({:?}) not found at render time", id));
1384 match plan {
1385 mz_compute_types::plan::GetPlan::PassArrangements => {
1386 assert!(
1388 keys.arranged
1389 .iter()
1390 .all(|(key, _, _)| collection.arranged.contains_key(key))
1391 );
1392 assert!(keys.raw <= collection.collection.is_some());
1393 collection.arranged.retain(|key, _value| {
1395 keys.arranged.iter().any(|(key2, _, _)| key2 == key)
1396 });
1397 collection
1398 }
1399 mz_compute_types::plan::GetPlan::Arrangement(key, row, mfp) => {
1400 let (oks, errs) = collection.as_collection_core(
1401 mfp,
1402 Some((key, row)),
1403 self.until.clone(),
1404 );
1405 CollectionBundle::from_edge(oks, errs)
1406 }
1407 mz_compute_types::plan::GetPlan::Collection(mfp) => {
1408 let (oks, errs) =
1409 collection.as_collection_core(mfp, None, self.until.clone());
1410 CollectionBundle::from_edge(oks, errs)
1411 }
1412 }
1413 }
1414 Mfp {
1415 input,
1416 mfp,
1417 input_key_val,
1418 } => {
1419 let input = expect_input(input);
1420 if mfp.is_identity() {
1422 input
1423 } else {
1424 let (oks, errs) =
1425 input.as_collection_core(mfp, input_key_val, self.until.clone());
1426 CollectionBundle::from_edge(oks, errs)
1427 }
1428 }
1429 FlatMap {
1430 input_key,
1431 input,
1432 exprs,
1433 func,
1434 mfp_after: mfp,
1435 } => {
1436 let input = expect_input(input);
1437 self.render_flat_map(input_key, input, exprs, func, mfp)
1438 }
1439 Join { inputs, plan } => {
1440 let inputs = inputs.into_iter().map(expect_input).collect();
1441 match plan {
1442 mz_compute_types::plan::join::JoinPlan::Linear(linear_plan) => {
1443 self.render_join(inputs, linear_plan)
1444 }
1445 mz_compute_types::plan::join::JoinPlan::Delta(delta_plan) => {
1446 self.render_delta_join(inputs, delta_plan)
1447 }
1448 }
1449 }
1450 Reduce {
1451 input_key,
1452 input,
1453 key_val_plan,
1454 plan,
1455 mfp_after,
1456 temporal_bucketing_strategy,
1457 } => {
1458 let input = expect_input(input);
1459 let mfp_option = (!mfp_after.is_identity()).then_some(mfp_after);
1460 self.render_reduce(
1461 input_key,
1462 input,
1463 key_val_plan,
1464 plan,
1465 mfp_option,
1466 temporal_bucketing_strategy,
1467 )
1468 }
1469 TopK {
1470 input,
1471 top_k_plan,
1472 temporal_bucketing_strategy,
1473 } => {
1474 let input = expect_input(input);
1475 self.render_topk(input, top_k_plan, temporal_bucketing_strategy)
1476 }
1477 Negate { input } => {
1478 let input = expect_input(input);
1479 let (oks, errs) = input
1480 .collection
1481 .clone()
1482 .expect("Negate input must be an unarranged collection");
1483 CollectionBundle::from_edge(columnar_negate(oks), errs)
1484 }
1485 Threshold {
1486 input,
1487 threshold_plan,
1488 } => {
1489 let input = expect_input(input);
1490 self.render_threshold(input, threshold_plan)
1491 }
1492 Union {
1493 inputs,
1494 consolidate_output,
1495 temporal_bucketing_strategies,
1496 } => {
1497 let mut oks = Vec::new();
1498 let mut errs = Vec::new();
1499 for (input, strategy) in inputs.into_iter().zip_eq(temporal_bucketing_strategies) {
1500 let (os, es) = expect_input(input)
1501 .collection
1502 .clone()
1503 .expect("Union input must be an unarranged collection");
1504 let os = if matches!(strategy, ArrangementStrategy::TemporalBucketing)
1508 && ENABLE_COMPUTE_TEMPORAL_BUCKETING.get(&self.config_set)
1509 {
1510 let summary: mz_repr::Timestamp = TEMPORAL_BUCKETING_SUMMARY
1511 .get(&self.config_set)
1512 .try_into()
1513 .expect("must fit");
1514 T::maybe_apply_temporal_bucketing(
1515 os.inner,
1516 self.as_of_frontier.clone(),
1517 summary,
1518 )
1519 } else {
1520 os
1521 };
1522 oks.push(os);
1523 errs.push(es);
1524 }
1525 let oks = concat_many(self.scope, oks);
1526 let oks = if consolidate_output {
1527 columnar_consolidate(oks, "UnionConsolidation")
1528 } else {
1529 oks
1530 };
1531 let errs = differential_dataflow::collection::concatenate(self.scope, errs);
1532 CollectionBundle::from_edge(oks, errs)
1533 }
1534 ArrangeBy {
1535 input_key,
1536 input,
1537 input_mfp,
1538 forms: keys,
1539 strategy,
1540 } => {
1541 let input = expect_input(input);
1542 input.ensure_collections(
1543 keys,
1544 input_key,
1545 input_mfp,
1546 self.as_of_frontier.clone(),
1547 self.until.clone(),
1548 &self.config_set,
1549 strategy,
1550 )
1551 }
1552 }
1553 }
1554
1555 fn log_dataflow_global_id(&self, dataflow_index: usize, global_id: GlobalId) {
1556 if let Some(logger) = &self.compute_logger {
1557 logger.log(&ComputeEvent::DataflowGlobal(DataflowGlobal {
1558 dataflow_index,
1559 global_id,
1560 }));
1561 }
1562 }
1563
1564 fn log_lir_mapping(&self, global_id: GlobalId, mapping: Vec<(LirId, LirMetadata)>) {
1565 if let Some(logger) = &self.compute_logger {
1566 logger.log(&ComputeEvent::LirMapping(LirMapping { global_id, mapping }));
1567 }
1568 }
1569
1570 fn log_operator_hydration(&self, bundle: &mut CollectionBundle<'scope, T>, lir_id: LirId) {
1571 match bundle.arranged.values_mut().next() {
1591 Some(arrangement) => {
1592 use ArrangementFlavor::*;
1593
1594 match arrangement {
1595 Local(a, _) => {
1596 a.stream = self.log_operator_hydration_inner(a.stream.clone(), lir_id);
1597 }
1598 Trace(_, a, _) => {
1599 a.stream = self.log_operator_hydration_inner(a.stream.clone(), lir_id);
1600 }
1601 }
1602 }
1603 None => {
1604 let (oks, _) = bundle
1605 .collection
1606 .as_mut()
1607 .expect("CollectionBundle invariant");
1608 let stream = self.log_operator_hydration_inner(oks.inner.clone(), lir_id);
1609 *oks = stream.as_collection();
1610 }
1611 }
1612 }
1613
1614 fn log_operator_hydration_inner<D>(
1615 &self,
1616 stream: Stream<'scope, T, D>,
1617 lir_id: LirId,
1618 ) -> Stream<'scope, T, D>
1619 where
1620 D: timely::Container + Clone + 'static,
1621 {
1622 let Some(logger) = self.compute_logger.clone() else {
1623 return stream.clone(); };
1625
1626 let export_ids = self.export_ids.clone();
1627
1628 let mut hydration_frontier = Antichain::new();
1636 for time in self.as_of_frontier.iter() {
1637 if let Some(time) = time.try_step_forward() {
1638 hydration_frontier.insert(Refines::to_inner(time));
1639 }
1640 }
1641
1642 let name = format!("LogOperatorHydration ({lir_id})");
1643 stream.unary_frontier(Pipeline, &name, |_cap, _info| {
1644 let mut hydrated = false;
1645
1646 for &export_id in &export_ids {
1647 logger.log(&ComputeEvent::OperatorHydration(OperatorHydration {
1648 export_id,
1649 lir_id,
1650 hydrated,
1651 }));
1652 }
1653
1654 move |(input, frontier), output| {
1655 input.for_each(|cap, data| {
1657 output.session(&cap).give_container(data);
1658 });
1659
1660 if hydrated {
1661 return;
1662 }
1663
1664 if PartialOrder::less_equal(&hydration_frontier.borrow(), &frontier.frontier()) {
1665 hydrated = true;
1666
1667 for &export_id in &export_ids {
1668 logger.log(&ComputeEvent::OperatorHydration(OperatorHydration {
1669 export_id,
1670 lir_id,
1671 hydrated,
1672 }));
1673 }
1674 }
1675 }
1676 })
1677 }
1678}
1679
1680#[allow(dead_code)] pub trait RenderTimestamp: MzTimestamp + Default + Refines<mz_repr::Timestamp> {
1683 fn system_time(&mut self) -> &mut mz_repr::Timestamp;
1688 fn system_delay(delay: mz_repr::Timestamp) -> <Self as Timestamp>::Summary;
1690 fn event_time(&self) -> mz_repr::Timestamp;
1692 fn event_time_mut(&mut self) -> &mut mz_repr::Timestamp;
1694 fn event_delay(delay: mz_repr::Timestamp) -> <Self as Timestamp>::Summary;
1696 fn step_back(&self) -> Self;
1699}
1700
1701pub trait MaybeBucketByTime: Timestamp + ColumnarData {
1708 fn maybe_apply_temporal_bucketing<'scope, D>(
1710 stream: Stream<'scope, Self, Column<(D, Self, Diff)>>,
1711 as_of: Antichain<mz_repr::Timestamp>,
1712 summary: mz_repr::Timestamp,
1713 ) -> Collection<'scope, Self, Column<(D, Self, Diff)>>
1714 where
1715 D: differential_dataflow::ExchangeData
1716 + crate::typedefs::MzData
1717 + differential_dataflow::Hashable
1718 + ColumnarData,
1719 for<'a> ::columnar::Ref<'a, D>: Copy + Ord + std::hash::Hash,
1720 for<'a> ::columnar::Ref<'a, Self>: Copy + Ord,
1721 for<'a> ::columnar::Ref<'a, Diff>: Ord,
1722 for<'a> <(D, Self, Diff) as ColumnarData>::Container:
1723 ColumnarPush<&'a (D, Self, Diff)> + ColumnarPush<::columnar::Ref<'a, (D, Self, Diff)>>;
1724
1725 fn maybe_apply_temporal_bucketing_vec<'scope, D>(
1731 stream: StreamVec<'scope, Self, (D, Self, Diff)>,
1732 as_of: Antichain<mz_repr::Timestamp>,
1733 summary: mz_repr::Timestamp,
1734 ) -> VecCollection<'scope, Self, D, Diff>
1735 where
1736 D: differential_dataflow::ExchangeData
1737 + crate::typedefs::MzData
1738 + differential_dataflow::Hashable
1739 + ColumnarData,
1740 for<'a> ::columnar::Ref<'a, D>: Copy + Ord + std::hash::Hash,
1741 for<'a> ::columnar::Ref<'a, Self>: Copy + Ord,
1742 for<'a> ::columnar::Ref<'a, Diff>: Ord,
1743 for<'a> <(D, Self, Diff) as ColumnarData>::Container:
1744 ColumnarPush<&'a (D, Self, Diff)> + ColumnarPush<::columnar::Ref<'a, (D, Self, Diff)>>;
1745}
1746
1747impl RenderTimestamp for mz_repr::Timestamp {
1748 fn system_time(&mut self) -> &mut mz_repr::Timestamp {
1749 self
1750 }
1751 fn system_delay(delay: mz_repr::Timestamp) -> <Self as Timestamp>::Summary {
1752 delay
1753 }
1754 fn event_time(&self) -> mz_repr::Timestamp {
1755 *self
1756 }
1757 fn event_time_mut(&mut self) -> &mut mz_repr::Timestamp {
1758 self
1759 }
1760 fn event_delay(delay: mz_repr::Timestamp) -> <Self as Timestamp>::Summary {
1761 delay
1762 }
1763 fn step_back(&self) -> Self {
1764 self.saturating_sub(1)
1765 }
1766}
1767
1768impl MaybeBucketByTime for mz_repr::Timestamp {
1769 fn maybe_apply_temporal_bucketing<'scope, D>(
1770 stream: Stream<'scope, Self, Column<(D, Self, Diff)>>,
1771 as_of: Antichain<mz_repr::Timestamp>,
1772 summary: mz_repr::Timestamp,
1773 ) -> Collection<'scope, Self, Column<(D, Self, Diff)>>
1774 where
1775 D: differential_dataflow::ExchangeData
1776 + crate::typedefs::MzData
1777 + differential_dataflow::Hashable
1778 + ColumnarData,
1779 for<'a> ::columnar::Ref<'a, D>: Copy + Ord + std::hash::Hash,
1780 for<'a> ::columnar::Ref<'a, Self>: Copy + Ord,
1781 for<'a> ::columnar::Ref<'a, Diff>: Ord,
1782 for<'a> <(D, Self, Diff) as ColumnarData>::Container:
1783 ColumnarPush<&'a (D, Self, Diff)> + ColumnarPush<::columnar::Ref<'a, (D, Self, Diff)>>,
1784 {
1785 stream.bucket(as_of, summary).as_collection()
1786 }
1787
1788 fn maybe_apply_temporal_bucketing_vec<'scope, D>(
1789 stream: StreamVec<'scope, Self, (D, Self, Diff)>,
1790 as_of: Antichain<mz_repr::Timestamp>,
1791 summary: mz_repr::Timestamp,
1792 ) -> VecCollection<'scope, Self, D, Diff>
1793 where
1794 D: differential_dataflow::ExchangeData
1795 + crate::typedefs::MzData
1796 + differential_dataflow::Hashable
1797 + ColumnarData,
1798 for<'a> ::columnar::Ref<'a, D>: Copy + Ord + std::hash::Hash,
1799 for<'a> ::columnar::Ref<'a, Self>: Copy + Ord,
1800 for<'a> ::columnar::Ref<'a, Diff>: Ord,
1801 for<'a> <(D, Self, Diff) as ColumnarData>::Container:
1802 ColumnarPush<&'a (D, Self, Diff)> + ColumnarPush<::columnar::Ref<'a, (D, Self, Diff)>>,
1803 {
1804 stream.bucket(as_of, summary).as_collection()
1805 }
1806}
1807
1808impl RenderTimestamp for Product<mz_repr::Timestamp, PointStamp<u64>> {
1809 fn system_time(&mut self) -> &mut mz_repr::Timestamp {
1810 &mut self.outer
1811 }
1812 fn system_delay(delay: mz_repr::Timestamp) -> <Self as Timestamp>::Summary {
1813 Product::new(delay, Default::default())
1814 }
1815 fn event_time(&self) -> mz_repr::Timestamp {
1816 self.outer
1817 }
1818 fn event_time_mut(&mut self) -> &mut mz_repr::Timestamp {
1819 &mut self.outer
1820 }
1821 fn event_delay(delay: mz_repr::Timestamp) -> <Self as Timestamp>::Summary {
1822 Product::new(delay, Default::default())
1823 }
1824 fn step_back(&self) -> Self {
1825 let inner = self.inner.clone();
1829 let mut vec = inner.into_inner();
1830 for item in vec.iter_mut() {
1831 *item = item.saturating_sub(1);
1832 }
1833 Product::new(self.outer.saturating_sub(1), PointStamp::new(vec))
1834 }
1835}
1836
1837impl MaybeBucketByTime for Product<mz_repr::Timestamp, PointStamp<u64>> {
1838 fn maybe_apply_temporal_bucketing<'scope, D>(
1839 stream: Stream<'scope, Self, Column<(D, Self, Diff)>>,
1840 _as_of: Antichain<mz_repr::Timestamp>,
1841 _summary: mz_repr::Timestamp,
1842 ) -> Collection<'scope, Self, Column<(D, Self, Diff)>>
1843 where
1844 D: differential_dataflow::ExchangeData
1845 + crate::typedefs::MzData
1846 + differential_dataflow::Hashable
1847 + ColumnarData,
1848 for<'a> ::columnar::Ref<'a, D>: Copy + Ord + std::hash::Hash,
1849 for<'a> ::columnar::Ref<'a, Self>: Copy + Ord,
1850 for<'a> ::columnar::Ref<'a, Diff>: Ord,
1851 for<'a> <(D, Self, Diff) as ColumnarData>::Container:
1852 ColumnarPush<&'a (D, Self, Diff)> + ColumnarPush<::columnar::Ref<'a, (D, Self, Diff)>>,
1853 {
1854 stream.as_collection()
1856 }
1857
1858 fn maybe_apply_temporal_bucketing_vec<'scope, D>(
1859 stream: StreamVec<'scope, Self, (D, Self, Diff)>,
1860 _as_of: Antichain<mz_repr::Timestamp>,
1861 _summary: mz_repr::Timestamp,
1862 ) -> VecCollection<'scope, Self, D, Diff>
1863 where
1864 D: differential_dataflow::ExchangeData
1865 + crate::typedefs::MzData
1866 + differential_dataflow::Hashable
1867 + ColumnarData,
1868 for<'a> ::columnar::Ref<'a, D>: Copy + Ord + std::hash::Hash,
1869 for<'a> ::columnar::Ref<'a, Self>: Copy + Ord,
1870 for<'a> ::columnar::Ref<'a, Diff>: Ord,
1871 for<'a> <(D, Self, Diff) as ColumnarData>::Container:
1872 ColumnarPush<&'a (D, Self, Diff)> + ColumnarPush<::columnar::Ref<'a, (D, Self, Diff)>>,
1873 {
1874 stream.as_collection()
1876 }
1877}
1878
1879#[derive(Clone)]
1889pub(crate) struct StartSignal {
1890 fut: futures::future::Shared<oneshot::Receiver<Infallible>>,
1895 token_ref: Weak<RefCell<Box<dyn Any>>>,
1897}
1898
1899impl StartSignal {
1900 pub fn new() -> (Self, Rc<dyn Any>) {
1903 let (tx, rx) = oneshot::channel::<Infallible>();
1904 let token: Rc<RefCell<Box<dyn Any>>> = Rc::new(RefCell::new(Box::new(tx)));
1905 let signal = Self {
1906 fut: rx.shared(),
1907 token_ref: Rc::downgrade(&token),
1908 };
1909 (signal, token)
1910 }
1911
1912 pub fn has_fired(&self) -> bool {
1913 self.token_ref.strong_count() == 0
1914 }
1915
1916 pub fn into_send_future(self) -> impl Future<Output = ()> + Send {
1921 use futures::FutureExt;
1922 self.fut.map(|_| ())
1923 }
1924
1925 pub fn drop_on_fire(&self, to_drop: Box<dyn Any>) {
1926 if let Some(token) = self.token_ref.upgrade() {
1927 let mut token = token.borrow_mut();
1928 let inner = std::mem::replace(&mut *token, Box::new(()));
1929 *token = Box::new((inner, to_drop));
1930 }
1931 }
1932}
1933
1934impl Future for StartSignal {
1935 type Output = ();
1936
1937 fn poll(mut self: Pin<&mut Self>, cx: &mut std::task::Context<'_>) -> Poll<Self::Output> {
1938 self.fut.poll_unpin(cx).map(|_| ())
1939 }
1940}
1941
1942pub(crate) trait WithStartSignal {
1944 fn with_start_signal(self, signal: StartSignal) -> Self;
1949}
1950
1951impl<'scope, Tr> WithStartSignal for Arranged<'scope, Tr>
1952where
1953 Tr: TraceReader<Time: RenderTimestamp> + Clone,
1954{
1955 fn with_start_signal(self, signal: StartSignal) -> Self {
1956 Arranged {
1957 stream: self.stream.with_start_signal(signal),
1958 trace: self.trace,
1959 }
1960 }
1961}
1962
1963impl<'scope, T: Timestamp, D> WithStartSignal for Stream<'scope, T, D>
1964where
1965 D: timely::Container + Clone + 'static,
1966{
1967 fn with_start_signal(self, signal: StartSignal) -> Self {
1968 let activations = self.scope().activations();
1969 self.unary(Pipeline, "StartSignal", |_cap, info| {
1970 let token = Box::new(ActivateOnDrop::new((), info.address, activations));
1971 signal.drop_on_fire(token);
1972
1973 let mut stash = Vec::new();
1974
1975 move |input, output| {
1976 if !signal.has_fired() {
1978 input.for_each(|cap, data| stash.push((cap, std::mem::take(data))));
1979 return;
1980 }
1981
1982 for (cap, mut data) in std::mem::take(&mut stash) {
1984 output.session(&cap).give_container(&mut data);
1985 }
1986
1987 input.for_each(|cap, data| {
1989 output.session(&cap).give_container(data);
1990 });
1991 }
1992 })
1993 }
1994}
1995
1996fn suppress_early_progress<'scope, T: Timestamp, D>(
2018 stream: Stream<'scope, T, D>,
2019 as_of: Antichain<T>,
2020) -> Stream<'scope, T, D>
2021where
2022 D: timely::Container + Clone,
2023{
2024 stream.unary_frontier(Pipeline, "SuppressEarlyProgress", |default_cap, _info| {
2025 let mut early_cap = Some(default_cap);
2026
2027 move |(input, frontier), output| {
2028 input.for_each_time(|data_cap, data| {
2029 if as_of.less_than(data_cap.time()) {
2030 let mut session = output.session(&data_cap);
2031 for data in data {
2032 session.give_container(data);
2033 }
2034 } else {
2035 let cap = early_cap.as_ref().expect("early_cap can't be dropped yet");
2036 let mut session = output.session(&cap);
2037 for data in data {
2038 session.give_container(data);
2039 }
2040 }
2041 });
2042
2043 if !PartialOrder::less_equal(&frontier.frontier(), &as_of.borrow()) {
2044 early_cap.take();
2045 }
2046 }
2047 })
2048}
2049
2050trait LimitProgress<T: Timestamp> {
2052 fn limit_progress(
2080 self,
2081 handle: MzProbeHandle<T>,
2082 slack_ms: u64,
2083 limit: Option<usize>,
2084 upper: Antichain<T>,
2085 name: String,
2086 ) -> Self;
2087}
2088
2089trait RecordTimes {
2091 fn for_each_time(&self, f: impl FnMut(mz_repr::Timestamp));
2093}
2094
2095impl<D, R> RecordTimes for Vec<(D, mz_repr::Timestamp, R)> {
2096 fn for_each_time(&self, mut f: impl FnMut(mz_repr::Timestamp)) {
2097 for (_, time, _) in self {
2098 f(*time);
2099 }
2100 }
2101}
2102
2103impl<D, R> RecordTimes for Column<(D, mz_repr::Timestamp, R)>
2104where
2105 D: ColumnarData,
2106 R: ColumnarData,
2107 (D, mz_repr::Timestamp, R): ColumnarData<
2108 Container = (
2109 D::Container,
2110 <mz_repr::Timestamp as ColumnarData>::Container,
2111 R::Container,
2112 ),
2113 >,
2114{
2115 fn for_each_time(&self, mut f: impl FnMut(mz_repr::Timestamp)) {
2116 for time in self.borrow().1.into_index_iter() {
2117 f(time);
2118 }
2119 }
2120}
2121
2122impl<'scope, C> LimitProgress<mz_repr::Timestamp> for Stream<'scope, mz_repr::Timestamp, C>
2125where
2126 C: timely::Container + Clone + RecordTimes,
2127{
2128 fn limit_progress(
2129 self,
2130 handle: MzProbeHandle<mz_repr::Timestamp>,
2131 slack_ms: u64,
2132 limit: Option<usize>,
2133 upper: Antichain<mz_repr::Timestamp>,
2134 name: String,
2135 ) -> Self {
2136 let scope = self.scope();
2137 let stream =
2138 self.unary_frontier(Pipeline, &format!("LimitProgress({name})"), |_cap, info| {
2139 let mut pending_times: BTreeSet<mz_repr::Timestamp> = BTreeSet::new();
2141 let mut retained_cap: Option<Capability<mz_repr::Timestamp>> = None;
2143
2144 let activator = scope.activator_for(info.address);
2145 handle.activate(activator.clone());
2146
2147 move |(input, frontier), output| {
2148 input.for_each(|cap, data| {
2149 data.for_each_time(|time| {
2150 let Some(time) = u64::from(time).checked_add(slack_ms) else {
2151 return;
2152 };
2153 let rounded_time = if slack_ms == 0 {
2157 time
2158 } else {
2159 (time / slack_ms).saturating_add(1).saturating_mul(slack_ms)
2160 };
2161 if !upper.less_than(&rounded_time.into()) {
2162 pending_times.insert(rounded_time.into());
2163 }
2164 });
2165 output.session(&cap).give_container(data);
2166 if retained_cap.as_ref().is_none_or(|c| {
2167 !c.time().less_than(cap.time()) && !upper.less_than(cap.time())
2168 }) {
2169 retained_cap = Some(cap.retain(0));
2170 }
2171 });
2172
2173 handle.with_frontier(|f| {
2174 while pending_times
2175 .first()
2176 .map_or(false, |retained_time| !f.less_than(&retained_time))
2177 {
2178 let _ = pending_times.pop_first();
2179 }
2180 });
2181
2182 while limit.map_or(false, |limit| pending_times.len() > limit) {
2183 let _ = pending_times.pop_first();
2184 }
2185
2186 match (retained_cap.as_mut(), pending_times.first()) {
2187 (Some(cap), Some(first)) => cap.downgrade(first),
2188 (_, None) => retained_cap = None,
2189 _ => {}
2190 }
2191
2192 if frontier.is_empty() {
2193 retained_cap = None;
2194 pending_times.clear();
2195 }
2196
2197 if !pending_times.is_empty() {
2198 tracing::debug!(
2199 name,
2200 info.global_id,
2201 pending_times = %PendingTimesDisplay(pending_times.iter().cloned()),
2202 frontier = ?frontier.frontier().get(0),
2203 probe = ?handle.with_frontier(|f| f.get(0).cloned()),
2204 ?upper,
2205 "pending times",
2206 );
2207 }
2208 }
2209 });
2210 stream
2211 }
2212}
2213
2214struct PendingTimesDisplay<T>(T);
2217
2218impl<T> std::fmt::Display for PendingTimesDisplay<T>
2219where
2220 T: IntoIterator<Item = mz_repr::Timestamp> + Clone,
2221{
2222 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
2223 let mut iter = self.0.clone().into_iter();
2224 write!(f, "[")?;
2225 if let Some(first) = iter.next() {
2226 write!(f, "{}", first)?;
2227 let mut last = u64::from(first);
2228 for time in iter {
2229 write!(f, ", +{}", u64::from(time) - last)?;
2230 last = u64::from(time);
2231 }
2232 }
2233 write!(f, "]")?;
2234 Ok(())
2235 }
2236}
2237
2238#[derive(Clone, Copy, Debug)]
2241struct Pairer {
2242 split_arity: usize,
2243}
2244
2245impl Pairer {
2246 fn new(split_arity: usize) -> Self {
2248 Self { split_arity }
2249 }
2250
2251 fn merge<'a, I1, I2>(&self, first: I1, second: I2) -> Row
2253 where
2254 I1: IntoIterator<Item = Datum<'a>>,
2255 I2: IntoIterator<Item = Datum<'a>>,
2256 {
2257 SharedRow::pack(first.into_iter().chain(second))
2258 }
2259
2260 fn split<'a>(&self, datum_iter: impl IntoIterator<Item = Datum<'a>>) -> (Row, Row) {
2262 let mut datum_iter = datum_iter.into_iter();
2263 let mut row_builder = SharedRow::get();
2264 let first = row_builder.pack_using(datum_iter.by_ref().take(self.split_arity));
2265 let second = row_builder.pack_using(datum_iter);
2266 (first, second)
2267 }
2268}