1use std::collections::BTreeMap;
15use std::iter;
16use std::sync::Arc;
17
18use differential_dataflow::{AsCollection, VecCollection};
19use mz_ore::cast::CastLossy;
20use mz_persist_client::operators::shard_source::SnapshotMode;
21use mz_repr::{Datum, Diff, GlobalId, Row, RowPacker};
22use mz_storage_operators::persist_source;
23use mz_storage_operators::persist_source::Subtime;
24use mz_storage_types::controller::CollectionMetadata;
25use mz_storage_types::dyncfgs;
26use mz_storage_types::errors::{
27 DataflowError, DecodeError, EnvelopeError, UpsertError, UpsertNullKeyError, UpsertValueError,
28};
29use mz_storage_types::parameters::StorageMaxInflightBytesConfig;
30use mz_storage_types::sources::envelope::{KeyEnvelope, NoneEnvelope, UpsertEnvelope, UpsertStyle};
31use mz_storage_types::sources::*;
32use mz_timely_util::builder_async::PressOnDropButton;
33use mz_timely_util::operator::CollectionExt;
34use mz_timely_util::order::refine_antichain;
35use serde::{Deserialize, Serialize};
36use timely::container::CapacityContainerBuilder;
37use timely::dataflow::StreamVec;
38use timely::dataflow::operators::vec::Map;
39use timely::dataflow::operators::{ConnectLoop, Feedback, Leave, OkErr};
40use timely::dataflow::scope::Scope;
41use timely::progress::{Antichain, Timestamp};
42
43use crate::decode::{render_decode_cdcv2, render_decode_delimited};
44use crate::healthcheck::{HealthStatusMessage, StatusNamespace};
45use crate::source::types::{DecodeResult, SourceOutput, SourceRender};
46use crate::source::{self, RawSourceCreationConfig, SourceExportCreationConfig};
47use crate::upsert::{UpsertKey, UpsertSourceTime, UpsertValue};
48
49pub fn render_source<'scope, 'root, C>(
61 scope: Scope<'scope, mz_repr::Timestamp>,
62 root_scope: Scope<'root, ()>,
63 dataflow_debug_name: &String,
64 connection: C,
65 description: IngestionDescription<CollectionMetadata>,
66 resume_stream: StreamVec<'scope, mz_repr::Timestamp, ()>,
67 storage_state: &crate::storage_state::StorageState,
68 base_source_config: RawSourceCreationConfig,
69) -> (
70 BTreeMap<
71 GlobalId,
72 (
73 VecCollection<'scope, mz_repr::Timestamp, Row, Diff>,
74 VecCollection<'scope, mz_repr::Timestamp, DataflowError, Diff>,
75 ),
76 >,
77 Vec<StreamVec<'root, (), HealthStatusMessage>>,
78 StreamVec<'scope, mz_repr::Timestamp, ()>,
79 Vec<PressOnDropButton>,
80)
81where
82 C: SourceConnection + SourceRender + 'static,
83 C::Time: UpsertSourceTime,
84{
85 let mut needed_tokens = Vec::new();
87
88 let (starter, mut start_signal) = tokio::sync::mpsc::channel::<()>(1);
101 let start_signal = async move {
102 let _ = start_signal.recv().await;
103 };
104
105 let (exports, health, remap_upper, source_tokens) = source::create_raw_source(
108 scope,
109 root_scope,
110 storage_state,
111 resume_stream,
112 &base_source_config,
113 connection,
114 start_signal,
115 );
116
117 needed_tokens.extend(source_tokens);
118
119 let mut health_streams = Vec::with_capacity(exports.len() + 1);
120 health_streams.push(health);
121
122 let mut outputs = BTreeMap::new();
123 for (export_id, export) in exports {
124 type CB<C> = CapacityContainerBuilder<C>;
125 let (ok_stream, err_stream) =
126 export.map_fallible::<CB<_>, CB<_>, _, _, _>("export-demux-ok-err", |r| r);
127
128 let mut error_collections = Vec::new();
132
133 let data_config = base_source_config.source_exports[&export_id]
134 .data_config
135 .clone();
136 let (ok, extra_tokens, health_stream) = render_source_stream(
137 scope,
138 dataflow_debug_name,
139 export_id,
140 ok_stream,
141 data_config,
142 &description,
143 &mut error_collections,
144 storage_state,
145 &base_source_config,
146 starter.clone(),
147 );
148 needed_tokens.extend(extra_tokens);
149
150 let err_collection = match error_collections.len() {
152 0 => err_stream,
153 _ => err_stream.concatenate(error_collections),
154 };
155
156 outputs.insert(export_id, (ok, err_collection));
157
158 health_streams.extend(health_stream.into_iter().map(|s| s.leave(root_scope)));
159 }
160 (outputs, health_streams, remap_upper, needed_tokens)
161}
162
163fn render_source_stream<'scope, FromTime>(
166 scope: Scope<'scope, mz_repr::Timestamp>,
167 dataflow_debug_name: &String,
168 export_id: GlobalId,
169 ok_source: VecCollection<'scope, mz_repr::Timestamp, SourceOutput<FromTime>, Diff>,
170 data_config: SourceExportDataConfig,
171 description: &IngestionDescription<CollectionMetadata>,
172 error_collections: &mut Vec<VecCollection<'scope, mz_repr::Timestamp, DataflowError, Diff>>,
173 storage_state: &crate::storage_state::StorageState,
174 base_source_config: &RawSourceCreationConfig,
175 rehydrated_token: impl std::any::Any + 'static,
176) -> (
177 VecCollection<'scope, mz_repr::Timestamp, Row, Diff>,
178 Vec<PressOnDropButton>,
179 Vec<StreamVec<'scope, mz_repr::Timestamp, HealthStatusMessage>>,
180)
181where
182 FromTime: Timestamp + Sync,
183 FromTime: UpsertSourceTime,
184{
185 let mut needed_tokens = vec![];
186
187 let SourceExportDataConfig { encoding, envelope } = data_config;
189
190 let SourceDesc {
191 connection: _,
192 timestamp_interval: _,
193 } = description.desc;
194
195 let (decoded_stream, decode_health) = match encoding {
196 None => (
197 ok_source.map(|r| DecodeResult {
198 key: Some(Ok(r.key)),
206 value: Some(Ok(r.value)),
207 metadata: r.metadata,
208 from_time: r.from_time,
209 }),
210 None,
211 ),
212 Some(encoding) => {
213 let (decoded_stream, decode_health) = render_decode_delimited(
214 ok_source,
215 encoding.key,
216 encoding.value,
217 dataflow_debug_name.clone(),
218 storage_state.metrics.decode_defs.clone(),
219 storage_state.storage_configuration.clone(),
220 );
221 (decoded_stream, Some(decode_health))
222 }
223 };
224
225 let (envelope_ok, envelope_health) = match &envelope {
227 SourceEnvelope::Upsert(upsert_envelope) => {
228 let upsert_input = upsert_commands(decoded_stream, upsert_envelope.clone());
229
230 let persist_clients = Arc::clone(&storage_state.persist_clients);
231 let resume_upper = base_source_config.resume_uppers[&export_id].clone();
233
234 let upper_ts = resume_upper
235 .as_option()
236 .expect("resuming an already finished ingestion")
237 .clone();
238 let outer_mz_scope = scope.clone();
239 let (upsert, health_update) = scope.scoped(
240 &format!("upsert_rehydration_backpressure({})", export_id),
241 |scope| {
242 let (
243 previous_ok,
244 previous_err,
245 previous_token,
246 feedback_handle,
247 backpressure_metrics,
248 ) = {
249 let as_of = Antichain::from_elem(upper_ts.saturating_sub(1));
250
251 let backpressure_max_inflight_bytes = get_backpressure_max_inflight_bytes(
252 &storage_state
253 .storage_configuration
254 .parameters
255 .storage_dataflow_max_inflight_bytes_config,
256 &storage_state.instance_context.cluster_memory_limit,
257 );
258
259 let (feedback_handle, flow_control, backpressure_metrics) =
260 if let Some(storage_dataflow_max_inflight_bytes) =
261 backpressure_max_inflight_bytes
262 {
263 tracing::info!(
264 ?backpressure_max_inflight_bytes,
265 "timely-{} using backpressure in upsert for source {}",
266 base_source_config.worker_id,
267 export_id
268 );
269 if !storage_state
270 .storage_configuration
271 .parameters
272 .storage_dataflow_max_inflight_bytes_config
273 .disk_only
274 || storage_state.instance_context.scratch_directory.is_some()
275 {
276 let (feedback_handle, feedback_data) =
277 scope.feedback(Default::default());
278
279 let backpressure_metrics = Some(
281 base_source_config
282 .metrics
283 .get_backpressure_metrics(export_id, scope.index()),
284 );
285
286 (
287 Some(feedback_handle),
288 Some(persist_source::FlowControl {
289 progress_stream: feedback_data,
290 max_inflight_bytes: storage_dataflow_max_inflight_bytes,
291 summary: (Default::default(), Subtime::least_summary()),
292 metrics: backpressure_metrics
293 .as_ref()
294 .map(|m| m.operator_metrics()),
295 }),
296 backpressure_metrics,
297 )
298 } else {
299 (None, None, None)
300 }
301 } else {
302 (None, None, None)
303 };
304
305 let storage_metadata = description.source_exports[&export_id]
306 .storage_metadata
307 .clone();
308
309 let error_handler =
310 storage_state.error_handler("upsert_rehydration", export_id);
311
312 let (ok_stream, err_stream, tok) = persist_source::persist_source_core(
313 outer_mz_scope,
314 scope,
315 export_id,
316 persist_clients,
317 storage_metadata,
318 None,
319 Some(as_of),
320 SnapshotMode::Include,
321 Antichain::new(),
322 None,
323 flow_control,
324 false.then_some(|| unreachable!()),
325 async {},
326 error_handler,
327 );
328 (
329 ok_stream.as_collection(),
330 err_stream.as_collection(),
331 Some(tok),
332 feedback_handle,
333 backpressure_metrics,
334 )
335 };
336
337 let export_statistics = storage_state
338 .aggregated_statistics
339 .get_source(&export_id)
340 .expect("statistics initialized")
341 .clone();
342 let export_config = SourceExportCreationConfig {
343 id: export_id,
344 worker_id: base_source_config.worker_id,
345 metrics: base_source_config.metrics.clone(),
346 source_statistics: export_statistics,
347 };
348 let (upsert, health_update, snapshot_progress, upsert_token) =
349 if dyncfgs::ENABLE_UPSERT_V2
350 .get(storage_state.storage_configuration.config_set())
351 {
352 let stash_flavor =
356 crate::upsert_continual_feedback_v2::UpsertStashFlavor::from_config(
357 storage_state.storage_configuration.config_set(),
358 );
359 crate::upsert::upsert_v2(
360 upsert_input.enter(scope),
361 upsert_envelope.clone(),
362 refine_antichain(&resume_upper),
363 previous_ok,
364 previous_err,
365 previous_token,
366 export_config,
367 backpressure_metrics,
368 stash_flavor,
369 )
370 } else {
371 crate::upsert::upsert(
372 upsert_input.enter(scope),
373 upsert_envelope.clone(),
374 refine_antichain(&resume_upper),
375 previous_ok,
376 previous_err,
377 previous_token,
378 export_config,
379 &storage_state.instance_context,
380 &storage_state.storage_configuration,
381 &storage_state.dataflow_parameters,
382 backpressure_metrics,
383 )
384 };
385
386 needed_tokens.push(upsert_token);
391
392 if dyncfgs::DELAY_SOURCES_PAST_REHYDRATION
396 .get(storage_state.storage_configuration.config_set())
397 {
398 crate::upsert::rehydration_finished(
399 scope.clone(),
400 base_source_config,
401 rehydrated_token,
402 refine_antichain(&resume_upper),
403 snapshot_progress.clone(),
404 );
405 } else {
406 drop(rehydrated_token)
407 };
408
409 if let Some(feedback_handle) = feedback_handle {
412 snapshot_progress.connect_loop(feedback_handle);
413 }
414
415 (
416 upsert.leave(outer_mz_scope),
417 health_update
418 .map(|(id, update)| HealthStatusMessage {
419 id,
420 namespace: StatusNamespace::Upsert,
421 update,
422 })
423 .leave(outer_mz_scope),
424 )
425 },
426 );
427
428 let (upsert_ok, upsert_err) = upsert.inner.ok_err(split_ok_err);
429 error_collections.push(upsert_err.as_collection());
430
431 (upsert_ok.as_collection(), Some(health_update))
432 }
433 SourceEnvelope::None(none_envelope) => {
434 let results = append_metadata_to_value(decoded_stream);
435
436 let flattened_stream = flatten_results_prepend_keys(none_envelope, results);
437
438 let (stream, errors) = flattened_stream.inner.ok_err(split_ok_err);
439
440 error_collections.push(errors.as_collection());
441 (stream.as_collection(), None)
442 }
443 SourceEnvelope::CdcV2 => {
444 let (oks, token) = render_decode_cdcv2(&decoded_stream);
445 needed_tokens.push(token);
446 (oks, None)
447 }
448 };
449
450 let health = decode_health.into_iter().chain(envelope_health).collect();
452 (envelope_ok, needed_tokens, health)
453}
454
455fn get_backpressure_max_inflight_bytes(
458 inflight_bytes_config: &StorageMaxInflightBytesConfig,
459 cluster_memory_limit: &Option<usize>,
460) -> Option<usize> {
461 let StorageMaxInflightBytesConfig {
462 max_inflight_bytes_default,
463 max_inflight_bytes_cluster_size_fraction,
464 disk_only: _,
465 } = inflight_bytes_config;
466
467 if max_inflight_bytes_default.is_some() {
469 let current_cluster_max_bytes_limit =
470 cluster_memory_limit.as_ref().and_then(|cluster_memory| {
471 max_inflight_bytes_cluster_size_fraction.map(|fraction| {
472 usize::cast_lossy(f64::cast_lossy(*cluster_memory) * fraction)
474 })
475 });
476 current_cluster_max_bytes_limit.or(*max_inflight_bytes_default)
477 } else {
478 None
479 }
480}
481
482fn split_ok_err<O, E, T, D>(x: (Result<O, E>, T, D)) -> Result<(O, T, D), (E, T, D)> {
485 match x {
486 (Ok(ok), ts, diff) => Ok((ok, ts, diff)),
487 (Err(err), ts, diff) => Err((err, ts, diff)),
488 }
489}
490
491#[derive(
493 Debug,
494 Clone,
495 Hash,
496 PartialEq,
497 Eq,
498 Ord,
499 PartialOrd,
500 Serialize,
501 Deserialize
502)]
503struct KV {
504 key: Option<Result<Row, DecodeError>>,
505 val: Option<Result<Row, DecodeError>>,
506}
507
508fn append_metadata_to_value<'scope, T: Timestamp, FromTime: Timestamp>(
509 results: VecCollection<'scope, T, DecodeResult<FromTime>, Diff>,
510) -> VecCollection<'scope, T, KV, Diff> {
511 results.map(move |res| {
512 let val = res.value.map(|val_result| {
513 val_result.map(|mut val| {
514 if !res.metadata.is_empty() {
515 RowPacker::for_existing_row(&mut val).extend_by_row(&res.metadata);
516 }
517 val
518 })
519 });
520
521 KV { val, key: res.key }
522 })
523}
524
525fn upsert_commands<'scope, T: Timestamp, FromTime: Timestamp>(
527 input: VecCollection<'scope, T, DecodeResult<FromTime>, Diff>,
528 upsert_envelope: UpsertEnvelope,
529) -> VecCollection<'scope, T, (UpsertKey, Option<UpsertValue>, FromTime), Diff> {
530 let mut row_buf = Row::default();
531 input.map(move |result| {
532 let from_time = result.from_time;
533
534 let key = match result.key {
535 Some(Ok(key)) => Ok(key),
536 None => Err(UpsertError::NullKey(UpsertNullKeyError)),
537 Some(Err(err)) => Err(UpsertError::KeyDecode(err)),
538 };
539
540 let key = match key {
542 Ok(key) => key,
543 Err(err) => match result.value {
544 Some(_) => {
545 return (
546 UpsertKey::from_key(Err(&err)),
547 Some(Err(Box::new(err))),
548 from_time,
549 );
550 }
551 None => return (UpsertKey::from_key(Err(&err)), None, from_time),
552 },
553 };
554
555 let key_row = match upsert_envelope.style {
557 UpsertStyle::Debezium { .. }
559 | UpsertStyle::Default(KeyEnvelope::Flattened)
560 | UpsertStyle::ValueErrInline {
561 key_envelope: KeyEnvelope::Flattened,
562 error_column: _,
563 } => key,
564 UpsertStyle::Default(KeyEnvelope::Named(_))
566 | UpsertStyle::ValueErrInline {
567 key_envelope: KeyEnvelope::Named(_),
568 error_column: _,
569 } => {
570 if key.iter().nth(1).is_none() {
571 key
572 } else {
573 row_buf.packer().push_list(key.iter());
574 row_buf.clone()
575 }
576 }
577 UpsertStyle::Default(KeyEnvelope::None)
578 | UpsertStyle::ValueErrInline {
579 key_envelope: KeyEnvelope::None,
580 error_column: _,
581 } => unreachable!(),
582 };
583
584 let key = UpsertKey::from_key(Ok(&key_row));
585
586 let metadata = result.metadata;
587
588 let value = match result.value {
589 Some(Ok(ref row)) => match upsert_envelope.style {
590 UpsertStyle::Debezium { after_idx } => match row.iter().nth(after_idx).unwrap() {
591 Datum::List(after) => {
592 let mut packer = row_buf.packer();
593 packer.extend(after.iter());
594 packer.extend_by_row(&metadata);
595 Some(Ok(row_buf.clone()))
596 }
597 Datum::Null => None,
598 d => panic!("type error: expected record, found {:?}", d),
599 },
600 UpsertStyle::Default(_) => {
601 let mut packer = row_buf.packer();
602 packer.extend_by_row(&key_row);
603 packer.extend_by_row(row);
604 packer.extend_by_row(&metadata);
605 Some(Ok(row_buf.clone()))
606 }
607 UpsertStyle::ValueErrInline { .. } => {
608 let mut packer = row_buf.packer();
609 packer.extend_by_row(&key_row);
610 packer.push(Datum::Null);
612 packer.extend_by_row(row);
613 packer.extend_by_row(&metadata);
614 Some(Ok(row_buf.clone()))
615 }
616 },
617 Some(Err(inner)) => {
618 match upsert_envelope.style {
619 UpsertStyle::ValueErrInline { .. } => {
620 let mut count = 0;
621 let err_string = inner.to_string();
623 let mut packer = row_buf.packer();
624 for datum in key_row.iter() {
625 packer.push(datum);
626 count += 1;
627 }
628 packer.push_list(iter::once(Datum::String(&err_string)));
630 count += 1;
631 let metadata_len = metadata.as_row_ref().iter().count();
632 packer.extend(
634 iter::repeat(Datum::Null)
635 .take(upsert_envelope.source_arity - count - metadata_len),
636 );
637 packer.extend_by_row(&metadata);
638 Some(Ok(row_buf.clone()))
639 }
640 _ => Some(Err(Box::new(UpsertError::Value(UpsertValueError {
641 for_key: key_row,
642 inner,
643 })))),
644 }
645 }
646 None => None,
647 };
648
649 (key, value, from_time)
650 })
651}
652
653fn flatten_results_prepend_keys<'scope, T: Timestamp>(
655 none_envelope: &NoneEnvelope,
656 results: VecCollection<'scope, T, KV, Diff>,
657) -> VecCollection<'scope, T, Result<Row, DataflowError>, Diff> {
658 let NoneEnvelope {
659 key_envelope,
660 key_arity,
661 } = none_envelope;
662
663 let null_key_columns = Row::pack_slice(&vec![Datum::Null; *key_arity]);
664
665 match key_envelope {
666 KeyEnvelope::None => {
667 results.flat_map(|KV { val, .. }| val.map(|result| result.map_err(Into::into)))
668 }
669 KeyEnvelope::Flattened => results
670 .flat_map(raise_key_value_errors)
671 .map(move |maybe_kv| {
672 maybe_kv.map(|(key, value)| {
673 let mut key = key.unwrap_or_else(|| null_key_columns.clone());
674 RowPacker::for_existing_row(&mut key).extend_by_row(&value);
675 key
676 })
677 }),
678 KeyEnvelope::Named(_) => {
679 results
680 .flat_map(raise_key_value_errors)
681 .map(move |maybe_kv| {
682 maybe_kv.map(|(key, value)| {
683 let mut key = key.unwrap_or_else(|| null_key_columns.clone());
684 let row = if key.iter().nth(1).is_none() {
687 RowPacker::for_existing_row(&mut key).extend_by_row(&value);
688 key
689 } else {
690 let mut new_row = Row::default();
691 let mut packer = new_row.packer();
692 packer.push_list(key.iter());
693 packer.extend_by_row(&value);
694 new_row
695 };
696 row
697 })
698 })
699 }
700 }
701}
702
703fn raise_key_value_errors(
705 KV { key, val }: KV,
706) -> Option<Result<(Option<Row>, Row), DataflowError>> {
707 match (key, val) {
708 (Some(Ok(key)), Some(Ok(value))) => Some(Ok((Some(key), value))),
709 (None, Some(Ok(value))) => Some(Ok((None, value))),
710 (_, Some(Err(e))) => Some(Err(e.into())),
712 (Some(Err(e)), _) => Some(Err(e.into())),
713 (None, None) => None,
714 _ => Some(Err(DataflowError::from(EnvelopeError::Flat(
716 "Value not present for message".into(),
717 )))),
718 }
719}
720
721#[cfg(test)]
722mod test {
723 use super::*;
724
725 #[mz_ore::test]
726 fn test_no_default() {
727 let config = StorageMaxInflightBytesConfig {
728 max_inflight_bytes_default: None,
729 max_inflight_bytes_cluster_size_fraction: Some(0.5),
730 disk_only: false,
731 };
732 let memory_limit = Some(1000);
733
734 let backpressure_inflight_bytes_limit =
735 get_backpressure_max_inflight_bytes(&config, &memory_limit);
736
737 assert_eq!(backpressure_inflight_bytes_limit, None)
738 }
739
740 #[mz_ore::test]
741 fn test_no_matching_size() {
742 let config = StorageMaxInflightBytesConfig {
743 max_inflight_bytes_default: Some(10000),
744 max_inflight_bytes_cluster_size_fraction: Some(0.5),
745 disk_only: false,
746 };
747
748 let backpressure_inflight_bytes_limit = get_backpressure_max_inflight_bytes(&config, &None);
749
750 assert_eq!(
751 backpressure_inflight_bytes_limit,
752 config.max_inflight_bytes_default
753 )
754 }
755
756 #[mz_ore::test]
757 fn test_calculated_cluster_limit() {
758 let config = StorageMaxInflightBytesConfig {
759 max_inflight_bytes_default: Some(10000),
760 max_inflight_bytes_cluster_size_fraction: Some(0.5),
761 disk_only: false,
762 };
763 let memory_limit = Some(2000);
764
765 let backpressure_inflight_bytes_limit =
766 get_backpressure_max_inflight_bytes(&config, &memory_limit);
767
768 assert_eq!(backpressure_inflight_bytes_limit, Some(1000));
770 }
771}