From 701add4161d2aec436be8b014ae8e61818e237ed Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Sat, 3 Oct 2026 15:12:21 +0000 Subject: [PATCH] feat(runtime): execute exact distinct identity counts --- .../src/query_physical_lowering.rs | 1 + .../src/operators/aggregate/mod.rs | 34 ++++++++ .../src/physical_planner/mod.rs | 7 ++ .../tests/blocking_resources.rs | 1 + .../tests/physical_semantics.rs | 87 +++++++++++++++++++ .../tests/physical_common/mod.rs | 64 ++++++++++++++ .../tests/sql_cardinality.rs | 46 ++++++++++ .../tests/sql_frequency_l2.rs | 71 +-------------- docs/develop_docs/planner-layering-status.md | 15 +++- 9 files changed, 258 insertions(+), 68 deletions(-) create mode 100644 crates/integration-tests/tests/sql_cardinality.rs diff --git a/crates/asap-aware-mapping/src/query_physical_lowering.rs b/crates/asap-aware-mapping/src/query_physical_lowering.rs index 996898a94..76fd2f8b3 100644 --- a/crates/asap-aware-mapping/src/query_physical_lowering.rs +++ b/crates/asap-aware-mapping/src/query_physical_lowering.rs @@ -1203,6 +1203,7 @@ fn supports_hash_aggregate( | AggIntent::PearsonCorr { .. } | AggIntent::Group | AggIntent::CountValues { .. } + | AggIntent::Cardinality { .. } | AggIntent::FrequencyL2 { .. } | AggIntent::FrequencyEntropy { .. } ) diff --git a/crates/asap-physical-operators/src/operators/aggregate/mod.rs b/crates/asap-physical-operators/src/operators/aggregate/mod.rs index c9715a41b..11b5d6049 100644 --- a/crates/asap-physical-operators/src/operators/aggregate/mod.rs +++ b/crates/asap-physical-operators/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) { @@ -134,6 +145,8 @@ 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), @@ -265,6 +278,27 @@ 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(); diff --git a/crates/asap-physical-operators/src/physical_planner/mod.rs b/crates/asap-physical-operators/src/physical_planner/mod.rs index 59d734d8e..9372589ca 100644 --- a/crates/asap-physical-operators/src/physical_planner/mod.rs +++ b/crates/asap-physical-operators/src/physical_planner/mod.rs @@ -826,6 +826,13 @@ fn bind_operation(node: &PostAsapDAGNode, 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, .. } => { diff --git a/crates/asap-physical-operators/tests/blocking_resources.rs b/crates/asap-physical-operators/tests/blocking_resources.rs index 579f7978d..873eef8d5 100644 --- a/crates/asap-physical-operators/tests/blocking_resources.rs +++ b/crates/asap-physical-operators/tests/blocking_resources.rs @@ -274,6 +274,7 @@ fn frequency_dictionary_enforces_memory_budget() { (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(); diff --git a/crates/asap-physical-operators/tests/physical_semantics.rs b/crates/asap-physical-operators/tests/physical_semantics.rs index 815108f1f..2ce2bda66 100644 --- a/crates/asap-physical-operators/tests/physical_semantics.rs +++ b/crates/asap-physical-operators/tests/physical_semantics.rs @@ -861,3 +861,90 @@ fn sql_sqrt_executes_numeric_and_null_arguments() { matches!(compiled.evaluate(&[Value::Float64(-1.0)]).unwrap(), Value::Float64(v) if v.is_nan()) ); } + +// 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_physical_operators::physical_planner::compile_node; + use planner_types::{ + post_asap::ExecutionDataState, + pre_asap::{AggIntent, GroupKeys, Reduction as PlanReduction}, + 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 = PostAsapDAGNode { + id: PostAsapNodeId(1), + payload: PostAsapOperatorPayload::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/integration-tests/tests/physical_common/mod.rs b/crates/integration-tests/tests/physical_common/mod.rs index 1612649f9..5184d836d 100644 --- a/crates/integration-tests/tests/physical_common/mod.rs +++ b/crates/integration-tests/tests/physical_common/mod.rs @@ -53,3 +53,67 @@ pub fn compile_post_asap_dag( )?; Ok(asap_types::ir::export::compile_post_asap_dag(&root)?) } + +// Execute a selected relational DAG against real raw connectors, including scan predicates. +#[allow(dead_code)] +pub fn execute_raw_rows( + root: &std::rc::Rc, + rows: Vec>, +) -> Vec> { + use asap_physical_operators::{ + physical_planner::bind_with_data_sources, + runtime::{Limits, RunContext}, + sources::{DataSources, MemorySource}, + }; + use asap_types::ir::export::{NonASAPOpKind, PostAsapOperatorPayload}; + use futures::{executor::block_on, StreamExt}; + use std::sync::Arc; + let wire = compile_post_asap_dag(root).unwrap(); + let scan = wire + .nodes + .iter() + .find(|node| { + matches!( + node.payload, + PostAsapOperatorPayload::Relational { + operator: NonASAPOpKind::Scan { .. } + } + ) + }) + .unwrap(); + let PostAsapOperatorPayload::Relational { + operator: NonASAPOpKind::Scan { source, .. }, + } = &scan.payload + else { + unreachable!(); + }; + let input = Arc::new(scan.output_schema.clone()); + let batch = Batch::try_new(input.clone(), rows).unwrap(); + let mut sources = DataSources::default(); + sources + .register( + source.clone(), + Arc::new(MemorySource::new(input, vec![batch]).unwrap()), + ) + .unwrap(); + let root_id = u64::from(wire.root.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/sql_cardinality.rs b/crates/integration-tests/tests/sql_cardinality.rs new file mode 100644 index 000000000..82ce6e597 --- /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_frontend_sql::{lower_sql, SqlCatalog}; +use asap_physical_operators::values::Value; +use asap_types::{ + pre_asap::{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/integration-tests/tests/sql_frequency_l2.rs b/crates/integration-tests/tests/sql_frequency_l2.rs index f7bea5760..5c79545b8 100644 --- a/crates/integration-tests/tests/sql_frequency_l2.rs +++ b/crates/integration-tests/tests/sql_frequency_l2.rs @@ -5,76 +5,13 @@ use asap_aware_mapping::{ SemanticEquivalentRewriteStrategy, }; use asap_frontend_sql::{lower_sql, SqlCatalog}; -use asap_physical_operators::{ - runtime::Scope, - values::{Batch, Value}, -}; +use asap_physical_operators::values::Value; use asap_types::{ - ir::{ - export::{NonASAPOpKind, PostAsapOperatorPayload}, - NonASAPOp, OperatorNode, - }, + ir::NonASAPOp, pre_asap::{DataType, Field, Schema}, types::AccuracyTarget, }; -use std::{collections::BTreeMap, rc::Rc, sync::Arc}; -fn run(root: &Rc, rows: Vec>) -> Vec> { - use asap_physical_operators::{ - physical_planner::bind_with_data_sources, - runtime::{Limits, RunContext}, - sources::{DataSources, MemorySource}, - }; - use futures::{executor::block_on, StreamExt}; - let wire = physical_common::compile_post_asap_dag(root).unwrap(); - let scan = wire - .nodes - .iter() - .find(|node| { - matches!( - node.payload, - PostAsapOperatorPayload::Relational { - operator: NonASAPOpKind::Scan { .. } - } - ) - }) - .unwrap(); - let PostAsapOperatorPayload::Relational { - operator: NonASAPOpKind::Scan { source, .. }, - } = &scan.payload - else { - unreachable!(); - }; - let input = Arc::new(scan.output_schema.clone()); - let batch = Batch::try_new(input.clone(), rows).unwrap(); - let mut sources = DataSources::default(); - sources - .register( - source.clone(), - Arc::new(MemorySource::new(input, vec![batch]).unwrap()), - ) - .unwrap(); - let root_id = u64::from(wire.root.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 -} // Original SQL and its logical alternative agree on filters and SQL's empty-input NULL. #[tokio::test] async fn sql_l2_original_and_rewrite_execute_equivalently() { @@ -104,8 +41,8 @@ async fn sql_l2_original_and_rewrite_execute_equivalently() { vec![Value::Utf8("discard".into()), Value::Bool(false)], ], ] { - let original = run(&root, rows.clone()); - let actual = run(rewritten, rows); + let original = physical_common::execute_raw_rows(&root, rows.clone()); + let actual = physical_common::execute_raw_rows(rewritten, rows); if original.iter().any(|row| !matches!(row[0], Value::Null)) { assert!( matches!(original[0][0], Value::Float64(v) if (v - 5.0_f64.sqrt()).abs() < 1e-12) diff --git a/docs/develop_docs/planner-layering-status.md b/docs/develop_docs/planner-layering-status.md index f2df5ec29..b0b6c7715 100644 --- a/docs/develop_docs/planner-layering-status.md +++ b/docs/develop_docs/planner-layering-status.md @@ -12,7 +12,7 @@ is a target contract, not a statement that its examples execute today. | Window composition | Mergeable state IR/native merge exists; physical pane compatibility and reuse cost helpers exist. | Automatic logical sliding/tumbling/EH alternatives over differing windows, boundary coverage and error proofs are absent. A merge kernel alone does not implement Examples 1/3. | | Physical materialization | Ephemeral/prepared/shared/continuously maintained lifecycle alternatives, costing, capabilities and latency checks exist. | Incremental query-time pane retention, historical backfill and the complete Example 4 matrix need executable implementations and explicit state/input contracts. | | Whole-workload selection | One unified selected DAG; shared states are interned and costed across their consumers. | `replacement.rs` documents its selection as non-exhaustive over interacting choices. The proposal's cheapest complete candidate guarantee and 54/156 inventories need a complete workload search/selection path. | -| Deployment inputs and execution | `PlanningModels` bundles cost, accuracy, evidence and capabilities; native typed UnivMon supports one build with three readouts. | At #557 native exact frequency L2/entropy fallback is absent. End-to-end SQL Example 2 is not established by the native UnivMon fixture. | +| Deployment inputs and execution | `PlanningModels` bundles cost, accuracy, evidence and capabilities; native typed UnivMon supports one build with three readouts. | At #557 native exact distinct/L2/entropy fallback is absent; the follow-ups supply those native bindings. End-to-end SQL Example 2 is not established by the native UnivMon fixture. | | Subtract/delete, parallelism, partitioning and resource planning | Some runtime memory/cancellation limits and maintenance capability flags exist. | These remain proposal TODOs; capability flags do not supply missing IR operators or a physical resource search. | ## Follow-up sequence @@ -61,3 +61,16 @@ product in Example 2 remains a gap because SQL overflow is observable. accuracy propagation and candidate retention; `integration-tests/tests/sql_frequency_l2.rs` executes both SQL and the rewrite through raw connectors and wire compilation. This step does not supply an L2 accuracy certificate or all sharing partitions. + +## Exact distinct follow-up acceptance + +The native binding executes `AggIntent::Cardinality` over one typed identity or +an ordered tuple, supports grouping, skips tuples containing NULL and returns +an Int64 zero for empty input. Typed key encoding preserves large integers, +normalizes signed zero and NaN payloads, and retains tuple boundaries. Dictionary +workspace is reserved and released under the normal execution limits. An empty +intent column list retains the existing sample-value convention. + +`integration-tests/tests/sql_cardinality.rs` lowers `COUNT(DISTINCT src_ip)` and +executes its exact path through raw scan predicates and native wire binding. +Filtered aggregate measures remain outside native binding's existing scope.