Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
25 changes: 5 additions & 20 deletions crates/integration-tests/tests/filtered_aggregates.rs
Original file line number Diff line number Diff line change
Expand Up @@ -411,8 +411,8 @@ async fn alternatives(sql: &str, target: AccuracyTarget) -> Vec<(String, Rc<Oper
/// Pass 1 offers a filtered single-measure aggregate the same alternatives
/// as the unfiltered one, each with `SummaryAgg.filter` set. All compose;
/// an alternative binds in the executor exactly when its unfiltered
/// counterpart does. The exact `Count` accumulator and HydraCms execute to
/// the exact plan's counts, `b` included.
/// counterpart does. Every alternative executes to the exact plan's
/// counts, `b` included.
#[tokio::test]
async fn pass1_offers_filtered_count_alternatives() {
// ε = 0.1 keeps the Hydra grid inside the default memory limit.
Expand Down Expand Up @@ -446,31 +446,16 @@ async fn pass1_offers_filtered_count_alternatives() {
})
.collect();
let labels: Vec<_> = 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)],
vec![s("c"), Value::Int64(1)],
]);
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}");
}
}

Expand Down
203 changes: 155 additions & 48 deletions crates/integration-tests/tests/pass1_sql_coverage.rs
Original file line number Diff line number Diff line change
@@ -1,14 +1,14 @@
//! 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;

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};
Expand Down Expand Up @@ -44,55 +44,146 @@ fn rows() -> Vec<Vec<Value>> {
.collect()
}

async fn inventory(sql: &str) -> LocalLogicalCandidates<usize> {
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<OperatorNode>) -> 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<Vec<Value>>); 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<T>(value: T) -> Evidence<T> {
Expand All @@ -103,18 +194,34 @@ fn declared<T>(value: T) -> Evidence<T> {
}
}

/// #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)));
}
Expand Down
77 changes: 76 additions & 1 deletion crates/logical-optimizer/src/pass1/logical_candidates.rs
Original file line number Diff line number Diff line change
Expand Up @@ -142,7 +142,17 @@ pub fn enumerate_local_logical_candidates<Id>(
}
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()],
Expand Down Expand Up @@ -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<OperatorNode>| {
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::<Vec<_>>()
};
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() {
Expand Down
Loading