Skip to main content

faucet_core/profiling/
sketch.rs

1//! Bounded-memory sketches behind column profiling (#708). Each is a few KB
2//! per column regardless of row count, so a 10M-row run profiles in the same
3//! memory as a 10-row one.
4//!
5//! - [`HyperLogLog`] — distinct-count estimate (2¹² registers, ~1.6 % error).
6//! - [`TopK`] — most-frequent values (Misra-Gries / Space-Saving counters).
7//! - [`Reservoir`] — a uniform sample for approximate quantiles (Algorithm R
8//!   with a deterministic xorshift generator, so a profile is reproducible).
9//! - [`Welford`] — streaming mean / variance / min / max.
10
11use std::collections::HashMap;
12use std::hash::{Hash, Hasher};
13
14/// 64-bit hash of a value, stable for the process (SipHash with fixed keys).
15pub fn hash64<T: Hash + ?Sized>(v: &T) -> u64 {
16    let mut h = std::hash::DefaultHasher::new();
17    v.hash(&mut h);
18    h.finish()
19}
20
21/// HyperLogLog distinct-count estimator with 2¹² registers (4 KiB).
22#[derive(Debug, Clone)]
23pub struct HyperLogLog {
24    registers: Vec<u8>,
25}
26
27const HLL_P: u32 = 12;
28const HLL_M: usize = 1 << HLL_P;
29
30impl Default for HyperLogLog {
31    fn default() -> Self {
32        Self::new()
33    }
34}
35
36impl HyperLogLog {
37    pub fn new() -> Self {
38        Self {
39            registers: vec![0; HLL_M],
40        }
41    }
42
43    /// Observe one already-hashed value.
44    pub fn insert_hash(&mut self, h: u64) {
45        let idx = (h >> (64 - HLL_P)) as usize;
46        let rest = h << HLL_P;
47        // Leading-zero rank of the remaining bits, 1-based; a zero remainder
48        // saturates at the width of the remainder.
49        let rank = if rest == 0 {
50            (64 - HLL_P) as u8 + 1
51        } else {
52            rest.leading_zeros() as u8 + 1
53        };
54        if rank > self.registers[idx] {
55            self.registers[idx] = rank;
56        }
57    }
58
59    /// Observe one value.
60    pub fn insert<T: Hash + ?Sized>(&mut self, v: &T) {
61        self.insert_hash(hash64(v));
62    }
63
64    /// Merge another sketch (register-wise max).
65    pub fn merge(&mut self, other: &HyperLogLog) {
66        for (a, b) in self.registers.iter_mut().zip(&other.registers) {
67            if *b > *a {
68                *a = *b;
69            }
70        }
71    }
72
73    /// The distinct-count estimate (rounded).
74    pub fn estimate(&self) -> u64 {
75        let m = HLL_M as f64;
76        let alpha = 0.7213 / (1.0 + 1.079 / m);
77        let sum: f64 = self.registers.iter().map(|&r| 2f64.powi(-(r as i32))).sum();
78        let raw = alpha * m * m / sum;
79        let zeros = self.registers.iter().filter(|&&r| r == 0).count();
80        let est = if raw <= 2.5 * m && zeros > 0 {
81            // Linear counting for the small range.
82            m * (m / zeros as f64).ln()
83        } else {
84            raw
85        };
86        est.round().max(0.0) as u64
87    }
88}
89
90/// Space-Saving top-k: at most `capacity` counters, each value's count is an
91/// upper bound whose error is bounded by `error`.
92#[derive(Debug, Clone)]
93pub struct TopK {
94    capacity: usize,
95    counters: HashMap<String, (u64, u64)>,
96}
97
98impl TopK {
99    /// A sketch that keeps `k` values with `4·k` counters (a common overhead
100    /// that keeps the top-`k` set exact in practice for skewed data).
101    pub fn new(k: usize) -> Self {
102        Self {
103            capacity: (k * 4).max(1),
104            counters: HashMap::new(),
105        }
106    }
107
108    pub fn insert(&mut self, value: &str) {
109        if let Some(c) = self.counters.get_mut(value) {
110            c.0 += 1;
111            return;
112        }
113        if self.counters.len() < self.capacity {
114            self.counters.insert(value.to_string(), (1, 0));
115            return;
116        }
117        // Evict the minimum counter and inherit its count as the error bound.
118        let (victim, min) = self
119            .counters
120            .iter()
121            .map(|(k, (c, _))| (k.clone(), *c))
122            .min_by(|a, b| a.1.cmp(&b.1).then_with(|| a.0.cmp(&b.0)))
123            .expect("capacity >= 1");
124        self.counters.remove(&victim);
125        self.counters.insert(value.to_string(), (min + 1, min));
126    }
127
128    /// The `k` most frequent values as `(value, count, error)`, count
129    /// descending then value ascending for determinism.
130    pub fn top(&self, k: usize) -> Vec<(String, u64, u64)> {
131        let mut all: Vec<(String, u64, u64)> = self
132            .counters
133            .iter()
134            .map(|(v, (c, e))| (v.clone(), *c, *e))
135            .collect();
136        all.sort_by(|a, b| b.1.cmp(&a.1).then_with(|| a.0.cmp(&b.0)));
137        all.truncate(k);
138        all
139    }
140}
141
142/// Uniform reservoir sample (Algorithm R) with a deterministic generator.
143#[derive(Debug, Clone)]
144pub struct Reservoir {
145    capacity: usize,
146    seen: u64,
147    values: Vec<f64>,
148    rng: u64,
149}
150
151impl Reservoir {
152    pub fn new(capacity: usize) -> Self {
153        Self {
154            capacity: capacity.max(1),
155            seen: 0,
156            values: Vec::new(),
157            rng: 0x9E37_79B9_7F4A_7C15,
158        }
159    }
160
161    fn next_u64(&mut self) -> u64 {
162        // xorshift64*
163        let mut x = self.rng;
164        x ^= x >> 12;
165        x ^= x << 25;
166        x ^= x >> 27;
167        self.rng = x;
168        x.wrapping_mul(0x2545_F491_4F6C_DD1D)
169    }
170
171    pub fn insert(&mut self, v: f64) {
172        self.seen += 1;
173        if self.values.len() < self.capacity {
174            self.values.push(v);
175            return;
176        }
177        let j = self.next_u64() % self.seen;
178        if (j as usize) < self.capacity {
179            self.values[j as usize] = v;
180        }
181    }
182
183    /// Approximate quantile over the sample; `None` while empty.
184    pub fn quantile(&self, q: f64) -> Option<f64> {
185        if self.values.is_empty() {
186            return None;
187        }
188        let mut sorted = self.values.clone();
189        sorted.sort_by(f64::total_cmp);
190        Some(crate::anomaly::quantile(&sorted, q))
191    }
192
193    pub fn len(&self) -> usize {
194        self.values.len()
195    }
196
197    pub fn is_empty(&self) -> bool {
198        self.values.is_empty()
199    }
200}
201
202/// Welford streaming moments plus exact min / max.
203#[derive(Debug, Clone, Default)]
204pub struct Welford {
205    pub count: u64,
206    mean: f64,
207    m2: f64,
208    pub min: f64,
209    pub max: f64,
210}
211
212impl Welford {
213    pub fn insert(&mut self, x: f64) {
214        if self.count == 0 {
215            self.min = x;
216            self.max = x;
217        } else {
218            self.min = self.min.min(x);
219            self.max = self.max.max(x);
220        }
221        self.count += 1;
222        let delta = x - self.mean;
223        self.mean += delta / self.count as f64;
224        self.m2 += delta * (x - self.mean);
225    }
226
227    pub fn mean(&self) -> Option<f64> {
228        (self.count > 0).then_some(self.mean)
229    }
230
231    /// Population standard deviation.
232    pub fn stddev(&self) -> Option<f64> {
233        (self.count > 0).then(|| (self.m2 / self.count as f64).sqrt())
234    }
235}
236
237#[cfg(test)]
238mod tests {
239    use super::*;
240
241    #[test]
242    fn hll_estimates_within_tolerance() {
243        for &n in &[1u64, 10, 100, 1_000, 50_000] {
244            let mut h = HyperLogLog::new();
245            for i in 0..n {
246                h.insert(&format!("value-{i}"));
247            }
248            let est = h.estimate() as f64;
249            let err = (est - n as f64).abs() / n as f64;
250            assert!(err < 0.05, "n {n} est {est} err {err}");
251        }
252        assert_eq!(HyperLogLog::new().estimate(), 0);
253    }
254
255    #[test]
256    fn hll_merge_is_register_max_and_duplicates_do_not_count() {
257        let mut a = HyperLogLog::new();
258        let mut b = HyperLogLog::new();
259        for i in 0..1000 {
260            a.insert(&i);
261            a.insert(&i);
262            b.insert(&(i + 500));
263        }
264        a.merge(&b);
265        let est = a.estimate() as f64;
266        assert!((est - 1500.0).abs() / 1500.0 < 0.05, "{est}");
267        let mut z = HyperLogLog::new();
268        z.insert_hash(0);
269        assert_eq!(z.registers.iter().filter(|&&r| r > 0).count(), 1);
270    }
271
272    #[test]
273    fn topk_keeps_the_heavy_hitters_with_error_bounds() {
274        let mut t = TopK::new(2);
275        for _ in 0..100 {
276            t.insert("a");
277        }
278        for _ in 0..50 {
279            t.insert("b");
280        }
281        for i in 0..40 {
282            t.insert(&format!("noise-{i}"));
283        }
284        let top = t.top(2);
285        assert_eq!(top[0].0, "a");
286        assert_eq!(top[0].1, 100);
287        assert_eq!(top[1].0, "b");
288        assert!(top[1].1 >= 50);
289        assert!(t.counters.len() <= 8);
290        // Evicted noise inherits the minimum as its error bound.
291        assert!(t.counters.values().any(|(_, e)| *e > 0));
292    }
293
294    #[test]
295    fn topk_orders_ties_by_value() {
296        let mut t = TopK::new(3);
297        for v in ["z", "y", "x"] {
298            t.insert(v);
299        }
300        let names: Vec<_> = t.top(3).into_iter().map(|(v, _, _)| v).collect();
301        assert_eq!(names, vec!["x", "y", "z"]);
302        assert!(TopK::new(0).top(1).is_empty());
303    }
304
305    #[test]
306    fn reservoir_is_uniform_enough_and_deterministic() {
307        let mut r = Reservoir::new(200);
308        for i in 0..10_000 {
309            r.insert(i as f64);
310        }
311        assert_eq!(r.len(), 200);
312        let median = r.quantile(0.5).unwrap();
313        assert!((3500.0..=6500.0).contains(&median), "{median}");
314        let mut r2 = Reservoir::new(200);
315        for i in 0..10_000 {
316            r2.insert(i as f64);
317        }
318        assert_eq!(r.quantile(0.9), r2.quantile(0.9));
319        assert!(Reservoir::new(4).quantile(0.5).is_none());
320        assert!(Reservoir::new(4).is_empty());
321    }
322
323    #[test]
324    fn welford_matches_closed_form() {
325        let mut w = Welford::default();
326        assert!(w.mean().is_none());
327        assert!(w.stddev().is_none());
328        for x in [2.0, 4.0, 4.0, 4.0, 5.0, 5.0, 7.0, 9.0] {
329            w.insert(x);
330        }
331        assert_eq!(w.count, 8);
332        assert_eq!(w.mean(), Some(5.0));
333        assert!((w.stddev().unwrap() - 2.0).abs() < 1e-12);
334        assert_eq!(w.min, 2.0);
335        assert_eq!(w.max, 9.0);
336    }
337}