diff --git a/crates/integration-tests/tests/filtered_aggregates.rs b/crates/integration-tests/tests/filtered_aggregates.rs index 578318d6b..e83a4e3dc 100644 --- a/crates/integration-tests/tests/filtered_aggregates.rs +++ b/crates/integration-tests/tests/filtered_aggregates.rs @@ -411,8 +411,8 @@ async fn alternatives(sql: &str, target: AccuracyTarget) -> Vec<(String, Rc = filtered.iter().map(|(label, _)| label.as_str()).collect(); - assert_eq!( - labels, - [ - "", - "ExactAggregate(Count, Count)", - "Cms", - "CountSketch", - "UnivMon", - "HydraCms" - ] - ); + assert_eq!(labels, ["", "ExactAggregate(Count, Count)", "HydraCms"]); let expected = printed([ vec![s("a"), Value::Int64(2)], vec![s("b"), Value::Int64(0)], @@ -464,13 +454,8 @@ async fn pass1_offers_filtered_count_alternatives() { ]); for ((label, root), plain) in filtered.iter().zip(&plain) { assert_eq!(binds(root).is_ok(), binds(plain).is_ok(), "{label}"); - if matches!( - label.as_str(), - "" | "ExactAggregate(Count, Count)" | "HydraCms" - ) { - // Few groups in a wide grid: Hydra's estimate is exact here. - assert_eq!(sorted(root), expected, "{label}"); - } + // Few groups in a wide grid: Hydra's estimate is exact here. + assert_eq!(sorted(root), expected, "{label}"); } } diff --git a/crates/integration-tests/tests/pass1_sql_coverage.rs b/crates/integration-tests/tests/pass1_sql_coverage.rs index 0b0daf97b..96fa5b311 100644 --- a/crates/integration-tests/tests/pass1_sql_coverage.rs +++ b/crates/integration-tests/tests/pass1_sql_coverage.rs @@ -1,6 +1,7 @@ //! Pass 1 alternatives over SQL row sources compose, compile and execute. mod executor_models; mod physical_common; +use asap_types::ir::NonASAPOp; use std::collections::BTreeMap; use std::rc::Rc; @@ -8,7 +9,6 @@ use asap_executor::values::Value; use asap_frontend_sql::{lower_sql, SqlCatalog}; use asap_logical_optimizer::pass1::logical_candidates::{ compose_logical_candidate, enumerate_choices, enumerate_local_logical_candidates, - LocalLogicalCandidates, }; use asap_plan_selection::plan_stages; use asap_types::ir::schema::{DataType, Field, Schema}; @@ -44,55 +44,146 @@ fn rows() -> Vec> { .collect() } -async fn inventory(sql: &str) -> LocalLogicalCandidates { - let root = lower_sql(sql, &catalog(), target()).await.unwrap(); - enumerate_local_logical_candidates(vec![(0, QueryRoot::Operator(root))], &BTreeMap::new()) - .unwrap() +/// The ε = 0.1 target: at 0.01 a HydraCms grid exceeds the default memory +/// limit, so Stage 3 would not price it. +fn coarse() -> AccuracyTarget { + AccuracyTarget::EpsilonDelta { + epsilon: 0.1, + delta: 0.01, + } } -fn has_estimate(root: &Rc) -> bool { - OperatorNode::reachable(root) - .iter() - .any(|n| matches!(n.operator, Operator::ASAP(ASAPOp::SummaryEstimate { .. }))) +/// SQL `COUNT(*)`, grouped, ungrouped and ranked, each with the rows it +/// returns from [`rows`]. +fn count_star_queries() -> [(&'static str, Vec>); 3] { + let count = |ip: &str, n| vec![Value::Utf8(ip.into()), Value::Int64(n)]; + [ + ( + "SELECT src_ip, COUNT(*) AS c FROM flows GROUP BY src_ip", + vec![count("a", 3), count("b", 2), count("c", 1)], + ), + ("SELECT COUNT(*) AS c FROM flows", vec![vec![Value::Int64(6)]]), + ( + "SELECT src_ip, COUNT(*) AS c FROM flows GROUP BY src_ip ORDER BY COUNT(*) DESC LIMIT 2", + vec![count("a", 3), count("b", 2)], + ), + ] } -/// `COUNT(*) … GROUP BY src_ip` (#509 Example 2's inner query) reads no -/// sample value. Every candidate composes and compiles; the exact ones, -/// the exact `Count` accumulator included, count rows per `src_ip`. +/// SQL `COUNT(*)` (#509 Example 2's inner query) reads no sample value, and +/// a group's rows all hash its key, so Pass 1 offers no per-group sketch. +/// Every candidate composes, compiles, binds and executes to the exact +/// counts: HydraCms too, whose few groups in a wide grid never collide. #[tokio::test] -async fn sql_count_star_group_by_candidates_compose_and_execute() { - let inventory = inventory("SELECT src_ip, COUNT(*) AS c FROM flows GROUP BY src_ip").await; - let choices = enumerate_choices(&inventory, usize::MAX); - assert!(choices.len() > 2, "exact and summary candidates"); - let mut executed = 0; - for choice in choices { - let roots = compose_logical_candidate(&inventory, &choice) - .unwrap_or_else(|e| panic!("{choice:?} composes: {e}")); - let QueryRoot::Operator(root) = &roots[0].1 else { - panic!("operator root") +async fn sql_count_star_candidates_compose_and_execute() { + // DataFusion 54 sorts the ranked query by the aggregate's own column, so + // `canonicalize` promotes it to `TopK` over the Count, as it does + // `ORDER BY c`. The executor has no native `TopK`, so only the unranked + // queries execute every choice here; the ranked query's priced + // candidates still bind (`priced_sql_count_star_candidates_bind`). + for (sql, expected) in count_star_queries().into_iter().take(2) { + let expected: Vec<_> = expected.iter().map(|row| format!("{row:?}")).collect(); + let root = lower_sql(sql, &catalog(), coarse()).await.unwrap(); + let inventory = enumerate_local_logical_candidates( + vec![(0, QueryRoot::Operator(root))], + &BTreeMap::new(), + ) + .unwrap(); + let choices = enumerate_choices(&inventory, usize::MAX); + assert!(choices.len() >= 2, "{sql}: pass-through and exact Count"); + for choice in choices { + let roots = compose_logical_candidate(&inventory, &choice) + .unwrap_or_else(|e| panic!("{sql}: {choice:?} composes: {e}")); + let QueryRoot::Operator(root) = &roots[0].1 else { + panic!("operator root") + }; + let mut rows: Vec<_> = physical_common::execute_raw_rows(root, rows()) + .iter() + .map(|row| format!("{row:?}")) + .collect(); + rows.sort(); + assert_eq!(rows, expected, "{sql}: {choice:?}"); + } + } +} + +/// Binds every root of `dag` in the executor, each scan reading an empty +/// in-memory source. +fn binds(dag: &asap_types::ir::physical_export::PhysicalASAPDAG) -> Result<(), String> { + use asap_executor::physical_planner::bind_with_data_sources; + use asap_executor::sources::{DataSources, MemorySource}; + use asap_types::ir::physical_export::PhysicalASAPOperatorPayload; + use std::sync::Arc; + let mut sources = DataSources::default(); + let mut registered = vec![]; + for node in &dag.nodes { + if let PhysicalASAPOperatorPayload::NonASAP(NonASAPOp::Scan { source, .. }) = &node.payload + { + if !registered.contains(source) { + let schema = Arc::new(node.output_schema.clone()); + sources + .register( + source.clone(), + Arc::new(MemorySource::new(schema, vec![]).unwrap()), + ) + .unwrap(); + registered.push(source.clone()); + } + } + } + let roots: Vec<_> = dag.roots.iter().map(|&root| root as u64).collect(); + bind_with_data_sources(dag, BTreeMap::new(), &roots, &sources) + .map(|_| ()) + .map_err(|e| e.to_string()) +} + +/// Every physical candidate Stage 3 prices binds in the executor: a priced +/// plan the executor cannot run could be selected. +async fn assert_priced_candidates_bind(queries: &[&str], target: AccuracyTarget) { + let mut roots = vec![]; + for (i, sql) in queries.iter().enumerate() { + let root = lower_sql(sql, &catalog(), target.clone()).await.unwrap(); + roots.push((i, QueryRoot::Operator(root))); + } + let demand = vec![ + RootDemand { + accuracy: Some(target), + recurrence: QueryRecurrence::OneTime { + invocations: 1, + execute_at: None, + }, + predictability: Predictability::default(), + latency_ms: None, }; - physical_common::compile_physical_asap_dag(root) - .unwrap_or_else(|e| panic!("{choice:?} compiles: {e}")); - if has_estimate(root) { - continue; + roots.len() + ]; + 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(roots, &demand, &data, executor_models(), 4096).unwrap(); + let enumeration = run.enumeration.unwrap(); + let mut priced = 0; + for physical in enumeration.candidates.iter().flat_map(|c| &c.physical) { + if enumeration.selection.costs.contains_key(&physical.id) { + binds(&physical.dag).unwrap_or_else(|e| panic!("{queries:?}: {}: {e}", physical.id)); + priced += 1; } - 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!("{choice:?}: unexpected row {other:?}"), - }) - .collect(); - assert_eq!( - counts, - [("a".into(), 3), ("b".into(), 2), ("c".into(), 1)], - "{choice:?}" - ); - executed += 1; } - assert_eq!(executed, 2, "pass-through and the exact Count accumulator"); + assert!( + priced > 0, + "{queries:?}: {:#?}", + enumeration.selection.rejected + ); +} + +#[tokio::test] +async fn priced_sql_count_star_candidates_bind() { + for (sql, _) in count_star_queries() { + assert_priced_candidates_bind(&[sql], coarse()).await; + } } fn declared(value: T) -> Evidence { @@ -103,18 +194,34 @@ fn declared(value: T) -> Evidence { } } +/// #509 Example 2's design queries, without their time window: the +/// runtime has no `now()` to bind it with. +const EXAMPLE2: [&str; 3] = [ + "SELECT COUNT(DISTINCT src_ip) FROM flows", + "SELECT -SUM(p * LN(p)) FROM (SELECT COUNT(*) * 1.0 / SUM(COUNT(*)) OVER () AS p FROM flows GROUP BY src_ip)", + "SELECT SQRT(SUM(c * c)) FROM (SELECT src_ip, COUNT(*) AS c FROM flows GROUP BY src_ip)", +]; + +/// Example 2 with Q3's floating product (as `planner_layering_example2`'s +/// `Q3_FLOAT`): the integer Q3's exact `Sum` over Int64 `c * c` is priced +/// but does not bind, since the executor sums only Float64 columns. +#[tokio::test] +async fn priced_example2_candidates_bind() { + let queries = [ + EXAMPLE2[0], + EXAMPLE2[1], + "SELECT SQRT(SUM(CAST(c AS DOUBLE) * CAST(c AS DOUBLE))) FROM (SELECT src_ip, COUNT(*) AS c FROM flows GROUP BY src_ip)", + ]; + assert_priced_candidates_bind(&queries, target()).await; +} + /// #509 Example 2's design queries (its integer Q3 keeps `COUNT(*) GROUP BY /// src_ip` as a target): every candidate builds through Stages 1 and 2. /// Stage 3 may still reject one, e.g. for accuracy. #[tokio::test] async fn example2_design_candidates_all_build() { - const QUERIES: [&str; 3] = [ - "SELECT COUNT(DISTINCT src_ip) FROM flows", - "SELECT -SUM(p * LN(p)) FROM (SELECT COUNT(*) * 1.0 / SUM(COUNT(*)) OVER () AS p FROM flows GROUP BY src_ip)", - "SELECT SQRT(SUM(c * c)) FROM (SELECT src_ip, COUNT(*) AS c FROM flows GROUP BY src_ip)", - ]; let mut roots = vec![]; - for (i, sql) in QUERIES.into_iter().enumerate() { + for (i, sql) in EXAMPLE2.into_iter().enumerate() { let root = lower_sql(sql, &catalog(), target()).await.unwrap(); roots.push((i, QueryRoot::Operator(root))); } diff --git a/crates/logical-optimizer/src/pass1/logical_candidates.rs b/crates/logical-optimizer/src/pass1/logical_candidates.rs index ba5d19391..1e57af85b 100644 --- a/crates/logical-optimizer/src/pass1/logical_candidates.rs +++ b/crates/logical-optimizer/src/pass1/logical_candidates.rs @@ -142,7 +142,17 @@ pub fn enumerate_local_logical_candidates( } if let Some(NonASAPOp::Aggregate { measures, .. }) = node.non_asap() { if let [intent] = measures.as_slice() { - let alternatives = local_realizations_for_intent(intent)?; + let mut alternatives = local_realizations_for_intent(intent)?; + // A SQL `COUNT(*)` group's rows all hash its key, so + // a per-group sketch is a counter with extra memory; + // only the shared Hydra grid (added below) helps. + if matches!(intent, AggIntent::Count { .. }) + && node.children().iter().all(|child| { + child.schema.closed && !child.schema.has_promql_series_identity() + }) + { + alternatives.retain(|a| !matches!(a, Realization::Sketch(_))); + } let mut target = LocalLogicalTarget { absorbs: vec![None; alternatives.len()], windows: vec![WindowForm::Whole; alternatives.len()], @@ -1125,6 +1135,71 @@ mod tests { assert!(hydra("count by (job) (m)", AccuracyTarget::Exact).is_empty()); } + /// A SQL `COUNT(*)` group's rows all hash its key, so Pass 1 offers no + /// per-group sketch: pass-through, the exact `Count` and HydraCms. A + /// PromQL count keeps Count-Min, Count Sketch and UnivMon. + #[test] + fn sql_count_star_offers_no_per_group_sketch() { + let approximate = AccuracyTarget::EpsilonDelta { + epsilon: 0.1, + delta: 0.01, + }; + let offered = |root: Rc| { + let inventory = enumerate_local_logical_candidates( + vec![(0, QueryRoot::Operator(root))], + &BTreeMap::new(), + ) + .unwrap(); + let target = &inventory.targets[0]; + target + .alternatives + .iter() + .zip(&target.groupings) + .map(|(a, g)| { + let family = match a { + Realization::PassThrough => "PassThrough".to_string(), + Realization::ExactAggregate { kind, .. } => format!("{kind:?}"), + Realization::Sketch(kind) => format!("{:?}", kind.algorithm()), + other => format!("{other:?}"), + }; + match *g == GroupingStrategy::default() { + true => family, + false => format!("Hydra{family}"), + } + }) + .collect::>() + }; + let mut schema = Schema::new(vec![ + asap_types::ir::schema::Field::plain("ts", DataType::Timestamp, false), + asap_types::ir::schema::Field::plain("src_ip", DataType::Utf8, false), + ]); + schema.closed = true; + let rows = crate::test_support::scan_from( + Source::Table { + table_ref: "flows".into(), + }, + schema, + ); + let count = AggIntent::Count { + accuracy: approximate.clone(), + }; + assert_eq!( + offered(crate::test_support::agg(vec![1], count, rows)), + ["PassThrough", "Count", "HydraCms"] + ); + for query in ["count by (job) (m)", "count_over_time(m[1m])"] { + let root = lower_promql(query, approximate.clone()); + let root = asap_types::ir::schema_support::with_promql_series_identity(&root).unwrap(); + let families = offered(root); + for sketch in ["Cms", "CountSketch", "UnivMon"] { + assert!( + families.iter().any(|f| f == sketch), + "{query}: {families:?}" + ); + } + } + } + /// Approximate requests must retain the exact execution alternative too. #[test] fn approximate_count_keeps_exact_and_universal_choices() { diff --git a/crates/planner/tests/stage_pipeline_selection.rs b/crates/planner/tests/stage_pipeline_selection.rs index bcec35982..322e698d3 100644 --- a/crates/planner/tests/stage_pipeline_selection.rs +++ b/crates/planner/tests/stage_pipeline_selection.rs @@ -430,8 +430,8 @@ async fn sql_hydra_count_dp_equals_exhaustive() { .targets .iter() .any(|t| t.groupings.iter().any(|g| *g != Default::default()))); - // (pass-through, Count acc, CMS, CountSketch, UnivMon, HydraCms) × (pass-through, KLL, DDSketch). - assert_dp_matches_exhaustive(&inventory, &workload, 18); + // (pass-through, Count acc, HydraCms) × (pass-through, KLL, DDSketch). + assert_dp_matches_exhaustive(&inventory, &workload, 9); } /// Filtered single-measure aggregates (`FILTER (WHERE …)`) get the same @@ -470,8 +470,8 @@ async fn sql_filtered_aggregates_dp_equals_exhaustive() { roots.push((index, QueryRoot::Operator(root))); } let inventory = stage1_logical_candidates(roots, &Default::default(), &[]).expect("Stage 1"); - // (pass-through, Count acc, CMS, CountSketch, UnivMon, HydraCms) × (pass-through, KLL, DDSketch). - assert_dp_matches_exhaustive(&inventory, &workload, 18); + // (pass-through, Count acc, HydraCms) × (pass-through, KLL, DDSketch). + assert_dp_matches_exhaustive(&inventory, &workload, 9); // Every filtered alternative builds through Stages 1 and 2. let exhaustive = select_exhaustive( &inventory,