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
21 changes: 15 additions & 6 deletions crates/devtools/src/bin/stage_pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -240,7 +240,9 @@ fn target_owners(inventory: &LocalLogicalCandidates<usize>) -> Vec<usize> {
.collect()
}

/// E.g. "Q1 exact · Q2 CMS+heap"; exact accumulators are listed in parentheses.
/// E.g. "Q1 exact · Q2 CMS+heap"; exact accumulators are listed in
/// parentheses. A sketch that absorbs the aggregate beneath it is
/// "whole-expression".
fn label(inventory: &LocalLogicalCandidates<usize>, owners: &[usize], choice: &[usize]) -> String {
(0..inventory.roots.len())
.map(|query| {
Expand All @@ -257,11 +259,18 @@ fn label(inventory: &LocalLogicalCandidates<usize>, owners: &[usize], choice: &[
Realization::ExactAggregate { kind, .. } => {
accumulators.push(format!("{kind:?} acc"))
}
Realization::Sketch(kind) => sketches.push(match kind.algorithm() {
SketchAlgorithm::CmsWithHeap => "CMS+heap".to_string(),
SketchAlgorithm::CountSketchWithHeap => "CountSketch+heap".to_string(),
other => format!("{other:?}"),
}),
Realization::Sketch(kind) => {
let name = match kind.algorithm() {
SketchAlgorithm::CmsWithHeap => "CMS+heap".to_string(),
SketchAlgorithm::CountSketchWithHeap => "CountSketch+heap".to_string(),
other => format!("{other:?}"),
};
// The sketch reads the inner aggregate's input and replaces it.
sketches.push(match target.absorbs[index] {
Some(_) => format!("whole-expression {name}"),
None => name,
});
}
other => sketches.push(format!("{other:?}")),
}
}
Expand Down
123 changes: 90 additions & 33 deletions crates/integration-tests/tests/planner_layering_example1.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,9 +9,10 @@
//! MVP scope: Stage 1 = Pass 1 + the identical-expression rule only (no
//! window-composition variants); Stage 2 = physical operator implementation
//! only (no materialization). Counts follow the planner's output (user
//! decision): 1 → 48 → 48 → 1, because Pass 1 also offers exact accumulators
//! (24 combinations) and Pass 2 adds a shared-input variant of each. The
//! doc's 1 → 54 → 156 → 1 needs window composition and materialization.
//! decision): 1 → 64 → 64 → 1, because Pass 1 also offers exact accumulators
//! and whole-expression top-k sketches (32 combinations) and Pass 2 adds a
//! shared-input variant of each. The doc's 1 → 54 → 156 → 1 needs window
//! composition and materialization.

use std::collections::{BTreeMap, BTreeSet, HashSet};

Expand Down Expand Up @@ -342,12 +343,15 @@ fn example1_workload() -> PlanningWorkload {

// ── DAG helpers ──────────────────────────────────────────────────────────

/// Q2's local option, read off its summary build node.
/// Q2's local option, read off its summary build node. A whole-expression
/// sketch reads the raw samples of the range, absorbing `sum_over_time`.
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
enum Q2Option {
Exact,
CountMinHeapPerJob,
CountSketchHeapPerJob,
WholeCountMinHeapPerJob,
WholeCountSketchHeapPerJob,
Hydra,
}

Expand Down Expand Up @@ -426,23 +430,36 @@ fn sketch_options(dag: &impl ExportedDag, nodes: &HashSet<NodeId>) -> BTreeSet<Q
Operator::ASAP(ASAPOp::SummaryAgg {
family: FieldDataType::Sketch(kind, grouping),
..
}) => Some(match grouping {
GroupingStrategy::SharedMultiSubpopulation {
kind: HydraKind::HydraCms,
..
} => Q2Option::Hydra,
GroupingStrategy::PerSubpopulationInstance
if *kind.algorithm() == SketchAlgorithm::CmsWithHeap =>
{
Q2Option::CountMinHeapPerJob
}
GroupingStrategy::PerSubpopulationInstance
if *kind.algorithm() == SketchAlgorithm::CountSketchWithHeap =>
{
Q2Option::CountSketchHeapPerJob
}
other => panic!("summary family outside Example 1: {kind:?} {other:?}"),
}),
}) => {
let whole = dag
.producers(id)
.iter()
.any(|&p| relational(payload(dag, p)).as_deref() == Some("time_range"));
let per_job = matches!(grouping, GroupingStrategy::PerSubpopulationInstance);
Some(match (grouping, kind.algorithm(), whole) {
(
GroupingStrategy::SharedMultiSubpopulation {
kind: HydraKind::HydraCms,
..
},
_,
_,
) => Q2Option::Hydra,
(_, SketchAlgorithm::CmsWithHeap, false) if per_job => {
Q2Option::CountMinHeapPerJob
}
(_, SketchAlgorithm::CmsWithHeap, true) if per_job => {
Q2Option::WholeCountMinHeapPerJob
}
(_, SketchAlgorithm::CountSketchWithHeap, false) if per_job => {
Q2Option::CountSketchHeapPerJob
}
(_, SketchAlgorithm::CountSketchWithHeap, true) if per_job => {
Q2Option::WholeCountSketchHeapPerJob
}
other => panic!("summary family outside Example 1: {other:?}"),
})
}
_ => None,
})
.collect()
Expand Down Expand Up @@ -524,9 +541,11 @@ fn choices(dag: &impl ExportedDag, query_roots: &[NodeId]) -> Choices {
)
}

