Skip to main content

faucet_core/profiling/
column.rs

1//! Per-run column profiling (#708): the accumulator that turns written pages
2//! into one [`RunProfile`] — a [`ColumnProfile`] per top-level column with
3//! null rate, type mix, distinct estimate, numeric / string summaries and
4//! the most frequent values. Nested objects and arrays count as one column
5//! (their canonical JSON is the value), matching schema-drift semantics.
6
7use super::sketch::{HyperLogLog, Reservoir, TopK, Welford, hash64};
8use super::spec::ProfilingSpec;
9use schemars::JsonSchema;
10use serde::{Deserialize, Serialize};
11use serde_json::Value;
12use std::collections::{BTreeMap, HashMap};
13
14/// Reservoir size for approximate quantiles.
15const RESERVOIR_CAPACITY: usize = 1024;
16/// Longest value text kept in the top-values sketch (bounds memory per
17/// counter; a longer value is truncated with a marker).
18const MAX_VALUE_TEXT: usize = 256;
19
20/// How many values of each JSON type the column held.
21#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
22#[serde(default)]
23pub struct TypeCounts {
24    pub null: u64,
25    pub boolean: u64,
26    pub integer: u64,
27    pub number: u64,
28    pub string: u64,
29    pub object: u64,
30    pub array: u64,
31}
32
33impl TypeCounts {
34    /// `(type name, count)` for every non-zero type, fixed order.
35    pub fn present(&self) -> Vec<(&'static str, u64)> {
36        [
37            ("null", self.null),
38            ("boolean", self.boolean),
39            ("integer", self.integer),
40            ("number", self.number),
41            ("string", self.string),
42            ("object", self.object),
43            ("array", self.array),
44        ]
45        .into_iter()
46        .filter(|(_, n)| *n > 0)
47        .collect()
48    }
49
50    /// The most common non-null type, if any value was non-null.
51    pub fn dominant(&self) -> Option<&'static str> {
52        self.present()
53            .into_iter()
54            .filter(|(t, _)| *t != "null")
55            .max_by_key(|(_, n)| *n)
56            .map(|(t, _)| t)
57    }
58}
59
60/// Summary of the numeric values a column held.
61#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
62pub struct NumericSummary {
63    pub count: u64,
64    pub min: f64,
65    pub max: f64,
66    pub mean: f64,
67    pub stddev: f64,
68    /// Approximate median (reservoir sample).
69    pub p50: f64,
70    /// Approximate 95th percentile (reservoir sample).
71    pub p95: f64,
72}
73
74/// Summary of the string values a column held (character lengths).
75#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
76pub struct StringSummary {
77    pub count: u64,
78    pub len_min: u64,
79    pub len_max: u64,
80    pub len_mean: f64,
81}
82
83/// One of a column's most frequent values.
84#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
85pub struct TopValue {
86    pub value: String,
87    /// Occurrences (an upper bound; `error` is the possible over-count).
88    pub count: u64,
89    #[serde(default, skip_serializing_if = "is_zero")]
90    pub error: u64,
91    /// `count / non-null values` in this run.
92    pub share: f64,
93}
94
95fn is_zero(n: &u64) -> bool {
96    *n == 0
97}
98
99/// The learned profile of one column over one run.
100#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
101pub struct ColumnProfile {
102    /// Records the run wrote (the denominator of `null_rate`).
103    pub rows: u64,
104    /// Records in which the column was present (including explicit nulls).
105    pub present: u64,
106    /// Records where the column was absent or null.
107    pub nulls: u64,
108    /// `nulls / rows` (0 when no rows).
109    pub null_rate: f64,
110    pub types: TypeCounts,
111    /// Estimated distinct non-null values (HyperLogLog, ~1.6 % error).
112    pub distinct: u64,
113    #[serde(default, skip_serializing_if = "Option::is_none")]
114    pub numeric: Option<NumericSummary>,
115    #[serde(default, skip_serializing_if = "Option::is_none")]
116    pub string: Option<StringSummary>,
117    /// Most frequent values, count descending. Absent for high-cardinality
118    /// columns and when `top_values: 0`.
119    #[serde(default, skip_serializing_if = "Option::is_none")]
120    pub top_values: Option<Vec<TopValue>>,
121    /// The estimated distinct count exceeded `categorical_max_distinct`, so
122    /// no values are published for this column.
123    #[serde(default)]
124    pub high_cardinality: bool,
125}
126
127impl ColumnProfile {
128    /// Non-null values observed.
129    pub fn non_null(&self) -> u64 {
130        self.rows.saturating_sub(self.nulls)
131    }
132}
133
134/// The profile of everything one run wrote.
135#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize, JsonSchema)]
136pub struct RunProfile {
137    /// Records the run wrote.
138    pub rows: u64,
139    /// Profiled columns, name-ordered.
140    pub columns: BTreeMap<String, ColumnProfile>,
141    /// Columns seen beyond `max_columns` and therefore not profiled.
142    #[serde(default, skip_serializing_if = "is_zero")]
143    pub skipped_columns: u64,
144}
145
146impl RunProfile {
147    /// Whether the column cap was hit.
148    pub fn truncated(&self) -> bool {
149        self.skipped_columns > 0
150    }
151}
152
153/// Per-column accumulator.
154#[derive(Debug)]
155struct ColumnAcc {
156    present: u64,
157    nulls: u64,
158    types: TypeCounts,
159    distinct: HyperLogLog,
160    top: Option<TopK>,
161    numeric: Welford,
162    reservoir: Reservoir,
163    string_len: Welford,
164}
165
166impl ColumnAcc {
167    fn new(top_values: usize) -> Self {
168        Self {
169            present: 0,
170            nulls: 0,
171            types: TypeCounts::default(),
172            distinct: HyperLogLog::new(),
173            top: (top_values > 0).then(|| TopK::new(top_values)),
174            numeric: Welford::default(),
175            reservoir: Reservoir::new(RESERVOIR_CAPACITY),
176            string_len: Welford::default(),
177        }
178    }
179
180    fn observe(&mut self, v: &Value) {
181        self.present += 1;
182        match v {
183            Value::Null => {
184                self.nulls += 1;
185                self.types.null += 1;
186                return;
187            }
188            Value::Bool(_) => self.types.boolean += 1,
189            Value::Number(n) => {
190                if n.is_i64() || n.is_u64() {
191                    self.types.integer += 1;
192                } else {
193                    self.types.number += 1;
194                }
195                if let Some(f) = n.as_f64()
196                    && f.is_finite()
197                {
198                    self.numeric.insert(f);
199                    self.reservoir.insert(f);
200                }
201            }
202            Value::String(s) => {
203                self.types.string += 1;
204                self.string_len.insert(s.chars().count() as f64);
205            }
206            Value::Object(_) => self.types.object += 1,
207            Value::Array(_) => self.types.array += 1,
208        }
209        let text = value_text(v);
210        // Distinctness is typed — `1` and `"1"` are two values — while the
211        // frequency sketch keys on the display text.
212        self.distinct
213            .insert_hash(hash64(&(type_tag(v), text.as_ref())));
214        if let Some(top) = &mut self.top {
215            top.insert(&truncate_text(&text));
216        }
217    }
218
219    fn finish(&self, rows: u64, spec: &ProfilingSpec) -> ColumnProfile {
220        let missing = rows.saturating_sub(self.present);
221        let nulls = self.nulls + missing;
222        let non_null = rows.saturating_sub(nulls);
223        let distinct = self.distinct.estimate();
224        let high_cardinality = distinct > spec.categorical_max_distinct;
225        let numeric = (self.numeric.count > 0).then(|| NumericSummary {
226            count: self.numeric.count,
227            min: self.numeric.min,
228            max: self.numeric.max,
229            mean: self.numeric.mean().unwrap_or(0.0),
230            stddev: self.numeric.stddev().unwrap_or(0.0),
231            p50: self.reservoir.quantile(0.5).unwrap_or(0.0),
232            p95: self.reservoir.quantile(0.95).unwrap_or(0.0),
233        });
234        let string = (self.string_len.count > 0).then(|| StringSummary {
235            count: self.string_len.count,
236            len_min: self.string_len.min as u64,
237            len_max: self.string_len.max as u64,
238            len_mean: self.string_len.mean().unwrap_or(0.0),
239        });
240        let top_values = match &self.top {
241            Some(top) if !high_cardinality && non_null > 0 => Some(
242                top.top(spec.top_values)
243                    .into_iter()
244                    .map(|(value, count, error)| TopValue {
245                        value,
246                        count,
247                        error,
248                        share: count as f64 / non_null as f64,
249                    })
250                    .collect(),
251            ),
252            _ => None,
253        };
254        ColumnProfile {
255            rows,
256            present: self.present,
257            nulls,
258            null_rate: if rows == 0 {
259                0.0
260            } else {
261                nulls as f64 / rows as f64
262            },
263            types: self.types,
264            distinct,
265            numeric,
266            string,
267            top_values,
268            high_cardinality,
269        }
270    }
271}
272
273/// The text a value is counted by: strings verbatim, everything else as
274/// canonical JSON (so `1` and `"1"` are distinct values).
275fn value_text(v: &Value) -> std::borrow::Cow<'_, str> {
276    match v {
277        Value::String(s) => std::borrow::Cow::Borrowed(s.as_str()),
278        other => std::borrow::Cow::Owned(other.to_string()),
279    }
280}
281
282fn type_tag(v: &Value) -> u8 {
283    match v {
284        Value::Null => 0,
285        Value::Bool(_) => 1,
286        Value::Number(_) => 2,
287        Value::String(_) => 3,
288        Value::Array(_) => 4,
289        Value::Object(_) => 5,
290    }
291}
292
293fn truncate_text(s: &str) -> String {
294    if s.chars().count() <= MAX_VALUE_TEXT {
295        return s.to_string();
296    }
297    let mut out: String = s.chars().take(MAX_VALUE_TEXT).collect();
298    out.push('…');
299    out
300}
301
302/// The run-level accumulator: hand it every page the sink wrote, then
303/// [`finish`](Profiler::finish).
304#[derive(Debug)]
305pub struct Profiler {
306    spec: ProfilingSpec,
307    rows: u64,
308    columns: HashMap<String, ColumnAcc>,
309    /// Column names seen and rejected by the cap (so the count is exact).
310    skipped: std::collections::HashSet<String>,
311}
312
313impl Profiler {
314    pub fn new(spec: ProfilingSpec) -> Self {
315        Self {
316            spec,
317            rows: 0,
318            columns: HashMap::new(),
319            skipped: std::collections::HashSet::new(),
320        }
321    }
322
323    pub fn spec(&self) -> &ProfilingSpec {
324        &self.spec
325    }
326
327    /// Records observed so far.
328    pub fn rows(&self) -> u64 {
329        self.rows
330    }
331
332    /// Fold one written page in. Non-object records count toward `rows` but
333    /// contribute no columns.
334    pub fn observe_page(&mut self, records: &[Value]) {
335        for record in records {
336            self.rows += 1;
337            let Some(obj) = record.as_object() else {
338                continue;
339            };
340            for (name, v) in obj {
341                if !self.spec.selects(name) {
342                    continue;
343                }
344                if let Some(acc) = self.columns.get_mut(name) {
345                    acc.observe(v);
346                    continue;
347                }
348                if self.columns.len() >= self.spec.max_columns {
349                    self.skipped.insert(name.clone());
350                    continue;
351                }
352                let mut acc = ColumnAcc::new(self.spec.top_values);
353                acc.observe(v);
354                self.columns.insert(name.clone(), acc);
355            }
356        }
357    }
358
359    /// The run's profile from everything observed so far.
360    pub fn finish(&self) -> RunProfile {
361        RunProfile {
362            rows: self.rows,
363            columns: self
364                .columns
365                .iter()
366                .map(|(name, acc)| (name.clone(), acc.finish(self.rows, &self.spec)))
367                .collect(),
368            skipped_columns: self.skipped.len() as u64,
369        }
370    }
371}
372
373#[cfg(test)]
374mod tests {
375    use super::*;
376    use serde_json::json;
377
378    fn profile(records: &[Value]) -> RunProfile {
379        let mut p = Profiler::new(ProfilingSpec::default());
380        p.observe_page(records);
381        p.finish()
382    }
383
384    #[test]
385    fn null_rate_counts_missing_and_explicit_nulls() {
386        let rp = profile(&[
387            json!({"a": 1}),
388            json!({"a": null}),
389            json!({}),
390            json!({"a": 2}),
391        ]);
392        let a = &rp.columns["a"];
393        assert_eq!(rp.rows, 4);
394        assert_eq!(a.present, 3);
395        assert_eq!(a.nulls, 2);
396        assert_eq!(a.null_rate, 0.5);
397        assert_eq!(a.non_null(), 2);
398        assert_eq!(a.types.null, 1);
399        assert_eq!(a.types.integer, 2);
400        assert_eq!(a.types.dominant(), Some("integer"));
401    }
402
403    #[test]
404    fn numeric_string_and_type_mix_summaries() {
405        let rp = profile(&[
406            json!({"n": 1, "s": "ab", "b": true, "o": {"x": 1}, "l": [1]}),
407            json!({"n": 2.5, "s": "abcd", "b": false, "o": {"x": 2}, "l": [1, 2]}),
408            json!({"n": 3, "s": "", "b": true, "o": {"x": 1}, "l": [1]}),
409        ]);
410        let n = rp.columns["n"].numeric.as_ref().unwrap();
411        assert_eq!(n.count, 3);
412        assert_eq!(n.min, 1.0);
413        assert_eq!(n.max, 3.0);
414        assert!((n.mean - 2.1666).abs() < 1e-3);
415        assert_eq!(n.p50, 2.5);
416        assert_eq!(rp.columns["n"].types.integer, 2);
417        assert_eq!(rp.columns["n"].types.number, 1);
418        let s = rp.columns["s"].string.as_ref().unwrap();
419        assert_eq!((s.len_min, s.len_max), (0, 4));
420        assert_eq!(s.len_mean, 2.0);
421        assert!(rp.columns["s"].numeric.is_none());
422        assert_eq!(rp.columns["b"].types.boolean, 3);
423        assert_eq!(rp.columns["o"].types.object, 3);
424        assert_eq!(rp.columns["o"].distinct, 2);
425        assert_eq!(rp.columns["l"].types.array, 3);
426        assert_eq!(
427            rp.columns["b"].types.present(),
428            vec![("boolean", 3)],
429            "present lists only non-zero types"
430        );
431    }
432
433    #[test]
434    fn top_values_with_shares_and_high_cardinality_cutoff() {
435        let mut rows = Vec::new();
436        for i in 0..100 {
437            let c = if i % 2 == 0 {
438                "eu"
439            } else if i % 3 == 0 {
440                "us"
441            } else {
442                "apac"
443            };
444            rows.push(json!({"region": c, "id": format!("id-{i}"), "n": null}));
445        }
446        let spec = ProfilingSpec {
447            categorical_max_distinct: 20,
448            ..Default::default()
449        };
450        let mut p = Profiler::new(spec);
451        p.observe_page(&rows);
452        let rp = p.finish();
453        let region = &rp.columns["region"];
454        let top = region.top_values.as_ref().unwrap();
455        assert_eq!(top[0].value, "eu");
456        assert_eq!(top[0].count, 50);
457        assert_eq!(top[0].share, 0.5);
458        assert!(!region.high_cardinality);
459        let id = &rp.columns["id"];
460        assert!(id.high_cardinality);
461        assert!(id.top_values.is_none());
462        assert!(id.distinct >= 95 && id.distinct <= 105, "{}", id.distinct);
463        // An all-null column has no values to rank.
464        assert!(rp.columns["n"].top_values.is_none());
465        assert_eq!(rp.columns["n"].null_rate, 1.0);
466    }
467
468    #[test]
469    fn top_values_zero_disables_the_sketch() {
470        let spec = ProfilingSpec {
471            top_values: 0,
472            ..Default::default()
473        };
474        let mut p = Profiler::new(spec);
475        p.observe_page(&[json!({"a": "x"})]);
476        assert!(p.finish().columns["a"].top_values.is_none());
477    }
478
479    #[test]
480    fn selection_cap_and_non_object_records() {
481        let mut spec = ProfilingSpec {
482            max_columns: 2,
483            ..Default::default()
484        };
485        spec.exclude.push("skip".into());
486        let mut p = Profiler::new(spec);
487        p.observe_page(&[
488            json!({"a": 1, "b": 2, "c": 3, "d": 4, "skip": 5, "_faucet_run_id": "r"}),
489            json!("scalar"),
490            json!({"a": 1, "c": 3, "e": 9}),
491        ]);
492        let rp = p.finish();
493        assert_eq!(rp.rows, 3);
494        assert_eq!(rp.columns.len(), 2);
495        assert!(rp.columns.contains_key("a") && rp.columns.contains_key("b"));
496        assert_eq!(
497            rp.skipped_columns, 3,
498            "c, d, e were capped; skip/_faucet_ excluded"
499        );
500        assert!(rp.truncated());
501        assert_eq!(p.rows(), 3);
502        assert_eq!(p.spec().max_columns, 2);
503    }
504
505    #[test]
506    fn strings_and_numbers_are_distinct_values_and_long_text_is_truncated() {
507        let long = "x".repeat(600);
508        let rp = profile(&[json!({"v": 1}), json!({"v": "1"}), json!({"v": long})]);
509        let v = &rp.columns["v"];
510        assert_eq!(v.distinct, 3);
511        let top = v.top_values.as_ref().unwrap();
512        let truncated = top.iter().find(|t| t.value.ends_with('…')).unwrap();
513        assert_eq!(truncated.value.chars().count(), MAX_VALUE_TEXT + 1);
514    }
515
516    #[test]
517    fn profile_round_trips_through_json() {
518        let rp = profile(&[json!({"a": 1.5, "b": "x"}), json!({"a": null})]);
519        let text = serde_json::to_string(&rp).unwrap();
520        let back: RunProfile = serde_json::from_str(&text).unwrap();
521        assert_eq!(back, rp);
522        assert!(!text.contains("skipped_columns"), "zero is elided: {text}");
523        assert!(!text.contains("\"error\""));
524    }
525
526    #[test]
527    fn empty_run_profile() {
528        let rp = profile(&[]);
529        assert_eq!(rp.rows, 0);
530        assert!(rp.columns.is_empty());
531        assert!(!rp.truncated());
532        assert_eq!(TypeCounts::default().dominant(), None);
533    }
534}