1use std::cell::RefCell;
11use std::cmp::Reverse;
12use std::convert::AsRef;
13use std::fmt::Debug;
14use std::hash::{Hash, Hasher};
15use std::path::PathBuf;
16use std::sync::Arc;
17
18use differential_dataflow::hashable::Hashable;
19use differential_dataflow::{AsCollection, VecCollection};
20use futures::StreamExt;
21use futures::future::FutureExt;
22use indexmap::map::Entry;
23use itertools::Itertools;
24use mz_ore::error::ErrorExt;
25use mz_repr::{Datum, DatumVec, Diff, GlobalId, Row};
26use mz_rocksdb::ValueIterator;
27use mz_sql_server_util::cdc::Lsn;
28use mz_storage_types::configuration::StorageConfiguration;
29use mz_storage_types::dyncfgs;
30use mz_storage_types::errors::{DataflowError, EnvelopeError, UpsertError};
31use mz_storage_types::sources::MzOffset;
32use mz_storage_types::sources::envelope::UpsertEnvelope;
33use mz_storage_types::sources::kafka::{KafkaTimestamp, RangeBound};
34use mz_storage_types::sources::mysql::GtidPartition;
35use mz_timely_util::builder_async::{
36 AsyncOutputHandle, Event as AsyncEvent, OperatorBuilder as AsyncOperatorBuilder,
37 PressOnDropButton,
38};
39use serde::{Deserialize, Serialize};
40use sha2::{Digest, Sha256};
41use timely::dataflow::channels::pact::Exchange;
42use timely::dataflow::operators::{Capability, InputCapability, Operator};
43use timely::dataflow::{Scope, StreamVec};
44use timely::order::{PartialOrder, TotalOrder};
45use timely::progress::timestamp::Refines;
46use timely::progress::{Antichain, Timestamp};
47
48use crate::healthcheck::HealthStatusUpdate;
49use crate::metrics::upsert::{UpsertBackpressureMetrics, UpsertMetrics};
50use crate::storage_state::StorageInstanceContext;
51use crate::{upsert_continual_feedback, upsert_continual_feedback_v2};
52use types::{
53 BincodeOpts, StateValue, UpsertState, UpsertStateBackend, consolidating_merge_function,
54 upsert_bincode_opts,
55};
56
57#[cfg(any(test, feature = "fuzzing"))]
58pub mod memory;
59pub(crate) mod rocksdb;
60pub(crate) mod types;
62
63pub type UpsertValue = Result<Row, Box<UpsertError>>;
64
65#[derive(
66 Copy,
67 Clone,
68 Hash,
69 PartialEq,
70 Eq,
71 PartialOrd,
72 Ord,
73 Serialize,
74 Deserialize,
75 bytemuck::AnyBitPattern,
76 bytemuck::NoUninit
77)]
78#[repr(transparent)]
79pub struct UpsertKey([u8; 32]);
80
81impl columnation::Columnation for UpsertKey {
82 type InnerRegion = columnation::CopyRegion<UpsertKey>;
83}
84
85mod columnar_upsert_key {
101 use super::UpsertKey;
102 use columnar::Columnar;
103 use mz_ore::cast::CastFrom;
104 use std::ops::Range;
105
106 #[derive(Clone, Copy, Default, Debug)]
108 pub struct UpsertKeys<T>(T);
109 impl<D, T: columnar::Push<D>> columnar::Push<D> for UpsertKeys<T> {
110 #[inline(always)]
111 fn push(&mut self, item: D) {
112 self.0.push(item)
113 }
114 }
115 impl<T: columnar::Clear> columnar::Clear for UpsertKeys<T> {
116 #[inline(always)]
117 fn clear(&mut self) {
118 self.0.clear()
119 }
120 }
121 impl<T: columnar::Len> columnar::Len for UpsertKeys<T> {
122 #[inline(always)]
123 fn len(&self) -> usize {
124 self.0.len()
125 }
126 }
127 impl<'a> columnar::Index for UpsertKeys<&'a [UpsertKey]> {
128 type Ref = &'a UpsertKey;
129
130 #[inline(always)]
131 fn get(&self, index: usize) -> Self::Ref {
132 &self.0[index]
133 }
134 }
135
136 impl Columnar for UpsertKey {
137 #[inline(always)]
138 fn into_owned<'a>(other: columnar::Ref<'a, Self>) -> Self {
139 *other
140 }
141 type Container = UpsertKeys<Vec<UpsertKey>>;
142 #[inline(always)]
143 fn reborrow<'b, 'a: 'b>(thing: columnar::Ref<'a, Self>) -> columnar::Ref<'b, Self>
144 where
145 Self: 'a,
146 {
147 thing
148 }
149 }
150
151 impl columnar::Borrow for UpsertKeys<Vec<UpsertKey>> {
152 type Ref<'a> = &'a UpsertKey;
153 type Borrowed<'a>
154 = UpsertKeys<&'a [UpsertKey]>
155 where
156 Self: 'a;
157 #[inline(always)]
158 fn borrow<'a>(&'a self) -> Self::Borrowed<'a> {
159 UpsertKeys(self.0.as_slice())
160 }
161 #[inline(always)]
162 fn reborrow<'b, 'a: 'b>(item: Self::Borrowed<'a>) -> Self::Borrowed<'b>
163 where
164 Self: 'a,
165 {
166 UpsertKeys(item.0)
167 }
168 #[inline(always)]
169 fn reborrow_ref<'b, 'a: 'b>(item: Self::Ref<'a>) -> Self::Ref<'b>
170 where
171 Self: 'a,
172 {
173 item
174 }
175 }
176
177 impl columnar::Container for UpsertKeys<Vec<UpsertKey>> {
178 #[inline(always)]
179 fn extend_from_self(&mut self, other: Self::Borrowed<'_>, range: Range<usize>) {
180 self.0.extend_from_self(other.0, range)
181 }
182 #[inline(always)]
183 fn reserve_for<'a, I>(&mut self, selves: I)
184 where
185 Self: 'a,
186 I: Iterator<Item = Self::Borrowed<'a>> + Clone,
187 {
188 self.0.reserve_for(selves.map(|s| s.0));
189 }
190 }
191
192 impl<'a> columnar::AsBytes<'a> for UpsertKeys<&'a [UpsertKey]> {
193 const SLICE_COUNT: usize = 1;
194 #[inline(always)]
195 fn get_byte_slice(&self, index: usize) -> (u64, &'a [u8]) {
196 mz_ore::soft_assert_no_log!(index < Self::SLICE_COUNT);
197 (
198 u64::cast_from(align_of::<UpsertKey>()),
199 bytemuck::cast_slice(self.0),
200 )
201 }
202 #[inline(always)]
203 fn as_bytes(&self) -> impl Iterator<Item = (u64, &'a [u8])> {
204 std::iter::once((
205 u64::cast_from(align_of::<UpsertKey>()),
206 bytemuck::cast_slice(self.0),
207 ))
208 }
209 }
210 impl<'a> columnar::FromBytes<'a> for UpsertKeys<&'a [UpsertKey]> {
211 const SLICE_COUNT: usize = 1;
212 #[inline(always)]
213 fn from_bytes(bytes: &mut impl Iterator<Item = &'a [u8]>) -> Self {
214 UpsertKeys(bytemuck::cast_slice(
215 bytes.next().expect("Iterator exhausted prematurely"),
216 ))
217 }
218 }
219}
220
221pub trait UpsertSourceTime
237where
238 for<'a> columnar::Ref<'a, Self::Order>: Ord,
239{
240 type Order: columnar::Columnar + Clone + Default + Ord + Send + Sync + 'static;
243 fn upsert_order(&self) -> Self::Order;
245}
246
247impl UpsertSourceTime for KafkaTimestamp {
248 type Order = (i64, u64);
254 fn upsert_order(&self) -> (i64, u64) {
255 let partition = match self.interval().lower {
256 RangeBound::NegInfinity => i64::MIN,
257 RangeBound::Elem(p, _) => i64::from(p),
258 RangeBound::PosInfinity => i64::MAX,
259 };
260 (partition, self.timestamp().offset)
261 }
262}
263
264impl UpsertSourceTime for MzOffset {
270 type Order = u64;
271 fn upsert_order(&self) -> u64 {
272 self.offset
273 }
274}
275
276macro_rules! upsert_source_time_unit {
283 ($($ty:ty),+ $(,)?) => {$(
284 impl UpsertSourceTime for $ty {
285 type Order = ();
286 fn upsert_order(&self) {
287 unreachable!(
288 "upsert source stash is not rendered for this source, but \
289 {} reached the projection",
290 std::any::type_name::<Self>(),
291 )
292 }
293 }
294 )+};
295}
296upsert_source_time_unit!(GtidPartition, Lsn);
297
298pub mod upsert_stash_spill {
312 pub fn set_enabled(enabled: bool) {
314 mz_timely_util::columnar::chunk::set_storage_spill_enabled(enabled);
315 }
316}
317
318pub mod upsert_stash_pager {
333 use std::sync::{LazyLock, RwLock};
334
335 use mz_timely_util::column_pager::{ColumnPager, shared_pager};
336
337 static PAGER: LazyLock<RwLock<ColumnPager>> =
340 LazyLock::new(|| RwLock::new(ColumnPager::disabled()));
341
342 pub fn set_enabled(enabled: bool) {
346 *PAGER.write().expect("upsert stash pager poisoned") = shared_pager(enabled);
347 }
348
349 pub fn pager() -> ColumnPager {
351 PAGER.read().expect("upsert stash pager poisoned").clone()
352 }
353}
354
355impl Debug for UpsertKey {
356 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
357 write!(f, "0x")?;
358 for byte in self.0 {
359 write!(f, "{:02x}", byte)?;
360 }
361 Ok(())
362 }
363}
364
365impl AsRef<[u8]> for UpsertKey {
366 #[inline(always)]
367 fn as_ref(&self) -> &[u8] {
371 &self.0
372 }
373}
374
375impl From<&[u8]> for UpsertKey {
376 fn from(bytes: &[u8]) -> Self {
377 UpsertKey(bytes.try_into().expect("invalid key length"))
378 }
379}
380
381type KeyHash = Sha256;
386
387impl UpsertKey {
388 pub fn from_key(key: Result<&Row, &UpsertError>) -> Self {
389 Self::from_iter(key.map(|r| r.iter()))
390 }
391
392 pub fn from_value(value: Result<&Row, &UpsertError>, key_indices: &[usize]) -> Self {
393 thread_local! {
394 static VALUE_DATUMS: RefCell<DatumVec> = RefCell::new(DatumVec::new());
396 }
397 VALUE_DATUMS.with(|value_datums| {
398 let mut value_datums = value_datums.borrow_mut();
399 let value = value.map(|v| value_datums.borrow_with(v));
400 let key = match value {
401 Ok(ref datums) => Ok(key_indices.iter().map(|&idx| datums[idx])),
402 Err(err) => Err(err),
403 };
404 Self::from_iter(key)
405 })
406 }
407
408 pub fn from_iter<'a, 'b>(
409 key: Result<impl Iterator<Item = Datum<'a>> + 'b, &UpsertError>,
410 ) -> Self {
411 thread_local! {
412 static KEY_DATUMS: RefCell<DatumVec> = RefCell::new(DatumVec::new());
414 }
415 KEY_DATUMS.with(|key_datums| {
416 let mut key_datums = key_datums.borrow_mut();
417 let mut key_datums = key_datums.borrow();
420 let key: Result<&[Datum], Datum> = match key {
421 Ok(key) => {
422 for datum in key {
423 key_datums.push(datum);
424 }
425 Ok(&*key_datums)
426 }
427 Err(UpsertError::Value(err)) => {
428 key_datums.extend(err.for_key.iter());
429 Ok(&*key_datums)
430 }
431 Err(UpsertError::KeyDecode(err)) => Err(Datum::Bytes(&err.raw)),
432 Err(UpsertError::NullKey(_)) => Err(Datum::Null),
433 };
434 let mut hasher = DigestHasher(KeyHash::new());
435 key.hash(&mut hasher);
436 Self(hasher.0.finalize().into())
437 })
438 }
439}
440
441struct DigestHasher<H: Digest>(H);
442
443impl<H: Digest> Hasher for DigestHasher<H> {
444 fn write(&mut self, bytes: &[u8]) {
445 self.0.update(bytes);
446 }
447
448 fn finish(&self) -> u64 {
449 panic!("digest wrapper used to produce a hash");
450 }
451}
452
453use std::convert::Infallible;
454use timely::container::CapacityContainerBuilder;
455use timely::dataflow::channels::pact::Pipeline;
456
457use self::types::ValueMetadata;
458
459pub fn rehydration_finished<'scope, T: Timestamp>(
463 scope: Scope<'scope, T>,
464 source_config: &crate::source::RawSourceCreationConfig,
465 token: impl std::any::Any + 'static,
467 resume_upper: Antichain<T>,
468 input: StreamVec<'scope, T, Infallible>,
469) {
470 let worker_id = source_config.worker_id;
471 let id = source_config.id;
472 let mut builder = AsyncOperatorBuilder::new(format!("rehydration_finished({id}"), scope);
473 let mut input = builder.new_disconnected_input(input, Pipeline);
474
475 builder.build(move |_capabilities| async move {
476 let mut input_upper = Antichain::from_elem(Timestamp::minimum());
477 while !PartialOrder::less_equal(&resume_upper, &input_upper) {
479 let Some(event) = input.next().await else {
480 break;
481 };
482 if let AsyncEvent::Progress(upper) = event {
483 input_upper = upper;
484 }
485 }
486 tracing::info!(
487 %worker_id,
488 source_id = %id,
489 "upsert source has downgraded past the resume upper ({resume_upper:?}) across all workers",
490 );
491 drop(token);
492 });
493}
494
495pub(crate) fn key_persist_feedback<'scope, T>(
501 ok: VecCollection<'scope, T, Row, Diff>,
502 err: VecCollection<'scope, T, DataflowError, Diff>,
503 key_indices: Vec<usize>,
504) -> VecCollection<'scope, T, (UpsertKey, UpsertValue), Diff>
505where
506 T: Timestamp,
507{
508 let keyed_ok = {
509 let key_indices = key_indices.clone();
510 ok.map(move |row| {
511 let key = UpsertKey::from_value(Ok(&row), &key_indices);
512 (key, Ok(row))
513 })
514 };
515 let keyed_err = err.flat_map(move |err| {
516 let err = match err {
517 DataflowError::EnvelopeError(err) => match *err {
518 EnvelopeError::Upsert(err) => Box::new(err),
519 EnvelopeError::Flat(_) => return None,
520 },
521 _ => return None,
522 };
523 let key = UpsertKey::from_value(Err(&err), &key_indices);
524 Some((key, Err(err)))
525 });
526 keyed_ok.concat(keyed_err)
527}
528
529pub(crate) fn upsert<'scope, T, FromTime>(
535 input: VecCollection<'scope, T, (UpsertKey, Option<UpsertValue>, FromTime), Diff>,
536 upsert_envelope: UpsertEnvelope,
537 resume_upper: Antichain<T>,
538 previous_ok: VecCollection<'scope, T, Row, Diff>,
539 previous_err: VecCollection<'scope, T, DataflowError, Diff>,
540 previous_token: Option<Vec<PressOnDropButton>>,
541 source_config: crate::source::SourceExportCreationConfig,
542 instance_context: &StorageInstanceContext,
543 storage_configuration: &StorageConfiguration,
544 dataflow_paramters: &crate::internal_control::DataflowParameters,
545 backpressure_metrics: Option<UpsertBackpressureMetrics>,
546) -> (
547 VecCollection<'scope, T, Result<Row, DataflowError>, Diff>,
548 StreamVec<'scope, T, (Option<GlobalId>, HealthStatusUpdate)>,
549 StreamVec<'scope, T, Infallible>,
550 PressOnDropButton,
551)
552where
553 T: Timestamp + TotalOrder + Sync,
554 T: Refines<mz_repr::Timestamp> + TotalOrder + Sync,
555 FromTime: Timestamp + Clone + Sync,
556{
557 let upsert_metrics = source_config.metrics.get_upsert_metrics(
558 source_config.id,
559 source_config.worker_id,
560 backpressure_metrics,
561 );
562
563 let rocksdb_cleanup_tries =
564 dyncfgs::STORAGE_ROCKSDB_CLEANUP_TRIES.get(storage_configuration.config_set());
565
566 let prevent_snapshot_buffering =
569 dyncfgs::STORAGE_UPSERT_PREVENT_SNAPSHOT_BUFFERING.get(storage_configuration.config_set());
570 let snapshot_buffering_max = dyncfgs::STORAGE_UPSERT_MAX_SNAPSHOT_BATCH_BUFFERING
572 .get(storage_configuration.config_set());
573
574 let rocksdb_use_native_merge_operator =
577 dyncfgs::STORAGE_ROCKSDB_USE_MERGE_OPERATOR.get(storage_configuration.config_set());
578
579 let upsert_config = UpsertConfig {
580 shrink_upsert_unused_buffers_by_ratio: storage_configuration
581 .parameters
582 .shrink_upsert_unused_buffers_by_ratio,
583 };
584
585 let thin_input = upsert_thinning(input);
586
587 let tuning = dataflow_paramters.upsert_rocksdb_tuning_config.clone();
588
589 let rocksdb_dir = instance_context
594 .scratch_directory
595 .clone()
596 .unwrap_or_else(|| PathBuf::from("/tmp"))
597 .join("storage")
598 .join("upsert")
599 .join(source_config.id.to_string())
600 .join(source_config.worker_id.to_string());
601
602 tracing::info!(
603 worker_id = %source_config.worker_id,
604 source_id = %source_config.id,
605 ?rocksdb_dir,
606 ?tuning,
607 ?rocksdb_use_native_merge_operator,
608 "rendering upsert source"
609 );
610
611 let rocksdb_shared_metrics = Arc::clone(&upsert_metrics.rocksdb_shared);
612 let rocksdb_instance_metrics = Arc::clone(&upsert_metrics.rocksdb_instance_metrics);
613
614 let env = instance_context
615 .rocksdb_env()
616 .expect("failed to create rocksdb env");
617
618 let rocksdb_init_fn = move || async move {
620 let merge_operator = if rocksdb_use_native_merge_operator {
621 Some((
622 "upsert_state_snapshot_merge_v1".to_string(),
623 |a: &[u8], b: ValueIterator<BincodeOpts, StateValue<T, FromTime>>| {
624 consolidating_merge_function::<T, FromTime>(a.into(), b)
625 },
626 ))
627 } else {
628 None
629 };
630 rocksdb::RocksDB::new(
631 mz_rocksdb::RocksDBInstance::new(
632 &rocksdb_dir,
633 mz_rocksdb::InstanceOptions::new(
634 env,
635 rocksdb_cleanup_tries,
636 merge_operator,
637 upsert_bincode_opts(),
640 ),
641 tuning,
642 rocksdb_shared_metrics,
643 rocksdb_instance_metrics,
644 )
645 .unwrap(),
646 )
647 };
648
649 upsert_operator(
650 thin_input,
651 upsert_envelope.key_indices,
652 resume_upper,
653 previous_ok,
654 previous_err,
655 previous_token,
656 upsert_metrics,
657 source_config,
658 rocksdb_init_fn,
659 upsert_config,
660 storage_configuration,
661 prevent_snapshot_buffering,
662 snapshot_buffering_max,
663 )
664}
665
666pub(crate) fn upsert_v2<'scope, T, FromTime>(
673 input: VecCollection<'scope, T, (UpsertKey, Option<UpsertValue>, FromTime), Diff>,
674 upsert_envelope: UpsertEnvelope,
675 resume_upper: Antichain<T>,
676 previous_ok: VecCollection<'scope, T, Row, Diff>,
677 previous_err: VecCollection<'scope, T, DataflowError, Diff>,
678 previous_token: Option<Vec<PressOnDropButton>>,
679 source_config: crate::source::SourceExportCreationConfig,
680 backpressure_metrics: Option<UpsertBackpressureMetrics>,
681 stash_flavor: upsert_continual_feedback_v2::UpsertStashFlavor,
682) -> (
683 VecCollection<'scope, T, Result<Row, DataflowError>, Diff>,
684 StreamVec<'scope, T, (Option<GlobalId>, HealthStatusUpdate)>,
685 StreamVec<'scope, T, Infallible>,
686 PressOnDropButton,
687)
688where
689 T: Timestamp + TotalOrder + Sync,
690 T: Refines<mz_repr::Timestamp> + differential_dataflow::lattice::Lattice,
691 T: columnation::Columnation,
692 T: columnar::Columnar + Default,
693 for<'a> columnar::Ref<'a, T>: Copy + Ord,
694 FromTime: Timestamp + Clone + Sync,
695 FromTime: UpsertSourceTime,
696{
697 let upsert_metrics = source_config.metrics.get_upsert_metrics(
698 source_config.id,
699 source_config.worker_id,
700 backpressure_metrics,
701 );
702
703 let thin_input = upsert_thinning(input);
704
705 tracing::info!(
706 worker_id = %source_config.worker_id,
707 source_id = %source_config.id,
708 ?stash_flavor,
709 "rendering upsert source (btreemap backend)"
710 );
711
712 upsert_continual_feedback_v2::upsert_inner(
713 stash_flavor,
714 thin_input,
715 upsert_envelope.key_indices,
716 resume_upper,
717 previous_ok,
718 previous_err,
719 previous_token,
720 upsert_metrics,
721 source_config,
722 )
723}
724
725fn upsert_operator<'scope, T, FromTime, F, Fut, US>(
728 input: VecCollection<'scope, T, (UpsertKey, Option<UpsertValue>, FromTime), Diff>,
729 key_indices: Vec<usize>,
730 resume_upper: Antichain<T>,
731 persist_ok: VecCollection<'scope, T, Row, Diff>,
732 persist_err: VecCollection<'scope, T, DataflowError, Diff>,
733 persist_token: Option<Vec<PressOnDropButton>>,
734 upsert_metrics: UpsertMetrics,
735 source_config: crate::source::SourceExportCreationConfig,
736 state: F,
737 upsert_config: UpsertConfig,
738 _storage_configuration: &StorageConfiguration,
739 prevent_snapshot_buffering: bool,
740 snapshot_buffering_max: Option<usize>,
741) -> (
742 VecCollection<'scope, T, Result<Row, DataflowError>, Diff>,
743 StreamVec<'scope, T, (Option<GlobalId>, HealthStatusUpdate)>,
744 StreamVec<'scope, T, Infallible>,
745 PressOnDropButton,
746)
747where
748 T: Timestamp + TotalOrder + Sync,
749 T: Refines<mz_repr::Timestamp> + TotalOrder + Sync,
750 F: FnOnce() -> Fut + 'static,
751 Fut: std::future::Future<Output = US>,
752 US: UpsertStateBackend<T, FromTime>,
753 FromTime: Debug + timely::ExchangeData + Clone + Ord + Sync,
754{
755 let use_continual_feedback_upsert = true;
759
760 tracing::info!(id = %source_config.id, %use_continual_feedback_upsert, "upsert operator implementation");
761
762 if use_continual_feedback_upsert {
763 upsert_continual_feedback::upsert_inner(
764 input,
765 key_indices,
766 resume_upper,
767 persist_ok,
768 persist_err,
769 persist_token,
770 upsert_metrics,
771 source_config,
772 state,
773 upsert_config,
774 prevent_snapshot_buffering,
775 snapshot_buffering_max,
776 )
777 } else {
778 upsert_classic(
779 input,
780 key_indices,
781 resume_upper,
782 persist_ok,
783 persist_err,
784 persist_token,
785 upsert_metrics,
786 source_config,
787 state,
788 upsert_config,
789 prevent_snapshot_buffering,
790 snapshot_buffering_max,
791 )
792 }
793}
794
795fn upsert_thinning<'scope, T, K, V, FromTime>(
800 input: VecCollection<'scope, T, (K, V, FromTime), Diff>,
801) -> VecCollection<'scope, T, (K, V, FromTime), Diff>
802where
803 T: Timestamp + TotalOrder,
804 K: timely::ExchangeData + Clone + Eq + Ord,
805 V: timely::ExchangeData + Clone,
806 FromTime: Timestamp,
807{
808 input
809 .inner
810 .unary(Pipeline, "UpsertThinning", |_, _| {
811 let mut capability: Option<InputCapability<T>> = None;
813 let mut updates = Vec::new();
815 move |input, output| {
816 input.for_each(|cap, data| {
817 assert!(
818 data.iter().all(|(_, _, diff)| diff.is_positive()),
819 "invalid upsert input"
820 );
821 updates.append(data);
822 match capability.as_mut() {
823 Some(capability) => {
824 if cap.time() <= capability.time() {
825 *capability = cap;
826 }
827 }
828 None => capability = Some(cap),
829 }
830 });
831 if let Some(capability) = capability.take() {
832 updates.sort_unstable_by(|a, b| {
835 let ((key1, _, from_time1), time1, _) = a;
836 let ((key2, _, from_time2), time2, _) = b;
837 Ord::cmp(
838 &(key1, time1, Reverse(from_time1)),
839 &(key2, time2, Reverse(from_time2)),
840 )
841 });
842 let mut session = output.session(&capability);
843 session.give_iterator(updates.drain(..).dedup_by(|a, b| {
844 let ((key1, _, _), time1, _) = a;
845 let ((key2, _, _), time2, _) = b;
846 (key1, time1) == (key2, time2)
847 }))
848 }
849 }
850 })
851 .as_collection()
852}
853
854fn stage_input<T, FromTime>(
857 stash: &mut Vec<(T, UpsertKey, Reverse<FromTime>, Option<UpsertValue>)>,
858 data: &mut Vec<((UpsertKey, Option<UpsertValue>, FromTime), T, Diff)>,
859 input_upper: &Antichain<T>,
860 resume_upper: &Antichain<T>,
861 storage_shrink_upsert_unused_buffers_by_ratio: usize,
862) where
863 T: PartialOrder,
864 FromTime: Ord,
865{
866 if PartialOrder::less_equal(input_upper, resume_upper) {
867 data.retain(|(_, ts, _)| resume_upper.less_equal(ts));
868 }
869
870 stash.extend(data.drain(..).map(|((key, value, order), time, diff)| {
871 assert!(diff.is_positive(), "invalid upsert input");
872 (time, key, Reverse(order), value)
873 }));
874
875 if storage_shrink_upsert_unused_buffers_by_ratio > 0 {
876 let reduced_capacity = stash.capacity() / storage_shrink_upsert_unused_buffers_by_ratio;
877 if reduced_capacity > stash.len() {
878 stash.shrink_to(reduced_capacity);
879 }
880 }
881}
882
883#[derive(Debug)]
886enum DrainStyle<'a, T> {
887 ToUpper(&'a Antichain<T>),
888 AtTime(T),
889}
890
891async fn drain_staged_input<S, T, FromTime, E>(
894 stash: &mut Vec<(T, UpsertKey, Reverse<FromTime>, Option<UpsertValue>)>,
895 commands_state: &mut indexmap::IndexMap<UpsertKey, types::UpsertValueAndSize<T, FromTime>>,
896 output_updates: &mut Vec<(UpsertValue, T, Diff)>,
897 multi_get_scratch: &mut Vec<UpsertKey>,
898 drain_style: DrainStyle<'_, T>,
899 error_emitter: &mut E,
900 state: &mut UpsertState<'_, S, T, FromTime>,
901 source_config: &crate::source::SourceExportCreationConfig,
902) where
903 S: UpsertStateBackend<T, FromTime>,
904 T: PartialOrder + Ord + Clone + Send + Sync + Serialize + Debug + 'static,
905 FromTime: timely::ExchangeData + Clone + Ord + Sync,
906 E: UpsertErrorEmitter<T>,
907{
908 stash.sort_unstable();
909
910 let idx = stash.partition_point(|(ts, _, _, _)| match &drain_style {
912 DrainStyle::ToUpper(upper) => !upper.less_equal(ts),
913 DrainStyle::AtTime(time) => ts <= time,
914 });
915
916 tracing::trace!(?drain_style, updates = idx, "draining stash in upsert");
917
918 commands_state.clear();
921 for (_, key, _, _) in stash.iter().take(idx) {
922 commands_state.entry(*key).or_default();
923 }
924
925 multi_get_scratch.clear();
928 multi_get_scratch.extend(commands_state.iter().map(|(k, _)| *k));
929 match state
930 .multi_get(multi_get_scratch.drain(..), commands_state.values_mut())
931 .await
932 {
933 Ok(_) => {}
934 Err(e) => {
935 error_emitter
936 .emit("Failed to fetch records from state".to_string(), e)
937 .await;
938 }
939 }
940
941 let mut commands = stash.drain(..idx).dedup_by(|a, b| {
945 let ((a_ts, a_key, _, _), (b_ts, b_key, _, _)) = (a, b);
946 a_ts == b_ts && a_key == b_key
947 });
948
949 let bincode_opts = types::upsert_bincode_opts();
950 while let Some((ts, key, from_time, value)) = commands.next() {
963 let mut command_state = if let Entry::Occupied(command_state) = commands_state.entry(key) {
964 command_state
965 } else {
966 panic!("key missing from commands_state");
967 };
968
969 let existing_value = &mut command_state.get_mut().value;
970
971 if let Some(cs) = existing_value.as_mut() {
972 cs.ensure_decoded(bincode_opts, source_config.id, Some(&key));
973 }
974
975 let existing_order = existing_value
979 .as_ref()
980 .and_then(|cs| cs.provisional_order(&ts));
981 if existing_order >= Some(&from_time.0) {
982 continue;
987 }
988
989 match value {
990 Some(value) => {
991 if let Some(old_value) =
992 existing_value.replace(StateValue::finalized_value(value.clone()))
993 {
994 if let Some(old_value) = old_value.into_decoded().finalized {
995 output_updates.push((old_value, ts.clone(), Diff::MINUS_ONE));
996 }
997 }
998 output_updates.push((value, ts, Diff::ONE));
999 }
1000 None => {
1001 if let Some(old_value) = existing_value.take() {
1002 if let Some(old_value) = old_value.into_decoded().finalized {
1003 output_updates.push((old_value, ts, Diff::MINUS_ONE));
1004 }
1005 }
1006
1007 *existing_value = Some(StateValue::tombstone());
1009 }
1010 }
1011 }
1012
1013 match state
1014 .multi_put(
1015 true, commands_state.drain(..).map(|(k, cv)| {
1017 (
1018 k,
1019 types::PutValue {
1020 value: cv.value.map(|cv| cv.into_decoded()),
1021 previous_value_metadata: cv.metadata.map(|v| ValueMetadata {
1022 size: v.size.try_into().expect("less than i64 size"),
1023 is_tombstone: v.is_tombstone,
1024 }),
1025 },
1026 )
1027 }),
1028 )
1029 .await
1030 {
1031 Ok(_) => {}
1032 Err(e) => {
1033 error_emitter
1034 .emit("Failed to update records in state".to_string(), e)
1035 .await;
1036 }
1037 }
1038}
1039
1040#[cfg(feature = "fuzzing")]
1044struct PanicErrorEmitter;
1045
1046#[cfg(feature = "fuzzing")]
1047#[async_trait::async_trait(?Send)]
1048impl<T> UpsertErrorEmitter<T> for PanicErrorEmitter {
1049 async fn emit(&mut self, context: String, e: anyhow::Error) {
1050 panic!("unexpected upsert state error during fuzzing: {context}: {e}");
1051 }
1052}
1053
1054#[cfg(feature = "fuzzing")]
1061pub async fn fuzz_drain_staged_input(
1062 parts: &types::FuzzUpsertParts,
1063 source_config: &crate::source::SourceExportCreationConfig,
1064 commands: Vec<(u64, UpsertKey, u64, Option<UpsertValue>)>,
1065 drain_to: u64,
1066 all_keys: &[UpsertKey],
1067) -> (Vec<(UpsertValue, u64, Diff)>, Vec<Option<UpsertValue>>) {
1068 let mut state = parts.state();
1069 let mut stash: Vec<(u64, UpsertKey, Reverse<u64>, Option<UpsertValue>)> = commands
1070 .into_iter()
1071 .map(|(ts, key, order, value)| (ts, key, Reverse(order), value))
1072 .collect();
1073 let mut commands_state = indexmap::IndexMap::new();
1074 let mut output = Vec::new();
1075 let mut multi_get_scratch = Vec::new();
1076 let mut emitter = PanicErrorEmitter;
1077
1078 drain_staged_input(
1079 &mut stash,
1080 &mut commands_state,
1081 &mut output,
1082 &mut multi_get_scratch,
1083 DrainStyle::ToUpper(&Antichain::from_elem(drain_to)),
1084 &mut emitter,
1085 &mut state,
1086 source_config,
1087 )
1088 .await;
1089
1090 let bincode_opts = types::upsert_bincode_opts();
1091 let mut results = vec![types::UpsertValueAndSize::default(); all_keys.len()];
1092 state
1093 .multi_get(all_keys.iter().copied(), results.iter_mut())
1094 .await
1095 .expect("multi_get in fuzz hook should not error");
1096 let final_state = results
1097 .into_iter()
1098 .map(|r| match r.value {
1099 None => None,
1100 Some(mut sv) => {
1101 sv.ensure_decoded(bincode_opts, GlobalId::User(0), None);
1102 sv.into_decoded().finalized
1103 }
1104 })
1105 .collect();
1106
1107 (output, final_state)
1108}
1109
1110pub(crate) struct UpsertConfig {
1113 pub shrink_upsert_unused_buffers_by_ratio: usize,
1114}
1115
1116fn upsert_classic<'scope, T, FromTime, F, Fut, US>(
1117 input: VecCollection<'scope, T, (UpsertKey, Option<UpsertValue>, FromTime), Diff>,
1118 key_indices: Vec<usize>,
1119 resume_upper: Antichain<T>,
1120 previous_ok: VecCollection<'scope, T, Row, Diff>,
1121 previous_err: VecCollection<'scope, T, DataflowError, Diff>,
1122 previous_token: Option<Vec<PressOnDropButton>>,
1123 upsert_metrics: UpsertMetrics,
1124 source_config: crate::source::SourceExportCreationConfig,
1125 state: F,
1126 upsert_config: UpsertConfig,
1127 prevent_snapshot_buffering: bool,
1128 snapshot_buffering_max: Option<usize>,
1129) -> (
1130 VecCollection<'scope, T, Result<Row, DataflowError>, Diff>,
1131 StreamVec<'scope, T, (Option<GlobalId>, HealthStatusUpdate)>,
1132 StreamVec<'scope, T, Infallible>,
1133 PressOnDropButton,
1134)
1135where
1136 T: Timestamp + TotalOrder + Sync,
1137 F: FnOnce() -> Fut + 'static,
1138 Fut: std::future::Future<Output = US>,
1139 US: UpsertStateBackend<T, FromTime>,
1140 FromTime: timely::ExchangeData + Clone + Ord + Sync,
1141{
1142 let mut builder = AsyncOperatorBuilder::new("Upsert".to_string(), input.scope());
1143
1144 let previous = key_persist_feedback(previous_ok, previous_err, key_indices);
1145 let (output_handle, output) = builder.new_output();
1146
1147 let (_snapshot_handle, snapshot_stream) =
1150 builder.new_output::<CapacityContainerBuilder<Vec<Infallible>>>();
1151
1152 let (mut health_output, health_stream) = builder.new_output();
1153 let mut input = builder.new_input_for(
1154 input.inner,
1155 Exchange::new(move |((key, _, _), _, _)| UpsertKey::hashed(key)),
1156 &output_handle,
1157 );
1158
1159 let mut previous = builder.new_input_for(
1160 previous.inner,
1161 Exchange::new(|((key, _), _, _)| UpsertKey::hashed(key)),
1162 &output_handle,
1163 );
1164
1165 let upsert_shared_metrics = Arc::clone(&upsert_metrics.shared);
1166 let shutdown_button = builder.build(move |caps| async move {
1167 let [mut output_cap, mut snapshot_cap, health_cap]: [_; 3] = caps.try_into().unwrap();
1168
1169 let mut state = UpsertState::<_, _, FromTime>::new(
1170 state().await,
1171 upsert_shared_metrics,
1172 &upsert_metrics,
1173 source_config.source_statistics.clone(),
1174 upsert_config.shrink_upsert_unused_buffers_by_ratio,
1175 );
1176 let mut events = vec![];
1177 let mut snapshot_upper = Antichain::from_elem(Timestamp::minimum());
1178
1179 let mut stash = vec![];
1180
1181 let mut error_emitter = (&mut health_output, &health_cap);
1182
1183 tracing::info!(
1184 ?resume_upper,
1185 ?snapshot_upper,
1186 "timely-{} upsert source {} starting rehydration",
1187 source_config.worker_id,
1188 source_config.id
1189 );
1190 while !PartialOrder::less_equal(&resume_upper, &snapshot_upper) {
1193 previous.ready().await;
1194 while let Some(event) = previous.next_sync() {
1195 match event {
1196 AsyncEvent::Data(_cap, data) => {
1197 events.extend(data.into_iter().filter_map(|((key, value), ts, diff)| {
1198 if !resume_upper.less_equal(&ts) {
1199 Some((key, value, diff))
1200 } else {
1201 None
1202 }
1203 }))
1204 }
1205 AsyncEvent::Progress(upper) => {
1206 snapshot_upper = upper;
1207 }
1208 };
1209 }
1210
1211 match state
1212 .consolidate_chunk(
1213 events.drain(..),
1214 PartialOrder::less_equal(&resume_upper, &snapshot_upper),
1215 )
1216 .await
1217 {
1218 Ok(_) => {
1219 if let Some(ts) = snapshot_upper.clone().into_option() {
1220 if !resume_upper.less_equal(&ts) {
1224 snapshot_cap.downgrade(&ts);
1225 output_cap.downgrade(&ts);
1226 }
1227 }
1228 }
1229 Err(e) => {
1230 UpsertErrorEmitter::<T>::emit(
1231 &mut error_emitter,
1232 "Failed to rehydrate state".to_string(),
1233 e,
1234 )
1235 .await;
1236 }
1237 }
1238 }
1239
1240 drop(events);
1241 drop(previous_token);
1242 drop(snapshot_cap);
1243
1244 while let Some(_event) = previous.next().await {}
1250
1251 if let Some(ts) = resume_upper.as_option() {
1253 output_cap.downgrade(ts);
1254 }
1255
1256 tracing::info!(
1257 "timely-{} upsert source {} finished rehydration",
1258 source_config.worker_id,
1259 source_config.id
1260 );
1261
1262 let mut commands_state: indexmap::IndexMap<_, types::UpsertValueAndSize<T, FromTime>> =
1265 indexmap::IndexMap::new();
1266 let mut multi_get_scratch = Vec::new();
1267
1268 let mut output_updates = vec![];
1270 let mut input_upper = Antichain::from_elem(Timestamp::minimum());
1271
1272 while let Some(event) = input.next().await {
1273 let events = [event]
1276 .into_iter()
1277 .chain(std::iter::from_fn(|| input.next().now_or_never().flatten()))
1278 .enumerate();
1279
1280 let mut partial_drain_time = None;
1281 for (i, event) in events {
1282 match event {
1283 AsyncEvent::Data(cap, mut data) => {
1284 tracing::trace!(
1285 time=?cap.time(),
1286 updates=%data.len(),
1287 "received data in upsert"
1288 );
1289 stage_input(
1290 &mut stash,
1291 &mut data,
1292 &input_upper,
1293 &resume_upper,
1294 upsert_config.shrink_upsert_unused_buffers_by_ratio,
1295 );
1296
1297 let event_time = cap.time();
1298 if prevent_snapshot_buffering && output_cap.time() == event_time {
1305 partial_drain_time = Some(event_time.clone());
1306 }
1307 }
1308 AsyncEvent::Progress(upper) => {
1309 tracing::trace!(?upper, "received progress in upsert");
1310 if PartialOrder::less_than(&upper, &resume_upper) {
1313 continue;
1314 }
1315
1316 partial_drain_time = None;
1319 drain_staged_input::<_, _, _, _>(
1320 &mut stash,
1321 &mut commands_state,
1322 &mut output_updates,
1323 &mut multi_get_scratch,
1324 DrainStyle::ToUpper(&upper),
1325 &mut error_emitter,
1326 &mut state,
1327 &source_config,
1328 )
1329 .await;
1330
1331 output_handle.give_container(&output_cap, &mut output_updates);
1332
1333 if let Some(ts) = upper.as_option() {
1334 output_cap.downgrade(ts);
1335 }
1336 input_upper = upper;
1337 }
1338 }
1339 let events_processed = i + 1;
1340 if let Some(max) = snapshot_buffering_max {
1341 if events_processed >= max {
1342 break;
1343 }
1344 }
1345 }
1346
1347 if let Some(partial_drain_time) = partial_drain_time {
1356 drain_staged_input::<_, _, _, _>(
1357 &mut stash,
1358 &mut commands_state,
1359 &mut output_updates,
1360 &mut multi_get_scratch,
1361 DrainStyle::AtTime(partial_drain_time),
1362 &mut error_emitter,
1363 &mut state,
1364 &source_config,
1365 )
1366 .await;
1367
1368 output_handle.give_container(&output_cap, &mut output_updates);
1369 }
1370 }
1371 });
1372
1373 (
1374 output.as_collection().map(|result| match result {
1375 Ok(ok) => Ok(ok),
1376 Err(err) => Err(DataflowError::from(EnvelopeError::Upsert(*err))),
1377 }),
1378 health_stream,
1379 snapshot_stream,
1380 shutdown_button.press_on_drop(),
1381 )
1382}
1383
1384#[async_trait::async_trait(?Send)]
1385pub(crate) trait UpsertErrorEmitter<T> {
1386 async fn emit(&mut self, context: String, e: anyhow::Error);
1387}
1388
1389#[async_trait::async_trait(?Send)]
1390impl<T: Timestamp> UpsertErrorEmitter<T>
1391 for (
1392 &mut AsyncOutputHandle<
1393 T,
1394 CapacityContainerBuilder<Vec<(Option<GlobalId>, HealthStatusUpdate)>>,
1395 >,
1396 &Capability<T>,
1397 )
1398{
1399 async fn emit(&mut self, context: String, e: anyhow::Error) {
1400 process_upsert_state_error::<T>(context, e, self.0, self.1).await
1401 }
1402}
1403
1404async fn process_upsert_state_error<T: Timestamp>(
1406 context: String,
1407 e: anyhow::Error,
1408 health_output: &AsyncOutputHandle<
1409 T,
1410 CapacityContainerBuilder<Vec<(Option<GlobalId>, HealthStatusUpdate)>>,
1411 >,
1412 health_cap: &Capability<T>,
1413) {
1414 let update = HealthStatusUpdate::halting(e.context(context).to_string_with_causes(), None);
1415 health_output.give(health_cap, (None, update));
1416 std::future::pending::<()>().await;
1417 unreachable!("pending future never returns");
1418}