/// The 24 Pass 1 combinations: Q1's rate and sum each raw or an exact
/// accumulator (4) × Q2's top-k exact, Count-Min + heap or CountSketch + heap
/// (3) × Q2's sum_over_time raw or an exact accumulator (2).
/// The 32 Pass 1 combinations: Q1's rate and sum each raw or an exact
/// accumulator (4) × (Q2's top-k exact, Count-Min + heap or CountSketch +
/// heap (3) × Q2's sum_over_time raw or an exact accumulator (2), or a
/// whole-expression Count-Min + heap or CountSketch + heap that absorbs
/// sum_over_time (2)).
fn expected_choices() -> BTreeSet<Choices> {
let q1 = [vec![], vec!["Rate"], vec!["Sum"], vec!["Rate", "Sum"]];
let q2 = [
Expand All @@ -542,6 +561,12 @@ fn expected_choices() -> BTreeSet<Choices> {
all.insert((owned(a), option, owned(&b)));
}
}
for option in [
Q2Option::WholeCountMinHeapPerJob,
Q2Option::WholeCountSketchHeapPerJob,
] {
all.insert((owned(a), option, vec![]));
}
}
all
}
Expand Down Expand Up @@ -612,13 +637,13 @@ fn stage0_queries_do_not_share_nodes() {

// ── Stage 1 ──────────────────────────────────────────────────────────────

/// Stage 1 outputs the 24 Pass 1 combinations twice: L1–L24 with separate
/// inputs, then L25–L48 with the shared input (Pass 2).
/// Stage 1 outputs the 32 Pass 1 combinations twice: L1–L32 with separate
/// inputs, then L33–L64 with the shared input (Pass 2).
#[test]
fn stage1_has_48_candidates_covering_every_combination_twice() {
fn stage1_has_64_candidates_covering_every_combination_twice() {
let (_, logical, _) = pipeline();
assert_eq!(logical.len(), 48);
for (half, shared) in [(&logical[..24], false), (&logical[24..], true)] {
assert_eq!(logical.len(), 64);
for (half, shared) in [(&logical[..32], false), (&logical[32..], true)] {
let found: BTreeSet<_> = half
.iter()
.map(|c| choices(&c.dag, &c.query_roots))
Expand Down Expand Up @@ -663,6 +688,8 @@ fn stage1_q2_summary_families_are_heap_sketches_and_hydra() {
BTreeSet::from([
Q2Option::CountMinHeapPerJob,
Q2Option::CountSketchHeapPerJob,
Q2Option::WholeCountMinHeapPerJob,
Q2Option::WholeCountSketchHeapPerJob,
Q2Option::Hydra
])
);
Expand All @@ -681,6 +708,8 @@ fn stage1_keeps_independent_and_shared_variants() {
Q2Option::Exact,
Q2Option::CountMinHeapPerJob,
Q2Option::CountSketchHeapPerJob,
Q2Option::WholeCountMinHeapPerJob,
Q2Option::WholeCountSketchHeapPerJob,
] {
assert!(
found.contains(&(option, false)),
Expand Down Expand Up @@ -724,11 +753,11 @@ fn stage1_candidates_are_valid_and_uniquely_named() {

// ── Stage 2 ──────────────────────────────────────────────────────────────

/// No candidate is discarded before Stage 3: Stage 2 maps the 48 logical candidates one-to-one.
/// No candidate is discarded before Stage 3: Stage 2 maps the 64 logical candidates one-to-one.
#[test]
fn stage2_keeps_every_logical_candidate() {
let (_, logical, physical) = pipeline();
assert_eq!(physical.len(), 48);
assert_eq!(physical.len(), 64);
let sources: BTreeSet<_> = physical.iter().map(|p| p.from_logical.as_str()).collect();
let logical_ids: BTreeSet<_> = logical.iter().map(|c| c.id.as_str()).collect();
assert_eq!(sources, logical_ids);
Expand Down Expand Up @@ -883,7 +912,7 @@ fn stage2_runtime_compiles_exactly_the_candidates_stage3_finds_valid() {
.filter(|r| !r.valid)
.map(|r| (r.id.as_str(), r.reason.as_str()))
.collect();
assert_eq!(invalid.len(), 16, "the Count-Min + heap candidates");
assert_eq!(invalid.len(), 24, "the Count-Min + heap candidates");
for p in &physical {
let compiled = compile_in_runtime(p);
match invalid.get(p.id.as_str()) {
Expand Down Expand Up @@ -1026,7 +1055,11 @@ fn stage3_shared_input_is_not_costlier() {
Some(((choices(&p.dag, &p.query_roots), shared), cost))
})
.collect();
for option in [Q2Option::Exact, Q2Option::CountSketchHeapPerJob] {
for option in [
Q2Option::Exact,
Q2Option::CountSketchHeapPerJob,
Q2Option::WholeCountSketchHeapPerJob,
] {
assert!(
by_combo.keys().any(|((_, o, _), _)| *o == option),
"{option:?} has no priced candidate"
Expand Down Expand Up @@ -1081,3 +1114,27 @@ fn stage3_selects_a_shared_input_plan() {
"{saving} vs {input_cost}"
);
}

/// Every Q2 realization, exact or sketch, whole-expression or not, returns
/// the same selected-rows schema (#579) in Stage 1, so consumers see one
/// shape. (Stage 2 exports exact top-k as sort → limit over the per-series
/// rows, whose schema still carries `ts`.)
#[test]
fn q2_roots_keep_one_schema_across_realizations() {
let (_, logical, _) = pipeline();
let schemas: Vec<_> = logical
.iter()
.map(|c| {
(
classify(&c.dag, &c.query_roots).0,
c.roots[1].schema.clone(),
)
})
.collect();
assert!(schemas
.iter()
.any(|(o, _)| *o == Q2Option::WholeCountSketchHeapPerJob));
for (option, schema) in &schemas {
assert_eq!(schema, &schemas[0].1, "{option:?}");
}
}
Loading