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.

21 changes: 19 additions & 2 deletions asap-planner-rs/src/bin/optimizer_cli.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,13 +3,15 @@

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,
};
use asap_planner::ControllerConfig;
use clap::Parser;
use rqe_optimizer::milp::Objective;
use rqe_optimizer::saturation::SaturationCurves;

#[derive(Parser, Debug)]
#[command(
Expand All @@ -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,
Expand All @@ -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<f64>,
Expand Down Expand Up @@ -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));
Expand All @@ -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());
Expand Down
50 changes: 35 additions & 15 deletions asap-planner-rs/src/optimizer/atomic_costs.rs
Original file line number Diff line number Diff line change
@@ -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<AtomicCostTable> {
let raw = std::fs::read_to_string(path)
.map_err(|e| anyhow::anyhow!("reading cost table {}: {e}", path.display()))?;
Expand Down Expand Up @@ -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::de::DeserializeOwned>() -> 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::*;
Expand All @@ -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();
Expand All @@ -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::<AtomicCostEntry>(&partial).is_err());
}
}
}
Loading
Loading