Skip to content
Merged
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
6 changes: 3 additions & 3 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

26 changes: 14 additions & 12 deletions asap-planner-rs/src/bin/optimizer_cli.rs
Original file line number Diff line number Diff line change
Expand Up @@ -188,33 +188,35 @@ fn run_milp(args: &Args, config: &ControllerConfig) -> anyhow::Result<()> {
tracing::debug!(?objective, cost_rows = costs.len(), "milp: inputs loaded");

let workload = build_milp_workload(config, &facts, args.data_ingestion_interval_ms)?;
let (deployments, solution) = solve_milp(&workload, &facts, &costs, objective)?;

let mut active: Vec<usize> = solution.mapping.clone();
active.sort();
active.dedup();
println!("=== Deployments: {} ===", active.len());
for &d in &active {
let dep = &deployments[d];
let solution = solve_milp(&workload, &facts, &costs, objective)?;

println!("=== Deployments: {} ===", solution.deployments.len());
for (d, planned) in solution.deployments.iter().enumerate() {
let dep = &planned.deployment;
println!(
" [{d}] {:?} {} config={} metric={} grouping={:?} window={}ms slide={}ms",
" [{d}] {:?} {} config={} metric={} grouping={:?} window={}ms slide={}ms retained={} key_tracker={}",
dep.capability,
dep.config.sketch,
dep.config.sketch_config,
dep.metric,
dep.grouping_labels,
dep.window_ms,
dep.slide_ms,
planned.retained_instance_count,
dep.key_tracker.is_some(),
);
}
println!("\n=== Raqes: {} ===", workload.raqes.len());
for ((raqe, &d), latency_ms) in workload
for ((raqe, planned), latency_ms) in workload
.raqes
.iter()
.zip(&solution.mapping)
.zip(&solution.raqes)
.zip(&solution.plan_cost.query_latency_ms)
{
println!(" {} -> [{d}] latency={latency_ms:.3e}ms", raqe.id);
println!(
" {} -> [{}] n={} latency={latency_ms:.3e}ms",
raqe.id, planned.deployment, planned.merged_instance_count
);
}
let cost = &solution.plan_cost;
println!(
Expand Down
151 changes: 113 additions & 38 deletions asap-planner-rs/src/optimizer/milp.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,11 +2,11 @@
//! workload config into `Raqe`s and solves for the cheapest deployments.

use promql_utilities::query_logics::enums::Statistic;
use rqe_optimizer::candidates::{build_all_candidates, eligible_deployments_for};
use rqe_optimizer::candidates::build_all_candidates;
use rqe_optimizer::enumerate::unservable;
use rqe_optimizer::milp::{minimize, MilpSolution, Objective};
use rqe_optimizer::{
validate_facts, AccuracyDirection, AtomicCostEntry, Capability, Deployment, LabelSet, Raqe,
WorkloadFacts,
validate_facts, AccuracyDirection, AtomicCostEntry, Capability, LabelSet, Raqe, WorkloadFacts,
};
use thiserror::Error;

Expand Down Expand Up @@ -37,34 +37,38 @@ pub enum MilpError {
query: String,
statistics: Vec<Statistic>,
},
#[error("query {query:?}: topk ranking (by value or by sample count) is unknown")]
TopkWeightingUnknown { query: String },
#[error("no eligible deployment for raqes: {0:?}")]
Unservable(Vec<String>),
#[error("MILP solve failed: {0}")]
Solver(String),
}

/// Candidate deployments and the MILP's choice among them. Only families in
/// sketch-bench's `DEPLOYABLE_FAMILIES` are candidates.
/// The cheapest plan. Only families in sketch-bench's `DEPLOYABLE_FAMILIES`
/// are candidates.
pub fn solve_milp(
workload: &MilpWorkload,
facts: &WorkloadFacts,
costs: &[AtomicCostEntry],
objective: Objective,
) -> Result<(Vec<Deployment>, MilpSolution), MilpError> {
) -> Result<MilpSolution, MilpError> {
let deployments = build_all_candidates(&workload.raqes, costs, facts, false);
tracing::debug!(
candidates = deployments.len(),
cost_rows = costs.len(),
"milp: built candidates"
);
let mut missing = Vec::new();
for raqe in &workload.raqes {
let eligible = eligible_deployments_for(raqe, &deployments).len();
tracing::debug!(raqe = %raqe.id, eligible, "milp: eligible deployments");
if eligible == 0 {
missing.push(raqe.id.clone());
}
}
// Occurrences of one item are identical Raqes (and contiguous), so
// checking the first of each is enough.
let one_per_item: Vec<Raqe> = workload
.raqes
.iter()
.enumerate()
.filter(|&(i, _)| i == 0 || workload.raqe_items[i] != workload.raqe_items[i - 1])
.map(|(_, raqe)| raqe.clone())
.collect();
let missing = unservable(&one_per_item, &deployments, facts);
if !missing.is_empty() {
return Err(MilpError::Unservable(missing));
}
Expand All @@ -74,10 +78,10 @@ pub fn solve_milp(
objective = objective.value(&solution.plan_cost),
cpu_secs_per_sec = solution.plan_cost.cpu_secs_per_sec(),
memory_bytes = solution.plan_cost.memory_bytes(),
mapping = ?solution.mapping,
deployments = solution.deployments.len(),
"milp: solved"
);
Ok((deployments, solution))
Ok(solution)
}

/// The MILP's view of a workload.
Expand Down Expand Up @@ -166,7 +170,9 @@ fn item_to_raqe(item: &OptimizerItem) -> Result<Raqe, MilpError> {
statistics: req.statistics.clone(),
});
};
let capability = capability(*statistic);
let Some(capability) = capability(*statistic, req.topk_count_events) else {
return Err(MilpError::TopkWeightingUnknown { query });
};
if !(0.0..=1.0).contains(&item.accuracy_sla) {
return Err(MilpError::AccuracySlaOutOfRange {
query,
Expand All @@ -179,7 +185,7 @@ fn item_to_raqe(item: &OptimizerItem) -> Result<Raqe, MilpError> {
// TopK keeps one heap per `topk by` bucket; its `grouping_labels` is the
// output label set (every label), which would cost one heap per series.
let grouping_labels: LabelSet = match capability {
Capability::TopK => req
Capability::TopKByValue | Capability::TopKByCount => req
.topk_by_labels
.as_ref()
.map(|labels| labels.labels.iter().cloned().collect())
Expand All @@ -202,16 +208,23 @@ fn item_to_raqe(item: &OptimizerItem) -> Result<Raqe, MilpError> {
})
}

fn capability(statistic: Statistic) -> Capability {
match statistic {
Statistic::Sum | Statistic::Count => Capability::SumOrCount,
/// `None` for a topk whose weighting is unknown: value-ranked
/// (`sum_over_time`) and count-ranked (`count_over_time`) top-k need separate
/// heaps.
fn capability(statistic: Statistic, topk_count_events: Option<bool>) -> Option<Capability> {
Some(match statistic {
Statistic::Sum => Capability::Sum,
Statistic::Count => Capability::Count,
Statistic::Rate | Statistic::Increase => Capability::RateOrIncrease,
Statistic::Min => Capability::Min,
Statistic::Max => Capability::Max,
Statistic::Quantile => Capability::Quantile,
Statistic::Cardinality => Capability::Cardinality,
Statistic::Topk => Capability::TopK,
}
Statistic::Topk => match topk_count_events? {
true => Capability::TopKByCount,
false => Capability::TopKByValue,
},
})
}

/// `accuracy_sla` is required accuracy: error metrics must stay within
Expand All @@ -222,13 +235,14 @@ fn accuracy_target(
) -> (&'static str, f64, AccuracyDirection) {
let max_error = 1.0 - accuracy_sla + SLA_EPSILON;
match capability {
Capability::SumOrCount
Capability::Sum
| Capability::Count
| Capability::Min
| Capability::Max
| Capability::RateOrIncrease
| Capability::Cardinality => (RELATIVE_ERROR, max_error, AccuracyDirection::LowerIsBetter),
Capability::Quantile => (MAX_RANK_ERROR, max_error, AccuracyDirection::LowerIsBetter),
Capability::TopK => (
Capability::TopKByValue | Capability::TopKByCount => (
PRECISION_AT_K,
accuracy_sla - SLA_EPSILON,
AccuracyDirection::HigherIsBetter,
Expand Down Expand Up @@ -352,10 +366,10 @@ metrics:
}

#[test]
fn sum_maps_to_sum_or_count_with_relative_error_ceiling() {
fn sum_maps_to_sum_with_relative_error_ceiling() {
let w = workload(&group("sum by (job) (http_requests_total)", 0.99));
let r = &w.raqes[0];
assert_eq!(r.capability, Capability::SumOrCount);
assert_eq!(r.capability, Capability::Sum);
assert_eq!(r.grouping_labels, labels(&["job"]));
assert_eq!(r.accuracy_metric, RELATIVE_ERROR);
assert_eq!(r.accuracy_direction, AccuracyDirection::LowerIsBetter);
Expand All @@ -364,6 +378,34 @@ metrics:
assert_eq!(r.latency_sla_ms, None);
}

#[test]
fn count_and_topk_weighting_map_to_their_own_capabilities() {
let cap = |query| workload(&group(query, 0.99)).raqes[0].capability;
assert_eq!(
cap("sum by (job) (count_over_time(http_requests_total[1m]))"),
Capability::Count
);
assert_eq!(
cap("topk(5, sum_over_time(http_requests_total[1m]))"),
Capability::TopKByValue
);
assert_eq!(
cap("topk(5, count_over_time(http_requests_total[1m]))"),
Capability::TopKByCount
);
}

#[test]
fn topk_with_unknown_weighting_is_an_error() {
let w = workload(&group("topk(5, http_requests_total)", 0.99));
let mut item = w.items[0].clone();
item.requirements.topk_count_events = None;
assert!(matches!(
item_to_raqe(&item),
Err(MilpError::TopkWeightingUnknown { .. })
));
}

#[test]
fn quantile_uses_max_rank_error() {
let w = workload(&group(
Expand All @@ -380,7 +422,7 @@ metrics:
fn topk_requires_precision_and_groups_by_its_buckets() {
let w = workload(&group("topk by (job) (5, http_requests_total)", 0.9));
let r = &w.raqes[0];
assert_eq!(r.capability, Capability::TopK);
assert_eq!(r.capability, Capability::TopKByValue);
assert_eq!(r.accuracy_metric, PRECISION_AT_K);
assert_eq!(r.accuracy_direction, AccuracyDirection::HigherIsBetter);
assert!((r.accuracy_sla - 0.9).abs() < 1e-6);
Expand Down Expand Up @@ -427,10 +469,13 @@ metrics:
let config = config(&(group(query, 0.99) + &group(query, 0.99)));
let facts = facts(&config);
let w = build_milp_workload(&config, &facts, SCRAPE_MS).unwrap();
let (deployments, solution) =
solve_milp(&w, &facts, &costs(), Objective::default()).unwrap();
assert_eq!(solution.mapping[0], solution.mapping[1]);
assert_eq!(deployments[solution.mapping[0]].config.sketch, "exact-sum");
let solution = solve_milp(&w, &facts, &costs(), Objective::default()).unwrap();
assert_eq!(solution.deployments.len(), 1);
assert_eq!(solution.raqes[0].deployment, solution.raqes[1].deployment);
assert_eq!(
solution.deployments[0].deployment.config.sketch,
"exact-sum"
);
}

#[test]
Expand All @@ -446,6 +491,17 @@ metrics:
assert!(matches!(err, MilpError::Unservable(ids) if ids.len() == 1));
}

#[test]
fn unservable_query_is_reported_once_per_item() {
let quantile = group("quantile_over_time(0.99, http_requests_total[5m])", 0.999);
let config = config(&(quantile.clone() + &quantile));
let facts = facts(&config);
let w = build_milp_workload(&config, &facts, SCRAPE_MS).unwrap();
assert_eq!(w.raqes.len(), 2);
let err = solve_milp(&w, &facts, &costs(), Objective::default()).unwrap_err();
assert!(matches!(err, MilpError::Unservable(ids) if ids == [w.raqes[0].id.clone()]));
}

#[test]
fn solve_picks_a_deployable_family_per_capability() {
let groups = group("sum by (job) (http_requests_total)", 0.99)
Expand All @@ -454,16 +510,35 @@ metrics:
let config = config(&groups);
let facts = facts(&config);
let w = build_milp_workload(&config, &facts, SCRAPE_MS).unwrap();
let (deployments, solution) =
solve_milp(&w, &facts, &costs(), Objective::default()).unwrap();
let solution = solve_milp(&w, &facts, &costs(), Objective::default()).unwrap();
let chosen: BTreeMap<Capability, &str> = w
.raqes
.iter()
.zip(&solution.mapping)
.map(|(r, &d)| (r.capability, deployments[d].config.sketch.as_str()))
.zip(&solution.raqes)
.map(|(r, p)| {
let sketch = &solution.deployments[p.deployment].deployment.config.sketch;
(r.capability, sketch.as_str())
})
.collect();
assert_eq!(chosen[&Capability::SumOrCount], "exact-sum");
assert_eq!(chosen[&Capability::Sum], "exact-sum");
assert_eq!(chosen[&Capability::Quantile], "kll-percall");
assert_eq!(chosen[&Capability::TopK], "cms-heap-topk-fastpath-vector2d");
assert_eq!(
chosen[&Capability::TopKByValue],
"cms-heap-topk-fastpath-vector2d"
);
}

#[test]
fn avg_sum_and_count_get_separate_deployments() {
// One exact-sum accumulator can't answer both halves of the avg
// rewrite; sharing would undercount ingest and memory.
let config = config(&group("avg by (job) (http_requests_total)", 0.99));
let facts = facts(&config);
let w = build_milp_workload(&config, &facts, SCRAPE_MS).unwrap();
let caps: Vec<_> = w.raqes.iter().map(|r| r.capability).collect();
assert_eq!(caps.len(), 2);
assert!(caps.contains(&Capability::Sum) && caps.contains(&Capability::Count));
let solution = solve_milp(&w, &facts, &costs(), Objective::default()).unwrap();
assert_eq!(solution.deployments.len(), 2);
}
}
Loading