From c4cd1d001d5a170ba72316d9e4cc0aca94eb74b2 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Sun, 4 Oct 2026 16:57:49 +0000 Subject: [PATCH] feat(executor): exact distinct, L2 and entropy, and typed UnivMon identities Port of the parked runtime PRs for #509 Example 2 onto asap-executor: - native UnivMon build and Cardinality / FrequencyL2 / FrequencyEntropy readouts (the old stack's 14c8ac8a, a prerequisite missing here); - #557: a UnivMon build accepts Utf8, Int64 and Bool identities as type-tagged keys; heap key bytes count toward state memory; - #559: exact FrequencyL2 / FrequencyEntropy reducers (bits, NULL skipped, 0 for an empty population, memory-accounted, cooperative); - #563: exact Cardinality reducer over typed tuples. The legacy raw-cost adapter lists the three intents as hash aggregates. Co-Authored-By: Claude Opus 5.5 --- crates/executor/src/capability.rs | 6 +- .../executor/src/operators/aggregate/mod.rs | 89 +++++++- crates/executor/src/operators/summary/mod.rs | 22 +- crates/executor/src/physical_planner/mod.rs | 13 ++ .../executor/src/summary_kernels/factory.rs | 20 ++ .../executor/src/summary_kernels/univmon.rs | 85 +++++++- crates/executor/src/values.rs | 16 ++ crates/executor/tests/blocking_resources.rs | 44 ++++ crates/executor/tests/physical_semantics.rs | 202 ++++++++++++++++++ crates/executor/tests/univmon_execution.rs | 182 ++++++++++++++++ .../tests/physical_common/mod.rs | 55 +++++ .../tests/precompute_raw_samples.rs | 18 +- .../tests/sql_cardinality.rs | 46 ++++ .../src/cost/query_physical_lowering.rs | 3 + 14 files changed, 779 insertions(+), 22 deletions(-) create mode 100644 crates/executor/tests/univmon_execution.rs create mode 100644 crates/integration-tests/tests/sql_cardinality.rs diff --git a/crates/executor/src/capability.rs b/crates/executor/src/capability.rs index 48c2395b..9d066730 100644 --- a/crates/executor/src/capability.rs +++ b/crates/executor/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/executor/src/operators/aggregate/mod.rs b/crates/executor/src/operators/aggregate/mod.rs index 94c2bb86..69e6c10d 100644 --- a/crates/executor/src/operators/aggregate/mod.rs +++ b/crates/executor/src/operators/aggregate/mod.rs @@ -13,6 +13,17 @@ impl Operator { for (name, reduction) in &measures { let (t, n) = match reduction { Reduction::Count => (DataType::Int64, false), + Reduction::Cardinality(columns) => { + if columns.is_empty() { + return Err(invalid( + "distinct aggregate requires at least one identity column", + )); + } + for column in columns { + plain(&input, *column)?; + } + (DataType::Int64, false) + } Reduction::Sum(i) | Reduction::Avg(i) => { let (t, _) = plain(&input, *i)?; if !matches!(t, DataType::Int64 | DataType::Float64) { @@ -34,6 +45,15 @@ impl Operator { } (DataType::Float64, false) } + Reduction::FrequencyL2(i) | Reduction::FrequencyEntropy(i) => { + if !matches!( + plain(&input, *i)?.0, + DataType::Bool | DataType::Int64 | DataType::Float64 | DataType::Utf8 + ) { + return Err(invalid("frequency aggregate requires a Boolean, Int64, Float64 or Utf8 identity")); + } + (DataType::Float64, false) + } Reduction::Min(i) | Reduction::Max(i) => { let (t, nullable) = plain(&input, *i)?; if !ordered(t) { @@ -125,10 +145,16 @@ impl Operator { #[derive(serde::Serialize, serde::Deserialize, Clone, Debug)] pub enum Reduction { Count, + /// Exact distinct tuple count; a tuple with any NULL component is skipped. + Cardinality(Vec), Sum(usize), Avg(usize), Min(usize), Max(usize), + /// L2 norm of unit-update frequencies; NULL identities are skipped. + FrequencyL2(usize), + /// Shannon entropy in bits; NULL identities are skipped. + FrequencyEntropy(usize), /// PromQL `quantile`: linear interpolation between closest ranks. Quantile { column: usize, @@ -205,7 +231,7 @@ async fn reduce( .map(|&i| rows[0][i].clone()) .collect::>(); for measure in measures { - result.push(reduce_one(&rows, measure, input, &mut work).await?); + result.push(reduce_one(&rows, measure, input, &mut work, context).await?); } workspace.grow(row_bytes(&result))?; output.push(result); @@ -243,6 +269,7 @@ async fn reduce_one( measure: &Reduction, input: &SchemaRef, work: &mut Cooperative, + context: &RunContext, ) -> Result { let column = match measure { Reduction::Count => { @@ -251,6 +278,66 @@ async fn reduce_one( )) } Reduction::Sum(i) | Reduction::Avg(i) | Reduction::Min(i) | Reduction::Max(i) => *i, + Reduction::Cardinality(columns) => { + let mut workspace = Workspace::new(context)?; + let mut identities = std::collections::BTreeSet::new(); + for row in rows { + work.checkpoint().await?; + if columns + .iter() + .any(|column| matches!(row[*column], Value::Null)) + { + continue; + } + let key = group_key(row, columns)?; + if !identities.contains(&key) { + workspace.grow(key_bytes(&key))?; + identities.insert(key); + } + } + return Ok(Value::Int64( + i64::try_from(identities.len()).map_err(|_| invalid("distinct count overflow"))?, + )); + } + Reduction::FrequencyL2(column) | Reduction::FrequencyEntropy(column) => { + let mut workspace = Workspace::new(context)?; + let mut counts = BTreeMap::, u64>::new(); + let mut total = 0_u64; + for row in rows { + work.checkpoint().await?; + let value = &row[*column]; + if matches!(value, Value::Null) { + continue; + } + if matches!(value, Value::Float64(v) if !v.is_finite()) { + return Err(invalid( + "frequency aggregate requires a finite floating identity", + )); + } + let key = value.key()?; + if !counts.contains_key(&key) { + workspace.grow( + 64 + std::mem::size_of::>() + + key.len() + + std::mem::size_of::(), + )?; + } + *counts.entry(key).or_default() += 1; + total += 1; + } + let mut result = 0.0_f64; + for count in counts.into_values() { + work.checkpoint().await?; + let count = count as f64; + if matches!(measure, Reduction::FrequencyL2(_)) { + result = result.hypot(count); + } else { + let probability = count / total as f64; + result -= probability * probability.log2(); + } + } + return Ok(Value::Float64(result)); + } Reduction::Quantile { column, q } => { let mut values = Vec::with_capacity(rows.len()); for row in rows { diff --git a/crates/executor/src/operators/summary/mod.rs b/crates/executor/src/operators/summary/mod.rs index dfec3e6d..9f4e9cf2 100644 --- a/crates/executor/src/operators/summary/mod.rs +++ b/crates/executor/src/operators/summary/mod.rs @@ -99,8 +99,14 @@ impl Operator { ) -> Result { crate::values::validate_family(&family)?; validate_groups(&input, &groups)?; - if plain(&input, value)?.0 != &DataType::Float64 { - return Err(invalid("summary numeric update requires Float64")); + let dtype = plain(&input, value)?.0; + let typed_frequency = matches!(&family, + SummaryFamilyType::Sketch(kind, _) if kind.algorithm() == &planner_types::ir::schema::SketchAlgorithm::UnivMon); + if *dtype != DataType::Float64 + && !(typed_frequency + && matches!(dtype, DataType::Utf8 | DataType::Int64 | DataType::Bool)) + { + return Err(invalid("summary update type is unsupported by its family")); } if let Some(time) = time { if plain(&input, time)? != (&DataType::Timestamp, false) { @@ -435,11 +441,10 @@ async fn build_summary( states.get_mut(&key).expect("inserted group"); // SQL aggregates ignore NULL samples while retaining the group. // A missing counter sample also contributes no observation. - let value = match row[value] { - Value::Float64(value) => value, - Value::Null => continue, - _ => return Err(invalid("summary update type")), - }; + let value = &row[value]; + if matches!(value, Value::Null) { + continue; + } let timestamp = if let Some(time) = time { let Value::Timestamp(time) = row[time] else { return Err(invalid("summary time type")); @@ -455,9 +460,8 @@ async fn build_summary( )); } updater - .validate_single_input(value) + .update_value(value, timestamp) .map_err(Error::Operator)?; - updater.update_single(value, timestamp); *previous = Some(timestamp); memory.resize(updater.memory_usage_bytes() + *overhead)?; } diff --git a/crates/executor/src/physical_planner/mod.rs b/crates/executor/src/physical_planner/mod.rs index f376bece..e74a944e 100644 --- a/crates/executor/src/physical_planner/mod.rs +++ b/crates/executor/src/physical_planner/mod.rs @@ -848,8 +848,21 @@ fn bind_operation(node: &PhysicalASAPDAGNode, inputs: &[SchemaRef]) -> Result Reduction::Count, + AggIntent::Cardinality { cols, .. } => { + Reduction::Cardinality(if cols.is_empty() { + vec![column(None)?] + } else { + cols.clone() + }) + } AggIntent::Sum { col } => Reduction::Sum(column(*col)?), AggIntent::Avg { col } => Reduction::Avg(column(*col)?), + AggIntent::FrequencyL2 { col, .. } => { + Reduction::FrequencyL2(column(*col)?) + } + AggIntent::FrequencyEntropy { col, .. } => { + Reduction::FrequencyEntropy(column(*col)?) + } AggIntent::Min { col } => Reduction::Min(column(*col)?), AggIntent::Max { col } => Reduction::Max(column(*col)?), _ => { diff --git a/crates/executor/src/summary_kernels/factory.rs b/crates/executor/src/summary_kernels/factory.rs index f6aca692..016c18c5 100644 --- a/crates/executor/src/summary_kernels/factory.rs +++ b/crates/executor/src/summary_kernels/factory.rs @@ -48,6 +48,21 @@ pub trait AccumulatorUpdater: Send { } } + /// Update from the typed native row. Numeric kernels retain their existing + /// Float64 contract; frequency kernels may accept nonnumeric identities. + fn update_value( + &mut self, + value: &crate::values::Value, + timestamp_ms: i64, + ) -> Result<(), String> { + let crate::values::Value::Float64(value) = value else { + return Err("summary update requires Float64".into()); + }; + self.validate_single_input(*value)?; + self.update_single(*value, timestamp_ms); + Ok(()) + } + /// Feed a single (value, timestamp_ms) pair — for SingleSubpopulation types. fn update_single(&mut self, value: f64, timestamp_ms: i64); @@ -690,6 +705,11 @@ impl AccumulatorUpdater for HllUpdater { } impl AccumulatorUpdater for UnivMonUpdater { + fn update_value(&mut self, value: &crate::values::Value, _: i64) -> Result<(), String> { + self.acc + .insert_value(value) + .map_err(|error| error.to_string()) + } fn is_keyed(&self) -> bool { false } diff --git a/crates/executor/src/summary_kernels/univmon.rs b/crates/executor/src/summary_kernels/univmon.rs index 54904282..9d679b5d 100644 --- a/crates/executor/src/summary_kernels/univmon.rs +++ b/crates/executor/src/summary_kernels/univmon.rs @@ -75,6 +75,25 @@ impl UnivMonAccumulator { Ok(()) } + /// SQL identities are kept in their original type: converting Int64 to + /// Float64 would collapse neighboring keys above 2^53. + pub fn insert_value(&mut self, value: &crate::values::Value) -> Result<(), Error> { + use crate::values::Value; + match value { + Value::Null => return Ok(()), + Value::Float64(number) if number.is_finite() => return self.insert_sample(*number), + Value::Bool(_) | Value::Int64(_) | Value::Utf8(_) => {} + _ => return Err("unsupported UnivMon identity type or nonfinite sample".into()), + } + let key = value.key()?; + self.inner + .bucket_size + .checked_add(1) + .ok_or("UnivMon count overflow")?; + self.inner.insert(&DataInput::Bytes(&key), 1); + Ok(()) + } + fn compatible(&self, other: &Self) -> bool { ( self.inner.heap_size, @@ -113,15 +132,25 @@ impl UnivMonAccumulator { impl AggregateCore for UnivMonAccumulator { fn approx_memory_bytes(&self) -> usize { - std::mem::size_of::().saturating_add( - self.inner.layer_size.saturating_mul( - self.inner - .sketch_row - .saturating_mul(self.inner.sketch_col) - .saturating_mul(16) - .saturating_add(self.inner.heap_size.saturating_mul(256)), - ), - ) + let key_bytes = (0..self.inner.layer_size) + .flat_map(|layer| self.inner.hh_layers[layer].heap()) + .map(|item| match &item.key { + asap_sketchlib::HeapItem::String(key) => key.capacity(), + asap_sketchlib::HeapItem::Bytes(key) => key.capacity(), + _ => 0, + }) + .fold(0usize, usize::saturating_add); + key_bytes + .saturating_add(std::mem::size_of::()) + .saturating_add( + self.inner.layer_size.saturating_mul( + self.inner + .sketch_row + .saturating_mul(self.inner.sketch_col) + .saturating_mul(16) + .saturating_add(self.inner.heap_size.saturating_mul(256)), + ), + ) } fn clone_boxed_core(&self) -> Box { Box::new(self.clone()) @@ -219,4 +248,42 @@ mod tests { assert!(UnivMonAccumulator::from_sketch(terminal).is_ok()); assert!(UnivMonAccumulator::from_sketch(UnivMon::init_univmon(4, 21, 16, 2)).is_err()); } + /// Typed keys remain distinct after merging and restoring persisted native state. + #[test] + fn typed_keys_merge_and_roundtrip() { + use crate::values::Value; + let mut left = UnivMonAccumulator::new(32, 5, 1024, 4).unwrap(); + let mut right = UnivMonAccumulator::new(32, 5, 1024, 4).unwrap(); + for _ in 0..2 { + left.insert_value(&Value::Utf8("192.0.2.1".into())).unwrap(); + right + .insert_value(&Value::Utf8("192.0.2.2".into())) + .unwrap(); + } + left.merge_in_place(&right).unwrap(); + let bytes = left.sketch().serialize_to_bytes().unwrap(); + let restored = + UnivMonAccumulator::from_sketch(UnivMon::deserialize_from_bytes(&bytes).unwrap()) + .unwrap(); + for (statistic, expected) in [ + (SketchStatistic::Cardinality, 2.0), + (SketchStatistic::FrequencyL2, 8.0_f64.sqrt()), + (SketchStatistic::FrequencyEntropy, 1.0), + ] { + assert!((restored.estimate(&statistic).unwrap() - expected).abs() < 0.01); + } + } + + /// Retained variable-length identities contribute to the runtime memory reservation. + #[test] + fn memory_accounts_for_string_identities() { + let mut state = UnivMonAccumulator::new(32, 5, 1024, 4).unwrap(); + let empty = state.approx_memory_bytes(); + state + .insert_value(&crate::values::Value::Utf8("x".repeat(4096).into())) + .unwrap(); + assert!(state.approx_memory_bytes() >= empty + 4096); + state.clear(); + assert_eq!(state.approx_memory_bytes(), empty); + } } diff --git a/crates/executor/src/values.rs b/crates/executor/src/values.rs index 9feb0fac..6f2bc76b 100644 --- a/crates/executor/src/values.rs +++ b/crates/executor/src/values.rs @@ -265,6 +265,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::ir::schema::SketchParams; validate_family(family)?; @@ -302,6 +303,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/executor/tests/blocking_resources.rs b/crates/executor/tests/blocking_resources.rs index a407f7e0..23848db2 100644 --- a/crates/executor/tests/blocking_resources.rs +++ b/crates/executor/tests/blocking_resources.rs @@ -247,3 +247,47 @@ fn weighted_summary_build_yields_within_a_batch() { drop(output); assert_eq!(run.retained_bytes(), 0); } + +// Frequency dictionaries count against the budget and release reservations on failure. +#[test] +fn frequency_dictionary_enforces_memory_budget() { + use asap_executor::operators::Reduction; + let input = schema(1); + let mut sources = PhysicalDAG::default(); + sources + .add( + 0, + vec![], + Operator::source( + input.clone(), + vec![Batch::try_new( + input.clone(), + (0..64).map(|i| vec![Value::Int64(i)]).collect(), + ) + .unwrap()], + ) + .unwrap(), + ) + .unwrap(); + for (reduction, succeeds) in [ + (Reduction::Count, true), + (Reduction::FrequencyL2(0), false), + (Reduction::FrequencyEntropy(0), false), + (Reduction::Cardinality(vec![0]), false), + ] { + let run = context(12_000); + let inputs = sources.execute(&[0], run.clone()).unwrap(); + let operator = + Operator::aggregate(input.clone(), vec![], vec![("result".into(), reduction)]).unwrap(); + let mut output = operator.start(inputs, run.clone()).unwrap(); + let result = block_on(output.next()).unwrap(); + if succeeds { + assert!(result.is_ok()); + } else { + assert!(matches!(result, Err(Error::MemoryLimit))); + } + drop(result); + drop(output); + assert_eq!(run.retained_bytes(), 0); + } +} diff --git a/crates/executor/tests/physical_semantics.rs b/crates/executor/tests/physical_semantics.rs index a078cfbd..4591a785 100644 --- a/crates/executor/tests/physical_semantics.rs +++ b/crates/executor/tests/physical_semantics.rs @@ -713,3 +713,205 @@ fn empty_exact_summary_extrema_agree_with_ordinary_aggregation() { assert!(matches!(rows[0][0], Value::Null)); } } + +// Exact frequency intents bind to native reducers without a sketch or numeric key conversion. +#[test] +fn exact_frequency_intents_execute_typed_keys_and_empty_input() { + use asap_executor::physical_planner::compile_node; + use planner_types::ir::operator::{AggIntent, GroupKeys, Reduction as PlanReduction}; + use planner_types::ir::properties::*; + use planner_types::types::AccuracyTarget; + for (dtype, values) in [ + ( + DataType::Utf8, + vec![Value::Utf8("a".into()), Value::Utf8("b".into())], + ), + ( + DataType::Int64, + vec![ + Value::Int64(9_007_199_254_740_992), + Value::Int64(9_007_199_254_740_993), + ], + ), + (DataType::Bool, vec![Value::Bool(false), Value::Bool(true)]), + ( + DataType::Float64, + vec![Value::Float64(-0.0), Value::Float64(1.0)], + ), + ] { + let input = schema(&[("key", dtype, true)]); + for (measure, name, expected) in [ + ( + AggIntent::FrequencyL2 { + col: Some(0), + accuracy: AccuracyTarget::Exact, + }, + "frequency_l2", + 8.0_f64.sqrt(), + ), + ( + AggIntent::FrequencyEntropy { + col: Some(0), + accuracy: AccuracyTarget::Exact, + }, + "frequency_entropy", + 1.0, + ), + ] { + let node = PhysicalASAPDAGNode { + coverage: None, + id: planner_types::ir::export::LogicalASAPNodeId(1), + payload: PhysicalASAPOperatorPayload::Relational { + operator: ValueOperation::Aggregate { + reduction: PlanReduction::Reduce(GroupKeys::none()), + measures: vec![measure], + output_names: vec![name.into()], + filters: vec![], + having: None, + }, + }, + output_state: ExecutionDataState::QUERY_ROWS, + output_schema: (*schema(&[(name, DataType::Float64, false)])).clone(), + guarantee: None, + }; + let operator = compile_node(&node, std::slice::from_ref(&input)) + .expect("exact frequency intent binds"); + let rows = values + .iter() + .flat_map(|v| [vec![v.clone()], vec![v.clone()]]) + .chain([vec![Value::Null]]) + .collect(); + let result = unary(input.clone(), vec![rows], operator.clone()); + assert!(matches!(result[0][0], Value::Float64(v) if (v - expected).abs() < 1e-12)); + for batches in [vec![], vec![vec![vec![Value::Null]]]] { + let result = unary(input.clone(), batches, operator.clone()); + assert!(matches!(result[0][0], Value::Float64(0.0))); + } + } + } +} + +// Each group gets its own frequency population, including one canonical signed-zero identity. +#[test] +fn exact_frequency_grouping_and_entropy_bits() { + let input = schema(&[ + ("group", DataType::Int64, false), + ("key", DataType::Float64, true), + ]); + let operator = Operator::aggregate( + input.clone(), + vec![0], + vec![ + ("l2".into(), Reduction::FrequencyL2(1)), + ("entropy".into(), Reduction::FrequencyEntropy(1)), + ], + ) + .unwrap(); + let rows = vec![ + vec![Value::Int64(1), Value::Float64(-0.0)], + vec![Value::Int64(1), Value::Float64(0.0)], + vec![Value::Int64(1), Value::Float64(0.0)], + vec![Value::Int64(1), Value::Float64(1.0)], + vec![Value::Int64(2), Value::Float64(2.0)], + vec![Value::Int64(2), Value::Null], + vec![Value::Int64(3), Value::Null], + ]; + let result = unary(input.clone(), vec![rows], operator.clone()); + assert_eq!(result.len(), 3); + let expected_entropy = -0.75_f64 * 0.75_f64.log2() - 0.25_f64 * 0.25_f64.log2(); + for (row, l2, entropy) in [ + (&result[0], 10.0_f64.sqrt(), expected_entropy), + (&result[1], 1.0, 0.0), + (&result[2], 0.0, 0.0), + ] { + assert!(matches!(row[1], Value::Float64(v) if (v - l2).abs() < 1e-12)); + assert!(matches!(row[2], Value::Float64(v) if (v - entropy).abs() < 1e-12)); + } + assert!(unary(input, vec![], operator).is_empty()); +} + +// Exact distinct binding preserves typed tuples, skips NULLs and returns zero on empty input. +#[test] +fn exact_cardinality_binds_and_executes_typed_tuples() { + use asap_executor::physical_planner::compile_node; + use planner_types::ir::operator::{AggIntent, GroupKeys, Reduction as PlanReduction}; + use planner_types::ir::properties::*; + use planner_types::types::AccuracyTarget; + let input = schema(&[ + ("key", DataType::Int64, true), + ("tag", DataType::Utf8, true), + ]); + for (cols, expected) in [(vec![0], 2), (vec![0, 1], 3)] { + let node = PhysicalASAPDAGNode { + coverage: None, + id: planner_types::ir::export::LogicalASAPNodeId(1), + payload: PhysicalASAPOperatorPayload::Relational { + operator: ValueOperation::Aggregate { + reduction: PlanReduction::Reduce(GroupKeys::none()), + measures: vec![AggIntent::Cardinality { + cols, + accuracy: AccuracyTarget::Exact, + }], + output_names: vec!["distinct".into()], + filters: vec![], + having: None, + }, + }, + output_state: ExecutionDataState::QUERY_ROWS, + output_schema: (*schema(&[("distinct", DataType::Int64, false)])).clone(), + guarantee: None, + }; + let operator = + compile_node(&node, std::slice::from_ref(&input)).expect("exact distinct intent binds"); + let rows = vec![ + vec![Value::Int64(9_007_199_254_740_992), Value::Utf8("a".into())], + vec![Value::Int64(9_007_199_254_740_992), Value::Utf8("a".into())], + vec![Value::Int64(9_007_199_254_740_992), Value::Utf8("b".into())], + vec![Value::Int64(9_007_199_254_740_993), Value::Utf8("a".into())], + vec![Value::Null, Value::Utf8("c".into())], + ]; + let result = unary(input.clone(), vec![rows], operator.clone()); + assert!(matches!(result[0][0], Value::Int64(v) if v == expected)); + for rows in [vec![], vec![vec![Value::Null, Value::Null]]] { + let result = unary(input.clone(), vec![rows], operator.clone()); + assert!(matches!(result[0][0], Value::Int64(0))); + } + } +} + +// Distinct uses grouped equality: signed zero and NaN payloads each form one identity. +#[test] +fn exact_cardinality_grouping_normalizes_float_identities() { + let input = schema(&[ + ("group", DataType::Int64, false), + ("key", DataType::Float64, true), + ]); + let operator = Operator::aggregate( + input.clone(), + vec![0], + vec![("distinct".into(), Reduction::Cardinality(vec![1]))], + ) + .unwrap(); + let result = unary( + input, + vec![vec![ + vec![Value::Int64(1), Value::Float64(0.0)], + vec![Value::Int64(1), Value::Float64(-0.0)], + vec![Value::Int64(1), Value::Float64(f64::NAN)], + vec![ + Value::Int64(1), + Value::Float64(f64::from_bits(f64::NAN.to_bits() + 1)), + ], + vec![Value::Int64(2), Value::Null], + ]], + operator, + ); + assert!(matches!( + result[0].as_slice(), + [Value::Int64(1), Value::Int64(2)] + )); + assert!(matches!( + result[1].as_slice(), + [Value::Int64(2), Value::Int64(0)] + )); +} diff --git a/crates/executor/tests/univmon_execution.rs b/crates/executor/tests/univmon_execution.rs new file mode 100644 index 00000000..b402d43a --- /dev/null +++ b/crates/executor/tests/univmon_execution.rs @@ -0,0 +1,182 @@ +//! One native UnivMon build supplies the three statistics in design Example 2. +use asap_executor::{ + operators::{Operator, SummaryEvaluation}, + plan::{PhysicalDAG, PhysicalOperator}, + runtime::{Limits, RunContext, Scope}, + values::{Batch, SchemaRef, Value}, + Error, +}; +use futures::{executor::block_on, StreamExt}; +use planner_types::ir::schema::{ + DataType, Field, FieldDataType, GroupingStrategy, SketchAlgorithm, SketchKind, SketchParams, + SketchStatistic, +}; +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::ir::schema::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(()) +} + +/// SQL source IPs and integer identifiers enter frequency state without numeric coercion. +#[test] +fn typed_frequency_keys_preserve_identity() -> Result<(), Error> { + for (dtype, keys) in [ + ( + DataType::Utf8, + vec![ + Value::Utf8("192.0.2.1".into()), + Value::Utf8("192.0.2.2".into()), + ], + ), + ( + DataType::Int64, + vec![ + Value::Int64(9_007_199_254_740_992), + Value::Int64(9_007_199_254_740_993), + ], + ), + (DataType::Bool, vec![Value::Bool(false), Value::Bool(true)]), + ] { + let input = Arc::new(planner_types::ir::schema::Schema::new(vec![Field::plain( + "src_ip", dtype, true, + )])); + let batch = Batch::try_new( + input.clone(), + vec![ + vec![keys[0].clone()], + vec![keys[0].clone()], + vec![keys[1].clone()], + vec![keys[1].clone()], + vec![Value::Null], + ], + )?; + let build = Operator::summary_build(input, family(), 0, None, vec![])?; + let output = build.output_schema(); + let mut dag = PhysicalDAG::default(); + dag.add( + 0, + vec![], + Operator::source(batch.schema().clone(), vec![batch])?, + )?; + dag.add(1, vec![0], build)?; + for (id, statistic) in [ + (2, SketchStatistic::Cardinality), + (3, SketchStatistic::FrequencyL2), + (4, SketchStatistic::FrequencyEntropy), + ] { + dag.add( + id, + vec![1], + Operator::evaluation(output.clone(), 0, SummaryEvaluation::Sketch(statistic))?, + )?; + } + 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::>()), + )); + for (stream, expected) in outputs.into_iter().zip([2.0, 8.0_f64.sqrt(), 1.0]) { + let batches = stream.into_iter().collect::, _>>()?; + assert!((value(&batches[0].rows()[0]) - expected).abs() < 0.01); + } + } + Ok(()) +} diff --git a/crates/integration-tests/tests/physical_common/mod.rs b/crates/integration-tests/tests/physical_common/mod.rs index 3403d27e..f87dd315 100644 --- a/crates/integration-tests/tests/physical_common/mod.rs +++ b/crates/integration-tests/tests/physical_common/mod.rs @@ -72,3 +72,58 @@ fn compile_with( asap_types::ir::apply_materialization_timings(root, assignment, &mut Default::default())?; Ok(asap_types::ir::export::compile_physical_asap_dag(&root)?) } + +/// Execute a relational DAG against a raw in-memory connector for its one +/// scan, including scan predicates; checks that no memory stays reserved. +#[allow(dead_code)] +pub fn execute_raw_rows( + root: &std::rc::Rc, + rows: Vec>, +) -> Vec> { + use asap_executor::{ + physical_planner::bind_with_data_sources, + sources::{DataSources, MemorySource}, + }; + use asap_types::ir::export::{NonASAPOpKind, PhysicalASAPOperatorPayload}; + use std::sync::Arc; + let wire = compile_physical_asap_dag(root).unwrap(); + let (source, schema) = wire + .nodes + .iter() + .find_map(|node| match &node.payload { + PhysicalASAPOperatorPayload::Relational { + operator: NonASAPOpKind::Scan { source, .. }, + } => Some((source.clone(), node.output_schema.clone())), + _ => None, + }) + .unwrap(); + let input = Arc::new(schema); + let batch = Batch::try_new(input.clone(), rows).unwrap(); + let mut sources = DataSources::default(); + sources + .register( + source, + Arc::new(MemorySource::new(input, vec![batch]).unwrap()), + ) + .unwrap(); + let root_id = u64::from(wire.roots[0].0); + let plan = bind_with_data_sources(&wire, BTreeMap::new(), &[root_id], &sources).unwrap(); + let context = RunContext::new( + Scope::Query { + evaluation_time_ms: 0, + revision: 1, + }, + Limits::default(), + ) + .unwrap(); + let result = block_on(async { + let mut output = plan.execute(&[root_id], context.clone()).unwrap().remove(0); + let mut rows = vec![]; + while let Some(batch) = output.next().await { + rows.extend_from_slice(batch.unwrap().rows()); + } + rows + }); + assert_eq!(context.retained_bytes(), 0); + result +} diff --git a/crates/integration-tests/tests/precompute_raw_samples.rs b/crates/integration-tests/tests/precompute_raw_samples.rs index d51358cb..50f80821 100644 --- a/crates/integration-tests/tests/precompute_raw_samples.rs +++ b/crates/integration-tests/tests/precompute_raw_samples.rs @@ -277,6 +277,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:?}"), } } @@ -322,7 +330,7 @@ fn check( let stored_only = matches!(family, FieldDataType::Sketch(kind, _) if kind.algorithm() == &asap_types::ir::schema::SketchAlgorithm::Cms); if stored_only || asap_executor::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()); @@ -335,7 +343,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(); diff --git a/crates/integration-tests/tests/sql_cardinality.rs b/crates/integration-tests/tests/sql_cardinality.rs new file mode 100644 index 00000000..dd1e9790 --- /dev/null +++ b/crates/integration-tests/tests/sql_cardinality.rs @@ -0,0 +1,46 @@ +//! The design's SQL distinct query executes its exact native fallback. +mod physical_common; +use asap_executor::values::Value; +use asap_frontend_sql::{lower_sql, SqlCatalog}; +use asap_types::{ + ir::schema::{DataType, Field, Schema}, + types::AccuracyTarget, +}; + +// COUNT(DISTINCT src_ip) skips NULL and preserves WHERE, Utf8 identities and zero on no input. +#[tokio::test] +async fn sql_distinct_executes_through_raw_scan_and_native_binding() { + let catalog = SqlCatalog::new().with_table( + "flows", + Schema::new(vec![ + Field::plain("src_ip", DataType::Utf8, true), + Field::plain("keep", DataType::Bool, false), + ]), + ); + let root = lower_sql( + "SELECT COUNT(DISTINCT src_ip) AS sources FROM flows WHERE keep", + &catalog, + AccuracyTarget::Exact, + ) + .await + .unwrap(); + for (rows, expected) in [ + (vec![], 0), + (vec![vec![Value::Null, Value::Bool(true)]], 0), + ( + vec![ + vec![Value::Utf8("a".into()), Value::Bool(true)], + vec![Value::Utf8("a".into()), Value::Bool(true)], + vec![Value::Utf8("b".into()), Value::Bool(true)], + vec![Value::Utf8("discard".into()), Value::Bool(false)], + vec![Value::Null, Value::Bool(true)], + ], + 2, + ), + ] { + let result = physical_common::execute_raw_rows(&root, rows); + assert!( + matches!(result.as_slice(), [row] if matches!(row.as_slice(), [Value::Int64(v)] if *v == expected)) + ); + } +} diff --git a/crates/plan-selection/src/cost/query_physical_lowering.rs b/crates/plan-selection/src/cost/query_physical_lowering.rs index 09693597..79073ba1 100644 --- a/crates/plan-selection/src/cost/query_physical_lowering.rs +++ b/crates/plan-selection/src/cost/query_physical_lowering.rs @@ -1204,6 +1204,9 @@ fn supports_hash_aggregate( | AggIntent::PearsonCorr { .. } | AggIntent::Group | AggIntent::CountValues { .. } + | AggIntent::Cardinality { .. } + | AggIntent::FrequencyL2 { .. } + | AggIntent::FrequencyEntropy { .. } ) }) }