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
72 changes: 70 additions & 2 deletions crates/logical-optimizer/src/accuracy/estimators/univmon.rs
Original file line number Diff line number Diff line change
@@ -1,17 +1,85 @@
//! UnivMon currently certifies only its exact unit-update total evaluation.
//! UnivMon certifies only its unit-update total, which the kernel reads
//! exactly (`calc_l1` returns the update count).
//!
//! Distinct count, L2 norm and entropy are deliberately left uncertified,
//! so Stage 3 rejects them unless a deployment supplies accuracy evidence:
//!
//! - Liu et al., "One Sketch to Rule Them All" (SIGCOMM 2016), rely on
//! Braverman and Ostrovsky's recursive sketch. With O(log n) layers, each of
//! which returns a (g, ε)-cover of its sampled substream (every g-heavy item
//! with a (1 ± ε) frequency), the recursive G-sum is a (1 ± ε)
//! approximation with probability 1 − δ. The layer sketch's size is only
//! stated asymptotically (O(ε⁻² log(1/δ)) per CountSketch, times
//! polylogarithmic factors). No constants are given that could be inverted
//! into a concrete `(heap_size, sketch_rows, sketch_cols, layers)`.
//! - The kernel (`asap_sketchlib::UnivMon::calc_g_sum_heuristic`, behind
//! `asap-executor`'s `UnivMonAccumulator`) keeps a fixed top-`heap_size`
//! heap per layer and estimates its items with the layer's CountSketch.
//! For distinct counts it also drops items below `L2 / sqrt(heap_size)`.
//! Nothing bounds the probability that a layer's heap is such a cover, so
//! the theorem's premise does not hold for these readouts. Even when every
//! heap holds every item, a CountSketch estimate may be at most zero for a
//! present item, so the distinct count is not certified in that case
//! either.
//! - A per-layer CountSketch does give a sound F₂ (AMS) bound, but `calc_l2`
//! reads the recursive sum, not that estimate.
//!
//! A sound guarantee needs either a readout with a proven bound (for
//! example, L2 from layer 0's CountSketch) or a per-layer cover guarantee for
//! the heap. Both are kernel changes.
use super::*;

pub(super) fn guarantee(query: &SketchStatistic) -> Option<ResultGuarantee> {
matches!(query, SketchStatistic::PointCount { value: None, .. })
.then(|| ResultGuarantee::exact("univmon_unit_update_total"))
}

/// One shape for every readout and requirement: no bound is inverted (see
/// the module docs), so every consumer of a key reads identical states.
pub(super) fn size_params() -> SketchParams {
// Baseline dimensions are candidates, not an inverted accuracy bound.
SketchParams::UnivMon {
heap_size: 256,
sketch_rows: 5,
sketch_cols: 1024,
layers: 16,
}
}

#[cfg(test)]
mod tests {
use super::*;
use asap_types::ir::scalar::ColumnRef;
use asap_types::ir::schema::{GroupingStrategy, SketchKind};

fn family() -> FieldDataType {
FieldDataType::Sketch(
SketchKind::new(SketchAlgorithm::UnivMon, size_params()),
GroupingStrategy::default(),
)
}

/// The exact total is certified; distinct count, L2 and entropy have no
/// sound bound for the shipped kernel and stay uncertified.
#[test]
fn only_the_total_is_certified() {
let total = SketchStatistic::PointCount {
key: ColumnRef::SampleValue,
value: None,
};
assert!(DefaultAccuracyModel
.local_guarantee(&family(), &total)
.is_some_and(|g| g.is_exact()));
for query in [
SketchStatistic::Cardinality,
SketchStatistic::FrequencyL2,
SketchStatistic::FrequencyEntropy,
] {
assert!(
DefaultAccuracyModel
.local_guarantee(&family(), &query)
.is_none(),
"{query:?}"
);
}
}
}
51 changes: 49 additions & 2 deletions crates/logical-optimizer/src/pass2/summary_capability.rs
Original file line number Diff line number Diff line change
Expand Up @@ -29,11 +29,14 @@ use crate::pass1::logical_candidates::{
};
use crate::pass1::replacement::{accuracy_budget, accuracy_target};

/// The estimates one summary serves: any quantile of one column from one
/// quantile summary (KLL, DDSketch).
/// The estimates one summary serves over one column.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Capability {
/// Any quantile, from one KLL or DDSketch.
Quantile,
/// Distinct count, L2 norm and entropy of the value frequencies, from one
/// UnivMon (#509 Example 2).
FrequencyMoments,
}

