diff --git a/Cargo.lock b/Cargo.lock index 568bf135..e1d15191 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#10415767b3366d5f56ff7bef5182a53cd055b6a7" +source = "git+https://github.com/ProjectASAP/sketch-bench?branch=main#30e725a49c89b1c0e4957132b1c82c9bc1198724" 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#10415767b3366d5f56ff7bef5182a53cd055b6a7" +source = "git+https://github.com/ProjectASAP/sketch-bench?branch=main#30e725a49c89b1c0e4957132b1c82c9bc1198724" 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#10415767b3366d5f56ff7bef5182a53cd055b6a7" +source = "git+https://github.com/ProjectASAP/sketch-bench?branch=main#30e725a49c89b1c0e4957132b1c82c9bc1198724" 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 d9db9193..91e083a2 100644 --- a/asap-planner-rs/src/bin/optimizer_cli.rs +++ b/asap-planner-rs/src/bin/optimizer_cli.rs @@ -3,6 +3,7 @@ use std::path::PathBuf; +use anyhow::Context; use asap_planner::optimizer::{ build_milp_workload, load_flat_atomic_cost_table, load_workload_facts, plan_to_planner_output, reject_avg_queries, solve_milp, MilpError, @@ -10,6 +11,7 @@ use asap_planner::optimizer::{ use asap_planner::ControllerConfig; use clap::Parser; use rqe_optimizer::milp::Objective; +use rqe_optimizer::saturation::SaturationCurves; #[derive(Parser, Debug)] #[command( @@ -26,7 +28,8 @@ struct Args { #[arg(long = "data-ingestion-interval-ms", value_parser = clap::value_parser!(u64).range(1..))] data_ingestion_interval_ms: u64, - /// The flat cost table `export_rqe_optimizer_costs.sh` writes + /// The flat cost table sketch-bench's + /// `study_saturation.py --phase optimizer-cost` writes /// (`rqe_atomic_costs.json`). #[arg(long = "atomic-costs")] atomic_costs: PathBuf, @@ -42,10 +45,17 @@ struct Args { allow_undeployable_families: bool, /// YAML workload facts: per metric, positive `value_range` and - /// `cardinality` per label set, including all labels (the series count). + /// `cardinality` per label set, including all labels (the series count), + /// plus the `shape` of each grouping sketches may serve. #[arg(long = "workload-facts")] workload_facts: PathBuf, + /// sketch-bench's saturation-study directory (`out_grid_1e7_cost/`, + /// `out_1e9/`): sketch accuracy is read off its error-vs-N curves at each + /// grouping's `shape`. + #[arg(long = "saturation-dir")] + saturation_dir: PathBuf, + /// Objective weight on CPU-sec/sec. Default: rqe-optimizer's. #[arg(long = "w-cpu", value_parser = parse_weight)] w_cpu: Option, @@ -90,6 +100,12 @@ fn run_milp(args: &Args, config: &ControllerConfig) -> anyhow::Result<()> { return Err(MilpError::MissingMetricHints.into()); }; let facts = load_workload_facts(&args.workload_facts, hints, args.data_ingestion_interval_ms)?; + let curves = SaturationCurves::load(&args.saturation_dir).with_context(|| { + format!( + "loading saturation curves from --saturation-dir {}", + args.saturation_dir.display() + ) + })?; let costs = load_flat_atomic_cost_table(&args.atomic_costs)?; let Objective::AUCCost { w_cpu, w_mem } = Objective::default(); let (w_cpu, w_mem) = (args.w_cpu.unwrap_or(w_cpu), args.w_mem.unwrap_or(w_mem)); @@ -112,6 +128,7 @@ fn run_milp(args: &Args, config: &ControllerConfig) -> anyhow::Result<()> { &costs, objective, args.allow_undeployable_families, + &|raqe, deployment| curves.accuracy(raqe, deployment, &facts), )?; println!("=== Deployments: {} ===", solution.deployments.len()); diff --git a/asap-planner-rs/src/optimizer/atomic_costs.rs b/asap-planner-rs/src/optimizer/atomic_costs.rs index 6b7f4046..77e4368f 100644 --- a/asap-planner-rs/src/optimizer/atomic_costs.rs +++ b/asap-planner-rs/src/optimizer/atomic_costs.rs @@ -1,12 +1,12 @@ -//! The flat atomic-cost table sketch-bench's `export_rqe_optimizer_costs.sh` -//! writes for the MILP. +//! The flat atomic-cost table sketch-bench's +//! `study_saturation.py --phase optimizer-cost` writes for the MILP. use std::path::Path; pub use rqe_optimizer::{AtomicCostEntry, AtomicCostTable}; -/// Load the flat cost table `export_rqe_optimizer_costs.sh` writes, for the -/// MILP. An invalid row is an error, not dropped. +/// Load the flat cost table `study_saturation.py --phase optimizer-cost` +/// writes, for the MILP. An invalid row is an error, not dropped. pub fn load_flat_atomic_cost_table(path: &Path) -> anyhow::Result { let raw = std::fs::read_to_string(path) .map_err(|e| anyhow::anyhow!("reading cost table {}: {e}", path.display()))?; @@ -37,6 +37,21 @@ fn valid_cost_entry(entry: &AtomicCostEntry) -> bool { .all(|cost| cost.is_finite() && *cost >= 0.0) } +/// A cost row's required `measured_at`, for test fixtures: the cost table's +/// shape (sketch-bench `study_saturation.py` COST_*). Generic so the +/// fixture needn't name `aqpbm_core::MeasuredAt`. +#[cfg(test)] +pub(crate) fn test_measured_at() -> T { + serde_json::from_value(serde_json::json!({ + "items_per_instance": 1_000_000, + "keys_per_instance": 10_000, + "value_range": null, + "merge_operand_items": null, + "distribution": null, + })) + .expect("MeasuredAt fixture") +} + #[cfg(test)] mod tests { use super::*; @@ -51,7 +66,7 @@ mod tests { fn flat_loader_rejects_non_finite_or_negative_costs() { let row = |insert: &str| { format!( - r#"{{"sketch":"hll","sketch_config":null,"mem_bytes_per_instance":1.0,"insert_cpu_secs":{insert},"merge_cpu_secs":1.0,"query_cpu_secs":1.0,"query_accuracy":{{}}}}"# + r#"{{"sketch":"hll","sketch_config":null,"mem_bytes_per_instance":1.0,"insert_cpu_secs":{insert},"merge_cpu_secs":1.0,"query_cpu_secs":1.0,"query_accuracy":{{"relative_error":0.0}},"accuracy_metric":"relative_error","measured_at":{{"items_per_instance":1000000,"keys_per_instance":10000,"value_range":null,"merge_operand_items":null,"distribution":null}}}}"# ) }; let file = tempfile::NamedTempFile::new().unwrap(); @@ -62,16 +77,21 @@ mod tests { assert!(err.to_string().contains("negative"), "{err}"); } - /// Rows with `measured_at` load, and older rows without it still do. + /// Rows name their accuracy metric and where they were measured; a row + /// without either is an old table, rejected rather than guessed. #[test] - fn atomic_cost_entry_accepts_optional_measured_at() { - let base = r#""sketch":"kll-percall","sketch_config":null,"mem_bytes_per_instance":1.0,"insert_cpu_secs":1.0,"merge_cpu_secs":1.0,"query_cpu_secs":1.0,"query_accuracy":{}"#; - let with = format!( - r#"{{{base},"measured_at":{{"items_per_instance":1000000,"keys_per_instance":100000,"value_range":[1.0,100000.0],"merge_operand_items":62500,"distribution":{{"kind":"zipf","skewness":1.1,"population_size":100000,"seed":42}}}}}}"# - ); - let entry: AtomicCostEntry = serde_json::from_str(&with).unwrap(); - assert_eq!(entry.measured_at.unwrap().items_per_instance, 1_000_000); - let without: AtomicCostEntry = serde_json::from_str(&format!("{{{base}}}")).unwrap(); - assert!(without.measured_at.is_none()); + fn atomic_cost_entry_requires_accuracy_metric_and_measured_at() { + let base = r#""sketch":"kll-percall","sketch_config":null,"mem_bytes_per_instance":1.0,"insert_cpu_secs":1.0,"merge_cpu_secs":1.0,"query_cpu_secs":1.0,"query_accuracy":{"mean_rank_err":0.01}"#; + let measured_at = r#""measured_at":{"items_per_instance":1000000,"keys_per_instance":null,"value_range":null,"merge_operand_items":62500,"distribution":{"kind":"pareto","alpha":2.0,"scale":1000.0,"seed":42}}"#; + let full = format!(r#"{{{base},"accuracy_metric":"mean_rank_err",{measured_at}}}"#); + let entry: AtomicCostEntry = serde_json::from_str(&full).unwrap(); + assert_eq!(entry.measured_at.items_per_instance, 1_000_000); + assert_eq!(entry.accuracy(), Some(0.01)); + for partial in [ + format!(r#"{{{base},{measured_at}}}"#), + format!(r#"{{{base},"accuracy_metric":"mean_rank_err"}}"#), + ] { + assert!(serde_json::from_str::(&partial).is_err()); + } } } diff --git a/asap-planner-rs/src/optimizer/milp.rs b/asap-planner-rs/src/optimizer/milp.rs index 9f2d4a9f..c87b144a 100644 --- a/asap-planner-rs/src/optimizer/milp.rs +++ b/asap-planner-rs/src/optimizer/milp.rs @@ -6,7 +6,8 @@ use rqe_optimizer::candidates::build_all_candidates; use rqe_optimizer::enumerate::unservable; use rqe_optimizer::milp::{minimize, MilpSolution, Objective}; use rqe_optimizer::{ - table_accuracy, validate_facts, AtomicCostEntry, Capability, LabelSet, Raqe, WorkloadFacts, + family_properties, validate_facts, Accuracy, AtomicCostEntry, Capability, LabelSet, Raqe, + WorkloadFacts, }; use thiserror::Error; @@ -41,6 +42,17 @@ pub enum MilpError { }, #[error("query {query:?}: topk ranking (by value or by sample count) is unknown")] TopkWeightingUnknown { query: String }, + #[error("topk query {query:?} has no literal k")] + TopkWithoutK { query: String }, + #[error( + "raqe {raqe:?}: sketches serve it, but workload facts give metric {metric:?} \ + grouping {grouping:?} no `shape`" + )] + MissingShape { + raqe: String, + metric: String, + grouping: LabelSet, + }, #[error("no eligible deployment for raqes: {0:?}")] Unservable(Vec), #[error("MILP solve failed: {0}")] @@ -49,20 +61,24 @@ pub enum MilpError { /// The cheapest plan. Only families in sketch-bench's `DEPLOYABLE_FAMILIES` /// are candidates unless `allow_undeployable_families`, whose plan can be -/// studied but not deployed. +/// studied but not deployed. `accuracy` is what a deployment achieves for a +/// Raqe: in production sketch-bench's saturation curves +/// (`SaturationCurves::accuracy`). pub fn solve_milp( workload: &MilpWorkload, facts: &WorkloadFacts, costs: &[AtomicCostEntry], objective: Objective, allow_undeployable_families: bool, + accuracy: &Accuracy, ) -> Result { + require_shapes(workload, facts, costs, allow_undeployable_families)?; let deployments = build_all_candidates( &workload.raqes, costs, facts, allow_undeployable_families, - &table_accuracy, + accuracy, ); tracing::debug!( candidates = deployments.len(), @@ -78,18 +94,12 @@ 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, &table_accuracy); + let missing = unservable(&one_per_item, &deployments, facts, accuracy); if !missing.is_empty() { return Err(MilpError::Unservable(missing)); } - let solution = minimize( - &workload.raqes, - &deployments, - facts, - objective, - &table_accuracy, - ) - .map_err(|e| MilpError::Solver(e.to_string()))?; + let solution = minimize(&workload.raqes, &deployments, facts, objective, 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(), @@ -100,6 +110,36 @@ pub fn solve_milp( Ok(solution) } +/// A Raqe a sketch family with cost rows may serve reads its accuracy off +/// the curves at its grouping's `shape`; without one it could only be +/// reported unservable. Families without rows are left to that path. +fn require_shapes( + workload: &MilpWorkload, + facts: &WorkloadFacts, + costs: &[AtomicCostEntry], + allow_undeployable_families: bool, +) -> Result<(), MilpError> { + for raqe in &workload.raqes { + let sketch_served = raqe + .capability + .candidate_families(allow_undeployable_families) + .any(|family| { + !family_properties(family).exact && costs.iter().any(|row| row.sketch == family) + }); + let has_shape = facts + .get(&raqe.metric) + .is_some_and(|m| m.data_shape.contains_key(&raqe.grouping_labels)); + if sketch_served && !has_shape { + return Err(MilpError::MissingShape { + raqe: raqe.id.clone(), + metric: raqe.metric.clone(), + grouping: raqe.grouping_labels.clone(), + }); + } + } + Ok(()) +} + /// The MILP's view of a workload. #[derive(Debug)] pub struct MilpWorkload { @@ -240,6 +280,21 @@ fn item_to_raqe(item: &OptimizerItem) -> Result { } let accuracy_sla = accuracy_target(capability, item.accuracy_sla); + // A top-k Raqe asks for its query's literal k (the largest, should an + // item carry several), which sizes and prices the heap. + let topk_k = match capability { + Capability::TopKByValue | Capability::TopKByCount => Some( + item.query_strings + .iter() + .map(|q| super::milp_output::topk_k(q)) + .try_fold(0, |max, k| Some(max.max(k?))) + .ok_or_else(|| MilpError::TopkWithoutK { + query: query.clone(), + })?, + ), + _ => None, + }; + // 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 { @@ -261,6 +316,7 @@ fn item_to_raqe(item: &OptimizerItem) -> Result { grouping_labels, accuracy_sla, latency_sla_ms: item.latency_sla_ms, + topk_k, }) } @@ -303,6 +359,8 @@ fn accuracy_target(capability: Capability, accuracy_sla: f64) -> f64 { pub(super) mod tests { use std::collections::BTreeMap; + pub(crate) use rqe_optimizer::table_accuracy; + use super::*; use crate::optimizer::workload_facts::parse_workload_facts; @@ -315,10 +373,13 @@ metrics: groups: - labels: [instance, job] cardinality: 100 + shape: {zipf_s: 1.1, distinct_keys: 10000, tail_index: 2.0} - labels: [job] cardinality: 10 + shape: {zipf_s: 1.1, distinct_keys: 10000, tail_index: 2.0} - labels: [] cardinality: 1 + shape: {zipf_s: 1.1, distinct_keys: 10000, tail_index: 2.0} "#; pub(crate) fn config(groups: &str) -> ControllerConfig { @@ -356,8 +417,8 @@ metrics: merge_cpu_secs: 1e-7, query_cpu_secs: 1e-7, query_accuracy: accuracy.iter().map(|(k, v)| (k.to_string(), *v)).collect(), - merge_accuracy: BTreeMap::new(), - measured_at: None, + accuracy_metric: rqe_optimizer::accuracy_key(sketch).0.to_string(), + measured_at: crate::optimizer::atomic_costs::test_measured_at(), } } @@ -365,13 +426,22 @@ metrics: vec![ 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)], - ), + topk_cost(32), + topk_cost(2048), ] } + /// A top-k row at one heap size; the cost table measures several + /// (sketch-bench `TOPK_HEAPS`), and a deployment needs `m · k`. + fn topk_cost(heap: u64) -> AtomicCostEntry { + let mut row = cost( + "cms-heap-topk-fastpath-vector2d", + &[("precision_at_k", 0.995)], + ); + row.sketch_config["params"] = serde_json::json!({"heap": heap}); + row + } + #[test] fn repeated_query_becomes_one_raqe_per_occurrence() { let query = "sum by (job) (http_requests_total)"; @@ -445,6 +515,23 @@ metrics: ); } + #[test] + fn topk_carries_its_literal_k_and_needs_one() { + let w = workload(&group( + "topk(7, sum_over_time(http_requests_total[1m]))", + 0.99, + )); + assert_eq!(w.raqes[0].topk_k, Some(7)); + let mut item = w.items[0].clone(); + item.query_strings = vec!["topk(scalar(up), http_requests_total)".into()]; + assert!(matches!( + item_to_raqe(&item), + Err(MilpError::TopkWithoutK { .. }) + )); + let sum = workload(&group("sum by (job) (http_requests_total)", 0.99)); + assert_eq!(sum.raqes[0].topk_k, None); + } + #[test] fn topk_with_unknown_weighting_is_an_error() { let w = workload(&group("topk(5, http_requests_total)", 0.99)); @@ -517,7 +604,15 @@ 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 solution = solve_milp(&w, &facts, &costs(), Objective::default(), false).unwrap(); + let solution = solve_milp( + &w, + &facts, + &costs(), + Objective::default(), + false, + &table_accuracy, + ) + .unwrap(); assert_eq!(solution.deployments.len(), 1); assert_eq!(solution.raqes[0].deployment, solution.raqes[1].deployment); assert_eq!( @@ -526,6 +621,62 @@ metrics: ); } + /// A sketch-served grouping without a `shape` fails by name, before + /// solving; an exact-only one needs none. + #[test] + fn a_sketch_served_grouping_without_a_shape_is_named() { + let quantile = "quantile_over_time(0.99, http_requests_total[5m])"; + let sum = "sum by (job) (http_requests_total)"; + let both = config(&(group(quantile, 0.99) + &group(sum, 0.99))); + let mut facts = facts(&both); + facts + .get_mut("http_requests_total") + .unwrap() + .data_shape + .clear(); + let w = build_milp_workload(&both, &facts, SCRAPE_MS).unwrap(); + let err = solve_milp( + &w, + &facts, + &costs(), + Objective::default(), + false, + &table_accuracy, + ) + .unwrap_err(); + assert!( + matches!(&err, MilpError::MissingShape { raqe, grouping, .. } + if raqe.contains(quantile) && *grouping == labels(&["instance", "job"])), + "{err:?}" + ); + // No sketch rows for quantiles: left to the Unservable path. + let exact_rows: Vec<_> = costs() + .into_iter() + .filter(|row| row.sketch == "exact-sum") + .collect(); + let err = solve_milp( + &w, + &facts, + &exact_rows, + Objective::default(), + false, + &table_accuracy, + ) + .unwrap_err(); + assert!(matches!(err, MilpError::Unservable(_)), "{err:?}"); + let exact_only = config(&group(sum, 0.99)); + let w = build_milp_workload(&exact_only, &facts, SCRAPE_MS).unwrap(); + assert!(solve_milp( + &w, + &facts, + &costs(), + Objective::default(), + false, + &table_accuracy + ) + .is_ok()); + } + #[test] fn solve_rejects_accuracy_no_sketch_meets() { // kll's max rank error 0.005 misses a 0.999 requirement (0.001). @@ -535,7 +686,15 @@ metrics: )); let facts = facts(&config); let w = build_milp_workload(&config, &facts, SCRAPE_MS).unwrap(); - let err = solve_milp(&w, &facts, &costs(), Objective::default(), false).unwrap_err(); + let err = solve_milp( + &w, + &facts, + &costs(), + Objective::default(), + false, + &table_accuracy, + ) + .unwrap_err(); assert!(matches!(err, MilpError::Unservable(ids) if ids.len() == 1)); } @@ -546,7 +705,15 @@ metrics: 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(), false).unwrap_err(); + let err = solve_milp( + &w, + &facts, + &costs(), + Objective::default(), + false, + &table_accuracy, + ) + .unwrap_err(); assert!(matches!(err, MilpError::Unservable(ids) if ids == [w.raqes[0].id.clone()])); } @@ -558,7 +725,15 @@ metrics: let config = config(&groups); let facts = facts(&config); let w = build_milp_workload(&config, &facts, SCRAPE_MS).unwrap(); - let solution = solve_milp(&w, &facts, &costs(), Objective::default(), false).unwrap(); + let solution = solve_milp( + &w, + &facts, + &costs(), + Objective::default(), + false, + &table_accuracy, + ) + .unwrap(); let chosen: BTreeMap = w .raqes .iter() @@ -586,7 +761,15 @@ metrics: 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(), false).unwrap(); + let solution = solve_milp( + &w, + &facts, + &costs(), + Objective::default(), + false, + &table_accuracy, + ) + .unwrap(); assert_eq!(solution.deployments.len(), 2); } diff --git a/asap-planner-rs/src/optimizer/milp_output.rs b/asap-planner-rs/src/optimizer/milp_output.rs index fbc517e5..715ea73b 100644 --- a/asap-planner-rs/src/optimizer/milp_output.rs +++ b/asap-planner-rs/src/optimizer/milp_output.rs @@ -41,8 +41,6 @@ pub enum MilpOutputError { }, #[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")] QueryOnTwoDeployments(String), #[error(transparent)] @@ -184,7 +182,7 @@ fn aggregation_config( vec![ ("depth", Value::from(integer_param("rows")?)), ("width", Value::from(integer_param("cols")?)), - ("heapsize", Value::from(max_topk_k(items)?)), + ("heapsize", Value::from(heap_size(deployment)?)), ], ), (_, capability) => { @@ -242,17 +240,18 @@ fn missing_param(variant: &str, param: &'static str) -> MilpOutputError { } } -/// One heap serves every topk query on the deployment, so it is sized for -/// the largest k. -fn max_topk_k(items: &[&OptimizerItem]) -> Result { - items - .iter() - .flat_map(|item| &item.query_strings) - .map(|query| topk_k(query).ok_or_else(|| MilpOutputError::TopkWithoutK(query.clone()))) - .try_fold(0, |max, k| Ok(max.max(k?))) +/// The heap the plan priced: sketch-bench's `heap` param, sized `m · k` for +/// the windows the query it was built for merges (or the smallest measured +/// heap above that). Another query sharing it may need more; the plan then +/// priced that query's accuracy as a lossy merge, so the heap is emitted as +/// priced. The engine takes each query's `k` from the query itself. +fn heap_size(deployment: &Deployment) -> Result { + rqe_optimizer::heap_capacity(&deployment.config) + .ok_or_else(|| missing_param(&deployment.config.sketch, "heap")) } -fn topk_k(query: &str) -> Option { +/// A `topk` query's literal `k`; `None` for anything else. +pub(crate) fn topk_k(query: &str) -> Option { let Ok(Expr::Aggregate(aggregate)) = promql_parser::parser::parse(query) else { return None; }; @@ -276,7 +275,7 @@ mod tests { use serde_json::json; use super::*; - use crate::optimizer::milp::tests::{config, costs, facts, group, SCRAPE_MS}; + use crate::optimizer::milp::tests::{config, costs, facts, group, table_accuracy, SCRAPE_MS}; use crate::optimizer::milp::{build_milp_workload, solve_milp}; fn plan(groups: &str) -> (ControllerConfig, MilpWorkload, MilpSolution) { @@ -287,11 +286,23 @@ mod tests { for cost in &mut costs { cost.sketch_config["params"] = match cost.sketch.as_str() { "kll-percall" => json!({"k": 200}), - "cms-heap-topk-fastpath-vector2d" => json!({"rows": 3, "cols": 2048}), + "cms-heap-topk-fastpath-vector2d" => json!({ + "rows": 3, + "cols": 2048, + "heap": cost.sketch_config["params"]["heap"], + }), _ => json!({}), }; } - let solution = solve_milp(&workload, &facts, &costs, Objective::default(), false).unwrap(); + let solution = solve_milp( + &workload, + &facts, + &costs, + Objective::default(), + false, + &table_accuracy, + ) + .unwrap(); (config, workload, solution) } @@ -322,8 +333,8 @@ mod tests { 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, + accuracy_metric: "mean_relative_value_error".into(), + measured_at: crate::optimizer::atomic_costs::test_measured_at(), } } @@ -368,7 +379,10 @@ mod tests { assert_eq!(a.aggregation_sub_type, "sum"); assert_eq!(a.parameters["depth"], json!(3)); assert_eq!(a.parameters["width"], json!(2048)); - assert_eq!(a.parameters["heapsize"], json!(5)); + assert_eq!( + a.parameters["heapsize"], + json!(planned_heap(&workload, &solution, topk_value)) + ); assert_eq!( aggregation_for(&output, topk_count).aggregation_sub_type, "count" @@ -382,21 +396,91 @@ mod tests { } } + /// The heap `query`'s deployment needs for it: `m · k` for the windows + /// its lookback merges and its own `k` (`Deployment::heap_needed`). + fn heap_needed(workload: &MilpWorkload, solution: &MilpSolution, query: &str) -> u64 { + let i = workload + .raqes + .iter() + .position(|r| r.id.contains(query)) + .unwrap_or_else(|| panic!("no raqe for {query}")); + let raqe = &workload.raqes[i]; + let deployment = &solution.deployments[solution.raqes[i].deployment].deployment; + deployment + .heap_needed(raqe.lookback_ms, raqe.topk_k()) + .expect("the lookback is whole windows") + } + + /// The heap the plan priced for `query`'s deployment: its `heap` param. + fn planned_heap(workload: &MilpWorkload, solution: &MilpSolution, query: &str) -> u64 { + let i = workload + .raqes + .iter() + .position(|r| r.id.contains(query)) + .unwrap_or_else(|| panic!("no raqe for {query}")); + let deployment = &solution.deployments[solution.raqes[i].deployment].deployment; + rqe_optimizer::heap_capacity(&deployment.config).expect("a heap top-k deployment") + } + + /// Each top-k query plans with its own literal k, and the engine's heap is + /// the one the plan priced: `m · k` for the windows it merges, or the + /// smallest measured heap (32) that already holds it. + #[test] + fn topk_plans_with_its_own_k_and_heap_m_times_k() { + for (k, measured_floor) in [(10, true), (100, false)] { + let query = format!("topk({k}, sum_over_time(http_requests_total[1m]))"); + let (config, workload, solution) = plan(&group(&query, 0.99)); + assert_eq!(workload.raqes[0].topk_k, Some(k)); + let deployment = &solution.deployments[0].deployment; + assert_eq!(rqe_optimizer::answered_k(&deployment.config), k); + let needed = heap_needed(&workload, &solution, &query); + let heap = planned_heap(&workload, &solution, &query); + if measured_floor { + assert_eq!(heap, rqe_optimizer::TOPK_K, "k = {k}"); + assert!(heap >= needed); + } else { + assert_eq!(heap, needed, "k = {k}: m · k"); + } + let output = plan_to_planner_output(&config, &workload, &solution).unwrap(); + assert_eq!( + aggregation_for(&output, &query).parameters["heapsize"], + json!(heap), + "k = {k}" + ); + } + } + + /// Queries sharing one heap get the heap the plan priced, which holds the + /// largest k among them. #[test] - fn shared_heap_is_sized_for_the_largest_k() { + fn a_shared_heap_is_the_planned_heap() { let small = "topk(5, sum_over_time(http_requests_total[1m]))"; let large = "topk(10, sum_over_time(http_requests_total[1m]))"; let (config, workload, solution) = plan(&(group(small, 0.99) + &group(large, 0.99))); assert_eq!(solution.deployments.len(), 1); + let heap = planned_heap(&workload, &solution, large); + assert!(heap >= heap_needed(&workload, &solution, large)); let output = plan_to_planner_output(&config, &workload, &solution).unwrap(); for query in [small, large] { assert_eq!( aggregation_for(&output, query).parameters["heapsize"], - json!(10) + json!(heap) ); } } + #[test] + fn a_heap_family_without_a_heap_is_a_missing_param() { + let query = "topk(5, sum_over_time(http_requests_total[1m]))"; + let (_, _, solution) = plan(&group(query, 0.99)); + let mut deployment = solution.deployments[0].deployment.clone(); + deployment.config.sketch = "cms-fastpath-vector2d".into(); + assert!(matches!( + heap_size(&deployment), + Err(MilpOutputError::MissingParam { param: "heap", .. }) + )); + } + #[test] fn ddsketch_plan_preserves_the_cost_row_alpha() { let query = "quantile_over_time(0.99, http_requests_total[5m])"; @@ -404,7 +488,15 @@ mod tests { 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 solution = solve_milp( + &workload, + &facts, + &costs, + Objective::default(), + false, + &table_accuracy, + ) + .unwrap(); let output = plan_to_planner_output(&config, &workload, &solution).unwrap(); let aggregation = aggregation_for(&output, query); @@ -419,7 +511,15 @@ mod tests { 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(); + let solution = solve_milp( + &workload, + &facts, + &costs, + Objective::default(), + false, + &table_accuracy, + ) + .unwrap(); assert!(matches!( plan_to_planner_output(&config, &workload, &solution), @@ -439,8 +539,15 @@ mod tests { 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(); + let solution = solve_milp( + &workload, + &facts, + &costs, + Objective::default(), + false, + &table_accuracy, + ) + .unwrap(); assert!(matches!( plan_to_planner_output(&config, &workload, &solution), Err(MilpOutputError::InvalidDdsAlpha { .. }) diff --git a/asap-planner-rs/src/optimizer/workload_facts.rs b/asap-planner-rs/src/optimizer/workload_facts.rs index c332781a..6e909803 100644 --- a/asap-planner-rs/src/optimizer/workload_facts.rs +++ b/asap-planner-rs/src/optimizer/workload_facts.rs @@ -1,11 +1,14 @@ //! Externally provided workload facts for the MILP planner: per metric, the -//! cardinality of each label set in use and its positive value range. The +//! cardinality of each label set in use, its positive value range, and, for +//! a grouping sketches serve, its fitted data shape (which keys +//! sketch-bench's saturation curves). 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}; +use rqe_optimizer::saturation::DataShape; use rqe_optimizer::{LabelSet, MetricFacts, Millis, WorkloadFacts}; use serde::Deserialize; use thiserror::Error; @@ -27,6 +30,11 @@ pub enum WorkloadFactsError { 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 {metric:?} labels {labels:?}: shape needs finite zipf_s >= 0, \ + distinct_keys >= 1 and tail_index > 0" + )] + InvalidShape { metric: String, labels: LabelSet }, #[error("metric {0:?} has facts but no `metrics:` hint giving its labels")] MetricWithoutHint(String), } @@ -50,6 +58,22 @@ struct MetricEntry { struct GroupEntry { labels: Vec, cardinality: u64, + /// The data one group's sketch sees, fitted by the caller as the worst + /// case over windows and groups. A grouping a sketch family with cost + /// rows would serve must carry one, or solving fails with `MissingShape`. + #[serde(default)] + shape: Option, +} + +/// [`DataShape`]'s fields: Zipf skew and distinct keys per group per window +/// (frequency, top-k, cardinality), and the values' Pareto tail index +/// (quantiles). +#[derive(Debug, Deserialize)] +#[serde(deny_unknown_fields)] +struct ShapeEntry { + zipf_s: f64, + distinct_keys: f64, + tail_index: f64, } pub fn load_workload_facts( @@ -86,6 +110,7 @@ pub fn parse_workload_facts( return Err(WorkloadFactsError::MetricWithoutHint(entry.metric)); }; let mut cardinality = BTreeMap::new(); + let mut data_shape = BTreeMap::new(); for group in entry.groups { let labels: LabelSet = group.labels.into_iter().collect(); if cardinality @@ -97,12 +122,37 @@ pub fn parse_workload_facts( labels, }); } + if let Some(shape) = group.shape { + let ShapeEntry { + zipf_s, + distinct_keys, + tail_index, + } = shape; + let finite = [zipf_s, distinct_keys, tail_index] + .iter() + .all(|x| x.is_finite()); + if !(finite && zipf_s >= 0.0 && distinct_keys >= 1.0 && tail_index > 0.0) { + return Err(WorkloadFactsError::InvalidShape { + metric: entry.metric, + labels, + }); + } + data_shape.insert( + labels, + DataShape { + zipf_s, + distinct_keys, + tail_index, + }, + ); + } } tracing::debug!( metric = %entry.metric, labels = ?hint.labels, scrape_interval_ms, cardinality = ?cardinality, + data_shape = ?data_shape, "workload facts: metric" ); let metric_facts = MetricFacts { @@ -110,7 +160,7 @@ pub fn parse_workload_facts( scrape_interval_ms, cardinality, value_range: Some((lo, hi)), - data_shape: BTreeMap::new(), + data_shape, }; if facts.insert(entry.metric.clone(), metric_facts).is_some() { return Err(WorkloadFactsError::DuplicateMetric(entry.metric)); @@ -134,6 +184,55 @@ mod tests { names.iter().map(|s| s.to_string()).collect() } + #[test] + fn a_group_with_a_shape_keys_the_curves() { + let yaml = r#" +metrics: + - metric: http_requests_total + value_range: [1.0, 1000.0] + groups: + - labels: [instance, job] + cardinality: 1200 + - labels: [job] + cardinality: 10 + shape: {zipf_s: 1.1, distinct_keys: 120, tail_index: 2.0} +"#; + let facts = parse_workload_facts(yaml, &hints(), 15_000).unwrap(); + let shapes = &facts["http_requests_total"].data_shape; + assert_eq!(shapes.len(), 1); + assert_eq!( + shapes[&labels(&["job"])], + DataShape { + zipf_s: 1.1, + distinct_keys: 120.0, + tail_index: 2.0, + } + ); + } + + #[test] + fn rejects_an_invalid_shape() { + for shape in [ + "{zipf_s: -0.1, distinct_keys: 120, tail_index: 2.0}", + "{zipf_s: 1.1, distinct_keys: 0.5, tail_index: 2.0}", + "{zipf_s: 1.1, distinct_keys: 120, tail_index: 0}", + "{zipf_s: .inf, distinct_keys: 120, tail_index: 2.0}", + ] { + let yaml = format!( + "metrics:\n - metric: http_requests_total\n value_range: [1.0, 1000.0]\n \ + groups:\n - labels: [job]\n cardinality: 10\n shape: {shape}\n" + ); + assert!( + matches!( + parse_workload_facts(&yaml, &hints(), 15_000), + Err(WorkloadFactsError::InvalidShape { ref metric, labels: ref got }) + if metric == "http_requests_total" && *got == labels(&["job"]) + ), + "{shape}" + ); + } + } + #[test] fn parses_cardinalities_and_takes_labels_from_hints() { let yaml = r#"