1use 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
42pub 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 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 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 pub traces: BTreeMap<LogVariant, TraceBundle>,
123 pub dataflow_index: usize,
125 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 fn step_logger(&self) -> Option<TimelyLogger> {
231 if let Some(logger) = self.worker.logging() {
232 return Some(logger);
235 }
236
237 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 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 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 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 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 activator.activate();
335 } else if logger.report_progress(*time) {
336 activator.activate();
337 }
338 },
339 )
340 }
341
342 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 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 activator.activate();
390 }
391 }
392 } else {
393 while let Some(container) = builder.finish() {
395 logger.publish_batch(std::mem::take(container));
396 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
410trait ExtractTimestamp: Clone + 'static {
412 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}