mz_storage/render.rs
1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Use of this software is governed by the Business Source License
4// included in the LICENSE file.
5//
6// As of the Change Date specified in that file, in accordance with
7// the Business Source License, use of this software will be governed
8// by the Apache License, Version 2.0.
9
10//! Renders ingestions and exports into timely dataflow
11//!
12//! ## Ingestions
13//!
14//! ### Overall structure
15//!
16//! Before describing any of the timely operators involved in ingesting a source it helps to
17//! understand the high level structure of the timely scopes involved. The reason for this
18//! structure is the fact that we ingest external sources with a source-specific, and source
19//! implementation defined, timestamp type which tracks progress in a way that the source
20//! implementation understands. Each source specific timestamp must be compatible with timely's
21//! `timely::progress::Timestamp` trait and so it's suitable to represent timely streams and by
22//! extension differential collections.
23//!
24//! On the other hand, Materialize expects a specific timestamp type for all its collections
25//! (currently `mz_repr::Timestamp`) so at some point the dataflow's timestamp must change. More
26//! generally, the ingestion dataflow starts with some timestamp type `FromTime` and ends with
27//! another timestamp type `IntoTime`.
28//!
29//! Here we run into a problem though because we want to start with a timely stream of type
30//! `Stream<G1: Scope<Timestamp=FromTime>, ..>` and end up using it in a scope `G2` whose timestamp
31//! type is `IntoTime`. Timely dataflows are organized in scopes where each scope has an associated
32//! timestamp type that must refine the timestamp type of its parent scope. What "refines" means is
33//! defined by the [`timely::progress::timestamp::Refines`] trait in timely. `FromTime` however
34//! does not refine `IntoTime` nor does `IntoTime` refine `FromTime`.
35//!
36//! In order to acomplish this we split ingestion dataflows in two scopes, both of which are
37//! children of the root timely scope. The first scope is timestamped with `FromTime` and the
38//! second one with `IntoTime`. To move timely streams from the one scope to the other we must do
39//! so manually. Each stream that needs to be transferred between scopes is first captured using
40//! [`timely::dataflow::operators::capture::capture::Capture`] into a tokio unbounded mpsc channel.
41//! The data in the channel record in full detail the worker-local view of the original stream and
42//! whoever controls the receiver can read in the events, in the standard way of consuming the
43//! async channel, and work with it. How the receiver is turned back into a timely stream in the
44//! destination scope is described in the next section.
45//!
46//! For now keep in mind the general structure of the dataflow:
47//!
48//!
49//! ```text
50//! +----------------RootScope(Timestamp=())------------------+
51//! | |
52//! | +---FromTime Scope---+ +---IntoTime Scope--+ | |
53//! | | | | | |
54//! | | *--+---------+--> | |
55//! | | | | | |
56//! | | <--+---------+--* | |
57//! | +--------------------+ ^ +-------------------+ |
58//! | | |
59//! | | |
60//! | data exchanged between |
61//! | scopes with capture/reclock |
62//! +---------------------------------------------------------+
63//! ```
64//!
65//! ### Detailed dataflow
66//!
67//! We are now ready to describe the detailed structure of the ingestion dataflow. The dataflow
68//! begins with the `source reader` dataflow fragment which is rendered in a `FromTime` timely
69//! scope. This scope's timestamp is controlled by the [`crate::source::types::SourceRender::Time`]
70//! associated type and can be anything the source implementation desires.
71//!
72//! Each source is free to render any arbitrary dataflow fragment in that scope as long as it
73//! produces the collections expected by the rest of the framework. The rendering is handled by the
74//! `[crate::source::types::SourceRender::render] method.
75//!
76//! When rendering a source dataflow we expect three outputs. First, a health output, which is how
77//! the source communicates status updates about its health. Second, a data output, which is the
78//! main output of a source and contains the data that will eventually be recorded in the persist
79//! shard. Finally, an optional upper frontier output, which tracks the overall upstream upper
80//! frontier. When a source doesn't provide a dedicated progress output the framework derives one
81//! by observing the progress of the data output. This output (derived or not) is what drives
82//! reclocking. When a source provides a dedicated upper output, it can manage it independently of
83//! the data output frontier. For example, it's possible that a source implementation queries the
84//! upstream system to learn what are the latest offsets for and set the upper output based on
85//! that, even before having started the actual ingestion, which would be presented as data and
86//! progress trickling in via the data output.
87//!
88//! ```text
89//! resume upper
90//! ,--------------------.
91//! / |
92//! health ,----+---. |
93//! output | source | |
94//! ,-----------| reader | |
95//! / +--,---.-+ |
96//! / / \ |
97//! +-----/----+ data / \ upper |
98//! | health | output/ \ output |
99//! | operator | | \ |
100//! +----------+ | | |
101//! FromTime | | |
102//! scope | | |
103//! -------------------------------------|-----------|---------------|---
104//! IntoTime | | |
105//! scope | ,----+-----. |
106//! | | remap | |
107//! | | operator | |
108//! | +---,------+ |
109//! | / |
110//! | / bindings |
111//! | / |
112//! ,-+-----+--. |
113//! | reclock | |
114//! | operator | |
115//! +-,--,---.-+ |
116//! ,----------´.-´ \ |
117//! _.-´ .-´ \ |
118//! _.-´ .-´ \ |
119//! .-´ ,´ \ |
120//! / / \ |
121//! ,----------. ,----------. ,----------. |
122//! | decode | | decode | .... | decode | |
123//! | output 0 | | output 1 | | output N | |
124//! +-----+----+ +-----+----+ +-----+----+ |
125//! | | | |
126//! | | | |
127//! ,-----+----. ,-----+----. ,-----+----. |
128//! | envelope | | envelope | .... | envelope | |
129//! | output 0 | | output 1 | | output N | |
130//! +----------+ +-----+----+ +-----+----+ |
131//! | | | |
132//! | | | |
133//! ,-----+----. ,-----+----. ,-----+----. |
134//! | persist | | persist | .... | persist | |
135//! | sink 0 | | sink 1 | | sink N | |
136//! +-----+----+ +-----+----+ +-----+----+ |
137//! \ \ / |
138//! `-. `, / |
139//! `-._ `-. / |
140//! `-._ `-. / |
141//! `---------. `-. / |
142//! +`---`---+---, |
143//! | resume | |
144//! | calculator | |
145//! +------+-----+ |
146//! \ |
147//! `-------------------´
148//! ```
149//!
150//! #### Reclocking
151//!
152//! Whenever a dataflow edge crosses the scope boundaries it must first be converted into a
153//! captured stream via the `[mz_timely_util::capture::UnboundedTokioCapture`] utility. This
154//! disassociates the stream and its progress information from the original timely scope and allows
155//! it to be read from a different place. The downside of this mechanism is that it's invisible to
156//! timely's progress tracking, but that seems like a necessary evil if we want to do reclocking.
157//!
158//! The two main ways these tokio-fied streams are turned back into normal timely streams in the
159//! destination scope are by the `reclock operator` and the `remap operator` which process the
160//! `data output` and `upper output` of the source reader respectively.
161//!
162//! The `remap operator` reads the `upper output`, which is composed only of frontiers, mints new
163//! bindings, and writes them into the remap shard. The final durable timestamp bindings are
164//! emitted as its output for consumption by the `reclock operator`.
165//!
166//! The `reclock operator` reads the `data output`, which contains both data and progress
167//! statements, and uses the bindings it receives from the `remap operator` to reclock each piece
168//! of data and each frontier statement into the target scope's timestamp and emit the reclocked
169//! stream in its output.
170//!
171//! #### Partitioning
172//!
173//! At this point we have a timely stream with correctly timestamped data in the mz time domain
174//! (`mz_repr::Timestamp`) which contains multiplexed messages for each of the potential subsources
175//! of this source. Each message selects the output it belongs to by setting the output field in
176//! [`crate::source::types::SourceMessage`]. By convention, the main source output is always output
177//! zero and subsources get the outputs from one onwards.
178//!
179//! However, regardless of whether the output is the main source or a subsource it is treated
180//! identically by the pipeline. Each output is demultiplexed into its own timely stream using
181//! [`timely::dataflow::operators::core::partition::Partition`] and the rest of the ingestion pipeline is
182//! rendered independently.
183//!
184//! #### Resumption frontier
185//!
186//! At the end of each per-output dataflow fragment is an instance of `persist_sink`, which is
187//! responsible for writing the final `Row` data into the corresponding output shard. The durable
188//! upper of each of the output shards is then recombined in a way that calculates the minimum
189//! upper frontier between them. This is what we refer to as the "resumption frontier" or "resume
190//! upper" and at this stage it is expressed in terms of `IntoTime` timestamps. As a final step,
191//! this resumption frontier is converted back into a `FromTime` timestamped frontier using
192//! `ReclockFollower::source_upper_at_frontier` and connected back to the source reader operator.
193//! This frontier is what drives the `OffsetCommiter` which informs the upstream system to release
194//! resources until the specified offsets.
195//!
196//! ## Exports
197//!
198//! Not yet documented
199
200use std::collections::BTreeMap;
201use std::rc::Rc;
202use std::sync::Arc;
203
204use mz_ore::error::ErrorExt;
205use mz_repr::{GlobalId, Row};
206use mz_storage_types::controller::CollectionMetadata;
207use mz_storage_types::dyncfgs;
208use mz_storage_types::oneshot_sources::{OneshotIngestionDescription, OneshotIngestionRequest};
209use mz_storage_types::sinks::StorageSinkDesc;
210use mz_storage_types::sources::{GenericSourceConnection, IngestionDescription, SourceConnection};
211use mz_timely_util::antichain::AntichainExt;
212use mz_timely_util::scope_label::ScopeExt;
213use timely::PartialOrder;
214use timely::dataflow::operators::vec::Map;
215use timely::dataflow::operators::{Concatenate, ConnectLoop, Feedback, Leave};
216use timely::progress::Antichain;
217use timely::worker::Worker as TimelyWorker;
218use tokio::sync::Semaphore;
219
220use crate::healthcheck::{HealthStatusMessage, HealthStatusUpdate, StatusNamespace};
221use crate::source::RawSourceCreationConfig;
222use crate::storage_state::StorageState;
223
224mod persist_sink;
225pub mod sinks;
226pub mod sources;
227
228/// Assemble the "ingestion" side of a dataflow, i.e. the sources.
229///
230/// This method creates a new dataflow to host the implementations of sources for the `dataflow`
231/// argument, and returns assets for each source that can import the results into a new dataflow.
232pub fn build_ingestion_dataflow(
233 timely_worker: &mut TimelyWorker,
234 storage_state: &mut StorageState,
235 primary_source_id: GlobalId,
236 description: IngestionDescription<CollectionMetadata>,
237 as_of: Antichain<mz_repr::Timestamp>,
238 resume_uppers: BTreeMap<GlobalId, Antichain<mz_repr::Timestamp>>,
239 source_resume_uppers: BTreeMap<GlobalId, Vec<Row>>,
240) {
241 let worker_id = timely_worker.index();
242 let worker_logging = timely_worker.logger_for("timely").map(Into::into);
243 let debug_name = primary_source_id.to_string();
244 let name = format!("Source dataflow: {debug_name}");
245 timely_worker.dataflow_core(&name, worker_logging, Box::new(()), |_, root_scope| {
246 let root_scope = root_scope.with_label();
247
248 // Here we need to create two scopes. One timestamped with `()`, which is the root scope,
249 // and one timestamped with `mz_repr::Timestamp` which is the final scope of the dataflow.
250 // Refer to the module documentation for an explanation of this structure.
251 // The scope.clone() occurs to allow import in the region.
252 root_scope.clone().scoped(&name, |mz_scope| {
253 let debug_name = format!("{debug_name}-sources");
254
255 let mut tokens = vec![];
256
257 let (feedback_handle, feedback) = mz_scope.feedback(Default::default());
258
259 let connection = description.desc.connection.clone();
260 tracing::info!(
261 id = %primary_source_id,
262 as_of = %as_of.pretty(),
263 resume_uppers = ?resume_uppers,
264 source_resume_uppers = ?source_resume_uppers,
265 "timely-{worker_id} building {} source pipeline", connection.name(),
266 );
267
268 let busy_signal = if dyncfgs::SUSPENDABLE_SOURCES
269 .get(storage_state.storage_configuration.config_set())
270 {
271 Arc::new(Semaphore::new(1))
272 } else {
273 Arc::new(Semaphore::new(Semaphore::MAX_PERMITS))
274 };
275
276 // Only Postgres keeps its remap upper still through a snapshot, so only its ceiling
277 // ends near where the frontier lands when the pin lifts. MySQL and SQL Server tick
278 // through theirs, which would hold the shard upper at the as_of for the whole replay.
279 // TODO: include them once their snapshots run concurrently with CDC.
280 let oltp_source = matches!(connection, GenericSourceConnection::Postgres(_));
281
282 let base_source_config = RawSourceCreationConfig {
283 name: format!("{}-{}", connection.name(), primary_source_id),
284 id: primary_source_id,
285 source_exports: description.source_exports.clone(),
286 timestamp_interval: description.desc.timestamp_interval,
287 worker_id: mz_scope.index(),
288 worker_count: mz_scope.peers(),
289 now_fn: storage_state.now.clone(),
290 metrics: storage_state.metrics.clone(),
291 as_of: as_of.clone(),
292 resume_uppers: resume_uppers.clone(),
293 source_resume_uppers,
294 remap_metadata: description.remap_metadata.clone(),
295 persist_clients: Arc::clone(&storage_state.persist_clients),
296 statistics: storage_state
297 .aggregated_statistics
298 .get_ingestion_stats(&primary_source_id),
299 shared_remap_upper: Rc::clone(
300 &storage_state.source_uppers[&description.remap_collection_id],
301 ),
302 // This might quite a large clone, but its just during rendering
303 config: storage_state.storage_configuration.clone(),
304 remap_collection_id: description.remap_collection_id,
305 busy_signal: Arc::clone(&busy_signal),
306 };
307
308 let (outputs, source_health, remap_upper, source_tokens) = match connection {
309 GenericSourceConnection::Kafka(c) => crate::render::sources::render_source(
310 mz_scope,
311 root_scope,
312 &debug_name,
313 c,
314 description.clone(),
315 feedback,
316 storage_state,
317 base_source_config,
318 ),
319 GenericSourceConnection::Postgres(c) => crate::render::sources::render_source(
320 mz_scope,
321 root_scope,
322 &debug_name,
323 c,
324 description.clone(),
325 feedback,
326 storage_state,
327 base_source_config,
328 ),
329 GenericSourceConnection::MySql(c) => crate::render::sources::render_source(
330 mz_scope,
331 root_scope,
332 &debug_name,
333 c,
334 description.clone(),
335 feedback,
336 storage_state,
337 base_source_config,
338 ),
339 GenericSourceConnection::SqlServer(c) => crate::render::sources::render_source(
340 mz_scope,
341 root_scope,
342 &debug_name,
343 c,
344 description.clone(),
345 feedback,
346 storage_state,
347 base_source_config,
348 ),
349 GenericSourceConnection::LoadGenerator(c) => crate::render::sources::render_source(
350 mz_scope,
351 root_scope,
352 &debug_name,
353 c,
354 description.clone(),
355 feedback,
356 storage_state,
357 base_source_config,
358 ),
359 };
360 tokens.extend(source_tokens);
361
362 let mut upper_streams = vec![];
363 let mut health_streams = Vec::with_capacity(source_health.len() + outputs.len());
364 health_streams.extend(source_health);
365 for (export_id, (ok, err)) in outputs {
366 let export = &description.source_exports[&export_id];
367 let source_data = ok.map(Ok).concat(err.map(Err));
368
369 let metrics = storage_state.metrics.get_source_persist_sink_metrics(
370 export_id,
371 primary_source_id,
372 worker_id,
373 &export.storage_metadata.data_shard,
374 );
375
376 tracing::info!(
377 id = %primary_source_id,
378 "timely-{worker_id}: persisting export {} of {}",
379 export_id,
380 primary_source_id
381 );
382
383 // An export snapshots when its resume upper is at or below the as_of. The
384 // controller uses the same test to hand the connector a minimum from-time resume
385 // upper.
386 let snapshotting = oltp_source
387 && resume_uppers
388 .get(&export_id)
389 .is_some_and(|upper| PartialOrder::less_equal(upper, &as_of));
390 let (upper_stream, errors, sink_tokens) = crate::render::persist_sink::render(
391 mz_scope,
392 export_id,
393 export.storage_metadata.clone(),
394 source_data,
395 storage_state,
396 metrics,
397 Arc::clone(&busy_signal),
398 snapshotting.then(|| as_of.clone()),
399 description.desc.timestamp_interval,
400 remap_upper.clone(),
401 );
402 upper_streams.push(upper_stream);
403 tokens.extend(sink_tokens);
404
405 let sink_health = errors.map(move |err: Rc<anyhow::Error>| {
406 let halt_status =
407 HealthStatusUpdate::halting(err.display_with_causes().to_string(), None);
408 HealthStatusMessage {
409 id: None,
410 namespace: StatusNamespace::Internal,
411 update: halt_status,
412 }
413 });
414 health_streams.push(sink_health.leave(root_scope));
415 }
416
417 mz_scope
418 .concatenate(upper_streams)
419 .connect_loop(feedback_handle);
420
421 let health_stream = root_scope.concatenate(health_streams);
422 let health_token = crate::healthcheck::health_operator(
423 root_scope,
424 storage_state.now.clone(),
425 resume_uppers
426 .iter()
427 .filter_map(|(id, frontier)| {
428 // If the collection isn't closed, then we will remark it as Starting as
429 // the dataflow comes up.
430 (!frontier.is_empty()).then_some(*id)
431 })
432 .collect(),
433 primary_source_id,
434 "source",
435 health_stream,
436 crate::healthcheck::DefaultWriter {
437 command_tx: storage_state.internal_cmd_tx.clone(),
438 updates: Rc::clone(&storage_state.shared_status_updates),
439 },
440 storage_state
441 .storage_configuration
442 .parameters
443 .record_namespaced_errors,
444 dyncfgs::STORAGE_SUSPEND_AND_RESTART_DELAY
445 .get(storage_state.storage_configuration.config_set()),
446 );
447 tokens.push(health_token);
448
449 storage_state
450 .source_tokens
451 .insert(primary_source_id, tokens);
452 })
453 });
454}
455
456/// do the export dataflow thing
457pub fn build_export_dataflow(
458 timely_worker: &mut TimelyWorker,
459 storage_state: &mut StorageState,
460 id: GlobalId,
461 description: StorageSinkDesc<CollectionMetadata, mz_repr::Timestamp>,
462) {
463 let worker_logging = timely_worker.logger_for("timely").map(Into::into);
464 let debug_name = id.to_string();
465 let name = format!("Source dataflow: {debug_name}");
466 timely_worker.dataflow_core(&name, worker_logging, Box::new(()), |_, scope| {
467 let scope = scope.with_label();
468
469 let mut tokens = vec![];
470 let (health_stream, sink_tokens) =
471 crate::render::sinks::render_sink(scope, storage_state, id, &description);
472 tokens.extend(sink_tokens);
473
474 // Note that sinks also have only 1 active worker, which simplifies the work that
475 // `health_operator` has to do internally.
476 let health_token = crate::healthcheck::health_operator(
477 scope,
478 storage_state.now.clone(),
479 [id].into_iter().collect(),
480 id,
481 "sink",
482 health_stream,
483 crate::healthcheck::DefaultWriter {
484 command_tx: storage_state.internal_cmd_tx.clone(),
485 updates: Rc::clone(&storage_state.shared_status_updates),
486 },
487 storage_state
488 .storage_configuration
489 .parameters
490 .record_namespaced_errors,
491 dyncfgs::STORAGE_SUSPEND_AND_RESTART_DELAY
492 .get(storage_state.storage_configuration.config_set()),
493 );
494 tokens.push(health_token);
495
496 storage_state.sink_tokens.insert(id, tokens);
497 });
498}
499
500pub(crate) fn build_oneshot_ingestion_dataflow(
501 timely_worker: &mut TimelyWorker,
502 storage_state: &mut StorageState,
503 ingestion_id: uuid::Uuid,
504 collection_id: GlobalId,
505 collection_meta: CollectionMetadata,
506 description: OneshotIngestionRequest,
507) {
508 let (results_tx, results_rx) = tokio::sync::mpsc::unbounded_channel();
509 let callback = move |result| {
510 // TODO(cf3): Do we care if the receiver has gone away?
511 //
512 // Persist is working on cleaning up leaked blobs, we could also use `OneshotReceiverExt`
513 // here, but that might run into the infamous async-Drop problem.
514 let _ = results_tx.send(result);
515 };
516 let connection_context = storage_state
517 .storage_configuration
518 .connection_context
519 .clone();
520 let enforce_external_addresses = mz_storage_types::dyncfgs::ENFORCE_EXTERNAL_ADDRESSES
521 .get(storage_state.storage_configuration.config_set());
522
523 let name = format!("Oneshot ingestion: {ingestion_id}");
524 let tokens = timely_worker.dataflow_named(&name, |scope| {
525 let scope = scope.with_label();
526 mz_storage_operators::oneshot_source::render(
527 scope,
528 Arc::clone(&storage_state.persist_clients),
529 connection_context,
530 collection_id,
531 collection_meta,
532 description,
533 enforce_external_addresses,
534 callback,
535 )
536 });
537 let ingestion_description = OneshotIngestionDescription {
538 tokens,
539 results: results_rx,
540 };
541
542 storage_state
543 .oneshot_ingestions
544 .insert(ingestion_id, ingestion_description);
545}