diff --git a/crates/devtools/src/bin/stage_pipeline.rs b/crates/devtools/src/bin/stage_pipeline.rs index 61e596d1..a125a4dd 100644 --- a/crates/devtools/src/bin/stage_pipeline.rs +++ b/crates/devtools/src/bin/stage_pipeline.rs @@ -44,7 +44,7 @@ use asap_logical_optimizer::Realization; use asap_plan_selection::PlanningModels; use asap_plan_selection::{plan_stages, Selection, Sharing, MAX_ENUMERATED_CANDIDATES}; use asap_types::ir::flat::{flatten, FlatDag}; -use asap_types::ir::schema::SketchAlgorithm; +use asap_types::ir::schema::{GroupingStrategy, SketchAlgorithm}; use asap_types::ir::schema_support::with_promql_series_identity; use asap_types::ir::{OperatorNode, QueryRoot}; use asap_types::types::AccuracyTarget; @@ -287,6 +287,12 @@ fn label(inventory: &LocalLogicalCandidates, owners: &[usize], choice: &[ SketchAlgorithm::CountSketchWithHeap => "CountSketch+heap".to_string(), other => format!("{other:?}"), }; + let name = match target.groupings[index] { + GroupingStrategy::PerSubpopulationInstance => name, + GroupingStrategy::SharedMultiSubpopulation { ref kind, .. } => { + format!("{kind:?}") + } + }; // The sketch reads the inner aggregate's input and replaces it. sketches.push(match target.absorbs[index] { Some(_) => format!("whole-expression {name}"), diff --git a/crates/integration-tests/tests/pass1_sql_coverage.rs b/crates/integration-tests/tests/pass1_sql_coverage.rs index c692a7ba..2d52ef9a 100644 --- a/crates/integration-tests/tests/pass1_sql_coverage.rs +++ b/crates/integration-tests/tests/pass1_sql_coverage.rs @@ -145,3 +145,137 @@ async fn example2_design_candidates_all_build() { assert!(enumeration.candidates.len() > 1); assert!(unbuilt.is_empty(), "{unbuilt:#?}"); } + +/// The HydraCms `SummaryAgg`s in `root` (#600's planner contract). +fn hydra_builds(root: &Rc) -> Vec> { + use asap_types::ir::schema::{GroupingStrategy, HydraKind}; + OperatorNode::reachable(root) + .into_iter() + .filter(|n| { + matches!( + &n.operator, + Operator::ASAP(ASAPOp::SummaryAgg { + grouping: GroupingStrategy::SharedMultiSubpopulation { + kind: HydraKind::HydraCms, + .. + }, + .. + }) + ) + }) + .collect() +} + +/// A grouped approximate count gets a HydraCms alternative (#580 W7): one +/// shared Count-Min grid for every `src_ip`, built from a unit weight and a +/// non-null item column. Stage 3 prices it (its grouping-aware guarantee +/// meets the target), and it compiles and executes. +#[tokio::test] +async fn grouped_count_offers_a_priced_executable_hydra_plan() { + use asap_types::ir::schema::{ + FieldDataType, GroupingStrategy, HydraParams, SketchParams, SummaryInputExpr, WeightDomain, + }; + let sql = "SELECT src_ip, COUNT(*) AS c FROM flows GROUP BY src_ip"; + // The grid holds shared_rows × shared_columns Count-Min cells, each as + // large as one per-group sketch: ε = 0.01 exceeds the default memory + // limit, so this uses ε = 0.1. + let target = AccuracyTarget::EpsilonDelta { + epsilon: 0.1, + delta: 0.01, + }; + let root = lower_sql(sql, &catalog(), target.clone()).await.unwrap(); + let demand = [RootDemand { + accuracy: Some(target), + recurrence: QueryRecurrence::OneTime { + invocations: 1, + execute_at: None, + }, + predictability: Predictability::default(), + latency_ms: None, + }]; + 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( + vec![(0, QueryRoot::Operator(root))], + &demand, + &data, + PlanningModels::builtin(), + 4096, + ) + .unwrap(); + let enumeration = run.enumeration.unwrap(); + let hydra: Vec<_> = enumeration + .candidates + .iter() + .filter_map(|c| { + let QueryRoot::Operator(root) = &c.logical.as_ref()?[0].1 else { + return None; + }; + (!hydra_builds(root).is_empty()).then(|| (c, root.clone())) + }) + .collect(); + assert!(!hydra.is_empty(), "a Hydra candidate is generated"); + for (candidate, root) in hydra { + // The all-query-time physical candidate comes first (#604). + let physical = candidate.physical.first().expect("Stage 2 builds it"); + assert!( + enumeration.selection.costs.contains_key(&physical.id), + "{} is priced: {:?}", + physical.id, + enumeration.selection.rejected + ); + for build in hydra_builds(&root) { + let Operator::ASAP(ASAPOp::SummaryAgg { + family: FieldDataType::Sketch(kind, family_grouping), + input, + grouping, + filter: None, + .. + }) = &build.operator + else { + panic!("HydraCms build") + }; + assert_eq!(family_grouping, grouping); + let GroupingStrategy::SharedMultiSubpopulation { + params: HydraParams::HydraCms { width, depth, .. }, + .. + } = grouping + else { + panic!("HydraCms params") + }; + assert_eq!( + kind.params(), + &SketchParams::Cms { + width: *width, + depth: *depth + } + ); + assert!(matches!(input.item, Some(SummaryInputExpr::Column(_)))); + assert_eq!(input.weight, SummaryInputExpr::Constant(1.0)); + assert!(matches!( + input.weight_domain, + WeightDomain::NonNegative { .. } + )); + } + 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!("unexpected row {other:?}"), + }) + .collect(); + // Few groups in a wide grid: no collisions, so the estimate is exact. + assert_eq!( + counts, + [("a".into(), 3), ("b".into(), 2), ("c".into(), 1)], + "{}", + physical.id + ); + } +} diff --git a/crates/logical-optimizer/src/accuracy/estimators/cms.rs b/crates/logical-optimizer/src/accuracy/estimators/cms.rs index 131dd1db..7a0c4347 100644 --- a/crates/logical-optimizer/src/accuracy/estimators/cms.rs +++ b/crates/logical-optimizer/src/accuracy/estimators/cms.rs @@ -64,6 +64,50 @@ mod tests { )); } + /// A shared HydraCms grid adds its collision term to the inner sketch's + /// error: sized like one per-group sketch for ε it misses ε, and sized + /// for ε/2 and δ/2 (Pass 1's split) it meets it. + #[test] + fn hydra_guarantee_adds_the_shared_grid_term() { + use crate::pass1::replacement::default_size_params; + use asap_types::ir::schema::{ + default_hydra_params, GroupingStrategy, HydraKind, SketchKind, + }; + let count = AggIntent::Count { + accuracy: AccuracyTarget::EpsilonDelta { + epsilon: 0.01, + delta: 0.01, + }, + }; + let target = AccuracyTarget::EpsilonDelta { + epsilon: 0.01, + delta: 0.01, + }; + let group_count = SketchStatistic::PointCount { + key: asap_types::ir::scalar::ColumnRef::SampleValue, + value: None, + }; + let guarantee = |epsilon: f64, delta: f64, hydra: bool| { + let params = default_size_params(SketchAlgorithm::Cms, &count, epsilon, delta); + let grouping = match hydra { + true => GroupingStrategy::SharedMultiSubpopulation { + kind: HydraKind::HydraCms, + params: default_hydra_params(HydraKind::HydraCms, ¶ms).unwrap(), + }, + false => GroupingStrategy::default(), + }; + DefaultAccuracyModel + .local_guarantee( + &FieldDataType::Sketch(SketchKind::new(SketchAlgorithm::Cms, params), grouping), + &group_count, + ) + .unwrap() + }; + assert!(DefaultAccuracyModel.satisfies(&guarantee(0.01, 0.01, false), &target)); + assert!(!DefaultAccuracyModel.satisfies(&guarantee(0.01, 0.01, true), &target)); + assert!(DefaultAccuracyModel.satisfies(&guarantee(0.005, 0.005, true), &target)); + } + #[test] fn heap_evaluation_retains_frequency_metric() { use asap_types::ir::schema::{GroupingStrategy, SketchKind}; diff --git a/crates/logical-optimizer/src/accuracy/estimators/mod.rs b/crates/logical-optimizer/src/accuracy/estimators/mod.rs index 6b094e30..5926e4b0 100644 --- a/crates/logical-optimizer/src/accuracy/estimators/mod.rs +++ b/crates/logical-optimizer/src/accuracy/estimators/mod.rs @@ -1,7 +1,7 @@ //! Dispatch committed estimator parameters to their accuracy models. use super::*; use asap_types::ir::operator::AggIntent; -use asap_types::ir::schema::GroupingStrategy; +use asap_types::ir::schema::{GroupingStrategy, HydraKind, HydraParams}; pub mod cardinality; pub mod cms; @@ -64,7 +64,37 @@ pub(super) fn local_guarantee( FieldDataType::ExactAggregate(kind, _) => { Some(ResultGuarantee::exact(format!("ExactAggregate({kind:?})"))) } - FieldDataType::Sketch(kind, _) => sketch_guarantee(kind.algorithm(), kind.params(), query), + FieldDataType::Sketch(kind, GroupingStrategy::PerSubpopulationInstance) => { + sketch_guarantee(kind.algorithm(), kind.params(), query) + } + FieldDataType::Sketch( + kind, + GroupingStrategy::SharedMultiSubpopulation { + kind: HydraKind::HydraCms, + params: + HydraParams::HydraCms { + shared_rows, + shared_columns, + .. + }, + }, + ) => { + // Groups that share a grid cell add their weight: at most + // e·N/shared_columns per row, and the minimum over rows exceeds + // it with probability e^-shared_rows. N is the whole input's + // weight, not one group's, so this bounds error relative to it. + let inner = sketch_guarantee(kind.algorithm(), kind.params(), query)?; + let stats = crate::accuracy::PropagationStats { + hydra_shared_grid_collision_bound: Some( + std::f64::consts::E / f64::from(*shared_columns), + ), + hydra_shared_grid_failure_probability: Some((-f64::from(*shared_rows)).exp()), + ..Default::default() + }; + Some(crate::pass1::grouping::hydra_guarantee(&inner, &stats)) + } + // No accuracy model for the other shared groupings. + FieldDataType::Sketch(..) => None, // No error model is registered for these families. FieldDataType::Sample(..) | FieldDataType::Wavelet(..) | FieldDataType::StatModel(..) => { None diff --git a/crates/logical-optimizer/src/pass1/grouping.rs b/crates/logical-optimizer/src/pass1/grouping.rs index a1c4ef66..3b5b1ca2 100644 --- a/crates/logical-optimizer/src/pass1/grouping.rs +++ b/crates/logical-optimizer/src/pass1/grouping.rs @@ -428,7 +428,10 @@ fn with_grouping( /// grid. The paper's collision term depends on deployment/data statistics; /// keeping those leaves symbolic makes the formula explicit while ensuring /// target satisfaction fails closed until a caller supplies them. -fn hydra_guarantee(inner: &ResultGuarantee, stats: &PropagationStats) -> ResultGuarantee { +pub(crate) fn hydra_guarantee( + inner: &ResultGuarantee, + stats: &PropagationStats, +) -> ResultGuarantee { let mut provenance = inner.provenance.clone(); provenance.extend(stats.evidence_provenance.clone()); provenance.push(GuaranteeSource::ChildGuarantee { diff --git a/crates/logical-optimizer/src/pass1/logical_candidates.rs b/crates/logical-optimizer/src/pass1/logical_candidates.rs index 778187b0..367c240c 100644 --- a/crates/logical-optimizer/src/pass1/logical_candidates.rs +++ b/crates/logical-optimizer/src/pass1/logical_candidates.rs @@ -11,9 +11,9 @@ use asap_types::ir::operator::{AggIntent, Reduction, Source}; use asap_types::ir::scalar::ColumnRef; use asap_types::ir::schema::Schema; use asap_types::ir::schema::{ - EntityIdentity, ExactKind, ExactParams, FieldDataType, GroupingStrategy, - NonNegativeWeightProof, SketchAlgorithm, SketchKind, SketchStatistic, SummaryInputExpr, - SummaryUpdate, WeightDomain, + default_hydra_params, DataType, EntityIdentity, ExactKind, ExactParams, FieldDataType, + GroupingStrategy, HydraKind, NonNegativeWeightProof, SketchAlgorithm, SketchKind, SketchParams, + SketchStatistic, SummaryInputExpr, SummaryUpdate, WeightDomain, }; use asap_types::ir::{ASAPOp, NonASAPOp, Operator, OperatorNode, QueryRoot, SchemaDerivationError}; use asap_types::types::AccuracyTarget; @@ -41,6 +41,9 @@ pub struct LocalLogicalTarget { /// /// [`add_window_forms`]: crate::pass2::window_composition::add_window_forms pub windows: Vec, + /// Per alternative, how a sketch's state serves the groups: one + /// instance per group, or one shared Hydra grid ([`add_hydra_alternatives`]). + pub groupings: Vec, /// Whether the target's input values are [`counter_samples`]. Absorbing /// alternatives read the input's own input, which then is too. pub counter_input: bool, @@ -140,16 +143,19 @@ 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)?; - targets.push(LocalLogicalTarget { + let mut target = LocalLogicalTarget { absorbs: vec![None; alternatives.len()], windows: vec![WindowForm::Whole; alternatives.len()], + groupings: vec![GroupingStrategy::default(); alternatives.len()], alternatives, counter_input: node .children() .iter() .all(|child| counter_samples(child, metric_types)), target: node, - }); + }; + add_hydra_alternatives(&mut target)?; + targets.push(target); } } } @@ -201,10 +207,98 @@ fn add_whole_expression_alternatives( target.alternatives.push(heap); target.absorbs.push(Some(inner)); target.windows.push(WindowForm::Whole); + target.groupings.push(GroupingStrategy::default()); } } } +/// Hydra (#580 W7, count first): a grouped approximate count may keep one +/// shared Count-Min grid for all its groups instead of one sketch per group. +/// The grid's collision term adds to the inner sketch's error, so each half +/// of the budget sizes one: the inner sketch and the grid get ε/2 and δ/2. +/// Offered when the target has groups (`by`, not `without`) and an item +/// column the kernel can hash ([`hydra_update`]). +pub fn add_hydra_alternatives( + target: &mut LocalLogicalTarget, +) -> Result<(), LogicalCandidateError> { + let Some(NonASAPOp::Aggregate { + child, + reduction, + measures, + .. + }) = target.target.non_asap() + else { + return Ok(()); + }; + let [intent @ AggIntent::Count { accuracy }] = measures.as_slice() else { + return Ok(()); + }; + if *accuracy == AccuracyTarget::Exact + || !crate::pass1::grouping::has_subpopulations(reduction) + || reduction.group_keys().is_some_and(|keys| keys.is_without()) + || hydra_update(reduction, &child.schema).is_none() + { + return Ok(()); + } + let (epsilon, delta) = accuracy_budget(accuracy); + let SketchParams::Cms { width, depth } = + default_size_params(SketchAlgorithm::Cms, intent, epsilon / 2.0, delta / 2.0) + else { + return Err(LogicalCandidateError::Unsupported("Count-Min sizing")); + }; + let kind = SketchKind::new(SketchAlgorithm::Cms, SketchParams::Cms { width, depth }); + let params = default_hydra_params(HydraKind::HydraCms, kind.params()) + .ok_or(LogicalCandidateError::Unsupported("HydraCms parameters"))?; + target.alternatives.push(Realization::Sketch(kind)); + target.absorbs.push(None); + target.windows.push(WindowForm::Whole); + target + .groupings + .push(GroupingStrategy::SharedMultiSubpopulation { + kind: HydraKind::HydraCms, + params, + }); + Ok(()) +} + +/// A HydraCms update (#600's contract): a unit weight per row, hashed by +/// a [`count_item`]. +fn hydra_update(reduction: &Reduction, child: &Schema) -> Option { + Some(SummaryUpdate { + item: Some(SummaryInputExpr::Column( + crate::pass1::replacement::column_ref(count_item(reduction, child)?), + )), + weight: SummaryInputExpr::Constant(1.0), + weight_domain: WeightDomain::NonNegative { + proof: NonNegativeWeightProof::UnitCount, + }, + }) +} + +/// The item a count sketch hashes per row: a group's count ignores its +/// value, so any non-null column the kernels hash (Utf8, Int64 or Bool) +/// serves, a grouping column first. PromQL labels are nullable; the series +/// identity is not. +fn count_item<'a>( + reduction: &Reduction, + child: &'a Schema, +) -> Option<&'a asap_types::ir::schema::Field> { + let keys = reduction + .group_keys() + .filter(|keys| !keys.is_without()) + .into_iter() + .flat_map(|keys| keys.iter().copied()); + keys.chain(0..child.fields.len()) + .filter_map(|index| child.fields.get(index)) + .find(|field| { + !field.nullable + && matches!( + field.dtype, + FieldDataType::Plain(DataType::Utf8 | DataType::Int64 | DataType::Bool) + ) + }) +} + /// The input and update of a whole-expression top-k over `target`'s inner /// aggregate, by the legacy keyed-additive rule, or `None` when it does not /// apply. Rows that carry the full series identity rank it as a column, as @@ -434,6 +528,7 @@ pub fn compose_logical_candidate( absorbs, counter_input: target.counter_input, window: target.windows[index], + grouping: &target.groupings[index], }, ) }) @@ -473,6 +568,7 @@ struct Chosen<'a> { /// The target's [`LocalLogicalTarget::counter_input`]. counter_input: bool, window: WindowForm, + grouping: &'a GroupingStrategy, } fn rewrite( @@ -511,6 +607,7 @@ fn realize( absorbs, counter_input, window, + grouping, } = *chosen; let Some(NonASAPOp::Aggregate { child, @@ -552,13 +649,17 @@ fn realize( None, ), Realization::Sketch(kind) => ( - FieldDataType::Sketch(kind.clone(), GroupingStrategy::default()), + FieldDataType::Sketch(kind.clone(), grouping.clone()), Some(statistic(intent)?), ), _ => return Err(LogicalCandidateError::Unsupported("summary family")), }; let mut input = match whole { Some((_, update)) => update, + None if *grouping != GroupingStrategy::PerSubpopulationInstance => { + hydra_update(reduction, &child.schema) + .ok_or(LogicalCandidateError::Unsupported("Hydra item column"))? + } None => summary_update(intent, &family, reduction, &child.schema)?, }; if counter_input @@ -577,7 +678,7 @@ fn realize( family: family.clone(), input: input.clone(), reduction: reduction.clone(), - grouping: GroupingStrategy::default(), + grouping: grouping.clone(), filter: None, }, ))?) @@ -721,27 +822,16 @@ fn summary_update( } /// SQL `COUNT(*)`: rows have no sample value, and every row counts, so each -/// adds a unit weight. A sketch hashes an item per row; its bare count -/// ignores the item's value, so any non-null column serves, a grouping -/// column first. +/// adds a unit weight. A sketch hashes a [`count_item`] per row. fn sql_row_count_update( sketch: bool, reduction: &Reduction, child: &Schema, ) -> Result { let item = if sketch { - let keys = reduction - .group_keys() - .filter(|keys| !keys.is_without()) - .into_iter() - .flat_map(|keys| keys.iter().copied()); - let column = keys - .chain(0..child.fields.len()) - .filter_map(|index| child.fields.get(index)) - .find(|field| field.is_plain() && !field.nullable) - .ok_or(LogicalCandidateError::Unsupported( - "a COUNT(*) sketch needs a non-null item column", - ))?; + let column = count_item(reduction, child).ok_or(LogicalCandidateError::Unsupported( + "a COUNT(*) sketch needs a non-null item column", + ))?; Some(SummaryInputExpr::Column( crate::pass1::replacement::column_ref(column), )) @@ -985,6 +1075,54 @@ mod tests { } } + /// Only an approximate count with `by` groups gets a HydraCms + /// alternative, sized for half the budget, with a matching grouping. + #[test] + fn grouped_approximate_count_offers_hydra() { + let approximate = AccuracyTarget::EpsilonDelta { + epsilon: 0.01, + delta: 0.01, + }; + let hydra = |query: &str, accuracy: AccuracyTarget| { + let root = lower_promql(query, accuracy); + let root = asap_types::ir::schema_support::with_promql_series_identity(&root).unwrap(); + let inventory = enumerate_local_logical_candidates( + vec![(0, QueryRoot::Operator(root))], + &BTreeMap::new(), + ) + .unwrap(); + let target = &inventory.targets[0]; + assert_eq!(target.groupings.len(), target.alternatives.len()); + target + .alternatives + .iter() + .zip(&target.groupings) + .filter(|(_, g)| **g != GroupingStrategy::default()) + .map(|(a, g)| (a.clone(), g.clone())) + .collect::>() + }; + let offered = hydra("count by (job) (m)", approximate.clone()); + let [(Realization::Sketch(kind), GroupingStrategy::SharedMultiSubpopulation { params, .. })] = + offered.as_slice() + else { + panic!("one HydraCms alternative: {offered:?}") + }; + assert_eq!(kind.algorithm(), &SketchAlgorithm::Cms); + // ⌈e/(ε/2)⌉ columns, ⌈ln(2/δ)⌉ rows, for the inner sketch and the grid. + assert_eq!( + *params, + asap_types::ir::schema::HydraParams::HydraCms { + width: 544, + depth: 6, + shared_rows: 6, + shared_columns: 544 + } + ); + assert!(hydra("count (m)", approximate.clone()).is_empty()); + assert!(hydra("count without (job) (m)", approximate).is_empty()); + assert!(hydra("count by (job) (m)", AccuracyTarget::Exact).is_empty()); + } + /// Approximate requests must retain the exact execution alternative too. #[test] fn approximate_count_keeps_exact_and_universal_choices() { diff --git a/crates/logical-optimizer/src/pass2/summary_capability.rs b/crates/logical-optimizer/src/pass2/summary_capability.rs index 8a52b796..0aa1d923 100644 --- a/crates/logical-optimizer/src/pass2/summary_capability.rs +++ b/crates/logical-optimizer/src/pass2/summary_capability.rs @@ -181,6 +181,8 @@ pub fn share_summary_capability( debug_assert!(out.targets[t].absorbs.iter().all(Option::is_none)); out.targets[t].absorbs = vec![None; alternatives.len()]; out.targets[t].windows = vec![Default::default(); alternatives.len()]; + // Count targets, the only ones with Hydra, are not keyed. + out.targets[t].groupings = vec![Default::default(); alternatives.len()]; out.targets[t].alternatives = alternatives; resized = true; } diff --git a/crates/logical-optimizer/src/pass2/window_composition.rs b/crates/logical-optimizer/src/pass2/window_composition.rs index ff62b1c2..be234fe3 100644 --- a/crates/logical-optimizer/src/pass2/window_composition.rs +++ b/crates/logical-optimizer/src/pass2/window_composition.rs @@ -206,32 +206,32 @@ pub fn add_window_forms(inventory: &mut LocalLogicalCandidates, demand: if !panes_tile_window(&window, pane_ms) { continue; } - let tumbling: Vec = target + let tumbling: Vec<(Realization, GroupingStrategy)> = target .alternatives .iter() .zip(&target.absorbs) - .filter(|(alternative, absorbs)| { - absorbs.is_none() && family(alternative).is_some_and(|f| f.family_merges()) + .zip(&target.groupings) + .filter(|((alternative, absorbs), grouping)| { + absorbs.is_none() + && family(alternative, grouping).is_some_and(|f| f.family_merges()) }) - .map(|(alternative, _)| alternative.clone()) + .map(|((alternative, _), grouping)| (alternative.clone(), grouping.clone())) .collect(); - for alternative in tumbling { + for (alternative, grouping) in tumbling { target.alternatives.push(alternative); target.absorbs.push(None); target.windows.push(WindowForm::Tumbling { pane_ms }); + target.groupings.push(grouping); } } } -fn family(realization: &Realization) -> Option { +fn family(realization: &Realization, grouping: &GroupingStrategy) -> Option { match realization { Realization::ExactAggregate { kind, params } => { Some(FieldDataType::ExactAggregate(kind.clone(), params.clone())) } - Realization::Sketch(kind) => Some(FieldDataType::Sketch( - kind.clone(), - GroupingStrategy::default(), - )), + Realization::Sketch(kind) => Some(FieldDataType::Sketch(kind.clone(), grouping.clone())), _ => None, } } @@ -508,7 +508,7 @@ mod tests { Ok(OperatorNode::new_shared(Operator::ASAP( ASAPOp::SummaryAgg { child: input, - family: family(&target.alternatives[1]).unwrap(), + family: family(&target.alternatives[1], &GroupingStrategy::default()).unwrap(), input: asap_types::ir::schema::SummaryUpdate::column( asap_types::ir::scalar::ColumnRef::SampleValue, ), diff --git a/crates/planner/tests/stage_pipeline_selection.rs b/crates/planner/tests/stage_pipeline_selection.rs index c8417193..38b5545e 100644 --- a/crates/planner/tests/stage_pipeline_selection.rs +++ b/crates/planner/tests/stage_pipeline_selection.rs @@ -390,6 +390,50 @@ async fn sql_dp_equals_exhaustive() { assert_dp_matches_exhaustive(&inventory, &workload, 15); } +/// A SQL grouped count, which also has a HydraCms alternative (#580 W7), +/// and a percentile select the exhaustive minimum. +#[tokio::test] +async fn sql_hydra_count_dp_equals_exhaustive() { + let accuracy = AccuracyTarget::Epsilon(0.1); + let queries = [ + "SELECT l_orderkey, COUNT(*) FROM lineitem GROUP BY l_orderkey", + "SELECT approx_percentile_cont(l_extendedprice, 0.99) FROM lineitem", + ]; + let workload = PlanningWorkload { + query_workload: QueryWorkload { + language: QueryLanguage::SQL(SqlDialect::DataFusionSQL), + query_batch: Some(queries.iter().map(|q| batch(q, accuracy.clone())).collect()), + repeating_queries: None, + }, + data_workload: Some(DataWorkload { + arrival: DataArrival::AtRest, + ..Default::default() + }), + }; + let catalog = SqlCatalog::new().with_table( + "lineitem", + Schema::new(vec![ + Field::plain("l_orderkey", DataType::Int64, false), + Field::plain("l_extendedprice", DataType::Float64, false), + ]), + ); + let mut roots = Vec::new(); + for (index, query) in queries.iter().enumerate() { + let root = lower_sql_dialect(query, &catalog, SqlDialect::DataFusionSQL, accuracy.clone()) + .await + .expect("lowers"); + roots.push((index, QueryRoot::Operator(root))); + } + let inventory = stage1_logical_candidates(roots, &Default::default(), &[]).expect("Stage 1"); + assert!(inventory[0] + .inventory + .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); +} + /// PromQL queries, each with its own ε (δ = 0.001). fn promql_with(queries: &[(&str, f64)]) -> PlanningWorkload { let mut workload = promql(&[], 1_000);