diff --git a/Cargo.lock b/Cargo.lock index c8b175c3..568bf135 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -112,7 +112,7 @@ checksum = "7f202df86484c868dbad7eaa557ef785d5c66295e41b460ef922eca0723b842c" [[package]] name = "aqpbm-core" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/sketch-bench?branch=main#3b31dfbdc8d0f3981312cfe3f8c3649787bfc381" +source = "git+https://github.com/ProjectASAP/sketch-bench?branch=main#10415767b3366d5f56ff7bef5182a53cd055b6a7" dependencies = [ "anyhow", "aqpbm-datagen", @@ -128,7 +128,7 @@ dependencies = [ [[package]] name = "aqpbm-datagen" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/sketch-bench?branch=main#3b31dfbdc8d0f3981312cfe3f8c3649787bfc381" +source = "git+https://github.com/ProjectASAP/sketch-bench?branch=main#10415767b3366d5f56ff7bef5182a53cd055b6a7" dependencies = [ "rand 0.9.4", "rand_distr", @@ -2388,7 +2388,7 @@ dependencies = [ [[package]] name = "rqe-optimizer" version = "0.1.0" -source = "git+https://github.com/ProjectASAP/sketch-bench?branch=main#3b31dfbdc8d0f3981312cfe3f8c3649787bfc381" +source = "git+https://github.com/ProjectASAP/sketch-bench?branch=main#10415767b3366d5f56ff7bef5182a53cd055b6a7" dependencies = [ "aqpbm-core", "good_lp", diff --git a/asap-planner-rs/src/bin/optimizer_cli.rs b/asap-planner-rs/src/bin/optimizer_cli.rs index 2ca22513..e2c7cd02 100644 --- a/asap-planner-rs/src/bin/optimizer_cli.rs +++ b/asap-planner-rs/src/bin/optimizer_cli.rs @@ -77,8 +77,8 @@ struct Args { )] allow_undeployable_families: bool, - /// MILP only. YAML workload facts: per metric, `cardinality` per label - /// set, including the set of all its labels (the series count). + /// MILP only. YAML workload facts: per metric, positive `value_range` + /// and `cardinality` per label set, including all labels (the series count). #[arg( long = "workload-facts", required_if_eq("milp", "true"), diff --git a/asap-planner-rs/src/optimizer/milp.rs b/asap-planner-rs/src/optimizer/milp.rs index 9aae3f64..7269ad3d 100644 --- a/asap-planner-rs/src/optimizer/milp.rs +++ b/asap-planner-rs/src/optimizer/milp.rs @@ -6,7 +6,7 @@ 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, LabelSet, Raqe, WorkloadFacts, + table_accuracy, validate_facts, AtomicCostEntry, Capability, LabelSet, Raqe, WorkloadFacts, }; use thiserror::Error; @@ -15,11 +15,6 @@ use crate::config::input::ControllerConfig; use super::pipeline::{extract_hinted_items, OptimizerPipelineError}; use super::solution::OptimizerItem; -/// `query_accuracy` keys in sketch-bench's cost export. Each is the worst case -/// its comparator reports. -const RELATIVE_ERROR: &str = "relative_error"; -const MAX_RANK_ERROR: &str = "max_rank_err"; -const PRECISION_AT_K: &str = "precision_at_k"; /// Slack on accuracy tolerances so `1 - sla` rounding (`1 - 0.9 = /// 0.0999...98`) doesn't reject a row measured exactly at the boundary. const SLA_EPSILON: f64 = 1e-9; @@ -55,8 +50,13 @@ pub fn solve_milp( objective: Objective, allow_undeployable_families: bool, ) -> Result { - let deployments = - build_all_candidates(&workload.raqes, costs, facts, allow_undeployable_families); + let deployments = build_all_candidates( + &workload.raqes, + costs, + facts, + allow_undeployable_families, + &table_accuracy, + ); tracing::debug!( candidates = deployments.len(), cost_rows = costs.len(), @@ -71,12 +71,18 @@ pub fn solve_milp( .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); + let missing = unservable(&one_per_item, &deployments, facts, &table_accuracy); if !missing.is_empty() { return Err(MilpError::Unservable(missing)); } - let solution = minimize(&workload.raqes, &deployments, facts, objective) - .map_err(|e| MilpError::Solver(e.to_string()))?; + let solution = minimize( + &workload.raqes, + &deployments, + facts, + objective, + &table_accuracy, + ) + .map_err(|e| MilpError::Solver(e.to_string()))?; tracing::debug!( objective = objective.value(&solution.plan_cost), cpu_secs_per_sec = solution.plan_cost.cpu_secs_per_sec(), @@ -129,9 +135,7 @@ pub fn build_milp_workload( grouping_labels = ?raqe.grouping_labels, lookback_ms = raqe.lookback_ms, interval_ms = raqe.interval_ms, - accuracy_metric = %raqe.accuracy_metric, accuracy_sla = raqe.accuracy_sla, - accuracy_direction = ?raqe.accuracy_direction, latency_sla_ms = ?raqe.latency_sla_ms, "milp inputs: item -> raqe" ); @@ -182,8 +186,7 @@ fn item_to_raqe(item: &OptimizerItem) -> Result { accuracy_sla: item.accuracy_sla, }); } - let (accuracy_metric, accuracy_sla, accuracy_direction) = - accuracy_target(capability, item.accuracy_sla); + let accuracy_sla = accuracy_target(capability, item.accuracy_sla); // 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. @@ -204,9 +207,7 @@ fn item_to_raqe(item: &OptimizerItem) -> Result { metric: req.metric.clone(), spatial_filter: req.spatial_filter_normalized.clone(), grouping_labels, - accuracy_metric: accuracy_metric.to_string(), accuracy_sla, - accuracy_direction, latency_sla_ms: item.latency_sla_ms, }) } @@ -232,10 +233,7 @@ fn capability(statistic: Statistic, topk_count_events: Option) -> Option (&'static str, f64, AccuracyDirection) { +fn accuracy_target(capability: Capability, accuracy_sla: f64) -> f64 { let max_error = 1.0 - accuracy_sla + SLA_EPSILON; match capability { Capability::Sum @@ -243,13 +241,9 @@ fn accuracy_target( | Capability::Min | Capability::Max | Capability::RateOrIncrease - | Capability::Cardinality => (RELATIVE_ERROR, max_error, AccuracyDirection::LowerIsBetter), - Capability::Quantile => (MAX_RANK_ERROR, max_error, AccuracyDirection::LowerIsBetter), - Capability::TopKByValue | Capability::TopKByCount => ( - PRECISION_AT_K, - accuracy_sla - SLA_EPSILON, - AccuracyDirection::HigherIsBetter, - ), + | Capability::Cardinality + | Capability::Quantile => max_error, + Capability::TopKByValue | Capability::TopKByCount => accuracy_sla - SLA_EPSILON, } } @@ -265,6 +259,7 @@ pub(super) mod tests { const FACTS: &str = r#" metrics: - metric: http_requests_total + value_range: [1.0, 1000.0] groups: - labels: [instance, job] cardinality: 100 @@ -316,11 +311,11 @@ metrics: pub(crate) fn costs() -> Vec { vec![ - cost("exact-sum", &[(RELATIVE_ERROR, 0.0)]), - cost("kll-percall", &[(MAX_RANK_ERROR, 0.005)]), + cost("exact-sum", &[("relative_error", 0.0)]), + cost("kll-percall", &[("mean_rank_err", 0.005)]), cost( "cms-heap-topk-fastpath-vector2d", - &[(PRECISION_AT_K, 0.995)], + &[("precision_at_k", 0.995)], ), ] } @@ -341,11 +336,13 @@ metrics: "quantile_over_time(0.99, http_requests_total[5m])", 0.9, )); - assert!(w.raqes[0].accuracy_ok(0.1)); - assert!(!w.raqes[0].accuracy_ok(0.1001)); + assert!(w.raqes[0].meets_sla("kll-percall", Some(0.1))); + assert!(w.raqes[0].meets_sla("dd", Some(0.1))); + assert!(!w.raqes[0].meets_sla("kll-percall", Some(0.1001))); + assert!(!w.raqes[0].meets_sla("dd", Some(0.1001))); let w = workload(&group("topk(5, http_requests_total)", 0.9)); - assert!(w.raqes[0].accuracy_ok(0.9)); - assert!(!w.raqes[0].accuracy_ok(0.8999)); + assert!(w.raqes[0].meets_sla("cms-heap-topk-fastpath-vector2d", Some(0.9))); + assert!(!w.raqes[0].meets_sla("cms-heap-topk-fastpath-vector2d", Some(0.8999))); } #[test] @@ -374,8 +371,6 @@ metrics: let r = &w.raqes[0]; 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); assert!((r.accuracy_sla - 0.01).abs() < 1e-6); assert_eq!(r.interval_ms, 60_000); assert_eq!(r.latency_sla_ms, None); @@ -410,14 +405,14 @@ metrics: } #[test] - fn quantile_uses_max_rank_error() { + fn quantile_uses_the_error_ceiling() { let w = workload(&group( "quantile_over_time(0.99, http_requests_total[5m])", 0.99, )); let r = &w.raqes[0]; assert_eq!(r.capability, Capability::Quantile); - assert_eq!(r.accuracy_metric, MAX_RANK_ERROR); + assert!((r.accuracy_sla - 0.01).abs() < 1e-6); assert_eq!(r.lookback_ms, 300_000); } @@ -426,8 +421,6 @@ metrics: let w = workload(&group("topk by (job) (5, http_requests_total)", 0.9)); let r = &w.raqes[0]; 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); // Not the all-labels output set, which would cost one heap per series. assert_eq!(r.grouping_labels, labels(&["job"])); diff --git a/asap-planner-rs/src/optimizer/milp_output.rs b/asap-planner-rs/src/optimizer/milp_output.rs index 0fa8e707..fbc517e5 100644 --- a/asap-planner-rs/src/optimizer/milp_output.rs +++ b/asap-planner-rs/src/optimizer/milp_output.rs @@ -39,6 +39,8 @@ pub enum MilpOutputError { variant: String, param: &'static str, }, + #[error("sketch-bench variant dd: alpha must be finite and in (0, 1), got {alpha}")] + InvalidDdsAlpha { alpha: f64 }, #[error("topk query {0:?} has no literal k")] TopkWithoutK(String), #[error("query {0:?} is served by two deployments; a query string can name only one")] @@ -130,13 +132,15 @@ fn aggregation_config( items: &[&OptimizerItem], ) -> Result { let variant = deployment.config.sketch.as_str(); - let param = |name: &'static str| { + let integer_param = |name: &'static str| { deployment.config.sketch_config["params"][name] .as_u64() - .ok_or_else(|| MilpOutputError::MissingParam { - variant: variant.to_string(), - param: name, - }) + .ok_or_else(|| missing_param(variant, name)) + }; + let float_param = |name: &'static str| { + deployment.config.sketch_config["params"][name] + .as_f64() + .ok_or_else(|| missing_param(variant, name)) }; let (aggregation_type, sub_type, parameters) = match (variant, deployment.capability) { ("exact-sum", Capability::Sum) => (AggregationType::MultipleSum, "sum", vec![]), @@ -149,12 +153,23 @@ fn aggregation_config( ("kll-percall", Capability::Quantile) => ( AggregationType::DatasketchesKLL, "", - vec![("K", param("k")?)], + vec![("K", Value::from(integer_param("k")?))], ), + ("dd", Capability::Quantile) => { + let alpha = float_param("alpha")?; + if !(alpha.is_finite() && alpha > 0.0 && alpha < 1.0) { + return Err(MilpOutputError::InvalidDdsAlpha { alpha }); + } + ( + AggregationType::DDSketch, + "", + vec![("alpha", Value::from(alpha))], + ) + } ("hll", Capability::Cardinality) => ( AggregationType::HLL, "", - vec![("precision", param("lg_k")?)], + vec![("precision", Value::from(integer_param("lg_k")?))], ), ( "cms-heap-topk-fastpath-vector2d", @@ -167,9 +182,9 @@ fn aggregation_config( "sum" }, vec![ - ("depth", param("rows")?), - ("width", param("cols")?), - ("heapsize", max_topk_k(items)?), + ("depth", Value::from(integer_param("rows")?)), + ("width", Value::from(integer_param("cols")?)), + ("heapsize", Value::from(max_topk_k(items)?)), ], ), (_, capability) => { @@ -212,7 +227,7 @@ fn aggregation_config( value_column: None, parameters: parameters .into_iter() - .map(|(name, value)| (name.to_string(), Value::from(value))) + .map(|(name, value)| (name.to_string(), value)) .collect(), rollup_labels: KeyByLabelNames::empty(), grouping_labels, @@ -220,6 +235,13 @@ fn aggregation_config( }) } +fn missing_param(variant: &str, param: &'static str) -> MilpOutputError { + MilpOutputError::MissingParam { + variant: variant.to_string(), + param, + } +} + /// One heap serves every topk query on the deployment, so it is sized for /// the largest k. fn max_topk_k(items: &[&OptimizerItem]) -> Result { @@ -245,9 +267,12 @@ fn topk_k(query: &str) -> Option { #[cfg(test)] mod tests { + use std::collections::BTreeMap; + use asap_types::aggregation_config::AggregationConfig; use asap_types::enums::QueryLanguage; use rqe_optimizer::milp::Objective; + use rqe_optimizer::AtomicCostEntry; use serde_json::json; use super::*; @@ -288,6 +313,20 @@ mod tests { .clone() } + fn dd_cost(params: serde_json::Value) -> AtomicCostEntry { + AtomicCostEntry { + sketch: "dd".into(), + sketch_config: json!({"algorithm": "dd", "params": params}), + mem_bytes_per_instance: 100.0, + insert_cpu_secs: 1e-7, + merge_cpu_secs: 1e-7, + query_cpu_secs: 1e-7, + query_accuracy: BTreeMap::from([("mean_relative_value_error".into(), 0.005)]), + merge_accuracy: BTreeMap::new(), + measured_at: None, + } + } + #[test] fn plan_round_trips_into_engine_configs() { let sum = "sum by (job) (http_requests_total)"; @@ -358,6 +397,57 @@ mod tests { } } + #[test] + fn ddsketch_plan_preserves_the_cost_row_alpha() { + let query = "quantile_over_time(0.99, http_requests_total[5m])"; + let config = config(&group(query, 0.99)); + let facts = facts(&config); + let workload = build_milp_workload(&config, &facts, SCRAPE_MS).unwrap(); + let costs = vec![dd_cost(json!({"alpha": 0.02}))]; + let solution = solve_milp(&workload, &facts, &costs, Objective::default(), false).unwrap(); + + let output = plan_to_planner_output(&config, &workload, &solution).unwrap(); + let aggregation = aggregation_for(&output, query); + assert_eq!(aggregation.aggregation_type, AggregationType::DDSketch); + assert_eq!(aggregation.parameters["alpha"], json!(0.02)); + } + + #[test] + fn ddsketch_plan_rejects_a_cost_row_without_alpha() { + let query = "quantile_over_time(0.99, http_requests_total[5m])"; + let config = config(&group(query, 0.99)); + let facts = facts(&config); + let workload = build_milp_workload(&config, &facts, SCRAPE_MS).unwrap(); + let costs = vec![dd_cost(json!({}))]; + let solution = solve_milp(&workload, &facts, &costs, Objective::default(), false).unwrap(); + + assert!(matches!( + plan_to_planner_output(&config, &workload, &solution), + Err(MilpOutputError::MissingParam { + variant, + param: "alpha", + }) if variant == "dd" + )); + } + + #[test] + fn ddsketch_plan_rejects_an_invalid_alpha() { + let query = "quantile_over_time(0.99, http_requests_total[5m])"; + let config = config(&group(query, 0.99)); + let facts = facts(&config); + let workload = build_milp_workload(&config, &facts, SCRAPE_MS).unwrap(); + + for alpha in [0.0, -0.01, 1.0] { + let costs = vec![dd_cost(json!({"alpha": alpha}))]; + let solution = + solve_milp(&workload, &facts, &costs, Objective::default(), false).unwrap(); + assert!(matches!( + plan_to_planner_output(&config, &workload, &solution), + Err(MilpOutputError::InvalidDdsAlpha { .. }) + )); + } + } + #[test] fn window_wider_than_slide_is_sliding() { let query = "quantile_over_time(0.99, http_requests_total[5m])"; diff --git a/asap-planner-rs/src/optimizer/workload_facts.rs b/asap-planner-rs/src/optimizer/workload_facts.rs index a5569926..c332781a 100644 --- a/asap-planner-rs/src/optimizer/workload_facts.rs +++ b/asap-planner-rs/src/optimizer/workload_facts.rs @@ -1,7 +1,7 @@ //! Externally provided workload facts for the MILP planner: per metric, the -//! cardinality of each label set in use. The entry for all of a metric's -//! labels is its series count. Labels come from the workload's `metrics:` -//! hints and the scrape interval from the caller. +//! cardinality of each label set in use and its positive value range. The +//! entry for all of a metric's labels is its series count. Labels come from +//! the workload's `metrics:` hints and the scrape interval from the caller. use std::collections::BTreeMap; use std::path::{Path, PathBuf}; @@ -25,6 +25,8 @@ pub enum WorkloadFactsError { DuplicateMetric(String), #[error("metric {metric:?}: duplicate cardinality for labels {labels:?}")] DuplicateLabels { metric: String, labels: LabelSet }, + #[error("metric {metric:?}: value range ({lo}, {hi}) needs 0 < lo <= hi < inf")] + InvalidValueRange { metric: String, lo: f64, hi: f64 }, #[error("metric {0:?} has facts but no `metrics:` hint giving its labels")] MetricWithoutHint(String), } @@ -39,6 +41,7 @@ struct FactsFile { #[serde(deny_unknown_fields)] struct MetricEntry { metric: String, + value_range: [f64; 2], groups: Vec, } @@ -71,6 +74,14 @@ pub fn parse_workload_facts( let file: FactsFile = serde_yaml::from_str(yaml)?; let mut facts = WorkloadFacts::new(); for entry in file.metrics { + let [lo, hi] = entry.value_range; + if !(lo > 0.0 && hi >= lo && hi.is_finite()) { + return Err(WorkloadFactsError::InvalidValueRange { + metric: entry.metric, + lo, + hi, + }); + } let Some(hint) = hints.iter().find(|h| h.metric == entry.metric) else { return Err(WorkloadFactsError::MetricWithoutHint(entry.metric)); }; @@ -98,8 +109,8 @@ pub fn parse_workload_facts( labels: hint.labels.iter().cloned().collect(), scrape_interval_ms, cardinality, - // ponytail: only sizes DDSketch, which ASAPQuery can't deploy yet. - value_range: None, + value_range: Some((lo, hi)), + data_shape: BTreeMap::new(), }; if facts.insert(entry.metric.clone(), metric_facts).is_some() { return Err(WorkloadFactsError::DuplicateMetric(entry.metric)); @@ -128,6 +139,7 @@ mod tests { let yaml = r#" metrics: - metric: http_requests_total + value_range: [1.0, 1000.0] groups: - labels: [instance, job] cardinality: 1200 @@ -138,13 +150,14 @@ metrics: let m = &facts["http_requests_total"]; assert_eq!(m.labels, labels(&["job", "instance"])); assert_eq!(m.scrape_interval_ms, 15_000); + assert_eq!(m.value_range, Some((1.0, 1000.0))); assert_eq!(m.cardinality[&labels(&["job", "instance"])], 1200); assert_eq!(m.cardinality[&labels(&["job"])], 10); } #[test] fn rejects_metric_without_hint() { - let yaml = "metrics:\n - metric: other\n groups: []\n"; + let yaml = "metrics:\n - metric: other\n value_range: [1.0, 1000.0]\n groups: []\n"; let err = parse_workload_facts(yaml, &hints(), 15_000).unwrap_err(); assert!(matches!(err, WorkloadFactsError::MetricWithoutHint(m) if m == "other")); } @@ -154,8 +167,10 @@ metrics: let yaml = r#" metrics: - metric: http_requests_total + value_range: [1.0, 1000.0] groups: [] - metric: http_requests_total + value_range: [1.0, 1000.0] groups: [] "#; let err = parse_workload_facts(yaml, &hints(), 15_000).unwrap_err(); @@ -167,6 +182,7 @@ metrics: let yaml = r#" metrics: - metric: http_requests_total + value_range: [1.0, 1000.0] groups: - labels: [job, instance] cardinality: 1200 @@ -185,4 +201,27 @@ metrics: WorkloadFactsError::Parse(_) )); } + + #[test] + fn requires_a_positive_value_range() { + let missing = r#" +metrics: + - metric: http_requests_total + groups: [] +"#; + assert!(matches!( + parse_workload_facts(missing, &hints(), 15_000), + Err(WorkloadFactsError::Parse(_)) + )); + + for value_range in ["[0.0, 1000.0]", "[2.0, 1.0]", "[1.0, .inf]"] { + let invalid = format!( + "metrics:\n - metric: http_requests_total\n value_range: {value_range}\n groups: []\n" + ); + assert!(matches!( + parse_workload_facts(&invalid, &hints(), 15_000), + Err(WorkloadFactsError::InvalidValueRange { .. }) + )); + } + } }