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
7 changes: 6 additions & 1 deletion docs/rqe_optimizer_TODO.md
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,11 @@ implemented, what is next, and what is intentionally deferred.
1-hour, 6-hour, and 1-day quantile RQEs to exercise those constraints.
- A regression test confirms the MILP matches exhaustive minimum-TCO search on
a tiny workload with shared deployment activation costs.
- Retained memory (`(x + max S) / y` instances per active deployment) is
scored, and `milp::minimize_cost` minimizes the hourly price on one EC2
machine family in fractional instances. Prices are a committed snapshot from
`scripts/fetch_ec2_pricing.py`; `small_problem --milp --machine-family NAME`
uses it.

## Next

Expand All @@ -43,7 +48,7 @@ implemented, what is next, and what is intentionally deferred.
## Explicitly deferred

- Query-result sharing across RQEs.
- Merge buffers, retained-storage capacity, and concurrent-query memory.
- Merge buffers and concurrent-query memory.
- RQE churn, replanning, and migration cost.
- Precomputed rollups; v1 merges selected base instances at query time.
- A policy for choosing one mapping from the reported Pareto frontier.
Expand Down
37 changes: 30 additions & 7 deletions docs/rqe_sketch_deployment_v1.md
Original file line number Diff line number Diff line change
Expand Up @@ -130,14 +130,14 @@ explicit on the RQE; it is never inferred from a metric name. A missing metric
does not pass.

Before mapping, prune a candidate only if another candidate can serve every
RQE it can and is no worse in query memory, ingest CPU, and latency for each
such RQE. This is safe because any mapping using the removed candidate can
RQE it can and is no worse in query memory, ingest CPU, latency, and retained
memory for each such RQE. This is safe because any mapping using the removed candidate can
substitute the remaining one without weakening a modeled objective.

Historical instances must be retained long enough to answer the RQEs assigned
to a deployment. Retention is an execution/storage detail in v1, not a
candidate parameter or a scored objective. After selecting a mapping, retain
the history needed by the largest assigned query window.
to a deployment. Retention is not a candidate parameter: a deployment retains
the history needed by the largest assigned query window. Its memory is scored
as retained memory (below) and priced by `minimize_cost`.

## Workload mapping

Expand Down Expand Up @@ -185,6 +185,20 @@ TCO_cpu = Σ_D ingest_cpu_D × u_D
latency_i = Σ_{D: (i,D)∈E} latency_{i,D} × z_{i,D}
```

`minimize_cost` instead prices the plan on one EC2 machine family `f` (vCPUs
`vcpu_f`, memory `gib_f`, hourly price `price_f`) in fractional instances
`n_f`, with a continuous `R_D ≥ 0` for each candidate's retained GiB:

```text
minimize price_f × n_f
R_D ≥ retained_gib_{i,D} × z_{i,D} for (i, D) ∈ E
vcpu_f × n_f ≥ TCO_cpu
gib_f × n_f ≥ Σ_D R_D
```

Prices come from `rqe-optimizer/data/ec2-pricing-<date>.json`, written by
`scripts/fetch_ec2_pricing.py`. Disk is not priced.

One solve needs a scalar objective, such as minimum `TCO_cpu` subject to
`M ≤ memory_budget` and optional `latency_i ≤ latency_budget_i`. To sample the
Pareto frontier, repeat solves over memory and latency budgets, or use a
Expand Down Expand Up @@ -216,6 +230,16 @@ The mapping-level memory objective is the worst individual query:
peak_query_memory = max_i query_memory_i
```

### Retained memory

A deployment holds `x / y` open instances plus the closed instances still
inside the longest lookback it serves:

```text
retained_memory = Σ_{D: u_D=1} card(D.labels) × mem_bytes(D.configuration)
× (D.x + max_{i: D(i)=D} S_i) / D.y
```

### CPU

