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
11 changes: 11 additions & 0 deletions .design_docs/optimizer-label-set-facts.example.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
# Label-set facts for optimizer-workload.example.yaml (asap-optimizer-cli --label-set-facts).
series:
- {metric: http_requests_total, spatial_filter: 'env="prod"', series_count: 200}
- {metric: latency_ms, spatial_filter: "", series_count: 50}
groups:
# sum by (job): output labels [job]; also the `topk by (job)` buckets.
- {metric: http_requests_total, spatial_filter: 'env="prod"', grouping_labels: [job], cardinality: 10}
# topk keeps every label: output labels are all of them.
- {metric: http_requests_total, spatial_filter: 'env="prod"', grouping_labels: [env, instance, job], cardinality: 200}
# quantile_over_time keeps every label.
- {metric: latency_ms, spatial_filter: "", grouping_labels: [instance], cardinality: 50}
121 changes: 82 additions & 39 deletions .design_docs/optimizer-v1-implementation-plan.md
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,12 @@ The optimizer output type is `OptimizerSolution`, which is then translated into
|---|---|
| `mergeable` | Two instances can be combined → supports Merge query method |
| `subtractable` | Element-wise difference defined → supports Subtract (tumbling only) |
| `subpopulation_aware` | One instance handles multiple label-group keys → N(s,g)=1 |
| `memory_grows_with_keys` | One instance is a keyed map whose memory grows per key (`Multiple*`), vs a fixed-size instance |

How many instances a config has comes from its `grouping_labels`, split the same way as the
legacy planner (`planner/labels.rs::set_subpopulation_labels`): keyed types keep the output
labels in `aggregated_labels` (one instance), per-group types in `grouping_labels` (one instance
per group), and `topk by (L)` gets one heap per value of `L`.

### Query method (not a free variable — derived from ingest_type × W vs range_a × algebra)

