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
8 changes: 7 additions & 1 deletion crates/devtools/src/bin/stage_pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,7 @@ use asap_plan_selection::{plan_stages, Selection, Sharing, MAX_ENUMERATED_CANDID
use asap_types::ir::export::{
compile_logical_asap_workload, LogicalASAPDAG, LogicalASAPDAGDocument,
};
use asap_types::ir::schema::SketchAlgorithm;
use asap_types::ir::schema::{GroupingStrategy, SketchAlgorithm};
use asap_types::ir::schema_support::with_promql_series_identity;
use asap_types::ir::{OperatorNode, QueryRoot};
use asap_types::types::AccuracyTarget;
Expand Down Expand Up @@ -287,6 +287,12 @@ fn label(inventory: &LocalLogicalCandidates<usize>, owners: &[usize], choice: &[
SketchAlgorithm::CountSketchWithHeap => "CountSketch+heap".to_string(),
other => format!("{other:?}"),
};
let name = match target.groupings[index] {
GroupingStrategy::PerSubpopulationInstance => name,
GroupingStrategy::SharedMultiSubpopulation { ref kind, .. } => {
format!("{kind:?}")
}
};
// The sketch reads the inner aggregate's input and replaces it.
sketches.push(match target.absorbs[index] {
Some(_) => format!("whole-expression {name}"),
Expand Down
134 changes: 134 additions & 0 deletions crates/integration-tests/tests/pass1_sql_coverage.rs
Original file line number Diff line number Diff line change
Expand Up @@ -145,3 +145,137 @@ async fn example2_design_candidates_all_build() {
assert!(enumeration.candidates.len() > 1);
assert!(unbuilt.is_empty(), "{unbuilt:#?}");
}

/// The HydraCms `SummaryAgg`s in `root` (#600's planner contract).
fn hydra_builds(root: &Rc<OperatorNode>) -> Vec<Rc<OperatorNode>> {
use asap_types::ir::schema::{GroupingStrategy, HydraKind};
OperatorNode::reachable(root)
.into_iter()
.filter(|n| {
matches!(
&n.operator,
Operator::ASAP(ASAPOp::SummaryAgg {
grouping: GroupingStrategy::SharedMultiSubpopulation {
kind: HydraKind::HydraCms,
..
},
..
})
)
})
.collect()
}

/// A grouped approximate count gets a HydraCms alternative (#580 W7): one
/// shared Count-Min grid for every `src_ip`, built from a unit weight and a
/// non-null item column. Stage 3 prices it (its grouping-aware guarantee
/// meets the target), and it compiles and executes.
#[tokio::test]
async fn grouped_count_offers_a_priced_executable_hydra_plan() {
use asap_types::ir::schema::{
FieldDataType, GroupingStrategy, HydraParams, SketchParams, SummaryInputExpr, WeightDomain,
};
let sql = "SELECT src_ip, COUNT(*) AS c FROM flows GROUP BY src_ip";
// The grid holds shared_rows × shared_columns Count-Min cells, each as
// large as one per-group sketch: ε = 0.01 exceeds the default memory
// limit, so this uses ε = 0.1.
let target = AccuracyTarget::EpsilonDelta {
epsilon: 0.1,
delta: 0.01,
};
let root = lower_sql(sql, &catalog(), target.clone()).await.unwrap();
let demand = [RootDemand {
accuracy: Some(target),
recurrence: QueryRecurrence::OneTime {
invocations: 1,
execute_at: None,
},
predictability: Predictability::default(),
latency_ms: None,
}];
let data = DataWorkload {
arrival: DataArrival::ContinuouslyIngesting,
ingestion_rate: declared(Rate(100_000.0)),
input_cardinality: declared(10_000_000),
..Default::default()
};
let run = plan_stages(
vec![(0, QueryRoot::Operator(root))],
&demand,
&data,
PlanningModels::builtin(),
4096,
)
.unwrap();
let enumeration = run.enumeration.unwrap();
let hydra: Vec<_> = enumeration
.candidates
.iter()
.filter_map(|c| {
let QueryRoot::Operator(root) = &c.logical.as_ref()?[0].1 else {
return None;
};
(!hydra_builds(root).is_empty()).then(|| (c, root.clone()))
})
.collect();
assert!(!hydra.is_empty(), "a Hydra candidate is generated");
for (candidate, root) in hydra {
// The all-query-time physical candidate comes first (#604).
let physical = candidate.physical.first().expect("Stage 2 builds it");
assert!(
enumeration.selection.costs.contains_key(&physical.id),
"{} is priced: {:?}",
physical.id,
enumeration.selection.rejected
);
for build in hydra_builds(&root) {
let Operator::ASAP(ASAPOp::SummaryAgg {
family: FieldDataType::Sketch(kind, family_grouping),
input,
grouping,
filter: None,
..
}) = &build.operator
else {
panic!("HydraCms build")
};
assert_eq!(family_grouping, grouping);
let GroupingStrategy::SharedMultiSubpopulation {
params: HydraParams::HydraCms { width, depth, .. },
..
} = grouping
else {
panic!("HydraCms params")
};
assert_eq!(
kind.params(),
&SketchParams::Cms {
width: *width,
depth: *depth
}
);
assert!(matches!(input.item, Some(SummaryInputExpr::Column(_))));
assert_eq!(input.weight, SummaryInputExpr::Constant(1.0));
assert!(matches!(
input.weight_domain,
WeightDomain::NonNegative { .. }
));
}
let mut rows = physical_common::execute_raw_rows(&root, rows());
rows.sort_by_key(|row| format!("{row:?}"));
let counts: Vec<_> = rows
.iter()
.map(|row| match (&row[0], &row[1]) {
(Value::Utf8(ip), Value::Int64(n)) => (ip.to_string(), *n),
other => panic!("unexpected row {other:?}"),
})
.collect();
// Few groups in a wide grid: no collisions, so the estimate is exact.
assert_eq!(
counts,
[("a".into(), 3), ("b".into(), 2), ("c".into(), 1)],
"{}",
physical.id
);
}
}
44 changes: 44 additions & 0 deletions crates/logical-optimizer/src/accuracy/estimators/cms.rs
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,50 @@ mod tests {
));
}

/// A shared HydraCms grid adds its collision term to the inner sketch's
/// error: sized like one per-group sketch for ε it misses ε, and sized
/// for ε/2 and δ/2 (Pass 1's split) it meets it.
#[test]
fn hydra_guarantee_adds_the_shared_grid_term() {
use crate::pass1::replacement::default_size_params;
use asap_types::ir::schema::{
default_hydra_params, GroupingStrategy, HydraKind, SketchKind,
};
let count = AggIntent::Count {
accuracy: AccuracyTarget::EpsilonDelta {
epsilon: 0.01,
delta: 0.01,
},
};
let target = AccuracyTarget::EpsilonDelta {
epsilon: 0.01,
delta: 0.01,
};
let group_count = SketchStatistic::PointCount {
key: asap_types::ir::scalar::ColumnRef::SampleValue,
value: None,
};
let guarantee = |epsilon: f64, delta: f64, hydra: bool| {
let params = default_size_params(SketchAlgorithm::Cms, &count, epsilon, delta);
let grouping = match hydra {
true => GroupingStrategy::SharedMultiSubpopulation {
kind: HydraKind::HydraCms,
params: default_hydra_params(HydraKind::HydraCms, &params).unwrap(),
},
false => GroupingStrategy::default(),
};
DefaultAccuracyModel
.local_guarantee(
&FieldDataType::Sketch(SketchKind::new(SketchAlgorithm::Cms, params), grouping),
&group_count,
)
.unwrap()
};
assert!(DefaultAccuracyModel.satisfies(&guarantee(0.01, 0.01, false), &target));
assert!(!DefaultAccuracyModel.satisfies(&guarantee(0.01, 0.01, true), &target));
assert!(DefaultAccuracyModel.satisfies(&guarantee(0.005, 0.005, true), &target));
}

#[test]
fn heap_evaluation_retains_frequency_metric() {
use asap_types::ir::schema::{GroupingStrategy, SketchKind};
Expand Down
34 changes: 32 additions & 2 deletions crates/logical-optimizer/src/accuracy/estimators/mod.rs
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
//! Dispatch committed estimator parameters to their accuracy models.
use super::*;
use asap_types::ir::operator::AggIntent;
use asap_types::ir::schema::GroupingStrategy;
use asap_types::ir::schema::{GroupingStrategy, HydraKind, HydraParams};

pub mod cardinality;
pub mod cms;
Expand Down Expand Up @@ -64,7 +64,37 @@ pub(super) fn local_guarantee(
FieldDataType::ExactAggregate(kind, _) => {
Some(ResultGuarantee::exact(format!("ExactAggregate({kind:?})")))
}
FieldDataType::Sketch(kind, _) => sketch_guarantee(kind.algorithm(), kind.params(), query),
FieldDataType::Sketch(kind, GroupingStrategy::PerSubpopulationInstance) => {
sketch_guarantee(kind.algorithm(), kind.params(), query)
}
FieldDataType::Sketch(
kind,
GroupingStrategy::SharedMultiSubpopulation {
kind: HydraKind::HydraCms,
params:
HydraParams::HydraCms {
shared_rows,
shared_columns,
..
},
},
) => {
// Groups that share a grid cell add their weight: at most
// e·N/shared_columns per row, and the minimum over rows exceeds
// it with probability e^-shared_rows. N is the whole input's
// weight, not one group's, so this bounds error relative to it.
let inner = sketch_guarantee(kind.algorithm(), kind.params(), query)?;
let stats = crate::accuracy::PropagationStats {
hydra_shared_grid_collision_bound: Some(
std::f64::consts::E / f64::from(*shared_columns),
),
hydra_shared_grid_failure_probability: Some((-f64::from(*shared_rows)).exp()),
..Default::default()
};
Some(crate::pass1::grouping::hydra_guarantee(&inner, &stats))
}
// No accuracy model for the other shared groupings.
FieldDataType::Sketch(..) => None,
// No error model is registered for these families.
FieldDataType::Sample(..) | FieldDataType::Wavelet(..) | FieldDataType::StatModel(..) => {
None
Expand Down
5 changes: 4 additions & 1 deletion crates/logical-optimizer/src/pass1/grouping.rs
Original file line number Diff line number Diff line change
Expand Up @@ -430,7 +430,10 @@ fn with_grouping(
/// grid. The paper's collision term depends on deployment/data statistics;
/// keeping those leaves symbolic makes the formula explicit while ensuring
/// target satisfaction fails closed until a caller supplies them.
fn hydra_guarantee(inner: &ResultGuarantee, stats: &PropagationStats) -> ResultGuarantee {
pub(crate) fn hydra_guarantee(
inner: &ResultGuarantee,
stats: &PropagationStats,
) -> ResultGuarantee {
let mut provenance = inner.provenance.clone();
provenance.extend(stats.evidence_provenance.clone());
provenance.push(GuaranteeSource::ChildGuarantee {
Expand Down
Loading