/// What a summary for a target would ingest and which estimates it would
Expand Down Expand Up @@ -77,6 +80,13 @@ fn key(node: &OperatorNode) -> Option<Key<'_>> {
}
let (capability, column) = match intent {
AggIntent::Quantile { col, .. } => (Capability::Quantile, *col),
// A distinct-tuple count has no UnivMon alternative.
AggIntent::Cardinality { cols, .. } if cols.len() <= 1 => {
(Capability::FrequencyMoments, cols.first().copied())
}
AggIntent::FrequencyL2 { col, .. } | AggIntent::FrequencyEntropy { col, .. } => {
(Capability::FrequencyMoments, *col)
}
_ => return None,
};
Some(Key {
Expand Down Expand Up @@ -259,6 +269,43 @@ mod tests {
assert!(!shared.resized);
}

/// Distinct count, entropy and L2 over one input form one key. UnivMon's
/// shape does not depend on the requirement, so all three alternatives
/// are the same state; the distinct count's other summaries are re-sized
/// for the strictest ε.
#[test]
fn frequency_moments_share_one_univmon() {
let base = inventory(&[
("distinct_over_time(src[1m])", 0.02),
("entropy_over_time(src[1m])", 0.05),
("l2_over_time(src[1m])", 0.01),
]);
let shared = share_summary_capability(&base).unwrap().expect("one key");
assert!(shared.resized);
let univmon: Vec<_> = shared
.inventory
.targets
.iter()
.map(|t| {
t.alternatives
.iter()
.find(|a| matches!(a, Realization::Sketch(kind) if *kind.algorithm() == SketchAlgorithm::UnivMon))
.cloned()
.expect("a UnivMon alternative")
})
.collect();
assert!(univmon.iter().all(|u| *u == univmon[0]));
let hll = |inv: &LocalLogicalCandidates<usize>| {
inv.targets[0].alternatives.iter().find_map(|a| match a {
Realization::Sketch(kind) if *kind.algorithm() == SketchAlgorithm::Hll => {
Some(kind.clone())
}
_ => None,
})
};
assert_ne!(hll(&shared.inventory), hll(&base));
}

/// A different window, selector or estimate family is a different key.
#[test]
fn different_summary_input_or_window_is_not_shared() {
Expand Down
58 changes: 56 additions & 2 deletions crates/planner/tests/summary_sharing.rs
Original file line number Diff line number Diff line change
Expand Up @@ -503,8 +503,9 @@ impl AccuracyModel for UnivMonEvidence {

/// Distinct count, entropy and L2 over one input, certified by an accuracy
/// model and selected by a cost model preferring UnivMon, read one UnivMon state: #515 sharing is the summary-capability rule
/// when the states are identical. The facade's stage pipeline does not plan
/// UnivMon sharing, so this runs the legacy search with the test model.
/// when the states are identical. This runs the legacy search; the stage
/// pipeline's counterpart is
/// `frequency_moments_share_one_univmon_in_the_stage_pipeline`.
#[test]
fn certified_frequency_evaluations_share_one_univmon_state() {
let queries = [
Expand Down Expand Up @@ -554,3 +555,56 @@ fn certified_frequency_evaluations_share_one_univmon_state() {
assert_eq!(states.len(), 3);
assert!(states.iter().all(|state| Rc::ptr_eq(state, &states[0])));
}

/// #509 Example 2 through the stage pipeline: distinct count, entropy and L2
/// of one input over one window, with requirements ε = 0.02, 0.05 and 0.01.
/// The summary-capability variant gives the three targets one UnivMon.
/// Certified by the synthetic model, that one state serves all three
/// queries. The built-in model has no sound bound for these readouts, so
/// Stage 3 rejects every UnivMon plan.
#[tokio::test]
async fn frequency_moments_share_one_univmon_in_the_stage_pipeline() {
let queries = [
("distinct_over_time(src[1m])", 0.02),
("entropy_over_time(src[1m])", 0.05),
("l2_over_time(src[1m])", 0.01),
];
let workload = promql_workload(&queries);
let plan = |models: PlanningModels<'static>| {
let workload = workload.clone();
async move {
let input = UserInput::new(
&workload,
FrontendInput::Promql {
now_ms: NOW_MS,
histograms: None,
},
models,
);
e2e_plan(input).await.expect("workload plans")
}
};
let univmon_states = |output: &PlanOutput| -> Vec<Rc<OperatorNode>> {
output
.plans
.iter()
.flat_map(plan_states)
.filter(|state| {
matches!(
&state.operator,
asap_types::ir::Operator::ASAP(ASAPOp::SummaryAgg { family: FieldDataType::Sketch(kind, _), .. })
if kind.algorithm() == &SketchAlgorithm::UnivMon
)
})
.collect()
};

let certified = plan(PlanningModels::builtin().with_accuracy(&UnivMonEvidence)).await;
let states = univmon_states(&certified);
assert_eq!(states.len(), 3, "{:?}", certified.selection);
assert!(states.iter().all(|state| Rc::ptr_eq(state, &states[0])));
assert_eq!(unique_deployments(&certified), 1);

let builtin = plan(PlanningModels::builtin()).await;
assert!(univmon_states(&builtin).is_empty());
}