Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 3 additions & 2 deletions crates/devtools/tests/stage_pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -253,8 +253,9 @@ fn selected_label(document: &Value) -> String {
}

/// Example 2: the three SQL statistics over `flows` lower through the SQL
/// frontend, are reported as SQL, and plan. The built-in models have no
/// UnivMon accuracy model, so the selected plan is the exact one with the
/// frontend, are reported as SQL, and plan. The built-in models certify
/// UnivMon only for the L2 norm (Q3), and a UnivMon sized for ε = 0.01 costs
/// more than exact counting, so the selected plan is the exact one with the
/// shared input (recorded, not required by the design).
#[test]
fn example2_plans_the_sql_workload() {
Expand Down
23 changes: 23 additions & 0 deletions crates/executor/src/physical_planner/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -987,6 +987,29 @@ fn summary_build(
groups(input, keys)?,
);
}
// UnivMon counts one occurrence of each row's item (SQL `COUNT(*) GROUP
// BY item`): its value-frequency build over the item column.
if let (
FieldDataType::Sketch(kind, _),
Some(SummaryInputExpr::Column(item)),
SummaryInputExpr::Constant(weight),
) = (family, &update.item, &update.weight)
{
if kind.algorithm() == &planner_types::ir::schema::SketchAlgorithm::UnivMon
&& *weight == 1.0
{
let PlannerReduction::Reduce(keys) = reduction else {
return Err(invalid("UnivMon requires explicit grouping columns"));
};
return Operator::summary_build(
input.clone(),
family.clone(),
named_column(input, item)?,
None,
groups(input, keys)?,
);
}
}
if let Some(item) = &update.item {
let PlannerReduction::Reduce(keys) = reduction else {
return Err(invalid("keyed summary requires explicit partitions"));
Expand Down
45 changes: 44 additions & 1 deletion crates/executor/src/summary_kernels/univmon.rs
Original file line number Diff line number Diff line change
Expand Up @@ -181,7 +181,9 @@ impl AggregateCore for UnivMonAccumulator {
value: None,
} => self.inner.calc_l1(),
SketchStatistic::Cardinality => self.inner.calc_card(),
SketchStatistic::FrequencyL2 => self.inner.calc_l2(),
// Layer 0 sees the whole stream; its row-median F₂ is the readout
// the planner certifies (`calc_l2` is the heavy-hitter G-sum).
SketchStatistic::FrequencyL2 => self.inner.l2_sketch_layers[0].get_l2(),
SketchStatistic::FrequencyEntropy => self.inner.calc_entropy(),
other => return Err(format!("UnivMon does not answer {other:?}").into()),
})
Expand Down Expand Up @@ -274,6 +276,47 @@ mod tests {
}
}

/// L2 reads layer 0's F₂ estimate, and at the planner's (0.01, 0.01)
/// sizing it is within 1% of a Zipf stream's true L2 norm.
#[test]
fn l2_is_layer0_f2_within_the_certified_bound() {
use asap_logical_optimizer::pass1::replacement::default_size_params;
use planner_types::ir::operator::agg_intent::default_cardinality;
use planner_types::ir::schema::{SketchAlgorithm, SketchParams};
let SketchParams::UnivMon {
heap_size,
sketch_rows,
sketch_cols,
layers,
} = default_size_params(SketchAlgorithm::UnivMon, &default_cardinality(), 0.01, 0.01)
else {
unreachable!()
};
let mut state = UnivMonAccumulator::new(
heap_size as usize,
sketch_rows as usize,
sketch_cols as usize,
usize::from(layers),
)
.unwrap();
// Zipf(1) frequencies over 5,000 keys: key i occurs ⌊10,000 / i⌋ times.
let mut f2 = 0.0;
for key in 1..=5_000u32 {
let count = 10_000 / key;
f2 += f64::from(count) * f64::from(count);
for _ in 0..count {
state.insert_sample(f64::from(key)).unwrap();
}
}
let l2 = state.estimate(&SketchStatistic::FrequencyL2).unwrap();
assert_eq!(l2, state.sketch().l2_sketch_layers[0].get_l2());
assert!(
(l2 - f2.sqrt()).abs() <= 0.01 * f2.sqrt(),
"{l2} vs {}",
f2.sqrt()
);
}

