From 84da65d6a6dcb1fe22b1683231cfccd98fafe622 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Mon, 5 Oct 2026 00:00:55 +0000 Subject: [PATCH 1/3] fix(pass1): offer no per-group sketch for SQL COUNT(*) A SQL COUNT(*) group's rows all hash its grouping column, so a per-group Count-Min, Count Sketch or UnivMon is a counter with extra memory, and the executor cannot bind them. Stage 3 priced the UnivMon one as valid, so it could select a plan that cannot run. Pass 1 now drops the per-group sketch alternatives of a count over SQL rows; HydraCms stays. PromQL counts are unchanged. Tests: every Pass 1 choice for grouped, ungrouped and ranked COUNT(*) executes to the exact counts, and every candidate Stage 3 prices for those queries and for Example 2 binds in the executor. Co-Authored-By: Claude Opus 5.5 --- .../tests/filtered_aggregates.rs | 25 +-- .../tests/pass1_sql_coverage.rs | 199 +++++++++++++----- .../src/pass1/logical_candidates.rs | 77 ++++++- .../planner/tests/stage_pipeline_selection.rs | 8 +- 4 files changed, 236 insertions(+), 73 deletions(-) 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..5b5b54a67 100644 --- a/crates/integration-tests/tests/pass1_sql_coverage.rs +++ b/crates/integration-tests/tests/pass1_sql_coverage.rs @@ -8,7 +8,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 +43,143 @@ 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() { + for (sql, expected) in count_star_queries() { + 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::export::PhysicalASAPDAG) -> Result<(), String> { + use asap_executor::physical_planner::bind_with_data_sources; + use asap_executor::sources::{DataSources, MemorySource}; + use asap_types::ir::export::{NonASAPOpKind, PhysicalASAPOperatorPayload}; + use std::sync::Arc; + let mut sources = DataSources::default(); + let mut registered = vec![]; + for node in &dag.nodes { + if let PhysicalASAPOperatorPayload::Relational { + operator: NonASAPOpKind::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| u64::from(root.0)).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 +190,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, From 3f8210813e11be931be5a044da4792921f24b277 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Mon, 5 Oct 2026 05:58:04 +0000 Subject: [PATCH 2/3] fix: integrate with #614 (DataFusion 54): execute only the unranked COUNT(*) queries DataFusion 54 plans ORDER BY COUNT(*) against the aggregate itself, so canonicalize promotes the ranked query to TopK over the Count (as it already did for ORDER BY c). The executor has no native TopK, so the compose-and-execute test runs every Pass 1 choice only for the grouped and ungrouped queries; the ranked one's priced candidates still bind. Co-Authored-By: Claude Opus 5.5 --- crates/integration-tests/tests/pass1_sql_coverage.rs | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/crates/integration-tests/tests/pass1_sql_coverage.rs b/crates/integration-tests/tests/pass1_sql_coverage.rs index 5b5b54a67..13398c416 100644 --- a/crates/integration-tests/tests/pass1_sql_coverage.rs +++ b/crates/integration-tests/tests/pass1_sql_coverage.rs @@ -75,7 +75,12 @@ fn count_star_queries() -> [(&'static str, Vec>); 3] { /// counts: HydraCms too, whose few groups in a wide grid never collide. #[tokio::test] async fn sql_count_star_candidates_compose_and_execute() { - for (sql, expected) in count_star_queries() { + // 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( From efbc3b10244f1d0855d79b3a87c6c507c61e326b Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Thu, 8 Oct 2026 15:04:06 +0000 Subject: [PATCH 3/3] test(integration): bind SQL COUNT(*) plans from Operator payloads Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01W7qG9aFyPij5uWsyAJCxDW --- crates/integration-tests/tests/pass1_sql_coverage.rs | 11 +++++------ 1 file changed, 5 insertions(+), 6 deletions(-) diff --git a/crates/integration-tests/tests/pass1_sql_coverage.rs b/crates/integration-tests/tests/pass1_sql_coverage.rs index 13398c416..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; @@ -108,17 +109,15 @@ async fn sql_count_star_candidates_compose_and_execute() { /// Binds every root of `dag` in the executor, each scan reading an empty /// in-memory source. -fn binds(dag: &asap_types::ir::export::PhysicalASAPDAG) -> Result<(), String> { +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::export::{NonASAPOpKind, PhysicalASAPOperatorPayload}; + 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::Relational { - operator: NonASAPOpKind::Scan { source, .. }, - } = &node.payload + if let PhysicalASAPOperatorPayload::NonASAP(NonASAPOp::Scan { source, .. }) = &node.payload { if !registered.contains(source) { let schema = Arc::new(node.output_schema.clone()); @@ -132,7 +131,7 @@ fn binds(dag: &asap_types::ir::export::PhysicalASAPDAG) -> Result<(), String> { } } } - let roots: Vec<_> = dag.roots.iter().map(|root| u64::from(root.0)).collect(); + 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())