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
10 changes: 6 additions & 4 deletions crates/asap-aware-mapping/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -44,8 +44,10 @@
//! coordinates logical choices and preserves shared nodes. Whether and when
//! a summary state is materialized is not decided here: every summary runs
//! at query time until Stage 2 materialization (#509) owns that choice.
//! - Run the whole pipeline through [`optimize`] with [`MajorPass`], which
//! performs the two steps above for every root of a parsed workload.
//! - Run the #509 stage pipeline through [`optimize`] with [`StagePipeline`]:
//! Stage 1 local alternatives, Stage 2 physical candidates, and Stage 3
//! selection, the only stage that prices plans. It does not use the
//! candidate search above.
//!
//! Models and evidence determine which choices the helpers can justify.
//! Physical operator binding, placement, storage, deployment, and execution
Expand Down Expand Up @@ -190,8 +192,8 @@ pub use explanation::{
};
pub use grouping::{has_subpopulations, HydraGroupingStrategy};
pub use pass::{
optimize, MajorPass, OptimizationInput, OptimizationInputError, OptimizationPass,
OptimizeError, PassNameConflict, PassRegistry, PlanOutput, PlanningModels, QueryPlan,
optimize, OptimizationInput, OptimizationInputError, OptimizationPass, OptimizeError,
PassNameConflict, PassRegistry, PlanOutput, PlanningModels, QueryPlan, StagePipeline,
};
pub use recurrence::{
evaluation_rate_of, total_cost, update_rate_from_data_workload, CostRate, EvaluationRate,
Expand Down
128 changes: 128 additions & 0 deletions crates/asap-aware-mapping/src/logical_candidates.rs
Original file line number Diff line number Diff line change
Expand Up @@ -130,6 +130,81 @@ pub fn enumerate_local_logical_candidates<Id>(
Ok(LocalLogicalCandidates { roots, targets })
}

/// Number of whole-workload candidates: one per choice of an alternative for
/// every target. Saturates rather than overflowing.
pub fn combination_count<Id>(inventory: &LocalLogicalCandidates<Id>) -> usize {
inventory
.targets
.iter()
.fold(1usize, |n, t| n.saturating_mul(t.alternatives.len()))
}

/// The first `max` choices in enumeration order: mixed radix, the last target
/// varying fastest. `choice[i]` indexes `inventory.targets[i].alternatives`.
pub fn enumerate_choices<Id>(
inventory: &LocalLogicalCandidates<Id>,
max: usize,
) -> Vec<Vec<usize>> {
let count = combination_count(inventory).min(max);
let mut choices = Vec::with_capacity(count);
let mut choice = vec![0; inventory.targets.len()];
for _ in 0..count {
choices.push(choice.clone());
for (digit, target) in choice.iter_mut().zip(&inventory.targets).rev() {
*digit += 1;
if *digit < target.alternatives.len() {
break;
}
*digit = 0;
}
}
choices
}

/// Position of `choice` in [`enumerate_choices`] order.
pub fn choice_index<Id>(inventory: &LocalLogicalCandidates<Id>, choice: &[usize]) -> usize {
inventory
.targets
.iter()
.zip(choice)
.fold(0usize, |index, (target, &digit)| {
index
.saturating_mul(target.alternatives.len())
.saturating_add(digit)
})
}

/// For each target, the targets directly beneath it: reachable from its input
/// without passing through another target.
pub fn nested_targets<Id>(inventory: &LocalLogicalCandidates<Id>) -> Vec<Vec<usize>> {
let position: HashMap<_, _> = inventory
.targets
.iter()
.enumerate()
.map(|(i, t)| (Rc::as_ptr(&t.target), i))
.collect();
inventory
.targets
.iter()
.map(|target| {
let mut found = Vec::new();
let mut seen = HashSet::new();
let mut stack: Vec<_> = target.target.children().into_iter().cloned().collect();
while let Some(node) = stack.pop() {
if !seen.insert(Rc::as_ptr(&node)) {
continue;
}
match position.get(&Rc::as_ptr(&node)) {
Some(&index) => found.push(index),
None => stack.extend(node.children().into_iter().cloned()),
}
}
found.sort_unstable();
found
})
.collect()
}

/// Build one whole-workload candidate (#509 Stage 1): `choice[i]` indexes
/// `inventory.targets[i].alternatives`. Each chosen non-pass-through target is
/// replaced by `SummaryAgg` followed by `SummaryEstimate` (sketch) or
Expand Down Expand Up @@ -381,6 +456,59 @@ fn statistic(intent: &AggIntent) -> Result<SketchStatistic, LogicalCandidateErro
#[cfg(test)]
mod tests {
use super::*;
use crate::test_support::lower_promql;

/// Every realization of a top-k returns the same selected-rows schema, so
/// an aggregate over a pass-through top-k composes like one over a sketch.
#[test]
fn aggregate_over_any_topk_realization_composes() {
let root = lower_promql(
"count(topk by (job) (10, sum_over_time(m[1m])))",
AccuracyTarget::EpsilonDelta {
epsilon: 0.01,
delta: 0.001,
},
);
let root = asap_types::ir::schema_support::with_promql_series_identity(&root).unwrap();
let inventory =
enumerate_local_logical_candidates(vec![(0, QueryRoot::Operator(root))]).unwrap();
let topk = inventory
.targets
.iter()
.position(|t| {
matches!(t.target.non_asap(), Some(NonASAPOp::Aggregate { measures, .. })
if matches!(measures.as_slice(), [AggIntent::TopK { .. }]))
})
.unwrap();
let mut schemas = Vec::new();
for (count, alternatives) in inventory.targets.iter().enumerate() {
if count == topk {
continue;
}
for outer in 0..alternatives.alternatives.len() {
for inner in 0..inventory.targets[topk].alternatives.len() {
let mut choice = vec![0; inventory.targets.len()];
choice[count] = outer;
choice[topk] = inner;
compose_logical_candidate(&inventory, &choice)
.unwrap_or_else(|e| panic!("{choice:?}: {e}"));
}
}
}
for inner in 0..inventory.targets[topk].alternatives.len() {
let mut choice = vec![0; inventory.targets.len()];
choice[topk] = inner;
let roots = compose_logical_candidate(&inventory, &choice).unwrap();
let QueryRoot::Operator(root) = &roots[0].1 else {
panic!("operator root")
};
let NonASAPOp::Aggregate { child, .. } = root.expect_non_asap() else {
panic!("count over top-k")
};
schemas.push(child.schema.clone());
}
assert!(schemas.windows(2).all(|w| w[0] == w[1]), "{schemas:#?}");
}
/// Approximate requests must retain the exact execution alternative too.
#[test]
fn approximate_count_keeps_exact_and_universal_choices() {
Expand Down
89 changes: 0 additions & 89 deletions crates/asap-aware-mapping/src/pass/major.rs

This file was deleted.

Loading