/// Retained variable-length identities contribute to the runtime memory reservation.
#[test]
fn memory_accounts_for_string_identities() {
Expand Down
16 changes: 10 additions & 6 deletions crates/frontend-promql/tests/univmon_candidates.rs
Original file line number Diff line number Diff line change
Expand Up @@ -68,11 +68,12 @@ fn candidate(query: &str, accuracy: AccuracyTarget) -> Rc<OperatorNode> {

#[test]
fn four_evaluations_share_one_value_frequency_state_and_keep_honest_guarantees() {
// Equal data, grouping and window produce one state independently of evaluation.
// Equal data, grouping, window and requirement produce one state
// independently of evaluation: UnivMon is sized for L2 whatever it reads.
let accuracy = AccuracyTarget::Epsilon(0.02);
let roots: Vec<_> = [
("distinct_over_time(m[5m])", accuracy.clone()),
("count_over_time(m[5m])", AccuracyTarget::Exact),
("count_over_time(m[5m])", accuracy.clone()),
("l2_over_time(m[5m])", accuracy.clone()),
("entropy_over_time(m[5m])", accuracy),
]
Expand Down Expand Up @@ -114,11 +115,14 @@ fn four_evaluations_share_one_value_frequency_state_and_keep_honest_guarantees()
let Operator::ASAP(ASAPOp::SummaryAgg { family, .. }) = &summary_input.operator else {
panic!()
};
assert!(
// Production certifies L2 from layer 0's F₂, but has no
// calibrated bound for distinct count or entropy.
assert_eq!(
DefaultAccuracyModel
.local_guarantee(family, query)
.is_none(),
"production has no calibrated error bound"
.map(|g| g.metric),
(*index == 2).then_some(ErrorMetric::RelativeValue),
"{query:?}"
);
}
post_asap_dag(root);
Expand All @@ -129,7 +133,7 @@ fn four_evaluations_share_one_value_frequency_state_and_keep_honest_guarantees()
fn uncalibrated_frequency_evaluations_do_not_bypass_accuracy_targets() {
// An unmeasured heuristic remains inspectable but is never certified or
// automatically selected for a caller-visible bounded-error result.
for query in ["entropy_over_time(m[5m])", "l2_over_time(m[5m])"] {
for query in ["entropy_over_time(m[5m])"] {
for target in [
AccuracyTarget::Exact,
AccuracyTarget::Epsilon(0.02),
Expand Down
8 changes: 7 additions & 1 deletion crates/integration-tests/tests/precompute_raw_samples.rs
Original file line number Diff line number Diff line change
Expand Up @@ -160,7 +160,13 @@ fn execute(
window_end_ms: 6000,
revision: 1,
},
Limits::default(),
// A UnivMon sized for L2 at ε = 0.02 holds 16 layers of 5 × 2^14
// counters (≈ 10 MiB) per population, past the default 64 MiB for
// this test's five series.
Limits {
max_bytes: 1 << 30,
..Limits::default()
},
)
.unwrap();
block_on(async {
Expand Down
7 changes: 3 additions & 4 deletions crates/logical-optimizer/src/accuracy/estimators/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,7 @@ pub(super) fn sketch_guarantee(
SketchParams::Kmv { .. } | SketchParams::Theta { .. } => {
cardinality::guarantee(algorithm, params, query)
}
SketchParams::UnivMon { .. } => univmon::guarantee(query),
SketchParams::UnivMon { .. } => univmon::guarantee(params, query),
}
}

Expand Down Expand Up @@ -108,9 +108,8 @@ pub(crate) fn size_params(
delta: f64,
) -> SketchParams {
match kind {
// Baseline dimensions are candidates, not an inverted error bound.
// Empirical models may size these; no theoretical guarantee is claimed.
SketchAlgorithm::UnivMon => univmon::size_params(),
// Sized for the L2 readout whatever the intent (see `univmon`).
SketchAlgorithm::UnivMon => univmon::size_params(eps, delta),
SketchAlgorithm::Kll => SketchParams::Kll { k: kll::kll_k(eps) },
SketchAlgorithm::Cms => SketchParams::Cms {
width: cms::cms_width(eps),
Expand Down
Loading