Skip to main content

mz_timely_util/columnar/chunk/
metrics.rs

1// Copyright Materialize, Inc. and contributors. All rights reserved.
2// Use of this software is governed by the Business Source License
3// included in the LICENSE file.
4
5//! Process-wide work counters for columnar chunk maintenance.
6//!
7//! Counts describe work attempted, including repeated processing of the same
8//! updates. Bytes are serialized uncompressed sizes, not extent or device
9//! writes, and only the commit stage records them. Row-size buckets are
10//! disjoint, indexed by ceil(log2(rows)) with zero and one sharing bucket zero
11//! and the last bucket absorbing everything larger. Counters at different
12//! stages must not be added together.
13
14use std::sync::atomic::{AtomicU64, Ordering};
15
16use mz_ore::cast::CastFrom;
17use mz_ore::metric;
18use mz_ore::metrics::{ComputedUIntGauge, MetricsRegistry};
19
20#[derive(Clone, Copy)]
21pub(super) enum Stage {
22    Merge,
23    Advance,
24    Commit,
25    Batch,
26}
27
28impl Stage {
29    const ALL: [(Self, &'static str); 4] = [
30        (Self::Merge, "merge"),
31        (Self::Advance, "advance"),
32        (Self::Commit, "commit"),
33        (Self::Batch, "batch"),
34    ];
35
36    fn counters(self) -> &'static Counters {
37        let index = match self {
38            Self::Merge => 0,
39            Self::Advance => 1,
40            Self::Commit => 2,
41            Self::Batch => 3,
42        };
43        &COUNTERS[index]
44    }
45}
46
47/// Buckets up to 2^24 rows, the last one absorbing anything larger. Chunks
48/// are bounded to 2 MiB, so only published batches can reach the last bucket.
49const BUCKETS: usize = 25;
50
51struct Counters {
52    calls: AtomicU64,
53    rows: AtomicU64,
54    bytes: AtomicU64,
55    buckets: [AtomicU64; BUCKETS],
56}
57
58static COUNTERS: [Counters; Stage::ALL.len()] = [const {
59    Counters {
60        calls: AtomicU64::new(0),
61        rows: AtomicU64::new(0),
62        bytes: AtomicU64::new(0),
63        buckets: [const { AtomicU64::new(0) }; BUCKETS],
64    }
65}; Stage::ALL.len()];
66
67fn bucket(rows: usize) -> usize {
68    let bits = usize::BITS - rows.saturating_sub(1).leading_zeros();
69    usize::cast_from(bits).min(BUCKETS - 1)
70}
71
72pub(super) fn record(stage: Stage, rows: usize, bytes: usize) {
73    let counters = stage.counters();
74    counters.calls.fetch_add(1, Ordering::Relaxed);
75    counters
76        .rows
77        .fetch_add(u64::cast_from(rows), Ordering::Relaxed);
78    if bytes > 0 {
79        counters
80            .bytes
81            .fetch_add(u64::cast_from(bytes), Ordering::Relaxed);
82    }
83    counters.buckets[bucket(rows)].fetch_add(1, Ordering::Relaxed);
84}
85
86/// Record one batch an upsert feedback arrangement published, including
87/// empty batches.
88///
89/// Storage records this, but the counters are registered with compute's
90/// pool metrics. They reach the scrape only because clusterd runs storage
91/// and compute in one process with one registry.
92pub fn record_batch(rows: usize) {
93    record(Stage::Batch, rows, 0);
94}
95
96pub(crate) fn register(registry: &MetricsRegistry) {
97    for (stage, name) in Stage::ALL {
98        let counters = stage.counters();
99        let _: ComputedUIntGauge = registry.register_computed_gauge(
100            metric!(name: "mz_column_chunk_work_calls_total", help: "Columnar work operations, by stage.", const_labels: {"stage" => name}),
101            move || counters.calls.load(Ordering::Relaxed),
102        );
103        let _: ComputedUIntGauge = registry.register_computed_gauge(
104            metric!(name: "mz_column_chunk_work_rows_total", help: "Rows processed by columnar work, including repeated visits.", const_labels: {"stage" => name}),
105            move || counters.rows.load(Ordering::Relaxed),
106        );
107        let _: ComputedUIntGauge = registry.register_computed_gauge(
108            metric!(name: "mz_column_chunk_work_bytes_total", help: "Uncompressed serialized bytes processed at instrumented columnar stages.", const_labels: {"stage" => name}),
109            move || counters.bytes.load(Ordering::Relaxed),
110        );
111        for (bucket, count) in counters.buckets.iter().enumerate() {
112            let _: ComputedUIntGauge = registry.register_computed_gauge(
113                metric!(name: "mz_column_chunk_work_size_total", help: "Disjoint columnar operation row-size buckets, indexed by ceil(log2(rows)), the last bucket unbounded.", const_labels: {"stage" => name, "log2_rows" => bucket.to_string()}),
114                move || count.load(Ordering::Relaxed),
115            );
116        }
117    }
118}
119
120#[cfg(test)]
121mod tests {
122    use super::*;
123
124    #[mz_ore::test]
125    fn buckets_are_ceil_log2() {
126        let cases = [
127            (0, 0),
128            (1, 0),
129            (2, 1),
130            (3, 2),
131            (4, 2),
132            (5, 3),
133            (1 << 23, 23),
134            ((1 << 23) + 1, 24),
135            (1 << 30, 24),
136            (usize::MAX, 24),
137        ];
138        for (rows, expected) in cases {
139            assert_eq!(bucket(rows), expected, "bucket for {rows} rows");
140        }
141    }
142
143    #[mz_ore::test]
144    fn all_stages_register_and_scrape() {
145        let registry = MetricsRegistry::new();
146        register(&registry);
147        let families = registry.gather();
148        assert_eq!(families.len(), 4);
149        assert_eq!(
150            families
151                .iter()
152                .map(|family| family.get_metric().len())
153                .sum::<usize>(),
154            Stage::ALL.len() * (3 + BUCKETS)
155        );
156    }
157}