Skip to main content

mz_compute/logging/
initialize.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//! Initialization of logging dataflows.
7
8use std::cell::RefCell;
9use std::collections::BTreeMap;
10use std::rc::Rc;
11use std::time::{Duration, Instant};
12
13use differential_dataflow::VecCollection;
14use differential_dataflow::dynamic::pointstamp::PointStamp;
15use differential_dataflow::logging::{DifferentialEvent, DifferentialEventBuilder};
16use mz_compute_client::logging::{LogVariant, LoggingConfig};
17use mz_dyncfg::ConfigSet;
18use mz_ore::metrics::MetricsRegistry;
19use mz_repr::{Diff, Timestamp};
20use mz_storage_operators::persist_source::Subtime;
21use mz_timely_util::columnar::Column;
22use mz_timely_util::columnar::builder::ColumnBuilder;
23use mz_timely_util::columnation::ColumnationChunker;
24use mz_timely_util::operator::CollectionExt;
25use mz_timely_util::scope_label::ScopeExt;
26use prometheus::IntCounter;
27use timely::ContainerBuilder;
28use timely::container::{ContainerBuilder as _, PushInto};
29use timely::logging::{StartStop, TimelyEvent, TimelyEventBuilder, TimelyLogger};
30use timely::logging_core::{Logger, Registry};
31use timely::order::Product;
32use timely::progress::reachability::logging::{TrackerEvent, TrackerEventBuilder};
33
34use crate::arrangement::manager::TraceBundle;
35use crate::extensions::arrange::{KeyCollection, MzArrange};
36use crate::logging::compute::{ComputeEvent, ComputeEventBuilder};
37use crate::logging::{BatchLogger, EventQueue, SharedLoggingState};
38use crate::metrics::LoggingMetrics;
39use crate::render::errors::DataflowErrorSer;
40use crate::typedefs::{ErrBatcher, ErrBuilder};
41
42/// Initialize logging dataflows.
43///
44/// Returns a logger for compute events, and for each `LogVariant` a trace bundle usable for
45/// retrieving logged records as well as the index of the exporting dataflow.
46pub fn initialize(
47    worker: &mut timely::worker::Worker,
48    config: &LoggingConfig,
49    metrics_registry: MetricsRegistry,
50    metrics: LoggingMetrics,
51    worker_config: Rc<ConfigSet>,
52    workers_per_process: usize,
53) -> LoggingTraces {
54    let interval_ms = std::cmp::max(1, config.interval.as_millis());
55
56    // Track time relative to the Unix epoch, rather than when the server
57    // started, so that the logging sources can be joined with tables and
58    // other real time sources for semi-sensible results.
59    let now = Instant::now();
60    let start_offset = std::time::SystemTime::now()
61        .duration_since(std::time::SystemTime::UNIX_EPOCH)
62        .expect("Failed to get duration since Unix epoch");
63
64    let mut context = LoggingContext {
65        worker,
66        config,
67        interval_ms,
68        now,
69        start_offset,
70        t_event_queue: EventQueue::new("t"),
71        r_event_queue: EventQueue::new("r"),
72        d_event_queue: EventQueue::new("d"),
73        c_event_queue: EventQueue::new("c"),
74        shared_state: Default::default(),
75        metrics_registry,
76        metrics,
77        worker_config,
78        workers_per_process,
79    };
80
81    // Depending on whether we should log the creation of the logging dataflows, we register the
82    // loggers with timely either before or after creating them.
83    let dataflow_index = context.worker.next_dataflow_index();
84    let traces = if config.log_logging {
85        context.register_loggers();
86        context.construct_dataflow()
87    } else {
88        let traces = context.construct_dataflow();
89        context.register_loggers();
90        traces
91    };
92
93    let compute_logger = worker.logger_for("materialize/compute").unwrap();
94    LoggingTraces {
95        traces,
96        dataflow_index,
97        compute_logger,
98    }
99}
100
101pub(super) type ReachabilityEvent = (usize, Vec<(usize, usize, bool, Timestamp, Diff)>);
102
103struct LoggingContext<'a> {
104    worker: &'a mut timely::worker::Worker,
105    config: &'a LoggingConfig,
106    interval_ms: u128,
107    now: Instant,
108    start_offset: Duration,
109    t_event_queue: EventQueue<Vec<(Duration, TimelyEvent)>>,
110    r_event_queue: EventQueue<Column<(Duration, ReachabilityEvent)>, 3>,
111    d_event_queue: EventQueue<Vec<(Duration, DifferentialEvent)>>,
112    c_event_queue: EventQueue<Column<(Duration, ComputeEvent)>>,
113    shared_state: Rc<RefCell<SharedLoggingState>>,
114    metrics_registry: MetricsRegistry,
115    metrics: LoggingMetrics,
116    worker_config: Rc<ConfigSet>,
117    workers_per_process: usize,
118}
119
120pub(crate) struct LoggingTraces {
121    /// Exported traces, by log variant.
122    pub traces: BTreeMap<LogVariant, TraceBundle>,
123    /// The index of the dataflow that exports the traces.
124    pub dataflow_index: usize,
125    /// The compute logger.
126    pub compute_logger: super::compute::Logger,
127}
128
129impl LoggingContext<'_> {
130    fn construct_dataflow(&mut self) -> BTreeMap<LogVariant, TraceBundle> {
131        let step_logger = self.step_logger();
132        self.worker.dataflow_core(
133            "Dataflow: logging",
134            step_logger,
135            Box::new(()),
136            |_, scope| {
137                let scope = scope.with_label();
138
139                let mut collections = BTreeMap::new();
140
141                let super::timely::Return {
142                    collections: timely_collections,
143                } = super::timely::construct(
144                    scope,
145                    self.config,
146                    self.t_event_queue.clone(),
147                    Rc::clone(&self.shared_state),
148                );
149                collections.extend(timely_collections);
150
151                let super::reachability::Return {
152                    collections: reachability_collections,
153                } = super::reachability::construct(scope, self.config, self.r_event_queue.clone());
154                collections.extend(reachability_collections);
155
156                let super::differential::Return {
157                    collections: differential_collections,
158                } = super::differential::construct(
159                    scope,
160                    self.config,
161                    self.d_event_queue.clone(),
162                    Rc::clone(&self.shared_state),
163                );
164                collections.extend(differential_collections);
165
166                let super::compute::Return {
167                    collections: compute_collections,
168                } = super::compute::construct(
169                    scope.clone(),
170                    scope.activations(),
171                    self.config,
172                    self.c_event_queue.clone(),
173                    Rc::clone(&self.shared_state),
174                );
175                collections.extend(compute_collections);
176
177                let super::prometheus::Return {
178                    collections: prometheus_collections,
179                } = super::prometheus::construct(
180                    scope,
181                    self.config,
182                    self.metrics_registry.clone(),
183                    self.now,
184                    self.start_offset,
185                    Rc::clone(&self.worker_config),
186                    self.workers_per_process,
187                );
188                collections.extend(prometheus_collections);
189
190                let super::resource_usage::Return {
191                    collections: resource_usage_collections,
192                } = super::resource_usage::construct(
193                    scope,
194                    self.config,
195                    self.now,
196                    self.start_offset,
197                    self.workers_per_process,
198                );
199                collections.extend(resource_usage_collections);
200
201                let errs = scope.scoped("logging errors", |scope| {
202                    let collection: KeyCollection<_, DataflowErrorSer, Diff> =
203                        VecCollection::empty(scope).into();
204                    collection
205                        .mz_arrange::<ColumnationChunker<_>, ErrBatcher<_, _>, ErrBuilder<_, _>, _>(
206                            "Arrange logging err",
207                        )
208                        .trace
209                });
210
211                let traces = collections
212                    .into_iter()
213                    .map(|(log, collection)| {
214                        let bundle = TraceBundle::new(collection.trace, errs.clone())
215                            .with_drop(collection.token);
216                        (log, bundle)
217                    })
218                    .collect();
219                traces
220            },
221        )
222    }
223
224    /// Construct the timely logger that the worker hands to the logging dataflow itself.
225    ///
226    /// The logging dataflow's operators log nothing unless `log_logging` is set, but the worker
227    /// logs every scheduling of the dataflow as a whole to this logger, which observes the
228    /// duration of each scheduling. With `log_logging` set, returns the worker's timely logger
229    /// instead and observes nothing.
230    fn step_logger(&self) -> Option<TimelyLogger> {
231        if let Some(logger) = self.worker.logging() {
232            // Forwarding events from a second logger would re-timestamp them at flush time, and
233            // `mz_scheduling_elapsed` already covers the logging dataflow in this mode.
234            return Some(logger);
235        }
236
237        // `dataflow_core` allocates the dataflow's identifier first.
238        let dataflow_id = self.worker.peek_identifier();
239        let step_duration_seconds = self.metrics.step_duration_seconds.clone();
240        let mut started = None;
241        let logger = Logger::<TimelyEventBuilder>::new(
242            self.now,
243            self.start_offset,
244            move |_time, data: &mut Option<Vec<(Duration, TimelyEvent)>>| {
245                let Some(data) = data else { return };
246                for (time, event) in data.drain(..) {
247                    if let TimelyEvent::Schedule(schedule) = event
248                        && schedule.id == dataflow_id
249                    {
250                        match schedule.start_stop {
251                            StartStop::Start => started = Some(time),
252                            StartStop::Stop => {
253                                if let Some(start) = started.take() {
254                                    let elapsed = time.saturating_sub(start);
255                                    step_duration_seconds.observe(elapsed.as_secs_f64());
256                                }
257                            }
258                        }
259                    }
260                }
261            },
262        );
263        // The worker flushes only registered loggers at the end of each step. Unregistered,
264        // observations would wait for the logger's buffer to fill.
265        let mut register = self.worker.log_register().expect("Logging must be enabled");
266        register.insert_logger("materialize/logging-step", logger.clone());
267        Some(logger.into())
268    }
269
270    /// Construct a new reachability logger for timestamp type `T`.
271    ///
272    /// Inserts a logger with the name `timely/reachability/{type_name::<T>()}`, following
273    /// Timely naming convention.
274    fn register_reachability_logger<T: ExtractTimestamp>(
275        &self,
276        registry: &mut Registry,
277        index: usize,
278    ) {
279        let logger = self.reachability_logger::<T>(index);
280        let type_name = std::any::type_name::<T>();
281        registry.insert_logger(&format!("timely/reachability/{type_name}"), logger);
282    }
283
284    /// Register all loggers with the timely worker.
285    ///
286    /// Registers the timely, differential, compute, and reachability loggers.
287    fn register_loggers(&self) {
288        let t_logger = self.simple_logger::<TimelyEventBuilder>(
289            self.t_event_queue.clone(),
290            self.metrics.timely_records_total.clone(),
291        );
292        let d_logger = self.simple_logger::<DifferentialEventBuilder>(
293            self.d_event_queue.clone(),
294            self.metrics.differential_records_total.clone(),
295        );
296        let c_logger = self.simple_logger::<ComputeEventBuilder>(
297            self.c_event_queue.clone(),
298            self.metrics.compute_records_total.clone(),
299        );
300
301        let mut register = self.worker.log_register().expect("Logging must be enabled");
302        register.insert_logger("timely", t_logger);
303        // Note that each reachability logger has a unique index, this is crucial to avoid dropping
304        // data because the event link structure is not multi-producer safe.
305        self.register_reachability_logger::<Timestamp>(&mut register, 0);
306        self.register_reachability_logger::<Product<Timestamp, PointStamp<u64>>>(&mut register, 1);
307        self.register_reachability_logger::<(Timestamp, Subtime)>(&mut register, 2);
308        register.insert_logger("differential/arrange", d_logger);
309        register.insert_logger("materialize/compute", c_logger.clone());
310
311        self.shared_state.borrow_mut().compute_logger = Some(c_logger);
312    }
313
314    fn simple_logger<CB: ContainerBuilder>(
315        &self,
316        event_queue: EventQueue<CB::Container>,
317        records_total: IntCounter,
318    ) -> Logger<CB> {
319        let [link] = event_queue.links;
320        let mut logger = BatchLogger::new(link, self.interval_ms, records_total);
321        let activator = event_queue.activator.clone();
322        Logger::new(
323            self.now,
324            self.start_offset,
325            move |time, data: &mut Option<CB::Container>| {
326                if let Some(data) = data.take() {
327                    logger.publish_batch(data);
328                    // Count every batch towards the replay's activation threshold, so the logging
329                    // dataflow drains events in bounded chunks. Without this, the replay only
330                    // wakes once per logging interval and processes the whole interval's events
331                    // in one uninterruptible call, stalling every other dataflow on the worker.
332                    // The activator is worker-local and never unparks the thread: a threshold
333                    // crossed while flushing before a park takes effect on the next wakeup.
334                    activator.activate();
335                } else if logger.report_progress(*time) {
336                    activator.activate();
337                }
338            },
339        )
340    }
341
342    /// Construct a reachability logger for timestamp type `T`. The index must
343    /// refer to a unique link in the reachability event queue.
344    fn reachability_logger<T>(&self, index: usize) -> Logger<TrackerEventBuilder<T>>
345    where
346        T: ExtractTimestamp,
347    {
348        let link = Rc::clone(&self.r_event_queue.links[index]);
349        let mut logger = BatchLogger::new(
350            link,
351            self.interval_ms,
352            self.metrics.reachability_records_total.clone(),
353        );
354        let mut massaged = Vec::new();
355        let mut builder = ColumnBuilder::default();
356        let activator = self.r_event_queue.activator.clone();
357
358        let action = move |batch_time: &Duration, data: &mut Option<Vec<_>>| {
359            if let Some(data) = data {
360                // Handle data
361                for (time, event) in data.drain(..) {
362                    match event {
363                        TrackerEvent::SourceUpdate(update) => {
364                            massaged.extend(update.updates.iter().map(
365                                |(node, port, time, diff)| {
366                                    let is_source = true;
367                                    (*node, *port, is_source, T::extract(time), Diff::from(*diff))
368                                },
369                            ));
370
371                            builder.push_into((time, (update.tracker_id, &massaged)));
372                            massaged.clear();
373                        }
374                        TrackerEvent::TargetUpdate(update) => {
375                            massaged.extend(update.updates.iter().map(
376                                |(node, port, time, diff)| {
377                                    let is_source = false;
378                                    (*node, *port, is_source, time.extract(), Diff::from(*diff))
379                                },
380                            ));
381
382                            builder.push_into((time, (update.tracker_id, &massaged)));
383                            massaged.clear();
384                        }
385                    }
386                    while let Some(container) = builder.extract() {
387                        logger.publish_batch(std::mem::take(container));
388                        // See `simple_logger`.
389                        activator.activate();
390                    }
391                }
392            } else {
393                // Handle a flush
394                while let Some(container) = builder.finish() {
395                    logger.publish_batch(std::mem::take(container));
396                    // See `simple_logger`.
397                    activator.activate();
398                }
399
400                if logger.report_progress(*batch_time) {
401                    activator.activate();
402                }
403            }
404        };
405
406        Logger::new(self.now, self.start_offset, action)
407    }
408}
409
410/// Helper trait to extract a timestamp from various types of timestamp used in rendering.
411trait ExtractTimestamp: Clone + 'static {
412    /// Extracts the timestamp from the type.
413    fn extract(&self) -> Timestamp;
414}
415
416impl ExtractTimestamp for Timestamp {
417    fn extract(&self) -> Timestamp {
418        *self
419    }
420}
421
422impl ExtractTimestamp for Product<Timestamp, PointStamp<u64>> {
423    fn extract(&self) -> Timestamp {
424        self.outer
425    }
426}
427
428impl ExtractTimestamp for (Timestamp, Subtime) {
429    fn extract(&self) -> Timestamp {
430        self.0
431    }
432}