Repository navigation
feat(planner): build rqe-optimizer Raqes and facts from the workload config #798
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
b0ce7a7
5561781
a7515f6
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -7,10 +7,13 @@ | |
| use std::path::PathBuf; | ||
|
|
||
| use asap_planner::optimizer::{ | ||
| load_optional_selected_atomic_cost_table, run_greedy_pipeline, AtomicCostTable, LabelSetFacts, | ||
| build_milp_workload, load_flat_atomic_cost_table, load_optional_selected_atomic_cost_table, | ||
| load_workload_facts, run_greedy_pipeline, solve_milp, AtomicCostTable, LabelSetFacts, | ||
| LabelSetFactsError, | ||
| }; | ||
| use asap_planner::ControllerConfig; | ||
| use clap::Parser; | ||
| use rqe_optimizer::milp::Objective; | ||
|
|
||
| #[derive(Parser, Debug)] | ||
| #[command( | ||
|
|
@@ -27,26 +30,57 @@ struct Args { | |
| #[arg(long = "data-ingestion-interval-ms", value_parser = clap::value_parser!(u64).range(1..))] | ||
| data_ingestion_interval_ms: u64, | ||
|
|
||
| /// YAML label-set facts: `series_count` per (metric, spatial filter) and | ||
| /// `cardinality` per (metric, spatial filter, grouping labels). | ||
| #[arg(long = "label-set-facts")] | ||
| label_set_facts: PathBuf, | ||
|
|
||
| /// Path to the versioned atomic-cost document sketch-bench's `atomic-costs` | ||
| /// subcommand exports. Requires --atomic-cost-workload to select exactly | ||
| /// one measured workload profile. Omitted: every | ||
| /// benchmarked-family candidate (CMS/HLL/KLL) is dropped, since there is | ||
| /// no data to cost it at — only trivial accumulators and EXACT remain | ||
| /// selectable. | ||
| #[arg(long = "atomic-costs")] | ||
| /// Greedy only. YAML label-set facts: `series_count` per (metric, spatial | ||
| /// filter) and `cardinality` per (metric, spatial filter, grouping labels). | ||
| #[arg( | ||
| long = "label-set-facts", | ||
| required_unless_present = "milp", | ||
| conflicts_with = "milp" | ||
| )] | ||
| label_set_facts: Option<PathBuf>, | ||
|
|
||
| /// Greedy: the versioned atomic-cost document sketch-bench's | ||
| /// `atomic-costs` subcommand exports; requires --atomic-cost-workload. | ||
| /// Omitted: every benchmarked-family candidate (CMS/HLL/KLL) is dropped, | ||
| /// leaving only trivial accumulators and EXACT. | ||
| /// MILP: the flat cost table `export_rqe_optimizer_costs.sh` writes | ||
| /// (`rqe_atomic_costs.json`); required. | ||
| #[arg(long = "atomic-costs", required_if_eq("milp", "true"))] | ||
| atomic_costs: Option<PathBuf>, | ||
|
|
||
| /// JSON `profiles[].workload` value copied from the sketch-bench atomic-cost | ||
| /// document. This makes the empirical workload profile explicit and avoids | ||
| /// mixing costs from different traces or time windows. | ||
| #[arg(long = "atomic-cost-workload", requires = "atomic_costs")] | ||
| /// Greedy only. JSON `profiles[].workload` value copied from the | ||
| /// sketch-bench atomic-cost document. This makes the empirical workload | ||
| /// profile explicit and avoids mixing costs from different traces or time | ||
| /// windows. | ||
| #[arg( | ||
| long = "atomic-cost-workload", | ||
| requires = "atomic_costs", | ||
| conflicts_with = "milp" | ||
| )] | ||
| atomic_cost_workload: Option<PathBuf>, | ||
|
|
||
| /// Plan with sketch-bench's rqe-optimizer MILP and print the plan; writes | ||
| /// no configs yet. | ||
| #[arg(long)] | ||
| milp: bool, | ||
|
|
||
| /// MILP only. YAML workload facts: per metric, `cardinality` per label | ||
| /// set, including the set of all its labels (the series count). | ||
| #[arg( | ||
| long = "workload-facts", | ||
| required_if_eq("milp", "true"), | ||
| requires = "milp" | ||
| )] | ||
| workload_facts: Option<PathBuf>, | ||
|
|
||
| /// MILP only. Objective weight on CPU-sec/sec. Default: rqe-optimizer's. | ||
| #[arg(long = "w-cpu", requires = "milp", value_parser = parse_weight)] | ||
| w_cpu: Option<f64>, | ||
|
|
||
| /// MILP only. Objective weight on memory GiB. Default: rqe-optimizer's. | ||
| #[arg(long = "w-mem", requires = "milp", value_parser = parse_weight)] | ||
| w_mem: Option<f64>, | ||
|
|
||
| #[arg(short, long, action = clap::ArgAction::Count)] | ||
| verbose: u8, | ||
| } | ||
|
|
@@ -64,7 +98,14 @@ fn main() -> anyhow::Result<()> { | |
|
|
||
| let yaml_str = std::fs::read_to_string(&args.input_config)?; | ||
| let config: ControllerConfig = serde_yaml::from_str(&yaml_str)?; | ||
| let facts = LabelSetFacts::from_path(&args.label_set_facts)?; | ||
| if args.milp { | ||
| return run_milp(&args, &config); | ||
| } | ||
| let facts = LabelSetFacts::from_path( | ||
| args.label_set_facts | ||
| .as_deref() | ||
| .expect("clap requires --label-set-facts without --milp"), | ||
| )?; | ||
|
|
||
| let atomic_cost_table = match load_optional_selected_atomic_cost_table( | ||
| args.atomic_costs.as_deref(), | ||
|
|
@@ -108,3 +149,105 @@ fn main() -> anyhow::Result<()> { | |
|
|
||
| Ok(()) | ||
| } | ||
|
|
||
| /// Objective weights must be finite and non-negative: a negative weight | ||
| /// rewards cost, and NaN poisons every coefficient. | ||
| fn parse_weight(s: &str) -> Result<f64, String> { | ||
|
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. An all-zero objective is accepted.
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Fixed in a7515f6. The CLI rejects |
||
| match s.parse::<f64>() { | ||
| Ok(w) if w.is_finite() && w >= 0.0 => Ok(w), | ||
| Ok(w) => Err(format!("must be finite and >= 0, got {w}")), | ||
| Err(e) => Err(e.to_string()), | ||
| } | ||
| } | ||
|
|
||
| fn run_milp(args: &Args, config: &ControllerConfig) -> anyhow::Result<()> { | ||
|
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Nit: workload facts load against
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Fixed in 5561781. A missing |
||
| config.warn_default_slas(); | ||
| let Some(hints) = config.metrics.as_deref() else { | ||
|
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. nit: this repeats the
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Leaving as is. The CLI guard has to run before |
||
| return Err(LabelSetFactsError::MissingMetricHints.into()); | ||
| }; | ||
| let facts = load_workload_facts( | ||
| args.workload_facts | ||
| .as_deref() | ||
| .expect("clap requires --workload-facts with --milp"), | ||
| hints, | ||
| args.data_ingestion_interval_ms, | ||
| )?; | ||
| let costs = load_flat_atomic_cost_table( | ||
| args.atomic_costs | ||
| .as_deref() | ||
| .expect("clap requires --atomic-costs with --milp"), | ||
| )?; | ||
| let Objective::AUCCost { w_cpu, w_mem } = Objective::default(); | ||
|
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Fixed in 5561781. |
||
| let (w_cpu, w_mem) = (args.w_cpu.unwrap_or(w_cpu), args.w_mem.unwrap_or(w_mem)); | ||
| // All-zero weights make every plan cost 0, so the solver's pick is arbitrary. | ||
| anyhow::ensure!( | ||
| w_cpu > 0.0 || w_mem > 0.0, | ||
| "--w-cpu and --w-mem are both 0; at least one must be positive" | ||
| ); | ||
| let objective = Objective::AUCCost { w_cpu, w_mem }; | ||
| 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]; | ||
| println!( | ||
| " [{d}] {:?} {} config={} metric={} grouping={:?} window={}ms slide={}ms", | ||
| dep.capability, | ||
| dep.config.sketch, | ||
| dep.config.sketch_config, | ||
| dep.metric, | ||
| dep.grouping_labels, | ||
| dep.window_ms, | ||
| dep.slide_ms, | ||
| ); | ||
| } | ||
| println!("\n=== Raqes: {} ===", workload.raqes.len()); | ||
| for ((raqe, &d), latency_ms) in workload | ||
| .raqes | ||
| .iter() | ||
| .zip(&solution.mapping) | ||
| .zip(&solution.plan_cost.query_latency_ms) | ||
| { | ||
| println!(" {} -> [{d}] latency={latency_ms:.3e}ms", raqe.id); | ||
| } | ||
| let cost = &solution.plan_cost; | ||
| println!( | ||
| "\nobjective={:.6e} cpu={:.6e} cpu-sec/sec memory={:.3} MB", | ||
| objective.value(cost), | ||
| cost.cpu_secs_per_sec(), | ||
| cost.memory_bytes() / 1e6, | ||
| ); | ||
| for (phase, c) in [ | ||
| ("ingest", &cost.ingest), | ||
| ("merge", &cost.merge), | ||
| ("query", &cost.query), | ||
| ("storage", &cost.storage), | ||
| ] { | ||
| println!( | ||
| " {phase}: cpu={:.6e} cpu-sec/sec memory={:.3} MB", | ||
| c.cpu_secs_per_sec, | ||
| c.memory_bytes / 1e6 | ||
| ); | ||
| } | ||
| Ok(()) | ||
| } | ||
|
|
||
| #[cfg(test)] | ||
| mod tests { | ||
| use super::parse_weight; | ||
|
|
||
| #[test] | ||
| fn weights_must_be_finite_and_non_negative() { | ||
| assert_eq!(parse_weight("0.5"), Ok(0.5)); | ||
| assert_eq!(parse_weight("0"), Ok(0.0)); | ||
| for bad in ["-1", "NaN", "inf", "x"] { | ||
| assert!(parse_weight(bad).is_err(), "{bad}"); | ||
| } | ||
| } | ||
| } | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -67,7 +67,8 @@ pub fn extract_aqes( | |
| metric_schema: &PromQLSchema, | ||
| scrape_interval_ms: u64, | ||
| ) -> Result<Vec<OptimizerItem>, OptimizerError> { | ||
| let mut acc: HashMap<OptimizerItemKey, (QueryRequirements, Vec<String>, f64)> = HashMap::new(); | ||
| let mut acc: HashMap<OptimizerItemKey, (QueryRequirements, Vec<String>, usize)> = | ||
| HashMap::new(); | ||
|
|
||
| for rqe in rqes { | ||
| if rqe.t_repeat_ms == 0 { | ||
|
|
@@ -82,13 +83,11 @@ pub fn extract_aqes( | |
| match extract_requirements(&leaf, metric_schema, scrape_interval_ms) { | ||
| Ok(req) => { | ||
| let key = OptimizerItemKey::from_rqe(&req, rqe); | ||
| let entry = acc.entry(key).or_insert_with(|| (req, Vec::new(), 0.0)); | ||
| let entry = acc.entry(key).or_insert_with(|| (req, Vec::new(), 0)); | ||
|
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Repeated leaves within one query are counted twice.
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Keeping two Raqes: the engine evaluates each arm separately. |
||
| if !entry.1.contains(&leaf) { | ||
| entry.1.push(leaf); | ||
| } | ||
| // query_frequency_hz must stay in Hz (queries per real second) | ||
| // regardless of t_repeat_ms's internal unit — 1000.0 / ms, not 1.0 / ms. | ||
| entry.2 += 1000.0 / rqe.t_repeat_ms as f64; | ||
| entry.2 += 1; | ||
|
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Correctness: |
||
| } | ||
| Err(reason) => { | ||
| return Err(OptimizerError::UnsupportedLeaf { | ||
|
|
@@ -104,10 +103,12 @@ pub fn extract_aqes( | |
| Ok(acc | ||
| .into_iter() | ||
| .map( | ||
| |(key, (requirements, query_strings, query_frequency_hz))| OptimizerItem { | ||
| |(key, (requirements, query_strings, occurrences))| OptimizerItem { | ||
| requirements, | ||
| query_strings, | ||
| query_frequency_hz, | ||
| // Hz (queries per real second): 1000.0 / ms, not 1.0 / ms. | ||
| query_frequency_hz: occurrences as f64 * 1000.0 / key.t_repeat_ms as f64, | ||
| occurrences, | ||
| t_repeat_ms: key.t_repeat_ms, | ||
| accuracy_sla: f64::from_bits(key.accuracy_sla_bits), | ||
| latency_sla_ms: key.latency_sla_ms_bits.map(f64::from_bits), | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Nit:
--label-set-factshas noconflicts_with = "milp", while--atomic-cost-workloaddoes, so it is silently ignored under--milp.There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Fixed in 5561781.
--label-set-factsnow hasconflicts_with = "milp".