Ingest CPU is paid once for every active deployment:
Expand Down Expand Up @@ -281,8 +305,7 @@ breakdown of `TCO_cpu`, together with the selected deployment mapping.
- **Latency SLAs:** v1 reports per-RQE latency but does not reject a mapping
for exceeding a target. Add optional per-RQE maximum latency as a hard
constraint when workloads supply targets.
- **Memory model:** merge buffers, retained-storage capacity, and concurrent
queries are not modeled.
- **Memory model:** merge buffers and concurrent queries are not modeled.
- **Static planning:** no RQE churn, replanning, or migration cost.
- **Rollups:** v1 does not precompute merged rollups. Queries merge their
selected base instances when they run.
Expand Down
34 changes: 34 additions & 0 deletions rqe-optimizer/data/ec2-pricing-2026-10-04.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
{
"source": "https://b0.p.awsstatic.com/pricing/2.0/meteredUnitMaps/ec2/USD/current/ec2-ondemand-without-sec-sel/US%20East%20(N.%20Virginia)/Linux/index.json",
"region": "US East (N. Virginia)",
"operating_system": "Linux",
"pricing": "on-demand",
"publication_date": "2026-09-25T17:45:21Z",
"fetched": "2026-10-04",
"families": [
{
"family": "compute_optimized",
"instance_type": "c7i.xlarge",
"vcpu": 4,
"memory_gib": 8.0,
"usd_per_hour": 0.1785,
"rate_code": "D8WVCVBE32N227KA.JRTCKXETXF.6YS6EN2CT7"
},
{
"family": "general_purpose",
"instance_type": "m7i.xlarge",
"vcpu": 4,
"memory_gib": 16.0,
"usd_per_hour": 0.2016,
"rate_code": "W8MD4QK5G3J3QYWR.JRTCKXETXF.6YS6EN2CT7"
},
{
"family": "memory_optimized",
"instance_type": "r7i.xlarge",
"vcpu": 4,
"memory_gib": 32.0,
"usd_per_hour": 0.2646,
"rate_code": "PRGSWU6294SJ3Z93.JRTCKXETXF.6YS6EN2CT7"
}
]
}
55 changes: 41 additions & 14 deletions rqe-optimizer/examples/small_problem.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,9 @@
//! mappings by default; pass `--progress-every N` to change that interval or
//! `--print-first N` to display example mappings. `--milp` solves the
//! minimum-TCO model without enumerating mappings. Repeat
//! `--latency-limit RQE_ID=SECONDS` to impose MILP latency bounds.
//! `--latency-limit RQE_ID=SECONDS` to impose MILP latency bounds. Add
//! `--machine-family NAME` (`compute_optimized`, `general_purpose`,
//! `memory_optimized`) to minimize that EC2 family's hourly price instead.
//! `--sample-mappings N` prints N feasible mappings and exits.