Expand Down Expand Up @@ -76,7 +81,7 @@ asap-planner-rs/src/optimizer/
├── translator.rs translate(&OptimizerSolution) → (StreamingConfig, InferenceConfig)
├── aqe_extractor.rs Rqe, extract_aqes()
├── pipeline.rs run_all_exact_pipeline(), run_greedy_pipeline()
├── dataset.rs CSV series inventory, schema validation, and N_g profiling [#693]
├── label_set_facts.rs external series counts / group cardinalities → per-AQE facts [#756]
├── sketch_properties.rs algebraic properties per AggregationType [Phase 2a, done]
├── candidate_gen.rs enumerate candidate configs per AQE [Phase 2b, done]
├── cost_model.rs ingest/query cost formulas [Phase 2c, done]
Expand Down Expand Up @@ -126,20 +131,23 @@ Each AQE gets its own best config independently (no cross-AQE sharing). MIP shar
#### 2a — `sketch_properties.rs` (done)

```rust
pub struct SketchProperties { pub mergeable: bool, pub subtractable: bool, pub subpopulation_aware: bool }
pub struct SketchProperties { pub mergeable: bool, pub subtractable: bool, pub memory_grows_with_keys: bool }
pub fn sketch_properties(t: AggregationType) -> SketchProperties
```

Implemented values — `CountMinSketchWithHeap` ended up `mergeable=false` (not `true`/unclear as
originally guessed): the heap top-k list doesn't compose across merged/subtracted windows even
though the underlying CMS cells would. Everything else matches the original guess:
- `CountMinSketch`: mergeable=T, subtractable=T, subpopulation_aware=T
- `CountMinSketchWithHeap`: mergeable=F, subtractable=F, subpopulation_aware=T
- `DatasketchesKLL`, `HydraKLL`: mergeable=T, subtractable=F, subpopulation_aware=F (true for HydraKLL too)
- `Sum`/`MultipleSum`: mergeable=T, subtractable=T, subpopulation_aware=F/T respectively
- `HLL`, `SetAggregator`, `DeltaSetAggregator`: mergeable=T, subtractable=F, subpopulation_aware=F
- `MinMax`/`MultipleMinMax`, `Increase`/`MultipleIncrease`: mergeable=T, subtractable=F, subpopulation_aware=F/T
- `SingleSubpopulation`/`MultipleSubpopulation` (legacy wrapper types): all false (unknown, treated conservatively)
though the underlying CMS cells would. `memory_grows_with_keys` is true only for the `Multiple*`
maps:
- `CountMinSketch`: mergeable=T, subtractable=T
- `CountMinSketchWithHeap`: mergeable=F, subtractable=F
- `DatasketchesKLL`, `HydraKLL`: mergeable=T, subtractable=F
- `Sum`/`MultipleSum`: mergeable=T, subtractable=T
- `HLL`, `SetAggregator`, `DeltaSetAggregator`: mergeable=T, subtractable=F
- `MinMax`/`MultipleMinMax`, `Increase`/`MultipleIncrease`: mergeable=T, subtractable=F

The optimizer never proposes single-group `Sum`/`MinMax`/`Increase`: with analytical per-key
memory they cost the same as their `Multiple*` forms (`OPTIMIZER_SKIPPED_AGG_TYPES`).

When a type is both `mergeable` and `subtractable` (e.g. `Sum`, `CountMinSketch`),
`candidate_gen.rs` prefers `Subtract` — it's O(1) regardless of `n`, strictly cheaper than `Merge`'s O(n).
Expand All @@ -151,11 +159,18 @@ pub struct CandidateConfig {
pub config: Option<AggregationConfig>, // None = EXACT fallback
pub query_method: QueryMethod,
pub n_windows: u64,
pub label_group_count: u64, // dataset-derived N_g
pub instance_count: u64, // from the config's grouping labels
pub output_group_count: u64, // cardinality of the AQE's output labels
pub key_config: Option<AggregationConfig>, // paired DeltaSet for CMS/HydraKLL
}
pub fn enumerate_candidates(aqe: &Aqe, scrape_interval_secs: u64) -> Vec<CandidateConfig>
pub fn enumerate_candidates_with_facts(aqe: &Aqe, scrape_interval_ms: u64, facts: &ItemFacts) -> Vec<CandidateConfig>
```

CMS and HydraKLL can't list their own keys, so each such candidate carries the same paired
DeltaSetAggregator the legacy planner deploys (`needs_key_aggregation`): same labels, Tumbling
at the value's slide, retaining `⌈range_a / slide⌉` panes. Greedy registers both and the
translator references both in the query config.

Enumeration: `compatible_agg_types(stat)` × param grid (hardcoded small grids per sketch type,
see `CMS_DEPTHS`/`CMS_WIDTHS`/etc. constants — replace with sketch-bench sweep results in Phase 3)
× window candidates × ingest type:
Expand All @@ -179,15 +194,23 @@ Label granularity: only proposes configs at the label granularity the AQE itself
pub struct AtomicCosts { mem_bytes_per_instance, insert_cpu_secs, merge_cpu_secs, subtract_cpu_secs,
query_cpu_secs, exact_query_cpu_secs: f64 }
pub struct CostWeights { ingest_mem, ingest_cpu, query_mem, query_cpu: f64 }
pub fn ingest_cost(candidate: &CandidateConfig, rho_g: f64, costs: &AtomicCosts, weights: &CostWeights) -> f64
pub fn ingest_cost(candidate: &CandidateConfig, arrival_rate_hz: f64, costs: &AtomicCosts, weights: &CostWeights) -> f64
pub fn query_cost(a: &Aqe, candidate: &CandidateConfig, costs: &AtomicCosts, weights: &CostWeights) -> f64
pub fn total_cost_rate(a: &Aqe, candidate: &CandidateConfig, rho_g: f64, costs: &AtomicCosts, weights: &CostWeights) -> f64
pub fn total_cost_rate(a: &Aqe, candidate: &CandidateConfig, arrival_rate_hz: f64, costs: &AtomicCosts, weights: &CostWeights) -> f64
```

Implements the design doc's formulas with `N(s,g) = 1` for subpopulation-aware sketches and
the dataset-derived `N_g` for per-label-group sketches. The CSV series inventory and exact
filter profiling are implemented in `dataset.rs` (#693); Prometheus-backed profiling remains
the follow-up in #525.
Group cardinality enters per candidate:

| | Multiplier |
|---|---|
| Memory, merge/subtract CPU | `output_group_count` for keyed maps (`Multiple*`), else `instance_count` |
| Query read CPU | `output_group_count` (one read per output group); top-k: `instance_count` (one heap read each) |
| Insert CPU | arrival rate `λ`, independent of cardinality |
| Paired key tracker | `output_group_count` key entries + one insert per item |

Trivial accumulators (`Sum`/`MinMax`/`Increase` and `Multiple*`) use an analytical per-group
entry, `(n_labels × 4 + value_bytes) × 8/7`: label values dictionary-encoded as `u32` codes
(dictionary amortized, not charged), hashbrown load-factor slack. CPU stays on the stub.
`AtomicCosts` added `exact_query_cpu_secs` (not in the original plan) — without a non-zero cost
for the EXACT fallback's query, `IngestCost=0, QueryCost=0` would make EXACT always win trivially.

Expand All @@ -212,18 +235,19 @@ it then, reusing `window_compatible`, `spatial_filter_compatible`, `topk_weighti
#### 2e — `greedy.rs` + `pipeline.rs` + `translator.rs` (done)

```rust
pub fn greedy_assign(aqes: Vec<Aqe>, scrape_interval_secs: u64, rho_g: f64,
costs: &AtomicCosts, weights: &CostWeights) -> OptimizerSolution
pub fn greedy_assign(aqes: Vec<Aqe>, scrape_interval_ms: u64, atomic_cost_table: &AtomicCostTable,
weights: &CostWeights, facts: &[ItemFacts]) -> OptimizerSolution
```

For each AQE independently: `argmin` over `enumerate_candidates(aqe, ...)` by `total_cost_rate`.
No feasibility filter needed (see 2d note) — every candidate from `enumerate_candidates` is valid
for that AQE by construction. Assigns sequential `aggregation_id`s to deployed configs.

`run_greedy_pipeline(config, dataset, scrape_interval_secs, rho_g) -> Result<(StreamingConfig, InferenceConfig), DatasetError>`
added to `pipeline.rs` alongside `run_all_exact_pipeline()`. The dataset is required for the
greedy path and supplies metric schemas plus the `N_g` profile for each AQE. `rho_g` is
currently a single placeholder value applied uniformly to every candidate — see open TODO below.
`run_greedy_pipeline(config, facts, scrape_interval_ms, atomic_cost_table) -> Result<(StreamingConfig, InferenceConfig), LabelSetFactsError>`
added to `pipeline.rs` alongside `run_all_exact_pipeline()`. The workload's `metrics:` hints
are required and supply the label schema; every workload metric must have one. `facts`
(`LabelSetFacts`) supplies each AQE's group cardinalities and series count; arrival rate is
derived as `series_count × 1000 / scrape_interval_ms`, assuming one sample per series per scrape.

`translator.rs::build_inference_config()` now populates `query_configs`: for each assignment with
a real `aggregation_id`, emits one `QueryConfig` per original query string, with
Expand Down Expand Up @@ -319,26 +343,44 @@ standalone `asap-optimizer-cli` binary (`asap-planner-rs/src/bin/optimizer_cli.r
```
cargo run -p asap_planner --bin asap-optimizer-cli -- \
--input_config <path/to/workload.yaml> \
--dataset <path/to/series-inventory.csv> \
--label-set-facts <path/to/label_set_facts.yaml> \
--data-ingestion-interval-ms 60000 \
[--rho 1.0] \
[--atomic-costs <path/to/atomic_costs.json> \
--atomic-cost-workload <path/to/workload-selector.json>]
```

Takes the same `ControllerConfig` YAML format as `asap-planner --input_config`.
The required CSV dataset is one row per unique metric series, with a fixed `metric`
column and label columns. It derives the metric schema and each AQE's distinct
group count; no live Prometheus connection is needed. A `metrics:` hints block, when
present, is checked against the dataset and mismatches fail loudly.
Prints deployed streaming configs and query configs to stdout. `--rho` is the
placeholder arrival rate (see TODOs below — not real yet). `--atomic-costs` is
Takes the same `ControllerConfig` YAML format as `asap-planner --input_config`; its
`metrics:` hints are required and must cover every workload metric, since they supply the
label schema. `--label-set-facts` is a YAML file of externally provided facts — the optimizer
never estimates them (example: `optimizer-label-set-facts.example.yaml`, for the workload
`optimizer-workload.example.yaml`; its top-k query needs `--atomic-costs` with a CMS-with-heap
reference row, otherwise the run fails with that item listed as unservable):

```yaml
series: # per (metric, spatial_filter)
- {metric: http_requests_total, spatial_filter: 'env="prod"', series_count: 200}
groups: # per (metric, spatial_filter, grouping_labels)
- {metric: http_requests_total, spatial_filter: 'env="prod"', grouping_labels: [job], cardinality: 10}
```

- `spatial_filter` and `grouping_labels` are required; `""` means unfiltered, `[]` means
full aggregation (cardinality must then be 1). Filters are normalized like the workload's.
- `grouping_labels` are the AQE's resolved output labels (all labels for `*_over_time` and
top-k; the `by` labels after resolving `without`). A `topk by (L)` query also needs a
`groups` row for `[L]`.
- `1 ≤ cardinality ≤ series_count`. Missing facts are an error listing every expected key;
unused facts only warn.
- Arrival rate is `series_count × 1000 / data-ingestion-interval-ms`.

No live Prometheus connection is needed.
Prints deployed streaming configs and query configs to stdout. `--atomic-costs` is
optional; when supplied it requires `--atomic-cost-workload`, a JSON file
containing the exact `profiles[].workload` value from that benchmark artifact.
The loader validates the document schema and rejects a selector that matches
zero or multiple profiles; it never mixes entries across workloads. Omit both
flags and ordinary unbenchmarked candidates use the flat stub, while
CMS-with-heap candidates warn and are dropped until a matching reference row is available.
CMS-with-heap candidates warn and are dropped until a matching reference row is available;
an item left with no candidate makes the run fail, listing every unservable item.

### Running with real sketch-bench costs

Expand All @@ -363,12 +405,13 @@ CMS-with-heap candidates warn and are dropped until a matching reference row is
# Run the actual optimizer:
cargo run -p asap_planner --bin asap-optimizer-cli -- \
--input_config workload.yaml --data-ingestion-interval-ms 60000 \
--dataset path/to/series-inventory.csv \
--label-set-facts path/to/label_set_facts.yaml \
--atomic-costs path/to/atomic_costs.json \
--atomic-cost-workload path/to/workload-selector.json
```

`candidate-gen-dump`'s output labels each resolved params row `[real]` or `[stub]`; candidates
`candidate-gen-dump`'s output labels each resolved params row `[real]`, `[stub]`, or
`[analytical mem, stub cpu]` (trivial accumulators); candidates
whose required cost row is missing are shown as `DROPPED`. CMS-with-heap uses its fixed-top-k
reference model when the `{rows,cols}` row is present and is dropped otherwise.

Expand All @@ -384,9 +427,9 @@ instead of `generator::generate_plan()`, likely behind an opt-in flag first.
|---|---|---|
| `aqe_extractor.rs` | `extract_requirements()` | Duplicates `build_query_requirements_promql` in `asap-query-engine/src/engines/simple_engine/promql.rs:614` — extract to shared free fn in `asap_types::query_requirements` |
| `capability_matching.rs` | `labels_compatible()`, line 86 | Relax to superset matching — do in Phase 3b |
| `cost_model.rs` | `ingest_cost()`, `query_cost()` | Prometheus-backed `N_g` profiling remains for #525; dataset-backed `N_g` is implemented in #693 |
| `cost_model.rs` | `key_tracker_ingest_cost()` | Paired DeltaSet uses analytical key bytes and stub insert CPU until sketch-bench measures it |
| `cost_model.rs` | `CostWeights::default()` | Self-consistent stub ratio (mem:cpu ≈ 1e-9:1), not real $/byte-sec vs $/cpu-sec calibration |
| `greedy.rs`, `pipeline.rs` | `rho_g` parameter | Single placeholder value applied uniformly; real per-config rates need Prometheus scrape-rate × active-series-count, not wired up |
| `label_set_facts.rs` | `resolve()` | Arrival rate assumes one sample per series per scrape; overestimates sparse or irregular series |
| `translator.rs` | `retention_count_for_assignment(Subtract)` | Returns hardcoded `1`; should be the actual checkpoint count needed to cover the full lookback |
| `promql/generator.rs` | `generate_plan()` doc comment | Flags that `Controller::generate()` still uses the hardcoded path, not the optimizer — see "Offline Testing" section above |
| — | Accuracy constraint | No `Error(a,g) ≤ ε_a` check exists anywhere; nothing stops picking an under-provisioned sketch (Phase 3d) |
Expand Down
12 changes: 12 additions & 0 deletions .design_docs/optimizer-workload.example.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,12 @@
# Example workload for asap-optimizer-cli; facts in optimizer-label-set-facts.example.yaml.
# The top-k query is only servable with --atomic-costs containing a CMS-with-heap
# reference row; without it the optimizer reports it as unservable.
query_groups:
- queries:
- 'sum by (job) (http_requests_total{env="prod"})'
- 'quantile_over_time(0.99, latency_ms[5m])'
- 'topk by (job) (5, sum_over_time(http_requests_total{env="prod"}[5m]))'
repetition_delay_ms: 60000
metrics:
- {metric: http_requests_total, labels: [job, instance, env]}
- {metric: latency_ms, labels: [instance]}
1 change: 0 additions & 1 deletion Cargo.lock

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
Expand Up @@ -86,12 +86,11 @@ pub fn does_precompute_operator_support_subpopulations(
// own top-k keys via CmsHeapItem. So it supports subpopulations the
// same way MultipleSum does: labels go in `aggregated`, not
// `grouping` (see `set_subpopulation_labels`).
AggregationType::CountMinSketchWithHeap if matches!(statistic, Statistic::Topk) => true,
AggregationType::CountMinSketchWithHeap => true,

AggregationType::HLL => false,

// Default: not supported
_ => panic!("Unexpected precompute operator: {}", precompute_operator),
AggregationType::HLL
| AggregationType::SetAggregator
| AggregationType::DeltaSetAggregator => false,
}
}

Expand Down
1 change: 0 additions & 1 deletion asap-planner-rs/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,6 @@ anyhow.workspace = true
tracing.workspace = true
tracing-subscriber.workspace = true
clap.workspace = true
csv = "1.4.0"
indexmap.workspace = true
chrono.workspace = true
promql-parser = "0.5.0"
Expand Down
24 changes: 20 additions & 4 deletions asap-planner-rs/src/bin/candidate_gen_dump.rs
Original file line number Diff line number Diff line change
Expand Up @@ -83,7 +83,11 @@ fn main() -> anyhow::Result<()> {
println!(" queries: {:?}", aqe.query_strings);

let candidates = enumerate_candidates(aqe, args.scrape_interval_ms);
print_candidates_grouped(&candidates, &atomic_cost_table);
print_candidates_grouped(
&candidates,
&atomic_cost_table,
aqe.requirements.grouping_labels.len(),
);
}

Ok(())
Expand All @@ -96,7 +100,11 @@ fn main() -> anyhow::Result<()> {
/// params (M) [× N windows = NM total]:
/// ...
/// EXACT is printed last as a single line.
fn print_candidates_grouped(candidates: &[CandidateConfig], atomic_cost_table: &AtomicCostTable) {
fn print_candidates_grouped(
candidates: &[CandidateConfig],
atomic_cost_table: &AtomicCostTable,
n_grouping_labels: usize,
) {
// Collect unique (agg_type_str, sub_type) keys in first-seen order.
let mut group_order: Vec<(String, String)> = Vec::new();
// (agg_type_str, sub_type) -> (agg_type, unique windows, unique params)
Expand Down Expand Up @@ -200,10 +208,18 @@ fn print_candidates_grouped(candidates: &[CandidateConfig], atomic_cost_table: &
} else {
let param_map: HashMap<String, Value> =
p.iter().map(|(k, v)| (k.clone(), v.clone())).collect();
match resolve_atomic_costs(atomic_cost_table, *agg_type, &param_map) {
match resolve_atomic_costs(
atomic_cost_table,
*agg_type,
&param_map,
n_grouping_labels,
) {
Some(costs) => {
let label = if costs == AtomicCosts::default() {
let stub = AtomicCosts::default();
let label = if costs == stub {
"stub"
} else if costs.insert_cpu_secs == stub.insert_cpu_secs {
"analytical mem, stub cpu"
} else {
"real"
};
Expand Down
Loading
Loading