Skip to content
Draft
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
13 changes: 11 additions & 2 deletions crates/integration-tests/tests/planner_layering_common/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@ use asap_physical_optimizer::implementation::physical_candidates::PhysicalCandid
#[path = "../executor_models/mod.rs"]
mod executor_models;

use asap_plan_selection::{plan_stages, Selection};
use asap_plan_selection::{plan_stages, PlanningModels, Selection};
use asap_types::ir::export::{
compile_logical_asap_workload, LogicalASAPDAG, LogicalASAPNodeId, LogicalASAPOperatorPayload,
LogicalASAPQueryRoot, NonASAPOpKind, PhysicalASAPDAG,
Expand Down Expand Up @@ -329,6 +329,15 @@ fn export(roots: &[QueryRoot]) -> (LogicalASAPDAG, Vec<LogicalASAPNodeId>) {
/// library's `plan_stages`, with each entry's demand (accuracy, recurrence,
/// predictability).
pub fn run_stages(workload: &PlanningWorkload, roots: Vec<QueryRoot>) -> Run {
run_stages_with(workload, roots, executor_models())
}

/// [`run_stages`] with the deployment inputs `models`.
pub fn run_stages_with(
workload: &PlanningWorkload,
roots: Vec<QueryRoot>,
models: PlanningModels<'_>,
) -> Run {
let (dag, query_roots) = export(&roots);
let stage0 = Logical {
id: "S0".into(),
Expand All @@ -345,7 +354,7 @@ pub fn run_stages(workload: &PlanningWorkload, roots: Vec<QueryRoot>) -> Run {
roots.into_iter().enumerate().collect(),
&demand,
workload.data_workload.as_ref().expect("data workload"),
executor_models(),
models,
MAX_CANDIDATES,
)
.expect("plans");
Expand Down
117 changes: 100 additions & 17 deletions crates/integration-tests/tests/planner_layering_example4.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,9 +15,12 @@ mod planner_layering_common;

use std::collections::{BTreeMap, BTreeSet};

use asap_plan_selection::{DeploymentCapabilities, PlanningModels};
use asap_types::ir::export::{LogicalASAPNodeId, LogicalASAPOperatorPayload};
use asap_types::ir::schema::SketchAlgorithm;
use asap_types::workload::{DataArrival, PlanningWorkload};
use asap_types::workload::{
DataArrival, DurationMs, PlanningWorkload, Query, Rate, RepeatedDemand, RepetitionInterval,
};
use planner_layering_common::*;

/// Pattern A as given: one ad hoc batch at T over mixed data.
Expand Down Expand Up @@ -326,26 +329,106 @@ fn stage2_b_sliding_kll_has_two_materialization_options() {
assert_eq!(found, BTreeSet::from([IngestionTime, QueryTimeKept]));
}

/// B2 rebuilds all five tumbling KLLs at every evaluation, so it costs at least B1 and B3.
// Which of B1 and B2 is cheaper depends on the workload and the deployment,
// not on the materialization alone (Q49). B1 keeps one KLL per series per
// pane for as long as the window; B2 keeps nothing but reads the window's raw
// samples at every evaluation, and pays to keep them when the deployment does
// not. So B1 wins when its panes are smaller than the raw samples they cover
// and those samples would otherwise be kept for the plan.

/// Pattern B with `series` series sampled every `interval_ms`, asking for the
/// p99 over the last `window_min` minutes every `every_min` minutes.
fn pattern_b_variant(
series: u64,
interval_ms: u64,
window_min: u64,
every_min: u64,
) -> PlanningWorkload {
let mut workload = pattern_b();
let data = workload.data_workload.as_mut().expect("data workload");
data.input_cardinality = declared(series);
data.data_ingestion_interval = declared(DurationMs(interval_ms));
data.ingestion_rate = declared(Rate(series as f64 * 1000.0 / interval_ms as f64));
let entry = &mut workload
.query_workload
.repeating_queries
.as_mut()
.expect("repeating")[0];
entry.query = Query(format!(
"quantile_over_time(0.99, latency_ms[{window_min}m])"
));
entry.demand =
RepeatedDemand::FixedInterval(RepetitionInterval((every_min * MINUTE_MS) as u32));
entry.time_selection.lookback = Some(DurationMs(window_min * MINUTE_MS));
workload
}

/// Plans `workload` on the executor with raw data kept or not by the deployment.
fn run_with_raw_data(workload: &PlanningWorkload, raw_data_retained: bool) -> Run {
let capabilities = DeploymentCapabilities {
raw_data_retained,
..asap_executor::capabilities()
};
let models = PlanningModels::builtin().with_capabilities(&capabilities);
run_stages_with(workload, lower_promql(workload), models)
}

/// The costs of B1 and B2 for the tumbling KLL of `every_min` minutes.
fn b1_b2(run: &Run, every_min: u64) -> (f64, f64) {
let form = Some(WindowForm::Tumbling {
length_ms: every_min * MINUTE_MS,
});
let options = options_of(run, form, 1);
(
cost(run, &options[&IngestionTime].0),
cost(run, &options[&NotMaterialized].0),
)
}

/// At Pattern B's 1M series, rebuilding the five panes at query time breaks
/// the 200 ms latency bound, so B2 is invalid and B1 is priced.
#[test]
#[ignore = "built-in model (#604, Q49): B2 rebuilds the panes in 310 ms, over the 200 ms latency bound, so it is not priced; B1 retains 6 panes of 1M per-series KLLs (6.1 GB) for 768.58 cost/s, and B2 costs 5.17 cost/s, or 45.17 when the deployment does not keep raw data (5 min of raw samples, 320 MB, 40.0 cost/s)"]
fn stage3_b_rebuilding_every_window_costs_most() {
fn stage3_b_rebuilding_a_million_series_breaks_the_latency_bound() {
let run = run_promql(&pattern_b());
let options = options_of(&run, tumbling(), 1);
let b2 = cost(&run, &options[&NotMaterialized].0);
for m in [IngestionTime, QueryTimeKept] {
assert!(cost(&run, &options[&m].0) <= b2, "{m:?} vs B2 {b2}");
}
let b2 = &options[&NotMaterialized].0;
let reason = run.invalid()[b2.as_str()];
assert!(reason.contains("latency bound"), "{b2}: {reason}");
assert!(run.cost(&options[&IngestionTime].0).is_some());
}

/// Repeating over arriving data, the built-in models pick B1 among the tumbling options.
/// When the deployment keeps raw data, B2 pays only to read it, while B1
/// keeps its panes: B2 is cheaper whenever it meets the latency bound.
#[test]
#[ignore = "built-in model (#604, Q49): B2 is over the 200 ms latency bound, so it is not priced; B1 costs 768.58 cost/s against B2's 5.17, or 45.17 when the deployment does not keep raw data"]
fn stage3_b_prefers_ingestion_time_tumbling_windows() {
let run = run_promql(&pattern_b());
let options = options_of(&run, tumbling(), 1);
let b1 = cost(&run, &options[&IngestionTime].0);
for (m, (id, _)) in &options {
assert!(b1 <= cost(&run, id), "B1 {b1} vs {m:?} {}", cost(&run, id));
}
fn stage3_b_kept_raw_data_makes_rebuilding_cheaper() {
let run = run_with_raw_data(&pattern_b_variant(10_000, 15_000, 5, 1), true);
let (b1, b2) = b1_b2(&run, 1);
assert!(b2 < b1, "B2 {b2} vs B1 {b1}");
}

/// Sampled every 15 s, a 1-min pane covers 4 samples per series, fewer bytes
/// than its KLL: B2 is cheaper even when it pays to keep the raw data.
#[test]
fn stage3_b_panes_larger_than_their_raw_data_lose() {
let run = run_with_raw_data(&pattern_b_variant(10_000, 15_000, 5, 1), false);
let (b1, b2) = b1_b2(&run, 1);
assert!(b2 < b1, "B2 {b2} vs B1 {b1}");
}

/// Sampled every second, a 10-min pane covers 600 samples per series, far
/// more bytes than its KLL. When the deployment does not keep raw data, B2
/// must keep the hour's samples, so B1 is cheaper, and the planner selects it.
#[test]
fn stage3_b_panes_smaller_than_their_raw_data_win() {
let run = run_with_raw_data(&pattern_b_variant(1_000, 1_000, 60, 10), false);
let (b1, b2) = b1_b2(&run, 10);
assert!(b1 < b2, "B1 {b1} vs B2 {b2}");
let options = options_of(
&run,
Some(WindowForm::Tumbling {
length_ms: 10 * MINUTE_MS,
}),
1,
);
assert_eq!(run.selection.selected, options[&IngestionTime].0);
}