1use std::fmt::Debug;
92
93use differential_dataflow::difference::{IsZero, Semigroup};
94use differential_dataflow::hashable::Hashable;
95use differential_dataflow::lattice::Lattice;
96use differential_dataflow::logging::Logger;
97use differential_dataflow::operators::arrange::agent::TraceAgent;
98use differential_dataflow::operators::arrange::arrangement::{Arranged, arrange_core};
99use differential_dataflow::trace::chunk::{ChunkBatch, ChunkBatcher, ChunkBuilder};
100use differential_dataflow::trace::{BatchReader, Batcher, Cursor, Description, TraceReader};
101use differential_dataflow::{AsCollection, VecCollection};
102use mz_dyncfg::ConfigSet;
103use mz_repr::{Datum, Diff, GlobalId, Row};
104#[cfg(feature = "fuzzing")]
106use mz_row_spine::DatumSeq;
107use mz_row_spine::{FundedValRowSpine, ValRowColPagedBuilder};
108use mz_storage_types::dyncfgs::ENABLE_UPSERT_CHUNKED_STASH;
109use mz_storage_types::errors::{DataflowError, EnvelopeError, UpsertError};
110use mz_timely_util::builder_async::{
111 AsyncOutputHandle, Event as AsyncEvent, OperatorBuilder as AsyncOperatorBuilder,
112 PressOnDropButton,
113};
114use mz_timely_util::columnar::batcher::ColumnChunker;
115use mz_timely_util::columnar::body::ColumnBody;
116use mz_timely_util::columnar::builder::ColumnBuilder;
117use mz_timely_util::columnar::chunk::{ChunkChunker, ColumnChunk};
118use mz_timely_util::columnar::merge_batcher::{ColumnMergeBatcher, PagedChunker};
119use mz_timely_util::columnar::unload::UnloadBatch;
120use mz_timely_util::columnar::{Col2ValPagedBatcher, Column};
121use mz_timely_util::containers::stack::FueledBuilder;
122use std::convert::Infallible;
123use timely::container::{CapacityContainerBuilder, PushInto};
124use timely::dataflow::channels::pact::{Exchange, Pipeline};
125use timely::dataflow::operators::generic::Operator;
126use timely::dataflow::operators::{Capability, CapabilitySet, Exchange as _};
127use timely::dataflow::{Stream, StreamVec};
128use timely::order::{PartialOrder, TotalOrder};
129use timely::progress::frontier::AntichainRef;
130use timely::progress::timestamp::Refines;
131use timely::progress::{Antichain, Timestamp};
132
133use crate::healthcheck::HealthStatusUpdate;
134use crate::metrics::upsert::UpsertMetrics;
135use crate::statistics::SourceStatistics;
136use crate::upsert::UpsertKey;
137use crate::upsert::UpsertSourceTime;
138use crate::upsert::UpsertValue;
139
140#[derive(Clone, Copy, Debug)]
145pub enum UpsertStashFlavor {
146 Paged,
150 Chunked,
153}
154
155impl UpsertStashFlavor {
156 pub fn from_config(config: &ConfigSet) -> Self {
160 if ENABLE_UPSERT_CHUNKED_STASH.get(config) {
161 Self::Chunked
162 } else {
163 Self::Paged
164 }
165 }
166}
167
168struct UpsertFeedbackBatcher<T: columnar::Columnar>(Col2ValPagedBatcher<UpsertKey, Row, T, Diff>);
180
181impl<T> Batcher for UpsertFeedbackBatcher<T>
182where
183 T: Timestamp + columnar::Columnar + Default + PartialOrder,
184 for<'a> columnar::Ref<'a, T>: Copy + Ord,
185{
186 type Output = Column<((UpsertKey, Row), T, Diff)>;
187 type Time = T;
188
189 fn new(logger: Option<Logger>, operator_id: usize) -> Self {
190 let mut batcher =
191 <Col2ValPagedBatcher<UpsertKey, Row, T, Diff> as Batcher>::new(logger, operator_id);
192 batcher.set_pager(crate::upsert::upsert_stash_pager::pager());
193 Self(batcher)
194 }
195
196 fn seal(&mut self, upper: Antichain<T>) -> (Vec<Self::Output>, Description<T>) {
197 self.0.seal(upper)
198 }
199
200 fn frontier(&mut self) -> AntichainRef<'_, T> {
201 self.0.frontier()
202 }
203}
204
205impl<T> PushInto<Column<((UpsertKey, Row), T, Diff)>> for UpsertFeedbackBatcher<T>
206where
207 T: Timestamp + columnar::Columnar + Default + PartialOrder,
208 for<'a> columnar::Ref<'a, T>: Copy + Ord,
209{
210 fn push_into(&mut self, chunk: Column<((UpsertKey, Row), T, Diff)>) {
211 self.0.push_into(chunk)
212 }
213}
214
215type FeedbackUpdate<T> = ((UpsertKey, Row), T, Diff);
218
219type FeedbackChunk<T> = ColumnChunk<(UpsertKey, Row), T, Diff>;
226
227type FeedbackSpine<T> =
230 mz_timely_util::funded_spine::Spine<std::rc::Rc<ChunkBatch<FeedbackChunk<T>>>>;
231
232#[derive(Clone, Debug, Default, columnar::Columnar)]
250#[columnar(derive(PartialEq, Eq, PartialOrd, Ord))]
251struct UpsertDiff<O> {
252 from_time: O,
253 value: Option<Row>,
254}
255
256impl<O> IsZero for UpsertDiff<O> {
257 fn is_zero(&self) -> bool {
258 false
259 }
260}
261
262impl<O: Ord + Clone> Semigroup for UpsertDiff<O> {
263 fn plus_equals(&mut self, rhs: &Self) {
264 if rhs.from_time > self.from_time {
265 *self = rhs.clone();
266 }
267 }
268}
269
270impl<'a, O> Semigroup<columnar::Ref<'a, UpsertDiff<O>>> for UpsertDiff<O>
277where
278 O: columnar::Columnar + Ord + Clone,
279{
280 fn plus_equals(&mut self, rhs: &columnar::Ref<'a, UpsertDiff<O>>) {
281 let rhs_from_time = <O as columnar::Columnar>::into_owned(rhs.from_time);
282 if rhs_from_time > self.from_time {
283 self.from_time = rhs_from_time;
284 self.value = <Option<Row> as columnar::Columnar>::into_owned(rhs.value);
285 }
286 }
287}
288
289type UpsertUpdate<T, O> = (UpsertKey, T, UpsertDiff<O>);
293
294type UpsertChunk<T, O> = ColumnChunk<UpsertKey, T, UpsertDiff<O>>;
297
298type UpsertChunkBatcher<T, O> = ChunkBatcher<UpsertChunk<T, O>>;
306
307type UpsertPagedBatcher<T, O> = ColumnMergeBatcher<UpsertKey, T, UpsertDiff<O>>;
312
313type UpsertChunker<T, O> = ColumnChunker<UpsertUpdate<T, O>>;
316
317type UpsertOutputHandle<T> =
321 AsyncOutputHandle<T, FueledBuilder<CapacityContainerBuilder<Vec<(UpsertValue, T, Diff)>>>>;
322
323#[doc(hidden)]
336pub fn upsert_value_to_row(value: &UpsertValue) -> Row {
337 let mut row = Row::default();
338 let mut packer = row.packer();
339 match value {
340 Ok(ok) => {
341 packer.push(Datum::UInt8(0));
342 packer.extend(ok.iter());
343 }
344 Err(err) => {
345 packer.push(Datum::UInt8(1));
346 let bytes =
347 bincode::serialize(err.as_ref()).expect("UpsertError is serializable via bincode");
348 packer.push(Datum::Bytes(&bytes));
349 }
350 }
351 row
352}
353
354fn upsert_value_byte_len(value: &UpsertValue) -> usize {
357 match value {
358 Ok(row) => row.byte_len(),
359 Err(err) => std::mem::size_of_val(err.as_ref()),
360 }
361}
362
363#[cfg(feature = "fuzzing")]
369#[doc(hidden)]
370pub fn datum_seq_to_upsert_value(seq: DatumSeq<'_>) -> UpsertValue {
371 decode_upsert_value(seq)
372}
373
374fn decode_upsert_value<'a>(mut iter: impl Iterator<Item = Datum<'a>>) -> UpsertValue {
377 let tag = match iter.next() {
378 Some(Datum::UInt8(tag)) => tag,
379 other => panic!("upsert value missing UInt8 tag, got {:?}", other),
380 };
381 match tag {
382 0 => {
383 let mut row = Row::default();
384 row.packer().extend(iter);
385 Ok(row)
386 }
387 1 => {
388 let bytes = match iter.next() {
389 Some(Datum::Bytes(b)) => b,
390 other => panic!("upsert error tag missing Bytes payload, got {:?}", other),
391 };
392 let err: UpsertError =
393 bincode::deserialize(bytes).expect("UpsertError bincode round-trip");
394 Err(Box::new(err))
395 }
396 tag => panic!("unknown upsert value tag {tag}"),
397 }
398}
399
400#[allow(clippy::disallowed_methods)]
415pub fn upsert_inner<'scope, T, FromTime>(
416 flavor: UpsertStashFlavor,
417 input: VecCollection<'scope, T, (UpsertKey, Option<UpsertValue>, FromTime), Diff>,
418 key_indices: Vec<usize>,
419 resume_upper: Antichain<T>,
420 persist_ok: VecCollection<'scope, T, Row, Diff>,
421 persist_err: VecCollection<'scope, T, DataflowError, Diff>,
422 persist_token: Option<Vec<PressOnDropButton>>,
423 upsert_metrics: UpsertMetrics,
424 source_config: crate::source::SourceExportCreationConfig,
425) -> (
426 VecCollection<'scope, T, Result<Row, DataflowError>, Diff>,
427 StreamVec<'scope, T, (Option<GlobalId>, HealthStatusUpdate)>,
428 StreamVec<'scope, T, Infallible>,
429 PressOnDropButton,
430)
431where
432 T: Timestamp + TotalOrder + Sync,
433 T: Refines<mz_repr::Timestamp> + differential_dataflow::lattice::Lattice,
434 T: columnation::Columnation,
435 T: columnar::Columnar + Default,
436 for<'a> columnar::Ref<'a, T>: Copy + Ord,
437 FromTime: Debug + timely::ExchangeData + Clone + Ord + Sync,
438 FromTime: UpsertSourceTime,
439{
440 let encoded = encode_feedback(
444 persist_ok,
445 persist_err,
446 key_indices,
447 source_config.source_statistics.clone(),
448 );
449 match flavor {
450 UpsertStashFlavor::Chunked => {
451 let persist_arranged = arrange_core::<
455 _,
456 _,
457 ChunkChunker<(UpsertKey, Row), T, Diff>,
458 ChunkBatcher<FeedbackChunk<T>>,
459 ChunkBuilder<FeedbackChunk<T>>,
460 FeedbackSpine<T>,
461 >(encoded, Pipeline, "Persist feedback");
462 build_upsert_operator::<ChunkedArm, _, _>(
463 input,
464 resume_upper,
465 persist_arranged,
466 persist_token,
467 upsert_metrics,
468 source_config,
469 )
470 }
471 UpsertStashFlavor::Paged => {
472 let persist_arranged = arrange_core::<
478 _,
479 _,
480 PagedChunker<((UpsertKey, Row), T, Diff)>,
481 UpsertFeedbackBatcher<T>,
482 ValRowColPagedBuilder<UpsertKey, T, Diff>,
483 FundedValRowSpine<UpsertKey, T, Diff>,
484 >(encoded, Pipeline, "Persist feedback");
485 build_upsert_operator::<PagedArm, _, _>(
486 input,
487 resume_upper,
488 persist_arranged,
489 persist_token,
490 upsert_metrics,
491 source_config,
492 )
493 }
494 }
495}
496
497fn encode_feedback<'scope, T>(
503 persist_ok: VecCollection<'scope, T, Row, Diff>,
504 persist_err: VecCollection<'scope, T, DataflowError, Diff>,
505 key_indices: Vec<usize>,
506 source_statistics: SourceStatistics,
507) -> Stream<'scope, T, Column<((UpsertKey, Row), T, Diff)>>
508where
509 T: Timestamp + TotalOrder + Sync,
510 T: Refines<mz_repr::Timestamp> + differential_dataflow::lattice::Lattice,
511 T: columnation::Columnation,
512 T: columnar::Columnar + Default,
513 for<'a> columnar::Ref<'a, T>: Copy + Ord,
514{
515 let persist_keyed = crate::upsert::key_persist_feedback(persist_ok, persist_err, key_indices);
516 let persist_keyed = persist_keyed
517 .inner
518 .exchange(move |((key, _), _, _)| UpsertKey::hashed(key))
521 .as_collection()
522 .inspect(move |((_, row), _, diff)| {
523 source_statistics.update_records_indexed_by(diff.into_inner());
524 source_statistics.update_bytes_indexed_by(
525 row.as_ref().map_or(0, |r| r.byte_len().try_into().unwrap()) * diff.into_inner(),
526 );
527 });
528 persist_keyed
529 .inner
530 .unary::<ColumnBuilder<((UpsertKey, Row), T, Diff)>, _, _, _>(
531 Pipeline,
532 "Persist feedback encode",
533 |_, _| {
534 move |input, output| {
535 input.for_each(|time, data| {
536 let mut session = output.session_with_builder(&time);
537 for ((key, value), ts, diff) in data.drain(..) {
538 let row = upsert_value_to_row(&value);
539 session.give(((&key, &row), &ts, &diff));
540 }
541 });
542 }
543 },
544 )
545}
546
547#[allow(clippy::disallowed_methods)]
555fn build_upsert_operator<'scope, A, T, FromTime>(
556 input: VecCollection<'scope, T, (UpsertKey, Option<UpsertValue>, FromTime), Diff>,
557 resume_upper: Antichain<T>,
558 persist_arranged: Arranged<'scope, TraceAgent<A::Spine>>,
559 persist_token: Option<Vec<PressOnDropButton>>,
560 upsert_metrics: UpsertMetrics,
561 source_config: crate::source::SourceExportCreationConfig,
562) -> (
563 VecCollection<'scope, T, Result<Row, DataflowError>, Diff>,
564 StreamVec<'scope, T, (Option<GlobalId>, HealthStatusUpdate)>,
565 StreamVec<'scope, T, Infallible>,
566 PressOnDropButton,
567)
568where
569 A: UpsertStashArm<T, FromTime::Order>,
570 T: Timestamp + TotalOrder + Sync,
571 T: Refines<mz_repr::Timestamp> + differential_dataflow::lattice::Lattice,
572 T: columnation::Columnation,
573 T: columnar::Columnar + Default,
574 for<'a> columnar::Ref<'a, T>: Copy + Ord,
575 FromTime: Debug + timely::ExchangeData + Clone + Ord + Sync,
576 FromTime: UpsertSourceTime,
577{
578 let mut persist_trace = persist_arranged.trace.clone();
579
580 use timely::dataflow::operators::Probe;
584 let (persist_probe, _persist_probe_stream) = persist_arranged.stream.probe();
585
586 let mut builder = AsyncOperatorBuilder::new("Upsert V2".to_string(), input.scope());
588
589 let (output_handle, output) = builder
590 .new_output::<FueledBuilder<CapacityContainerBuilder<Vec<(UpsertValue, T, Diff)>>>>();
591 let (_snapshot_handle, snapshot_stream) =
592 builder.new_output::<CapacityContainerBuilder<Vec<Infallible>>>();
593 let (_health_output, health_stream) = builder
594 .new_output::<CapacityContainerBuilder<Vec<(Option<GlobalId>, HealthStatusUpdate)>>>();
595
596 let mut input = builder.new_input_for(
597 input.inner,
598 Exchange::new(move |((key, _, _), _, _)| UpsertKey::hashed(key)),
599 &output_handle,
600 );
601
602 let mut persist_wakeup = builder.new_disconnected_input(_persist_probe_stream, Pipeline);
606
607 let shutdown_button = builder.build(move |caps| async move {
608 let _persist_token = persist_token;
611
612 let [output_cap, snapshot_cap, _health_cap]: [_; 3] = caps.try_into().unwrap();
613 drop(output_cap);
614 let mut snapshot_cap = CapabilitySet::from_elem(snapshot_cap);
615
616 let mut hydrating = true;
617
618 let mut batcher = A::new_batcher();
624 let mut chunker: UpsertChunker<T, FromTime::Order> = Default::default();
627 let mut push_buffer: Vec<UpsertUpdate<T, FromTime::Order>> = Vec::new();
630
631 let mut stash_cap: Option<Capability<T>> = None;
635 let mut input_upper = Antichain::from_elem(Timestamp::minimum());
636
637 let snapshot_start = std::time::Instant::now();
638 let mut prev_persist_upper = Antichain::from_elem(Timestamp::minimum());
639
640 let mut rehydration_total: u64 = 0;
642 let mut rehydration_updates: u64 = 0;
643
644 loop {
650 tokio::select! {
652 _ = input.ready() => {}
653 _ = persist_wakeup.ready() => {
654 while let Some(event) = persist_wakeup.next_sync() {
655 if let AsyncEvent::Data(_, batches) = event {
656 for batch in batches {
657 mz_timely_util::columnar::chunk::metrics::record_batch(batch.len());
658 }
659 }
660 }
661 }
662 }
663
664 while let Some(event) = input.next_sync() {
669 match event {
670 AsyncEvent::Data(cap, data) => {
671 let mut pushed_any = false;
672 for ((key, value, from_time), ts, diff) in data {
673 assert!(diff.is_positive(), "invalid upsert input");
674 if PartialOrder::less_equal(&input_upper, &resume_upper)
675 && !resume_upper.less_equal(&ts)
676 {
677 continue;
678 }
679 let value = value.as_ref().map(upsert_value_to_row);
680 let from_time = from_time.upsert_order();
681 push_buffer.push((key, ts, UpsertDiff { from_time, value }));
682 pushed_any = true;
683 }
684 if pushed_any {
687 stash_cap = Some(match stash_cap {
688 Some(prev) if cap.time() < prev.time() => cap,
689 Some(prev) => prev,
690 None => cap,
691 });
692 }
693 }
694 AsyncEvent::Progress(upper) => {
695 if PartialOrder::less_than(&upper, &resume_upper) {
696 continue;
697 }
698 input_upper = upper;
699 }
700 }
701 }
702
703 A::flush(&mut push_buffer, &mut chunker, &mut batcher);
707
708 let persist_upper = persist_probe.with_frontier(|f| f.to_owned());
717
718 if persist_upper != prev_persist_upper {
719 let last_rehydration_chunk =
720 hydrating && PartialOrder::less_equal(&resume_upper, &persist_upper);
721
722 if last_rehydration_chunk {
723 hydrating = false;
724 upsert_metrics
725 .rehydration_latency
726 .set(snapshot_start.elapsed().as_secs_f64());
727 upsert_metrics.rehydration_total.set(rehydration_total);
728 upsert_metrics.rehydration_updates.set(rehydration_updates);
729 tracing::info!(
730 worker_id = %source_config.worker_id,
731 source_id = %source_config.id,
732 "upsert finished rehydration",
733 );
734 snapshot_cap.downgrade(&[]);
735 }
736
737 let _ = snapshot_cap.try_downgrade(persist_upper.iter());
738
739 persist_trace.set_logical_compaction(persist_upper.borrow());
741 persist_trace.set_physical_compaction(persist_upper.borrow());
742
743 prev_persist_upper = persist_upper.clone();
744 }
745
746 if let Some(cap) = stash_cap.as_mut()
783 && !persist_upper.less_than(cap.time())
784 && PartialOrder::less_than(&persist_upper, &input_upper)
785 {
786 let (sealed, _description) = batcher.seal(input_upper.clone());
790 let remaining_frontier = batcher.frontier().to_owned();
792
793 let mut ineligible = Vec::new();
794 let drain_stats = A::drain(
798 sealed,
799 &mut ineligible,
800 &output_handle,
801 &*cap,
802 &persist_upper,
803 &mut persist_trace,
804 source_config.worker_id,
805 source_config.id,
806 )
807 .await;
808
809 upsert_metrics.multi_get_size.inc_by(drain_stats.eligible);
810 upsert_metrics
811 .multi_get_result_count
812 .inc_by(drain_stats.result_count);
813 upsert_metrics
814 .multi_put_size
815 .inc_by(drain_stats.output_count);
816 upsert_metrics.upsert_inserts.inc_by(drain_stats.inserts);
817 upsert_metrics.upsert_updates.inc_by(drain_stats.updates);
818 upsert_metrics.upsert_deletes.inc_by(drain_stats.deletes);
819
820 if hydrating {
821 rehydration_total += drain_stats.inserts;
822 rehydration_updates += drain_stats.eligible;
823 }
824
825 let min_ineligible_ts = ineligible.iter().map(|(_, ts, _)| ts).min().cloned();
830 A::flush(&mut ineligible, &mut chunker, &mut batcher);
831
832 let min_ts = remaining_frontier
835 .elements()
836 .first()
837 .into_iter()
838 .chain(min_ineligible_ts.as_ref())
839 .min();
840 match min_ts {
841 Some(min_ts) => cap.downgrade(min_ts),
842 None => stash_cap = None,
845 }
846 }
847
848 if input_upper.is_empty() {
849 break;
850 }
851 }
852 });
853
854 (
855 output
856 .as_collection()
857 .map(|result: UpsertValue| match result {
858 Ok(ok) => Ok(ok),
859 Err(err) => Err(DataflowError::from(EnvelopeError::Upsert(*err))),
860 }),
861 health_stream,
862 snapshot_stream,
863 shutdown_button.press_on_drop(),
864 )
865}
866
867trait UpsertStashArm<T, O>
872where
873 T: Timestamp + Lattice + columnar::Columnar + Default,
874 for<'a> columnar::Ref<'a, T>: Copy + Ord,
875 O: columnar::Columnar + Default + Ord + Clone + Send + Sync + 'static,
876 for<'a> columnar::Ref<'a, O>: Ord + Copy,
877{
878 type Spine: TraceReader<Time = T> + 'static;
881 type Batcher: Batcher<Time = T> + 'static;
884
885 fn new_batcher() -> Self::Batcher;
887
888 fn push_chunk(batcher: &mut Self::Batcher, chunk: ColumnBody<UpsertUpdate<T, O>>);
891
892 fn flush(
897 updates: &mut Vec<UpsertUpdate<T, O>>,
898 chunker: &mut UpsertChunker<T, O>,
899 batcher: &mut Self::Batcher,
900 ) {
901 use timely::container::{ContainerBuilder as _, PushInto as _};
902 if updates.is_empty() {
903 return;
904 }
905 let mut raw: Column<UpsertUpdate<T, O>> = Default::default();
906 for update in updates.drain(..) {
907 raw.push_into(&update);
908 }
909 chunker.push_into(&mut raw);
910 while let Some(chunk) = chunker.extract() {
911 Self::push_chunk(batcher, std::mem::take(chunk));
912 }
913 }
914
915 async fn drain(
918 sealed: Vec<<Self::Batcher as Batcher>::Output>,
919 ineligible: &mut Vec<UpsertUpdate<T, O>>,
920 output_handle: &UpsertOutputHandle<T>,
921 output_cap: &Capability<T>,
922 persist_upper: &Antichain<T>,
923 trace: &mut TraceAgent<Self::Spine>,
924 worker_id: usize,
925 source_id: GlobalId,
926 ) -> DrainStats;
927}
928
929struct ChunkedArm;
941
942impl<T, O> UpsertStashArm<T, O> for ChunkedArm
943where
944 T: Timestamp + TotalOrder + Lattice + Sync,
945 T: columnation::Columnation + columnar::Columnar + Default,
946 for<'a> columnar::Ref<'a, T>: Copy + Ord,
947 O: columnar::Columnar + Default + Ord + Clone + Send + Sync + 'static,
948 for<'a> columnar::Ref<'a, O>: Ord + Copy,
949{
950 type Spine = FeedbackSpine<T>;
951 type Batcher = UpsertChunkBatcher<T, O>;
952
953 fn new_batcher() -> Self::Batcher {
954 Batcher::new(None, 0)
955 }
956
957 fn push_chunk(batcher: &mut Self::Batcher, chunk: ColumnBody<UpsertUpdate<T, O>>) {
958 batcher.push_into(ColumnChunk::from_body(chunk));
959 }
960
961 async fn drain(
962 sealed: Vec<UpsertChunk<T, O>>,
963 ineligible: &mut Vec<UpsertUpdate<T, O>>,
964 output_handle: &UpsertOutputHandle<T>,
965 output_cap: &Capability<T>,
966 persist_upper: &Antichain<T>,
967 trace: &mut TraceAgent<Self::Spine>,
968 worker_id: usize,
969 source_id: GlobalId,
970 ) -> DrainStats {
971 drain_sealed_input_chunked(
972 sealed.into_iter().map(ColumnChunk::into_body),
973 ineligible,
974 output_handle,
975 output_cap,
976 persist_upper,
977 trace,
978 worker_id,
979 source_id,
980 )
981 .await
982 }
983}
984
985struct PagedArm;
989
990impl<T, O> UpsertStashArm<T, O> for PagedArm
991where
992 T: Timestamp + TotalOrder + Lattice + Sync,
993 T: columnation::Columnation + columnar::Columnar + Default,
994 for<'a> columnar::Ref<'a, T>: Copy + Ord,
995 O: columnar::Columnar + Default + Ord + Clone + Send + Sync + 'static,
996 for<'a> columnar::Ref<'a, O>: Ord + Copy,
997{
998 type Spine = FundedValRowSpine<UpsertKey, T, Diff>;
999 type Batcher = UpsertPagedBatcher<T, O>;
1000
1001 fn new_batcher() -> Self::Batcher {
1002 let mut batcher: UpsertPagedBatcher<T, O> = Batcher::new(None, 0);
1003 batcher.set_pager(crate::upsert::upsert_stash_pager::pager());
1004 batcher
1005 }
1006
1007 fn push_chunk(batcher: &mut Self::Batcher, chunk: ColumnBody<UpsertUpdate<T, O>>) {
1008 batcher.push_into(Column::from(chunk));
1010 }
1011
1012 async fn drain(
1013 sealed: Vec<Column<UpsertUpdate<T, O>>>,
1014 ineligible: &mut Vec<UpsertUpdate<T, O>>,
1015 output_handle: &UpsertOutputHandle<T>,
1016 output_cap: &Capability<T>,
1017 persist_upper: &Antichain<T>,
1018 trace: &mut TraceAgent<Self::Spine>,
1019 worker_id: usize,
1020 source_id: GlobalId,
1021 ) -> DrainStats {
1022 drain_sealed_input_paged(
1023 sealed,
1024 ineligible,
1025 output_handle,
1026 output_cap,
1027 persist_upper,
1028 trace,
1029 worker_id,
1030 source_id,
1031 )
1032 .await
1033 }
1034}
1035
1036enum TimeClass {
1038 AlreadyPersisted,
1040 Eligible,
1042 Ineligible,
1044}
1045
1046fn classify_time<T: PartialOrder>(persist_upper: &Antichain<T>, ts: &T) -> TimeClass {
1053 if !persist_upper.less_equal(ts) {
1054 TimeClass::AlreadyPersisted
1055 } else if persist_upper.less_than(ts) {
1056 TimeClass::Ineligible
1057 } else {
1058 TimeClass::Eligible
1059 }
1060}
1061
1062struct DrainStats {
1065 eligible: u64,
1067 result_count: u64,
1069 inserts: u64,
1071 updates: u64,
1073 deletes: u64,
1075 output_count: u64,
1077}
1078
1079async fn drain_sealed_input_chunked<T, O>(
1105 sealed: impl Iterator<Item = ColumnBody<UpsertUpdate<T, O>>>,
1106 ineligible: &mut Vec<UpsertUpdate<T, O>>,
1107 output_handle: &UpsertOutputHandle<T>,
1108 output_cap: &Capability<T>,
1109 persist_upper: &Antichain<T>,
1110 trace: &mut TraceAgent<FeedbackSpine<T>>,
1111 worker_id: usize,
1112 source_id: GlobalId,
1113) -> DrainStats
1114where
1115 T: Timestamp + TotalOrder + Lattice + Sync,
1116 T: columnation::Columnation + columnar::Columnar + Default,
1117 for<'a> columnar::Ref<'a, T>: Copy + Ord,
1118 O: columnar::Columnar,
1119{
1120 let mut eligible_count: u64 = 0;
1121 let mut result_count: u64 = 0;
1122 let mut output_count: u64 = 0;
1123 let mut inserts: u64 = 0;
1124 let mut updates: u64 = 0;
1125 let mut deletes: u64 = 0;
1126
1127 let batches = trace
1131 .batches_through(Antichain::new().borrow())
1132 .expect("complete batch set for the feedback trace; is it closed?");
1133
1134 const PROBE_WINDOW: usize = 1024;
1139
1140 for chunk in sealed {
1141 use columnar::{Index, Len};
1142 let view = chunk.borrow();
1143 let total = view.len();
1144 let mut start = 0;
1145 while start < total {
1146 let mut probe_col = <UpsertKey as columnar::Columnar>::Container::default();
1152 let mut probe_count = 0usize;
1153 let mut end = total;
1154 {
1155 use columnar::Push;
1156 let mut last_probe: Option<&UpsertKey> = None;
1157 for index in start..total {
1158 let (key, ts, _diff) = view.get(index);
1159 let ts = <T as columnar::Columnar>::into_owned(ts);
1160 if matches!(classify_time(persist_upper, &ts), TimeClass::Eligible) {
1161 if last_probe != Some(key) {
1162 if probe_count == PROBE_WINDOW {
1163 end = index;
1164 break;
1165 }
1166 probe_col.push(key);
1167 probe_count += 1;
1168 last_probe = Some(key);
1169 }
1170 }
1171 }
1172 }
1173
1174 let mut old_values: std::collections::BTreeMap<UpsertKey, UpsertValue> =
1180 std::collections::BTreeMap::new();
1181 if probe_count > 0 {
1182 use columnar::Borrow;
1183 let mut staging = <FeedbackUpdate<T> as columnar::Columnar>::Container::default();
1184 for batch in &batches {
1185 batch.extract_into(probe_col.borrow(), &mut staging);
1186 }
1187 let staged = staging.borrow();
1188 let mut hits: Vec<_> = (0..staged.len())
1189 .map(|i| {
1190 let ((key, val), _time, diff) = staged.get(i);
1191 (key, val, <Diff as columnar::Columnar>::into_owned(diff))
1192 })
1193 .collect();
1194 hits.sort_by(|a, b| (a.0, a.1).cmp(&(b.0, b.1)));
1195 let mut i = 0;
1196 while i < hits.len() {
1197 let (key, val, _) = hits[i];
1198 let mut count = Diff::ZERO;
1199 let mut j = i;
1200 while j < hits.len() && hits[j].0 == key && hits[j].1 == val {
1201 count += hits[j].2;
1202 j += 1;
1203 }
1204 if count.is_positive() {
1205 assert!(
1206 count == 1.into(),
1207 "unexpected multiple entries for the same key in persist trace"
1208 );
1209 let prev = old_values.insert(*key, decode_upsert_value(val.iter()));
1210 assert!(
1211 prev.is_none(),
1212 "unexpected multiple values for the same key in persist trace"
1213 );
1214 }
1215 i = j;
1216 }
1217 }
1218
1219 for index in start..end {
1221 let (key, ts, diff) = view.get(index);
1222 let ts = <T as columnar::Columnar>::into_owned(ts);
1223 match classify_time(persist_upper, &ts) {
1224 TimeClass::AlreadyPersisted => continue,
1225 TimeClass::Ineligible => {
1226 ineligible.push((
1228 *key,
1229 ts,
1230 <UpsertDiff<O> as columnar::Columnar>::into_owned(diff),
1231 ));
1232 continue;
1233 }
1234 TimeClass::Eligible => {}
1235 }
1236
1237 eligible_count += 1;
1241 let old_value = old_values.remove(key);
1242
1243 if old_value.is_some() {
1244 result_count += 1;
1245 }
1246
1247 match diff.value {
1248 Some(row) => {
1249 if let Some(old_val) = old_value {
1250 let size = upsert_value_byte_len(&old_val);
1251 output_handle
1252 .give_fueled(
1253 output_cap,
1254 (old_val, ts.clone(), Diff::MINUS_ONE),
1255 size,
1256 )
1257 .await;
1258 output_count += 1;
1259 updates += 1;
1260 } else {
1261 inserts += 1;
1262 }
1263 let new_val = decode_upsert_value(row.iter());
1264 let size = upsert_value_byte_len(&new_val);
1265 output_handle
1266 .give_fueled(output_cap, (new_val, ts, Diff::ONE), size)
1267 .await;
1268 output_count += 1;
1269 }
1270 None => {
1271 if let Some(old_val) = old_value {
1272 let size = upsert_value_byte_len(&old_val);
1273 output_handle
1274 .give_fueled(output_cap, (old_val, ts, Diff::MINUS_ONE), size)
1275 .await;
1276 output_count += 1;
1277 deletes += 1;
1278 }
1279 }
1280 }
1281 }
1282 start = end;
1283 }
1284 }
1285
1286 tracing::debug!(
1287 worker_id = %worker_id,
1288 source_id = %source_id,
1289 ineligible = ineligible.len(),
1290 eligible = eligible_count,
1291 "drained stash",
1292 );
1293
1294 DrainStats {
1295 eligible: eligible_count,
1296 result_count,
1297 inserts,
1298 updates,
1299 deletes,
1300 output_count,
1301 }
1302}
1303
1304async fn drain_sealed_input_paged<T, O>(
1315 sealed: Vec<Column<UpsertUpdate<T, O>>>,
1316 ineligible: &mut Vec<UpsertUpdate<T, O>>,
1317 output_handle: &UpsertOutputHandle<T>,
1318 output_cap: &Capability<T>,
1319 persist_upper: &Antichain<T>,
1320 trace: &mut TraceAgent<FundedValRowSpine<UpsertKey, T, Diff>>,
1321 worker_id: usize,
1322 source_id: GlobalId,
1323) -> DrainStats
1324where
1325 T: Timestamp + TotalOrder + Lattice + Sync,
1326 T: columnation::Columnation + columnar::Columnar,
1327 O: columnar::Columnar,
1328{
1329 use columnar::Index as _;
1330
1331 let mut eligible_count: u64 = 0;
1346 let mut result_count: u64 = 0;
1347 let mut output_count: u64 = 0;
1348 let mut inserts: u64 = 0;
1349 let mut updates: u64 = 0;
1350 let mut deletes: u64 = 0;
1351
1352 let (mut cursor, storage) = trace.cursor();
1353
1354 for chunk in &sealed {
1355 for (key, ts, diff) in chunk.borrow().into_index_iter() {
1356 let ts = <T as columnar::Columnar>::into_owned(ts);
1357 match classify_time(persist_upper, &ts) {
1358 TimeClass::AlreadyPersisted => continue,
1359 TimeClass::Ineligible => {
1360 ineligible.push((
1362 *key,
1363 ts,
1364 <UpsertDiff<O> as columnar::Columnar>::into_owned(diff),
1365 ));
1366 continue;
1367 }
1368 TimeClass::Eligible => {}
1369 }
1370
1371 eligible_count += 1;
1376 cursor.seek_key(&storage, key);
1377 let old_value = match cursor.get_key(&storage) {
1378 Some(found) if found == key => {
1379 let mut result = None;
1380 while let Some(val) = cursor.get_val(&storage) {
1381 let mut count = Diff::ZERO;
1382 cursor.map_times(&storage, |_time, d| {
1383 count += d.clone();
1384 });
1385 if count.is_positive() {
1386 assert!(
1387 count == 1.into(),
1388 "unexpected multiple entries for the same key in persist trace"
1389 );
1390 assert!(
1391 result.is_none(),
1392 "unexpected multiple values for the same key in persist trace"
1393 );
1394 result = Some(decode_upsert_value(val));
1395 }
1396 cursor.step_val(&storage);
1397 }
1398 result
1399 }
1400 _ => None,
1401 };
1402
1403 if old_value.is_some() {
1404 result_count += 1;
1405 }
1406
1407 match diff.value {
1408 Some(row) => {
1409 if let Some(old_val) = old_value {
1410 let size = upsert_value_byte_len(&old_val);
1411 output_handle
1412 .give_fueled(output_cap, (old_val, ts.clone(), Diff::MINUS_ONE), size)
1413 .await;
1414 output_count += 1;
1415 updates += 1;
1416 } else {
1417 inserts += 1;
1418 }
1419 let new_val = decode_upsert_value(row.iter());
1420 let size = upsert_value_byte_len(&new_val);
1421 output_handle
1422 .give_fueled(output_cap, (new_val, ts, Diff::ONE), size)
1423 .await;
1424 output_count += 1;
1425 }
1426 None => {
1427 if let Some(old_val) = old_value {
1428 let size = upsert_value_byte_len(&old_val);
1429 output_handle
1430 .give_fueled(output_cap, (old_val, ts, Diff::MINUS_ONE), size)
1431 .await;
1432 output_count += 1;
1433 deletes += 1;
1434 }
1435 }
1436 }
1437 }
1438 }
1439
1440 tracing::debug!(
1441 worker_id = %worker_id,
1442 source_id = %source_id,
1443 ineligible = ineligible.len(),
1444 eligible = eligible_count,
1445 "drained stash",
1446 );
1447
1448 DrainStats {
1449 eligible: eligible_count,
1450 result_count,
1451 inserts,
1452 updates,
1453 deletes,
1454 output_count,
1455 }
1456}
1457
1458#[cfg(test)]
1459mod test {
1460 use mz_ore::metrics::MetricsRegistry;
1464 use mz_persist_types::ShardId;
1465 use mz_repr::{Datum, Timestamp as MzTimestamp};
1466 use mz_storage_operators::persist_source::Subtime;
1467 use mz_storage_types::sources::SourceEnvelope;
1468 use mz_storage_types::sources::envelope::{KeyEnvelope, UpsertEnvelope, UpsertStyle};
1469 use timely::dataflow::operators::capture::Extract;
1470 use timely::dataflow::operators::{Capture, Input};
1471 use timely::progress::Timestamp;
1472
1473 use crate::metrics::StorageMetrics;
1474 use crate::metrics::upsert::UpsertMetricDefs;
1475 use crate::source::SourceExportCreationConfig;
1476 use crate::statistics::{SourceStatistics, SourceStatisticsMetricDefs};
1477
1478 use super::*;
1479
1480 impl UpsertSourceTime for i32 {
1483 type Order = i32;
1484 fn upsert_order(&self) -> i32 {
1485 *self
1486 }
1487 }
1488
1489 type Ts = (MzTimestamp, Subtime);
1490
1491 fn new_ts(ts: u64) -> Ts {
1492 (MzTimestamp::new(ts), Subtime::minimum())
1493 }
1494
1495 fn key(k: i64) -> UpsertKey {
1496 UpsertKey::from_key(Ok(&Row::pack_slice(&[Datum::Int64(k)])))
1497 }
1498
1499 fn row(k: i64, v: i64) -> Row {
1500 Row::pack_slice(&[Datum::Int64(k), Datum::Int64(v)])
1501 }
1502
1503 macro_rules! upsert_test {
1508 (|$input:ident, $persist:ident, $worker:ident| $body:block) => {{
1509 let run = |flavor: UpsertStashFlavor| {
1510 let output_handle = timely::execute_directly(move |$worker| {
1511 let (mut $input, mut $persist, output_handle) = $worker
1512 .dataflow::<MzTimestamp, _, _>(|scope| {
1513 scope.scoped::<Ts, _, _>("upsert", |scope| {
1514 let (input_handle, input) = scope.new_input();
1515 let (persist_handle, persist_input) = scope.new_input();
1516 let (_persist_err_handle, persist_err_input) = scope.new_input();
1517 let source_id = GlobalId::User(0);
1518
1519 let reg = MetricsRegistry::new();
1520 let upsert_defs = UpsertMetricDefs::register_with(®);
1521 let upsert_metrics =
1522 UpsertMetrics::new(&upsert_defs, source_id, 0, None);
1523
1524 let reg2 = MetricsRegistry::new();
1525 let storage_metrics = StorageMetrics::register_with(®2);
1526
1527 let reg3 = MetricsRegistry::new();
1528 let stats_defs =
1529 SourceStatisticsMetricDefs::register_with(®3);
1530 let envelope = SourceEnvelope::Upsert(UpsertEnvelope {
1531 source_arity: 2,
1532 style: UpsertStyle::Default(KeyEnvelope::Flattened),
1533 key_indices: vec![0],
1534 });
1535 let source_statistics = SourceStatistics::new(
1536 source_id, 0, &stats_defs, source_id, &ShardId::new(),
1537 envelope, Antichain::from_elem(Timestamp::minimum()),
1538 );
1539 let source_config = SourceExportCreationConfig {
1540 id: source_id,
1541 worker_id: 0,
1542 metrics: storage_metrics,
1543 source_statistics,
1544 };
1545
1546 let (output, _, _, button) = upsert_inner(
1547 flavor,
1548 input.as_collection(),
1549 vec![0],
1550 Antichain::from_elem(Timestamp::minimum()),
1551 persist_input.as_collection(),
1552 persist_err_input.as_collection(),
1553 None,
1554 upsert_metrics,
1555 source_config,
1556 );
1557 std::mem::forget(button);
1558 (input_handle, persist_handle, output.inner.capture())
1559 })
1560 });
1561
1562 $body
1563
1564 output_handle
1565 });
1566
1567 let mut actual: Vec<_> = output_handle
1568 .extract()
1569 .into_iter()
1570 .flat_map(|(_cap, container)| container)
1571 .collect();
1572 differential_dataflow::consolidation::consolidate_updates(&mut actual);
1573 actual
1574 };
1575
1576 let paged = run(UpsertStashFlavor::Paged);
1577 let chunked = run(UpsertStashFlavor::Chunked);
1578 assert_eq!(paged, chunked, "stash flavors must produce equal output");
1579 chunked
1580 }};
1581 }
1582
1583 #[mz_ore::test]
1584 #[cfg_attr(miri, ignore)]
1585 fn gh_9160_repro() {
1586 let actual = upsert_test!(|input, persist, worker| {
1587 let key0 = key(0);
1588 let key1 = key(1);
1589 let value1 = row(0, 0);
1590 let value3 = row(0, 1);
1591 let value4 = row(0, 2);
1592
1593 input.send(((key0, Some(Ok(value1.clone())), 1), new_ts(0), Diff::ONE));
1594 input.advance_to(new_ts(2));
1595 worker.step();
1596
1597 persist.send((value1, new_ts(0), Diff::ONE));
1598 persist.advance_to(new_ts(1));
1599 worker.step();
1600
1601 input.send_batch(&mut vec![
1602 ((key1, None, 2), new_ts(2), Diff::ONE),
1603 ((key0, Some(Ok(value3)), 3), new_ts(3), Diff::ONE),
1604 ]);
1605 input.advance_to(new_ts(3));
1606 input.send_batch(&mut vec![(
1607 (key0, Some(Ok(value4)), 4),
1608 new_ts(3),
1609 Diff::ONE,
1610 )]);
1611 input.advance_to(new_ts(4));
1612 worker.step();
1613
1614 persist.advance_to(new_ts(3));
1615 worker.step();
1616 });
1617
1618 let value1 = row(0, 0);
1619 let value4 = row(0, 2);
1620 let expected: Vec<(Result<Row, DataflowError>, _, _)> = vec![
1621 (Ok(value1.clone()), new_ts(0), Diff::ONE),
1622 (Ok(value1), new_ts(3), Diff::MINUS_ONE),
1623 (Ok(value4), new_ts(3), Diff::ONE),
1624 ];
1625 assert_eq!(actual, expected);
1626 }
1627
1628 #[mz_ore::test]
1629 #[cfg_attr(miri, ignore)]
1630 fn out_of_order_keys_across_timestamps() {
1631 let actual = upsert_test!(|input, persist, worker| {
1632 let key_high = key(99);
1633 let key_low = key(1);
1634 let val_a = row(99, 1);
1635 let val_b = row(1, 2);
1636
1637 input.send(((key_high, Some(Ok(val_a.clone())), 1), new_ts(0), Diff::ONE));
1638 input.advance_to(new_ts(1));
1639 worker.step();
1640 persist.send((val_a.clone(), new_ts(0), Diff::ONE));
1641 persist.advance_to(new_ts(1));
1642 worker.step();
1643
1644 input.send(((key_low, Some(Ok(val_b.clone())), 2), new_ts(1), Diff::ONE));
1645 input.advance_to(new_ts(2));
1646 worker.step();
1647 persist.send((val_b.clone(), new_ts(1), Diff::ONE));
1648 persist.advance_to(new_ts(2));
1649 worker.step();
1650
1651 let val_a2 = row(99, 10);
1652 let val_b2 = row(1, 20);
1653 input.send_batch(&mut vec![
1654 (
1655 (key_high, Some(Ok(val_a2.clone())), 3),
1656 new_ts(2),
1657 Diff::ONE,
1658 ),
1659 ((key_low, Some(Ok(val_b2.clone())), 4), new_ts(2), Diff::ONE),
1660 ]);
1661 input.advance_to(new_ts(3));
1662 worker.step();
1663 persist.advance_to(new_ts(3));
1664 worker.step();
1665 });
1666
1667 let val_a = row(99, 1);
1668 let val_b = row(1, 2);
1669 let val_a2 = row(99, 10);
1670 let val_b2 = row(1, 20);
1671 let expected: Vec<(Result<Row, DataflowError>, _, _)> = vec![
1672 (Ok(val_b.clone()), new_ts(1), Diff::ONE),
1673 (Ok(val_b), new_ts(2), Diff::MINUS_ONE),
1674 (Ok(val_b2), new_ts(2), Diff::ONE),
1675 (Ok(val_a.clone()), new_ts(0), Diff::ONE),
1676 (Ok(val_a), new_ts(2), Diff::MINUS_ONE),
1677 (Ok(val_a2), new_ts(2), Diff::ONE),
1678 ];
1679 let mut actual_sorted = actual;
1680 let mut expected_sorted = expected;
1681 actual_sorted.sort();
1682 expected_sorted.sort();
1683 assert_eq!(actual_sorted, expected_sorted);
1684 }
1685
1686 #[mz_ore::test]
1687 #[cfg_attr(miri, ignore)]
1688 fn rehydration_then_update() {
1689 let actual = upsert_test!(|input, persist, worker| {
1690 let k = key(42);
1691 let old_val = row(42, 100);
1692 let new_val = row(42, 200);
1693
1694 persist.send((old_val, new_ts(0), Diff::ONE));
1695 persist.advance_to(new_ts(1));
1696 worker.step();
1697
1698 input.send(((k, Some(Ok(new_val)), 1), new_ts(1), Diff::ONE));
1699 input.advance_to(new_ts(2));
1700 worker.step();
1701 persist.advance_to(new_ts(2));
1702 worker.step();
1703 });
1704
1705 let old_val = row(42, 100);
1706 let new_val = row(42, 200);
1707 let expected: Vec<(Result<Row, DataflowError>, _, _)> = vec![
1708 (Ok(old_val), new_ts(1), Diff::MINUS_ONE),
1709 (Ok(new_val), new_ts(1), Diff::ONE),
1710 ];
1711 assert_eq!(actual, expected);
1712 }
1713
1714 #[mz_ore::test]
1715 #[cfg_attr(miri, ignore)]
1716 fn drain_crosses_probe_window() {
1717 const KEYS: i64 = 1500;
1722 let actual = upsert_test!(|input, persist, worker| {
1723 for k in 0..KEYS {
1724 persist.send((row(k, k), new_ts(0), Diff::ONE));
1725 }
1726 persist.advance_to(new_ts(1));
1727 worker.step();
1728
1729 for k in 0..KEYS {
1730 input.send(((key(k), Some(Ok(row(k, k + 1))), 1), new_ts(1), Diff::ONE));
1731 }
1732 input.advance_to(new_ts(2));
1733 worker.step();
1734 persist.advance_to(new_ts(2));
1735 worker.step();
1736 });
1737
1738 let mut expected: Vec<(Result<Row, DataflowError>, _, _)> = Vec::new();
1739 for k in 0..KEYS {
1740 expected.push((Ok(row(k, k)), new_ts(1), Diff::MINUS_ONE));
1741 expected.push((Ok(row(k, k + 1)), new_ts(1), Diff::ONE));
1742 }
1743 let mut actual_sorted = actual;
1744 actual_sorted.sort();
1745 expected.sort();
1746 assert_eq!(actual_sorted, expected);
1747 }
1748
1749 #[mz_ore::test]
1758 #[cfg_attr(miri, ignore)]
1759 fn drain_reads_spilled_chunks() {
1760 use mz_ore::pool::Pool;
1761 use mz_timely_util::columnar::chunk::set_spill_override;
1762
1763 let pool = Pool::new().expect("pool creation");
1764 set_spill_override(Some(pool.clone()));
1765
1766 const KEYS: i64 = 1500;
1767 let actual = upsert_test!(|input, persist, worker| {
1768 for k in 0..KEYS {
1769 persist.send((row(k, k), new_ts(0), Diff::ONE));
1770 }
1771 persist.advance_to(new_ts(1));
1772 worker.step();
1773
1774 for k in 0..KEYS {
1775 input.send(((key(k), Some(Ok(row(k, k + 1))), 1), new_ts(1), Diff::ONE));
1776 }
1777 input.advance_to(new_ts(2));
1778 worker.step();
1779 persist.advance_to(new_ts(2));
1780 worker.step();
1781 });
1782
1783 set_spill_override(None);
1784 assert!(
1785 pool.stats().inserts > 0,
1786 "chunks should have spilled through the pool"
1787 );
1788
1789 let mut expected: Vec<(Result<Row, DataflowError>, _, _)> = Vec::new();
1790 for k in 0..KEYS {
1791 expected.push((Ok(row(k, k)), new_ts(1), Diff::MINUS_ONE));
1792 expected.push((Ok(row(k, k + 1)), new_ts(1), Diff::ONE));
1793 }
1794 let mut actual_sorted = actual;
1795 actual_sorted.sort();
1796 expected.sort();
1797 assert_eq!(actual_sorted, expected);
1798 }
1799
1800 #[mz_ore::test]
1801 #[cfg_attr(miri, ignore)]
1802 fn delete_existing_key() {
1803 let actual = upsert_test!(|input, persist, worker| {
1804 let k = key(7);
1805 let val = row(7, 77);
1806
1807 input.send(((k, Some(Ok(val.clone())), 1), new_ts(0), Diff::ONE));
1808 input.advance_to(new_ts(1));
1809 worker.step();
1810 persist.send((val, new_ts(0), Diff::ONE));
1811 persist.advance_to(new_ts(1));
1812 worker.step();
1813
1814 input.send(((k, None, 2), new_ts(1), Diff::ONE));
1815 input.advance_to(new_ts(2));
1816 worker.step();
1817 persist.advance_to(new_ts(2));
1818 worker.step();
1819 });
1820
1821 let val = row(7, 77);
1822 let expected: Vec<(Result<Row, DataflowError>, _, _)> = vec![
1823 (Ok(val.clone()), new_ts(0), Diff::ONE),
1824 (Ok(val), new_ts(1), Diff::MINUS_ONE),
1825 ];
1826 assert_eq!(actual, expected);
1827 }
1828
1829 #[mz_ore::test]
1830 #[cfg_attr(miri, ignore)]
1831 fn multi_batch_rehydration() {
1832 let actual = upsert_test!(|input, persist, worker| {
1833 let k = key(5);
1834 let old_val = row(5, 10);
1835 let new_val = row(5, 20);
1836 let updated_val = row(5, 30);
1837
1838 persist.send((old_val.clone(), new_ts(0), Diff::ONE));
1839 persist.send((old_val, new_ts(0), Diff::MINUS_ONE));
1840 persist.send((new_val, new_ts(0), Diff::ONE));
1841 persist.advance_to(new_ts(1));
1842 worker.step();
1843
1844 input.send(((k, Some(Ok(updated_val)), 1), new_ts(1), Diff::ONE));
1845 input.advance_to(new_ts(2));
1846 worker.step();
1847 persist.advance_to(new_ts(2));
1848 worker.step();
1849 });
1850
1851 let new_val = row(5, 20);
1852 let updated_val = row(5, 30);
1853 let expected: Vec<(Result<Row, DataflowError>, _, _)> = vec![
1854 (Ok(new_val), new_ts(1), Diff::MINUS_ONE),
1855 (Ok(updated_val), new_ts(1), Diff::ONE),
1856 ];
1857 assert_eq!(actual, expected);
1858 }
1859
1860 #[mz_ore::test]
1861 #[cfg_attr(miri, ignore)]
1862 fn delete_nonexistent_key() {
1863 let actual = upsert_test!(|input, persist, worker| {
1864 let k = key(99);
1865
1866 persist.advance_to(new_ts(1));
1867 worker.step();
1868
1869 input.send(((k, None, 1), new_ts(1), Diff::ONE));
1870 input.advance_to(new_ts(2));
1871 worker.step();
1872 persist.advance_to(new_ts(2));
1873 worker.step();
1874 });
1875
1876 assert!(actual.is_empty(), "expected empty output, got: {actual:?}");
1877 }
1878
1879 #[mz_ore::test]
1880 #[cfg_attr(miri, ignore)]
1881 fn reinsert_after_delete() {
1882 let actual = upsert_test!(|input, persist, worker| {
1883 let k = key(3);
1884 let val_a = row(3, 10);
1885 let val_b = row(3, 20);
1886
1887 input.send(((k, Some(Ok(val_a.clone())), 1), new_ts(0), Diff::ONE));
1888 input.advance_to(new_ts(1));
1889 worker.step();
1890 persist.send((val_a.clone(), new_ts(0), Diff::ONE));
1891 persist.advance_to(new_ts(1));
1892 worker.step();
1893
1894 input.send(((k, None, 2), new_ts(1), Diff::ONE));
1895 input.advance_to(new_ts(2));
1896 worker.step();
1897 persist.send((val_a, new_ts(1), Diff::MINUS_ONE));
1898 persist.advance_to(new_ts(2));
1899 worker.step();
1900
1901 input.send(((k, Some(Ok(val_b.clone())), 3), new_ts(2), Diff::ONE));
1902 input.advance_to(new_ts(3));
1903 worker.step();
1904 persist.advance_to(new_ts(3));
1905 worker.step();
1906 });
1907
1908 let val_a = row(3, 10);
1909 let val_b = row(3, 20);
1910 let mut expected: Vec<(Result<Row, DataflowError>, _, _)> = vec![
1911 (Ok(val_a.clone()), new_ts(0), Diff::ONE),
1912 (Ok(val_a), new_ts(1), Diff::MINUS_ONE),
1913 (Ok(val_b), new_ts(2), Diff::ONE),
1914 ];
1915 expected.sort();
1916 let mut actual = actual;
1917 actual.sort();
1918 assert_eq!(actual, expected);
1919 }
1920
1921 #[mz_ore::test]
1922 #[cfg_attr(miri, ignore)]
1923 fn idempotent_update() {
1924 let actual = upsert_test!(|input, persist, worker| {
1925 let k = key(11);
1926 let val = row(11, 50);
1927
1928 input.send(((k, Some(Ok(val.clone())), 1), new_ts(0), Diff::ONE));
1929 input.advance_to(new_ts(1));
1930 worker.step();
1931 persist.send((val.clone(), new_ts(0), Diff::ONE));
1932 persist.advance_to(new_ts(1));
1933 worker.step();
1934
1935 input.send(((k, Some(Ok(val.clone())), 2), new_ts(1), Diff::ONE));
1936 input.advance_to(new_ts(2));
1937 worker.step();
1938 persist.advance_to(new_ts(2));
1939 worker.step();
1940 });
1941
1942 let val = row(11, 50);
1943 let expected: Vec<(Result<Row, DataflowError>, _, _)> =
1944 vec![(Ok(val), new_ts(0), Diff::ONE)];
1945 assert_eq!(actual, expected);
1946 }
1947
1948 #[mz_ore::test]
1965 #[cfg_attr(miri, ignore)]
1966 fn lagging_replacement_below_upper_strands_data() {
1967 for flavor in [UpsertStashFlavor::Paged, UpsertStashFlavor::Chunked] {
1968 let (frontier, emitted) = run_below_upper_scenario_v2(flavor);
1969
1970 assert!(
1974 emitted.is_empty(),
1975 "below-upper data should be dropped, not emitted; got {emitted:?} ({flavor:?})"
1976 );
1977 assert_eq!(
1978 frontier,
1979 vec![new_ts(11)],
1980 "v2 output frontier should advance to the input upper, not pin below \
1981 persist_upper ({flavor:?})"
1982 );
1983 assert!(
1984 frontier[0] >= new_ts(10),
1985 "v2 output frontier {frontier:?} should reach at least persist_upper (10) \
1986 ({flavor:?})"
1987 );
1988 }
1989 }
1990
1991 fn run_below_upper_scenario_v2(
1994 flavor: UpsertStashFlavor,
1995 ) -> (Vec<Ts>, Vec<(Result<Row, DataflowError>, Ts, Diff)>) {
1996 use timely::dataflow::operators::Probe;
1997
1998 let (frontier, capture) = timely::execute_directly(move |worker| {
1999 let (mut input, mut persist, probe, capture) =
2000 worker.dataflow::<MzTimestamp, _, _>(|scope| {
2001 scope.scoped::<Ts, _, _>("upsert", |scope| {
2002 let (input_handle, input) = scope.new_input();
2003 let (persist_handle, persist_input) = scope.new_input();
2004 let (_persist_err_handle, persist_err_input) = scope.new_input();
2005 let source_id = GlobalId::User(0);
2006
2007 let reg = MetricsRegistry::new();
2008 let upsert_defs = UpsertMetricDefs::register_with(®);
2009 let upsert_metrics = UpsertMetrics::new(&upsert_defs, source_id, 0, None);
2010
2011 let reg2 = MetricsRegistry::new();
2012 let storage_metrics = StorageMetrics::register_with(®2);
2013
2014 let reg3 = MetricsRegistry::new();
2015 let stats_defs = SourceStatisticsMetricDefs::register_with(®3);
2016 let envelope = SourceEnvelope::Upsert(UpsertEnvelope {
2017 source_arity: 2,
2018 style: UpsertStyle::Default(KeyEnvelope::Flattened),
2019 key_indices: vec![0],
2020 });
2021 let source_statistics = SourceStatistics::new(
2022 source_id,
2023 0,
2024 &stats_defs,
2025 source_id,
2026 &ShardId::new(),
2027 envelope,
2028 Antichain::from_elem(Timestamp::minimum()),
2029 );
2030 let source_config = SourceExportCreationConfig {
2031 id: source_id,
2032 worker_id: 0,
2033 metrics: storage_metrics,
2034 source_statistics,
2035 };
2036
2037 let (output, _, _, button) = upsert_inner(
2038 flavor,
2039 input.as_collection(),
2040 vec![0],
2041 Antichain::from_elem(Timestamp::minimum()),
2042 persist_input.as_collection(),
2043 persist_err_input.as_collection(),
2044 None,
2045 upsert_metrics,
2046 source_config,
2047 );
2048 std::mem::forget(button);
2049 let (probe, stream) = output.inner.probe();
2050 (input_handle, persist_handle, probe, stream.capture())
2051 })
2052 });
2053
2054 persist.advance_to(new_ts(10));
2057 for _ in 0..20 {
2058 worker.step();
2059 }
2060
2061 input.send(((key(0), Some(Ok(row(0, 1))), 1), new_ts(5), Diff::ONE));
2064 input.send(((key(1), Some(Ok(row(1, 2))), 2), new_ts(7), Diff::ONE));
2065 input.advance_to(new_ts(11));
2066 for _ in 0..20 {
2067 worker.step();
2068 }
2069
2070 (probe.with_frontier(|f| f.to_vec()), capture)
2071 });
2072
2073 let mut emitted: Vec<_> = capture
2074 .extract()
2075 .into_iter()
2076 .flat_map(|(_cap, c)| c)
2077 .collect();
2078 differential_dataflow::consolidation::consolidate_updates(&mut emitted);
2079 (frontier, emitted)
2080 }
2081}