Skip to main content

mz_compute/
logging.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//! Logging dataflows for events generated by various subsystems.
11
12pub mod compute;
13mod differential;
14pub(super) mod initialize;
15mod prometheus;
16mod reachability;
17mod resource_usage;
18mod timely;
19
20use std::any::Any;
21use std::collections::BTreeMap;
22use std::marker::PhantomData;
23use std::rc::Rc;
24use std::time::{Duration, Instant};
25
26use ::prometheus::IntCounter;
27use ::timely::container::{CapacityContainerBuilder, PushInto};
28use ::timely::dataflow::Stream;
29use ::timely::dataflow::channels::pact::Pipeline;
30use ::timely::dataflow::operators::capture::{Event, EventLink, EventPusher};
31use ::timely::dataflow::operators::generic::Session;
32use ::timely::dataflow::operators::{Capability, CapabilityTrait, InputCapability, Operator};
33use ::timely::progress::Timestamp as TimelyTimestamp;
34use ::timely::scheduling::Activator;
35use ::timely::{Container, ContainerBuilder};
36use differential_dataflow::trace::Batcher;
37use mz_compute_client::logging::{ComputeLog, DifferentialLog, LogVariant, TimelyLog};
38use mz_expr::{MirScalarExpr, permutation_for_arrangement};
39use mz_repr::{Datum, Diff, Row, RowPacker, RowRef, Timestamp};
40use mz_timely_util::activator::RcActivator;
41use mz_timely_util::columnar::builder::ColumnBuilder;
42use mz_timely_util::operator::consolidate_pact;
43
44use crate::logging::compute::Logger as ComputeLogger;
45use crate::typedefs::RowRowAgent;
46
47pub use crate::logging::initialize::initialize;
48
49/// An update of value `D` at a time and with a diff.
50pub(super) type Update<D> = (D, Timestamp, Diff);
51/// A pusher for containers `C`.
52/// An output session for the specified container builder.
53pub(super) type OutputSession<'a, 'b, CB> =
54    Session<'a, 'b, Timestamp, CB, InputCapability<Timestamp>>;
55/// An output session for vector-based containers of updates `D`, using a capacity container builder.
56pub(super) type OutputSessionVec<'a, 'b, D> =
57    OutputSession<'a, 'b, CapacityContainerBuilder<Vec<D>>>;
58/// An output session for columnar containers of updates `D`, using a column builder.
59pub(super) type OutputSessionColumnar<'a, 'b, D> = OutputSession<'a, 'b, ColumnBuilder<D>>;
60
61/// Logs events as a timely stream, with progress statements.
62struct BatchLogger<C, P>
63where
64    P: EventPusher<Timestamp, C>,
65{
66    /// Time in milliseconds of the current expressed capability.
67    time_ms: Timestamp,
68    /// Pushes events to the logging dataflow.
69    event_pusher: P,
70    /// Each time is advanced to the strictly next millisecond that is a multiple of this interval.
71    /// This means we should be able to perform the same action on timestamp capabilities, and only
72    /// flush buffers when this timestamp advances.
73    interval_ms: u128,
74    /// Counts the records of published batches.
75    records_total: IntCounter,
76    _marker: PhantomData<C>,
77}
78
79impl<C, P> BatchLogger<C, P>
80where
81    P: EventPusher<Timestamp, C>,
82{
83    /// Creates a new batch logger.
84    fn new(event_pusher: P, interval_ms: u128, records_total: IntCounter) -> Self {
85        BatchLogger {
86            time_ms: Timestamp::minimum(),
87            event_pusher,
88            interval_ms,
89            records_total,
90            _marker: PhantomData,
91        }
92    }
93}
94
95impl<C, P> BatchLogger<C, P>
96where
97    P: EventPusher<Timestamp, C>,
98    C: Container,
99{
100    /// Publishes a batch of logged events.
101    fn publish_batch(&mut self, data: C) {
102        let records = u64::try_from(data.record_count()).expect("record count is non-negative");
103        self.records_total.inc_by(records);
104        self.event_pusher.push(Event::Messages(self.time_ms, data));
105    }
106
107    /// Indicate progress up to `time`, advances the capability.
108    ///
109    /// Returns `true` if the capability was advanced.
110    fn report_progress(&mut self, time: Duration) -> bool {
111        let time_ms = ((time.as_millis() / self.interval_ms) + 1) * self.interval_ms;
112        let new_time_ms: Timestamp = time_ms.try_into().expect("must fit");
113        if self.time_ms < new_time_ms {
114            self.event_pusher
115                .push(Event::Progress(vec![(new_time_ms, 1), (self.time_ms, -1)]));
116            self.time_ms = new_time_ms;
117            true
118        } else {
119            false
120        }
121    }
122}
123
124impl<C, P> Drop for BatchLogger<C, P>
125where
126    P: EventPusher<Timestamp, C>,
127{
128    fn drop(&mut self) {
129        self.event_pusher
130            .push(Event::Progress(vec![(self.time_ms, -1)]));
131    }
132}
133
134/// Parts to connect a logging dataflows the timely runtime.
135///
136/// This is just a bundle-type intended to make passing around its contents in the logging
137/// initialization code more convenient.
138///
139/// The `N` type parameter specifies the number of links to create for the event queue. We need
140/// separate links for queues that feed from multiple loggers because the `EventLink` type is not
141/// multi-producer safe (it is a linked-list, and multiple writers would blindly append, replacing
142/// existing new data, and cutting off other writers).
143#[derive(Clone)]
144struct EventQueue<C, const N: usize = 1> {
145    links: [Rc<EventLink<Timestamp, C>>; N],
146    activator: RcActivator,
147}
148
149impl<C, const N: usize> EventQueue<C, N> {
150    fn new(name: &str) -> Self {
151        let activator_name = format!("{name}_activator");
152        let activate_after = 128;
153        Self {
154            links: [(); N].map(|_| Rc::new(EventLink::new())),
155            activator: RcActivator::new(activator_name, activate_after),
156        }
157    }
158}
159
160/// State shared between different logging dataflow fragments.
161#[derive(Default)]
162struct SharedLoggingState {
163    /// Activators for arrangement heap size operators.
164    arrangement_size_activators: BTreeMap<usize, Activator>,
165    /// Shared compute logger.
166    compute_logger: Option<ComputeLogger>,
167}
168
169/// Helper to pack collections of [`Datum`]s into key and value row.
170pub(crate) struct PermutedRowPacker {
171    key: Vec<usize>,
172    value: Vec<usize>,
173    key_row: Row,
174    value_row: Row,
175}
176
177impl PermutedRowPacker {
178    /// Construct based on the information within the log variant.
179    pub(crate) fn new<V: Into<LogVariant>>(variant: V) -> Self {
180        let variant = variant.into();
181        let key = variant.index_by();
182        let (_, value) = permutation_for_arrangement(
183            &key.iter()
184                .cloned()
185                .map(MirScalarExpr::column)
186                .collect::<Vec<_>>(),
187            variant.desc().arity(),
188        );
189        Self {
190            key,
191            value,
192            key_row: Row::default(),
193            value_row: Row::default(),
194        }
195    }
196
197    /// Pack a slice of datums suitable for the key columns in the log variant.
198    pub(crate) fn pack_slice(&mut self, datums: &[Datum]) -> (&RowRef, &RowRef) {
199        self.pack_by_index(|packer, index| packer.push(datums[index]))
200    }
201
202    /// Pack using a callback suitable for the key columns in the log variant.
203    pub(crate) fn pack_by_index<F: Fn(&mut RowPacker, usize)>(
204        &mut self,
205        logic: F,
206    ) -> (&RowRef, &RowRef) {
207        let mut packer = self.key_row.packer();
208        for index in &self.key {
209            logic(&mut packer, *index);
210        }
211
212        let mut packer = self.value_row.packer();
213        for index in &self.value {
214            logic(&mut packer, *index);
215        }
216
217        (&self.key_row, &self.value_row)
218    }
219}
220
221/// Downgrade `cap` to the next logging-interval boundary and schedule the operator's next
222/// activation there. Returns the time the capability now holds.
223///
224/// `now` and `start_offset` must be the ones the logging dataflow was constructed with, so that
225/// every collection in it reports on the same boundaries. Scheduling off the boundary rather than
226/// off a fixed delay keeps the output frontier progressing at the logging rate without drifting
227/// from wall-clock elapsed time.
228///
229/// NOTE: downgrading the capability asserts the collection is complete up to the new time, so an
230/// operator that samples less often than the logging interval publishes a stale value rather than
231/// withholding it.
232pub(super) fn downgrade_to_interval_boundary(
233    cap: &mut Capability<Timestamp>,
234    activator: &Activator,
235    now: Instant,
236    start_offset: Duration,
237    interval_ms: u128,
238) -> Timestamp {
239    let elapsed = now.elapsed().as_millis();
240    let time_ms: u128 = ((elapsed + start_offset.as_millis()) / interval_ms + 1) * interval_ms;
241    let ts: Timestamp = time_ms.try_into().expect("must fit");
242    cap.downgrade(&ts);
243
244    let next_boundary_ms = time_ms - start_offset.as_millis();
245    let next_activation =
246        now + Duration::from_millis(next_boundary_ms.try_into().expect("must fit"));
247    activator.activate_after(next_activation.saturating_duration_since(Instant::now()));
248
249    ts
250}
251
252/// Emit the difference between two snapshots of a sampled source as updates at `ts`.
253///
254/// A sampled source reports its whole state on every read, so a changed value has to be expressed
255/// as a retraction of the previous one paired with an insertion of the new one. A key absent from
256/// `current` is retracted without a replacement, so a source that stops being readable drops out of
257/// the collection rather than lingering at its last value.
258///
259/// `pack` takes the packer as an argument instead of closing over it, because it hands back rows
260/// borrowed from it.
261pub(super) fn emit_snapshot_diff<K, V, CB, P, F>(
262    session: &mut Session<'_, '_, Timestamp, CB, P>,
263    packer: &mut PermutedRowPacker,
264    prev: &BTreeMap<K, V>,
265    current: &BTreeMap<K, V>,
266    ts: Timestamp,
267    pack: F,
268) where
269    K: Ord,
270    V: PartialEq,
271    CB: ContainerBuilder + for<'a> PushInto<((&'a RowRef, &'a RowRef), Timestamp, Diff)>,
272    P: CapabilityTrait<Timestamp>,
273    F: for<'a> Fn(&'a mut PermutedRowPacker, &K, &V) -> (&'a RowRef, &'a RowRef),
274{
275    for (key, value) in prev {
276        if current.get(key) != Some(value) {
277            let row = pack(packer, key, value);
278            session.give((row, ts, Diff::MINUS_ONE));
279        }
280    }
281    for (key, value) in current {
282        if prev.get(key) != Some(value) {
283            let row = pack(packer, key, value);
284            session.give((row, ts, Diff::ONE));
285        }
286    }
287}
288
289/// Information about a collection exported from a logging dataflow.
290struct LogCollection {
291    /// Trace handle providing access to the logged records.
292    trace: RowRowAgent<Timestamp, Diff>,
293    /// Token that should be dropped to drop this collection.
294    token: Rc<dyn Any>,
295}
296
297/// A single-purpose function to consolidate and pack updates for log collection.
298///
299/// The function first consolidates worker-local updates using the [`Pipeline`] pact, then converts
300/// the updates into `(Row, Row)` pairs using the provided logic function. It is crucial that the
301/// data is not exchanged between workers, as the consolidation would not function as desired
302/// otherwise.
303pub(super) fn consolidate_and_pack<'scope, Chu, B, CB, L, F, C>(
304    input: Stream<'scope, Timestamp, C>,
305    log: L,
306    mut logic: F,
307) -> Stream<'scope, Timestamp, CB::Container>
308where
309    B: Batcher<Time = Timestamp> + 'static,
310    Chu: ContainerBuilder<Container = B::Output> + for<'a> PushInto<&'a mut C> + 'static,
311    C: Container + Clone + 'static,
312    B::Output: Clone,
313    CB: ContainerBuilder,
314    L: Into<LogVariant>,
315    F: FnMut(B::Output, &mut PermutedRowPacker, &mut OutputSession<CB>) + 'static,
316{
317    let log = log.into();
318    // TODO: Use something other than the debug representation of the log variant as a name.
319    let c_name = &format!("Consolidate {log:?}");
320    let u_name = &format!("ToRow {log:?}");
321    let mut packer = PermutedRowPacker::new(log);
322    let consolidated = consolidate_pact::<Chu, B, _, _>(input, Pipeline, c_name);
323    consolidated.unary::<CB, _, _, _>(Pipeline, u_name, |_, _| {
324        move |input, output| {
325            input.for_each_time(|time, data| {
326                let mut session = output.session_with_builder(&time);
327                for item in data.flatten().flat_map(|data| data.drain(..)) {
328                    logic(item, &mut packer, &mut session);
329                }
330            });
331        }
332    })
333}