diff --git a/crates/asap-physical-operators/src/capability.rs b/crates/asap-physical-operators/src/capability.rs index 8bdb1891e..9e3d84249 100644 --- a/crates/asap-physical-operators/src/capability.rs +++ b/crates/asap-physical-operators/src/capability.rs @@ -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(), @@ -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, diff --git a/crates/asap-physical-operators/src/values.rs b/crates/asap-physical-operators/src/values.rs index 41640b7b6..6bf84ef92 100644 --- a/crates/asap-physical-operators/src/values.rs +++ b/crates/asap-physical-operators/src/values.rs @@ -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)?; @@ -303,6 +304,21 @@ fn validate_state(family: &SummaryFamilyType, state: &dyn AggregateCore) -> Resu .as_any() .downcast_ref::() .is_some_and(|s| s.inner.precision == u32::from(*precision)), + SketchParams::UnivMon { + heap_size, + sketch_rows, + sketch_cols, + layers, + } => state + .as_any() + .downcast_ref::() + .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::() diff --git a/crates/asap-physical-operators/tests/univmon_execution.rs b/crates/asap-physical-operators/tests/univmon_execution.rs new file mode 100644 index 000000000..6ddc1f0e7 --- /dev/null +++ b/crates/asap-physical-operators/tests/univmon_execution.rs @@ -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::>()), + )); + let batches: Vec<_> = outputs + .into_iter() + .map(|stream| stream.into_iter().collect::, _>>()) + .collect::>()?; + 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(()) +} diff --git a/crates/integration-tests/tests/precompute_raw_samples.rs b/crates/integration-tests/tests/precompute_raw_samples.rs index 4e119c185..40c3cf7aa 100644 --- a/crates/integration-tests/tests/precompute_raw_samples.rs +++ b/crates/integration-tests/tests/precompute_raw_samples.rs @@ -276,6 +276,14 @@ fn evaluations(state: &dyn AggregateCore, family: &FieldDataType) -> Vec { .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:?}"), } } @@ -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()); @@ -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::>::new();