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
6 changes: 5 additions & 1 deletion crates/asap-physical-operators/src/capability.rs
Original file line number Diff line number Diff line change
Expand Up @@ -167,7 +167,7 @@ pub fn validate_native_family(family: &SummaryFamilyType) -> Result<(), Error> {
match family {
SummaryFamilyType::ExactAggregate(..) => {}
SummaryFamilyType::Sketch(kind, _)
if matches!(kind.algorithm(), A::Kll | A::DDSketch | A::Hll) => {}
if matches!(kind.algorithm(), A::Kll | A::DDSketch | A::Hll | A::UnivMon) => {}
_ => {
return Err(Error::Invalid(
"summary family has no native DAG state implementation".into(),
Expand Down Expand Up @@ -207,6 +207,10 @@ pub fn validate_sketch_evaluation(
(A::DDSketch, _) => bare_count,
(A::Hll, SketchStatistic::Cardinality) => true,
(A::Hll, _) => bare_count,
(A::UnivMon, SketchStatistic::PointCount { value: None, .. })
| (A::UnivMon, SketchStatistic::Cardinality)
| (A::UnivMon, SketchStatistic::FrequencyL2)
| (A::UnivMon, SketchStatistic::FrequencyEntropy) => true,
// Only count intents read a Count-Min bare count, and their
// updates have unit weight; the evaluation is typed Int64 on that basis.
(A::Cms, _) => bare_count,
Expand Down
16 changes: 16 additions & 0 deletions crates/asap-physical-operators/src/values.rs
Original file line number Diff line number Diff line change
Expand Up @@ -266,6 +266,7 @@ fn validate_state(family: &SummaryFamilyType, state: &dyn AggregateCore) -> Resu
use crate::summary_kernels::{
count_min_sketch::CountMinSketchAccumulator, datasketches_kll::DatasketchesKLLAccumulator,
dd_sketch::DDSketchAccumulator, exact::ExactAccumulator, hll_sketch::HllSketchAccumulator,
univmon::UnivMonAccumulator,
};
use planner_types::post_asap::SketchParams;
validate_family(family)?;
Expand Down Expand Up @@ -303,6 +304,21 @@ fn validate_state(family: &SummaryFamilyType, state: &dyn AggregateCore) -> Resu
.as_any()
.downcast_ref::<HllSketchAccumulator>()
.is_some_and(|s| s.inner.precision == u32::from(*precision)),
SketchParams::UnivMon {
heap_size,
sketch_rows,
sketch_cols,
layers,
} => state
.as_any()
.downcast_ref::<UnivMonAccumulator>()
.is_some_and(|s| {
let sketch = s.sketch();
sketch.heap_size == *heap_size as usize
&& sketch.sketch_row == *sketch_rows as usize
&& sketch.sketch_col == *sketch_cols as usize
&& sketch.layer_size == *layers as usize
}),
SketchParams::Cms { width, depth } => state
.as_any()
.downcast_ref::<CountMinSketchAccumulator>()
Expand Down
112 changes: 112 additions & 0 deletions crates/asap-physical-operators/tests/univmon_execution.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,112 @@
//! One native UnivMon build supplies the three statistics in design Example 2.
use asap_physical_operators::{
operators::{Operator, SummaryEvaluation},
plan::{PhysicalDAG, PhysicalOperator},
runtime::{Limits, RunContext, Scope},
values::{Batch, SchemaRef, Value},
Error,
};
use futures::{executor::block_on, StreamExt};
use planner_types::{
post_asap::{
Field, FieldDataType, GroupingStrategy, SketchAlgorithm, SketchKind, SketchParams,
SketchStatistic,
},
pre_asap::DataType,
};
use std::sync::Arc;

fn family() -> FieldDataType {
FieldDataType::Sketch(
SketchKind::new(
SketchAlgorithm::UnivMon,
SketchParams::UnivMon {
heap_size: 64,
sketch_rows: 5,
sketch_cols: 128,
layers: 4,
},
),
GroupingStrategy::default(),
)
}

fn value(row: &[Value]) -> f64 {
match row {
[Value::Float64(value)] => *value,
other => panic!("expected one floating point result, got {other:?}"),
}
}

#[test]
fn one_univmon_state_answers_distinct_l2_and_entropy() -> Result<(), Error> {
let input: SchemaRef = Arc::new(planner_types::pre_asap::Schema::new(vec![Field::plain(
"value",
DataType::Float64,
false,
)]));
let batch = Batch::try_new(
input.clone(),
[1.0, 1.0, 2.0, 2.0, 2.0]
.into_iter()
.map(|value| vec![Value::Float64(value)])
.collect(),
)?;
let state = Operator::summary_build(input, family(), 0, None, vec![])?;
let state_schema = state.output_schema();
let mut dag = PhysicalDAG::default();
dag.add(
0,
vec![],
Operator::source(batch.schema().clone(), vec![batch])?,
)?;
dag.add(1, vec![0], state)?;
dag.add(
2,
vec![1],
Operator::evaluation(
state_schema.clone(),
0,
SummaryEvaluation::Sketch(SketchStatistic::Cardinality),
)?,
)?;
dag.add(
3,
vec![1],
Operator::evaluation(
state_schema.clone(),
0,
SummaryEvaluation::Sketch(SketchStatistic::FrequencyL2),
)?,
)?;
dag.add(
4,
vec![1],
Operator::evaluation(
state_schema,
0,
SummaryEvaluation::Sketch(SketchStatistic::FrequencyEntropy),
)?,
)?;

let context = RunContext::new(
Scope::Query {
evaluation_time_ms: 0,
revision: 1,
},
Limits::default(),
)?;
let outputs = block_on(futures::future::join_all(
dag.execute(&[2, 3, 4], context)?
.into_iter()
.map(|stream| stream.collect::<Vec<_>>()),
));
let batches: Vec<_> = outputs
.into_iter()
.map(|stream| stream.into_iter().collect::<Result<Vec<_>, _>>())
.collect::<Result<_, _>>()?;
assert!((value(batches[0][0].rows().first().unwrap()) - 2.0).abs() < 1.0);
assert!((value(batches[1][0].rows().first().unwrap()) - 13.0_f64.sqrt()).abs() < 1.0);
assert!(value(batches[2][0].rows().first().unwrap()) > 0.0);
Ok(())
}
18 changes: 16 additions & 2 deletions crates/integration-tests/tests/precompute_raw_samples.rs
Original file line number Diff line number Diff line change
Expand Up @@ -276,6 +276,14 @@ fn evaluations(state: &dyn AggregateCore, family: &FieldDataType) -> Vec<f64> {
.map(|q| state.estimate(&SketchStatistic::Quantile { q }).unwrap())
.collect(),
SketchAlgorithm::Hll => vec![state.estimate(&SketchStatistic::Cardinality).unwrap()],
SketchAlgorithm::UnivMon => [
SketchStatistic::Cardinality,
SketchStatistic::FrequencyL2,
SketchStatistic::FrequencyEntropy,
]
.iter()
.map(|statistic| state.estimate(statistic).unwrap())
.collect(),
other => panic!("unexpected unkeyed sketch {other:?}"),
}
}
Expand Down Expand Up @@ -321,7 +329,7 @@ fn check(
let stored_only = matches!(family, FieldDataType::Sketch(kind, _)
if kind.algorithm() == &asap_types::post_asap::SketchAlgorithm::Cms);
if stored_only || asap_physical_operators::capability::validate_native_family(family).is_err() {
// Families without a native state (e.g. UnivMon), or with native
// Families without a native state, or with native
// stored state only (plain CMS), are outside precompute execution;
// their compile must fail.
assert!(precompute::compile(dag, &[source], &[root]).is_err());
Expand All @@ -334,7 +342,13 @@ fn check(
other => format!("{other:?}"),
};
if let FieldDataType::Sketch(kind, _) = family {
if let (Some(keyed), false) = (&input.item, kind.algorithm() == &SketchAlgorithm::Hll) {
if let (Some(keyed), false) = (
&input.item,
matches!(
kind.algorithm(),
SketchAlgorithm::Hll | SketchAlgorithm::UnivMon
),
) {
// Keyed heaps: every item's estimated weight is its exact
// total at this scale (no collisions in the fixture).
let mut expected = BTreeMap::<Series, BTreeMap<String, f64>>::new();
Expand Down
Loading