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
160 changes: 63 additions & 97 deletions crates/devtools/src/bin/sketch_coverage.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,36 +2,24 @@
//
// Lowers every query in every corpus we have (mirrors `variant_coverage`'s
// corpus list exactly, so the two reports are directly comparable) with an
// *approximate* `AccuracyTarget`, runs `asap_logical_optimizer::explain_replacements`
// over each corpus as one workload, and reports the MVP demo's query-coverage
// metric: of the queries that lowered successfully, what fraction got
//
// - a `SketchApproximation` candidate (a genuine sketch alternative was
// found for at least one aggregate in the query — the KLL-vs-DDSketch
// kind of degree of freedom), and/or
// - a `CommonSubexpressionReuse` candidate (the query shares a sub-DAG,
// inside itself or with another query in the same corpus, that a
// build-once-and-share candidate was found for).
// *approximate* `AccuracyTarget`, runs Stage 1's Pass 1
// (`enumerate_local_logical_candidates`) over each query, and reports which
// sketch alternatives Pass 1 offers for it, plus the fraction of lowered
// queries with at least one.
//
// `--epsilon <f64>` (default 0.01) sets the `AccuracyTarget` every query in
// every corpus lowers with. Without an approximate target,
// `ASAPStrategies` never has a genuine sketch alternative to
// report.
//
// This is a workload-level count (`explain_replacements` runs
// `search_workload` once per corpus, over every query in it together), not
// just a per-query re-run of the single-target path — so cross-query CSE
// reuse inside one corpus shows up here the same way it would in the
// dag-viewer's Union mode.
// every corpus lowers with. Without an approximate target, Pass 1 offers no
// sketch alternative.

use asap_devtools::lower_promql_with_data_ingestion_interval;
use asap_frontend_sql::{lower_sql_dialect, SqlCatalog};
use asap_logical_optimizer::{explain_replacements, ExplanationKind};
use asap_types::ir::schema::{DataType, Field, Schema};
use asap_types::ir::OperatorNode;
use asap_logical_optimizer::pass1::logical_candidates::enumerate_local_logical_candidates;
use asap_logical_optimizer::Realization;
use asap_types::ir::schema::{DataType, Field, GroupingStrategy, Schema};
use asap_types::ir::{OperatorNode, QueryRoot};
use asap_types::types::AccuracyTarget;
use asap_types::workload::SqlDialect;
use std::collections::BTreeSet;
use std::collections::{BTreeMap, BTreeSet};
use std::rc::Rc;

/// Line-based `#`/`--` comment stripping, then split on `;` — the shape every
Expand Down Expand Up @@ -123,9 +111,10 @@ struct CorpusCoverage {
name: &'static str,
lowered: usize,
failed: usize,
sketch_covered: usize,
cse_covered: usize,
either_covered: usize,
/// Lowered queries Pass 1 rejected.
rejected: usize,
/// Per lowered query Pass 1 accepted, the sketch alternatives it offers.
sketches: Vec<(String, BTreeSet<String>)>,
}

fn pct(n: usize, total: usize) -> String {
Expand All @@ -136,86 +125,69 @@ fn pct(n: usize, total: usize) -> String {
}
}

/// `explain_replacements`' `location` is a comma-joined list of breadcrumbs
/// (`collect_locations` in `explanation.rs`), one per path from a workload
/// root to the target — e.g. `root "q3"` or `root "q3" > lhs`. A location
/// "covers" `label` (a bare `root "qN"` breadcrumb) if `label` is exactly one
/// of those comma-separated entries or the prefix of one that goes deeper —
/// i.e. the explanation's target is reachable from that query's root at all.
fn covers(location: &str, label: &str) -> bool {
location
.split(", ")
.any(|loc| loc == label || loc.starts_with(&format!("{label} > ")))
}

fn root_label(id: &str) -> String {
format!("root {id:?}")
}