use std::collections::BTreeSet;
Expand All @@ -32,8 +34,8 @@ use rqe_optimizer::candidates::{
build_all_candidates, build_all_candidates_unpruned, eligible_deployments_for,
};
use rqe_optimizer::enumerate::{brute_force, for_each_mapping, for_each_mapping_while, unservable};
use rqe_optimizer::milp::{minimize_tco, MilpBounds};
use rqe_optimizer::objectives::score;
use rqe_optimizer::milp::{minimize_cost, minimize_tco, MilpBounds};
use rqe_optimizer::objectives::{score, MachineFamily};
use rqe_optimizer::pareto::{pareto_front, ParetoFront};
use rqe_optimizer::{
AccuracyDirection, AtomicCostTable, Capability, LabelSet, LabelSetInfo, LabelSetTable, Rqe,
Expand Down Expand Up @@ -61,6 +63,22 @@ const CARDINALITY_ERR: &str = "relative_error";
const TOPK_PRECISION: &str = "precision_at_k";

const COST_TABLE_PATH: &str = "out/rqe_atomic_costs.json";
const EC2_PRICING: &str = include_str!("../data/ec2-pricing-2026-10-04.json");

fn machine_family() -> Option<MachineFamily> {
let args: Vec<_> = std::env::args().collect();
let name = args
.windows(2)
.find(|pair| pair[0] == "--machine-family")
.map(|pair| pair[1].clone())?;
let families = MachineFamily::from_pricing_json(EC2_PRICING).expect("committed snapshot");
Some(
families
.into_iter()
.find(|family| family.family == name)
.unwrap_or_else(|| panic!("unknown --machine-family: {name}")),
)
}

fn label_set(names: &[&str]) -> LabelSet {
names.iter().map(|s| s.to_string()).collect()
Expand Down Expand Up @@ -329,19 +347,28 @@ fn main() {
}

if std::env::args().any(|arg| arg == "--milp") {
let latency_bounds = latency_bounds(&rqes);
let solution = minimize_tco(
&rqes,
&deployments,
&label_sets,
&MilpBounds {
max_peak_query_memory_bytes: None,
max_query_latency_secs: latency_bounds,
},
)
let bounds = MilpBounds {
max_peak_query_memory_bytes: None,
max_query_latency_secs: latency_bounds(&rqes),
};
let family = machine_family();
let solution = match &family {
Some(family) => minimize_cost(&rqes, &deployments, &label_sets, &bounds, family),
None => minimize_tco(&rqes, &deployments, &label_sets, &bounds),
}
.expect("small_problem MILP should be feasible");
if let Some(family) = &family {
println!(
"MILP minimum-cost solution on {}: ${:.4}/hour, {:.4} instances, \
retained_mem={:.0}MB",
family.family,
family.usd_per_hour(&solution.objectives),
family.instances(&solution.objectives),
solution.objectives.retained_memory_bytes / 1e6,
);
}
println!(
"MILP minimum-TCO solution: peak_query_mem={:.0}MB, ingest={:.3e}, \
"MILP solution: peak_query_mem={:.0}MB, ingest={:.3e}, \
merge={:.3e}, query={:.3e}, total={:.3e} cpu-sec/sec",
solution.objectives.peak_query_memory_bytes / 1e6,
solution.objectives.ingest_cpu_secs_per_sec,
Expand Down
53 changes: 49 additions & 4 deletions rqe-optimizer/src/candidates.rs
Original file line number Diff line number Diff line change
Expand Up @@ -101,8 +101,8 @@ pub fn build_all_candidates_unpruned(rqes: &[Rqe], costs: &[AtomicCostEntry]) ->
/// This comparison is deliberately local to a capability/label-set group.
/// Within such a group, label cardinality and arrival rate are common
/// multipliers, so comparing per-instance query memory and
/// `active_instances * insert_cost` is sufficient. Query latency is checked
/// for each RQE the dominated candidate can serve.
/// `active_instances * insert_cost` is sufficient. Query latency and retained
/// memory are checked for each RQE the dominated candidate can serve.
pub fn prune_dominated_candidates(rqes: &[Rqe], candidates: Vec<Deployment>) -> Vec<Deployment> {
let eligibility: Vec<Vec<bool>> = candidates
.iter()
Expand Down Expand Up @@ -166,8 +166,10 @@ fn candidate_dominates(
.enumerate()
.all(|(rqe_index, &serves)| {
!serves
|| query_latency(replacement, &rqes[rqe_index])
|| (query_latency(replacement, &rqes[rqe_index])
<= query_latency(original, &rqes[rqe_index])
&& retained_memory(replacement, &rqes[rqe_index])
<= retained_memory(original, &rqes[rqe_index]))
})
}

Expand All @@ -193,8 +195,10 @@ fn strictly_better(
.enumerate()
.any(|(rqe_index, &serves)| {
serves
&& query_latency(replacement, &rqes[rqe_index])
&& (query_latency(replacement, &rqes[rqe_index])
< query_latency(original, &rqes[rqe_index])
|| retained_memory(replacement, &rqes[rqe_index])
< retained_memory(original, &rqes[rqe_index]))
})
}

Expand All @@ -206,6 +210,15 @@ fn query_latency(deployment: &Deployment, rqe: &Rqe) -> f64 {
+ instances.saturating_sub(1) as f64 * deployment.config.merge_cpu_secs
}

/// Retained bytes per label group when serving `rqe`. Label cardinality is a
/// common multiplier within a group, so it is left out.
fn retained_memory(deployment: &Deployment, rqe: &Rqe) -> f64 {
deployment
.retained_instance_count(rqe.lookback_secs)
.expect("candidate coverage only contains exactly tiled RQEs") as f64
* deployment.config.mem_bytes_per_instance
}

pub fn eligible_deployments_for(r: &Rqe, deployments: &[Deployment]) -> Vec<usize> {
deployments
.iter()
Expand Down Expand Up @@ -334,4 +347,36 @@ mod tests {

assert_eq!(retained, vec![large_window]);
}

#[test]
fn retains_candidate_with_lower_retained_memory() {
let r = rqe("r", 60, 60);
// Holds (60 + 60) / 60 = 2 instances of 10 bytes.
let whole_window = Deployment {
capability: Capability::Freq,
labels: LabelSet::new(),
config: AtomicCostEntry {
mem_bytes_per_instance: 10.0,
..cost()
},
window_secs: 60,
slide_secs: 60,
};
// Smaller, faster sketches, but (20 + 60) / 20 = 4 of them: 24 bytes.
let panes = Deployment {
config: AtomicCostEntry {
mem_bytes_per_instance: 6.0,
query_cpu_secs: 0.5,
merge_cpu_secs: 0.1,
..cost()
},
window_secs: 20,
slide_secs: 20,
..whole_window.clone()
};

let retained = prune_dominated_candidates(&[r], vec![whole_window.clone(), panes.clone()]);

assert_eq!(retained, vec![whole_window, panes]);
}
}
8 changes: 8 additions & 0 deletions rqe-optimizer/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -152,6 +152,14 @@ impl Deployment {
}
Some(lookback_secs / self.window_secs)
}

/// Instances held to serve a lookback: `x / y` open ones plus the closed
/// ones still inside the lookback, `(x + S) / y`.
pub fn retained_instance_count(&self, lookback_secs: Seconds) -> Option<u64> {
self.query_instance_count(lookback_secs)?;
self.active_instance_count()?;
Some((self.window_secs + lookback_secs) / self.slide_secs)
}
}

/// A full mapping: one deployment index (into the `deployments` slice passed
Expand Down
Loading
Loading