/// Run `explain_replacements` over one corpus's already-lowered roots as one
/// workload, then attribute each finding back to the query root(s) it's
/// reachable from.
/// Run Pass 1 over each of one corpus's already-lowered queries and collect
/// the sketch alternatives it offers for any of the query's targets. A Hydra
/// alternative is listed as `Hydra(<algorithm>)`.
fn analyze_corpus(
name: &'static str,
roots: Vec<(String, Rc<OperatorNode>)>,
failed: usize,
) -> CorpusCoverage {
let lowered = roots.len();
let labels: Vec<String> = roots.iter().map(|(id, _)| root_label(id)).collect();
let explanations = explain_replacements(roots);

let mut sketch_covered: BTreeSet<usize> = BTreeSet::new();
let mut cse_covered: BTreeSet<usize> = BTreeSet::new();
for explanation in &explanations {
for (i, label) in labels.iter().enumerate() {
if !covers(&explanation.location, label) {
continue;
}
match explanation.kind {
ExplanationKind::SketchApproximation => {
sketch_covered.insert(i);
}
ExplanationKind::CommonSubexpressionReuse => {
cse_covered.insert(i);
let mut rejected = 0;
let mut sketches = Vec::new();
for (id, root) in roots {
let Ok(inventory) = enumerate_local_logical_candidates(
vec![(id.clone(), QueryRoot::Operator(root))],
&BTreeMap::new(),
) else {
rejected += 1;
continue;
};
let offered = inventory
.targets
.iter()
.flat_map(|target| target.alternatives.iter().zip(&target.groupings))
.filter_map(|(alternative, grouping)| match alternative {
Realization::Sketch(kind) if *grouping == GroupingStrategy::default() => {
Some(format!("{:?}", kind.algorithm()))
}
// `#[non_exhaustive]`: a future kind just doesn't count
// toward either bucket here until this tool is taught about it.
_ => {}
}
}
Realization::Sketch(kind) => Some(format!("Hydra({:?})", kind.algorithm())),
_ => None,
})
.collect();
sketches.push((id, offered));
}
let either_covered = sketch_covered.union(&cse_covered).count();

CorpusCoverage {
name,
lowered,
failed,
sketch_covered: sketch_covered.len(),
cse_covered: cse_covered.len(),
either_covered,
rejected,
sketches,
}
}

fn covered(r: &CorpusCoverage) -> usize {
r.sketches.iter().filter(|(_, s)| !s.is_empty()).count()
}

fn report(r: &CorpusCoverage) {
println!("--- {} ---", r.name);
println!("lowered: {}, failed: {}", r.lowered, r.failed);
println!(
"sketch-approximable: {}/{} ({})",
r.sketch_covered,
r.lowered,
pct(r.sketch_covered, r.lowered)
);
println!(
"CSE-shareable: {}/{} ({})",
r.cse_covered,
r.lowered,
pct(r.cse_covered, r.lowered)
"lowered: {}, failed: {}, rejected by Pass 1: {}",
r.lowered, r.failed, r.rejected
);
for (id, offered) in &r.sketches {
if !offered.is_empty() {
let offered: Vec<_> = offered.iter().map(String::as_str).collect();
println!(" {id}: {}", offered.join(", "));
}
}
println!(
"either (coverage): {}/{} ({})",
r.either_covered,
"sketch-approximable: {}/{} ({})",
covered(r),
r.lowered,
pct(r.either_covered, r.lowered)
pct(covered(r), r.lowered)
);
println!();
}
Expand Down Expand Up @@ -332,22 +304,16 @@ async fn main() {

let total_lowered: usize = results.iter().map(|r| r.lowered).sum();
let total_failed: usize = results.iter().map(|r| r.failed).sum();
let total_sketch: usize = results.iter().map(|r| r.sketch_covered).sum();
let total_cse: usize = results.iter().map(|r| r.cse_covered).sum();
let total_either: usize = results.iter().map(|r| r.either_covered).sum();
let total_rejected: usize = results.iter().map(|r| r.rejected).sum();
let total_sketch: usize = results.iter().map(covered).sum();

println!("=== global (epsilon = {epsilon}) ===");
println!("total lowered: {total_lowered}, total failed: {total_failed}");
println!(
"sketch-approximable: {total_sketch}/{total_lowered} ({})",
pct(total_sketch, total_lowered)
"total lowered: {total_lowered}, total failed: {total_failed}, \
rejected by Pass 1: {total_rejected}"
);
println!(
"CSE-shareable: {total_cse}/{total_lowered} ({})",
pct(total_cse, total_lowered)
);
println!(
"either (coverage): {total_either}/{total_lowered} ({})",
pct(total_either, total_lowered)
"sketch-approximable: {total_sketch}/{total_lowered} ({})",
pct(total_sketch, total_lowered)
);
}
27 changes: 27 additions & 0 deletions crates/executor/tests/common/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -61,3 +61,30 @@ pub fn selected_dag(
};
root.clone()
}

/// Every whole-query candidate Stage 1's Pass 1 lists for `root` (up to
/// 4096), composed; choices that do not compose are skipped.
pub fn stage1_candidates(root: &Rc<OperatorNode>) -> Vec<Rc<OperatorNode>> {
use asap_logical_optimizer::pass1::logical_candidates::{
compose_logical_candidate, enumerate_choices, enumerate_local_logical_candidates,
};
use planner_types::ir::QueryRoot;
let inventory = enumerate_local_logical_candidates(
vec![(0, QueryRoot::Operator(Rc::clone(root)))],
&Default::default(),
)
.unwrap();
enumerate_choices(&inventory, 4096)
.iter()
.filter_map(|choice| {
match compose_logical_candidate(&inventory, choice)
.ok()?
.remove(0)
.1
{
QueryRoot::Operator(node) => Some(node),
QueryRoot::Scalar(_) => None,
}
})
.collect()
}
42 changes: 18 additions & 24 deletions crates/executor/tests/deployment_computation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,12 +7,11 @@ use asap_executor::{
runtime::{Limits, RunContext, Scope},
values::{Batch, Value},
};
use common::{compile_physical_asap_dag, selected_dag};
use common::{compile_physical_asap_dag, selected_dag, stage1_candidates};
use futures::{executor::block_on, StreamExt};
use planner_types::ir::physical_export::{PhysicalASAPDAG, PhysicalASAPOperatorPayload};
use planner_types::ir::schema::*;
use planner_types::ir::ASAPOp;
use planner_types::{types::AccuracyTarget, workload::*};
use planner_types::{ir::schema::*, types::AccuracyTarget, workload::*};
use std::{collections::BTreeMap, rc::Rc, sync::Arc};

fn lower(query: &str) -> Rc<planner_types::ir::OperatorNode> {
Expand Down Expand Up @@ -611,31 +610,26 @@ fn population_sums_and_averages_are_compensated() {
#[test]
fn stored_count_min_bare_count_compiles_to_a_evaluation() {
use asap_executor::summary_kernels::CountMinSketchAccumulator;
use asap_logical_optimizer::{Replacement, ReplacementStrategy, TargetSubDAG};
let root = lower_with("count(up)", AccuracyTarget::Epsilon(0.02));
let dag = asap_logical_optimizer::ASAPStrategies::default()
.replacements(&TargetSubDAG::new(&root))
let dag = stage1_candidates(&root)
.into_iter()
.find_map(|candidate| match candidate.replacement {
Replacement::SubDAG(node) => {
let dag = compile_physical_asap_dag(&node).ok()?;
let bare_count = dag.nodes.iter().any(|n| {
matches!(
&n.payload,
PhysicalASAPOperatorPayload::ASAP(ASAPOp::SummaryEstimate {
query: SketchStatistic::PointCount { value: None, .. },
..
})
)
});
let count_min = dag.nodes.iter().any(|n| {
matches!(&n.payload, PhysicalASAPOperatorPayload::ASAP(ASAPOp::SummaryAgg {
.find_map(|node| {
let dag = compile_physical_asap_dag(&node).ok()?;
let bare_count = dag.nodes.iter().any(|n| {
matches!(
&n.payload,
PhysicalASAPOperatorPayload::ASAP(ASAPOp::SummaryEstimate {
query: SketchStatistic::PointCount { value: None, .. },
..
})
)
});
let count_min = dag.nodes.iter().any(|n| {
matches!(&n.payload, PhysicalASAPOperatorPayload::ASAP(ASAPOp::SummaryAgg {
family: FieldDataType::Sketch(kind, _), ..
}) if kind.algorithm() == &SketchAlgorithm::Cms)
});
(bare_count && count_min).then_some(dag)
}
_ => None,
});
(bare_count && count_min).then_some(dag)
})
.expect("Planner lists a Count-Min candidate for count(up)");
let state = dag
Expand Down
Loading