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
18 changes: 6 additions & 12 deletions crates/integration-tests/tests/planner_layering_example3.rs
Original file line number Diff line number Diff line change
Expand Up @@ -498,10 +498,12 @@ fn stage3_b_selects_cheapest_valid() {

/// Merged KLL and DDSketch panes keep the family's guarantee, so Stage 3
/// finds no tumbling candidate inaccurate. Rebuilding all five panes at
/// every evaluation takes 310 ms of query-time work, over the 200 ms latency
/// bound (S6); with the panes maintained at ingestion time, only the merge
/// and the estimate remain and the candidate is valid and priced. Panes
/// kept at query time (B3) need a capability the executor lacks.
/// every evaluation takes 130 ms of query-time work, within the 200 ms
/// latency bound (S6), now that the panes' time shifts are free and each
/// 1-min range pays only for its own rows (Q66); with the panes maintained at
/// ingestion time, only the merge and the estimate remain. Both are valid
/// and priced. Panes kept at query time (B3) need a capability the executor
/// lacks.
#[test]
fn stage3_b_tumbling_candidates_are_valid() {
let run = run_b();
Expand All @@ -526,14 +528,6 @@ fn stage3_b_tumbling_candidates_are_valid() {
p.id,
invalid.get(p.id.as_str())
),
"" => assert!(
invalid
.get(p.id.as_str())
.is_some_and(|r| r.contains("310.0 ms") && r.contains("200 ms latency")),
"{}: {:?}",
p.id,
invalid.get(p.id.as_str())
),
_ => {
assert!(
!invalid.contains_key(p.id.as_str()),
Expand Down
14 changes: 8 additions & 6 deletions crates/integration-tests/tests/planner_layering_example4.rs
Original file line number Diff line number Diff line change
Expand Up @@ -406,16 +406,18 @@ fn b1_b2(run: &Run, every_min: u64) -> (f64, f64) {
)
}

/// 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.
/// At Pattern B's 1M series, rebuilding the five panes at query time takes
/// 130 ms, within the 200 ms latency bound: the panes' time shifts are free
/// and each 1-min range pays only for its own rows (Q66). So B2 and B1 are
/// both valid and priced, and B2, which keeps nothing, is cheaper.
#[test]
fn stage3_b_rebuilding_a_million_series_breaks_the_latency_bound() {
fn stage3_b_rebuilding_a_million_series_meets_the_latency_bound() {
let run = run_promql(&pattern_b());
let options = options_of(&run, tumbling(), 1);
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());
assert!(!run.invalid().contains_key(b2.as_str()), "{b2}");
let (b1, b2) = b1_b2(&run, 1);
assert!(b2 < b1, "B2 {b2} vs B1 {b1}");
}

/// When the deployment keeps raw data, B2 pays only to read it, while B1
Expand Down
93 changes: 82 additions & 11 deletions crates/plan-selection/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1738,22 +1738,32 @@ fn price_nodes(
(out, estimate, format!("limit to {} rows", out.rows))
}
// At query time a range keeps only its own span of a longer
// scan: a filter on the timestamp.
// scan. Rows are ordered by time, so it seeks to that span
// and is charged for the rows it keeps, not those it skips
// (Q66).
NonASAPOp::TimeRange { range, .. } if !ingestion => {
let rows = (rows_per_ms * range.as_millis() as f64).round() as u64;
let out = edge(input.rows.min(rows.max(1)));
let estimate = estimate_operator(
PhysicalOperator::Filter {
predicate_operations_per_row: 1,
},
OperatorStatistics::Filter { edges: unary(out) },
);
(
out,
estimate,
format!("time range {range:?}: pass {} rows", out.rows),
Ok(ResourceEstimate::new(out.rows as f64, width, 0)),
format!(
"time range {range:?}: keep {} of {} rows",
out.rows, input.rows
),
)
}
// A shift only re-labels time: the executor folds it into
// the time bounds of the read below it, so it does no
// per-row work (Q66).
NonASAPOp::TimeShift { shift, .. } => (
input,
Ok(ResourceEstimate::new(0.0, 0, 0)),
format!(
"time shift {} ms: re-label {} rows, free",
shift.offset_ms, input.rows
),
),
other => {
let estimate = estimate_operator(
PhysicalOperator::PassThrough,
Expand Down Expand Up @@ -2837,8 +2847,8 @@ mod tests {
.filter(|n| !n.output_state.timing.is_query_time() && cost.per_node[&n.id].cost > 0.0)
.map(|n| cost.per_node[&n.id].detail.clone())
.collect();
// The scan, one shift, one range and the newest pane.
assert_eq!(ingestion.len(), 4, "{ingestion:#?}");
// The scan, one range and the newest pane; the shift is free (Q66).
assert_eq!(ingestion.len(), 3, "{ingestion:#?}");
}

/// Kept at query time (B3, Q59), each evaluation builds the newest pane
Expand Down Expand Up @@ -3196,6 +3206,67 @@ mod tests {
assert!(shared <= separate, "shared {shared} vs separate {separate}");
}

/// Q66: over a shared 5-year scan, a time shift is free and each time
/// range is charged for the rows it keeps: the 1-year range costs a
/// fifth of the 5-year one, not the same.
#[test]
fn time_shift_is_free_and_time_range_pays_for_the_rows_it_keeps() {
let queries = [
"quantile_over_time(0.99, latency_ms[5y])",
"quantile_over_time(0.99, latency_ms[1y] offset 2y)",
];
let roots = queries
.iter()
.enumerate()
.map(|(i, query)| {
let root = lower_promql(query, AccuracyTarget::Exact);
let root =
asap_types::ir::schema_support::with_promql_series_identity(&root).unwrap();
(i, QueryRoot::Operator(root))
})
.collect();
let stage1 = stage1_logical_candidates(roots, &Default::default(), &[]).unwrap();
let demand = vec![every_10s(None); queries.len()];
let data = data();
let variant = *variants(&stage1)
.iter()
.find(|v| v.sharing == Sharing::IdenticalExpressions)
.unwrap();
let raw = vec![0; variant.inventory.targets.len()];
let models = PlanningModels::builtin();
let (_, mut stage2) = realize(variant, &raw, &demand, &data, &models).unwrap();
let candidate = stage2.candidates.swap_remove(0);
let cost = assess(&candidate, &demand, &data, &models).unwrap();
let of = |pick: fn(&NonASAPOp<PhysicalASAPNodeId>) -> bool| -> Vec<&NodeCost> {
candidate
.dag
.nodes
.iter()
.filter(|n| matches!(&n.payload, Payload::NonASAP(operator) if pick(operator)))
.map(|n| &cost.per_node[&n.id])
.collect()
};
let shifts = of(|op| matches!(op, NonASAPOp::TimeShift { .. }));
assert_eq!(shifts.len(), 1);
assert_eq!(shifts[0].cost, 0.0, "{}", shifts[0].detail);
let scan = of(|op| matches!(op, NonASAPOp::Scan { .. }))[0].rows;
let mut ranges = of(|op| matches!(op, NonASAPOp::TimeRange { .. }));
ranges.sort_by_key(|n| n.rows);
let [year, five_years] = ranges.as_slice() else {
panic!("two time ranges");
};
assert_eq!(five_years.rows, scan);
assert!((five_years.rows as f64 / year.rows as f64 - 5.0).abs() < 0.01);
// Charged per kept row: the same price per row for both.
let per_row = |n: &NodeCost| n.cost / n.rows as f64;
assert!((per_row(year) / per_row(five_years) - 1.0).abs() < 1e-9);
assert!(
year.detail.contains(&format!("of {scan} rows")),
"{}",
year.detail
);
}

// ── Deployment capabilities (C3, Q48) ────────────────────────────────

fn reasons(selection: &Selection) -> BTreeMap<&str, &str> {
Expand Down
16 changes: 14 additions & 2 deletions docs/design_docs/proposals/stage3-cost-model.md
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,18 @@ the longest range plus its offset. Each range then passes only its own
span. For example, `x[1y] offset 2y` scans 3 years, and a 1-year range over a
shared 5-year scan passes 1 year of rows.

**Time shifts and ranges (Q66).** A time shift only re-labels time: the
executor folds the offset into the time bounds of the read below it, so a
shift is free at either timing. A query-time time range is charged one
operation per row it keeps, not per row of the scan below it: rows are
ordered by time, so it seeks to its span. The 1-year range above costs 1
year of rows, not 5. (At ingestion time a range keeps every arriving row, so
both counts agree.) Without this, every shifted window over a shared long
scan paid for the whole scan twice, and Example 3a's shared 1-year segments
(18 104 per second; scan 12 264, five ranges and five builds of 584 each)
lost to the shared-input plan (25 112), which reads the same scan with fewer
shifts and ranges.

**Workload cost** is the sum over the candidate's DAG nodes:
`cost(P) = Σ_n cost(n)`, in cost per second.

Expand Down Expand Up @@ -292,8 +304,8 @@ candidate scales by the same factor, so their ranking is unchanged.

**With materialization.** Stage 2 also offers Q2's exact `sum_over_time` in
six 10-s panes maintained at ingestion time. The cheapest such candidate costs
47.407 per second: the panes' ingestion work is small (scan 0.387, shift and
range 0.133, newest-pane build 0.067), but the newest pane retains seven
47.340 per second: the panes' ingestion work is small (scan 0.387, range
0.067, the shift free, newest-pane build 0.067), but the newest pane retains seven
panes of 1 000 000 per-series sums, 336 MB, for 42.0 per second of memory.
The selection is unchanged. Q2's 100 ms latency bound rejects 40 candidates,
every Count-Sketch + heap among them.
Expand Down
6 changes: 3 additions & 3 deletions tools/dag-viewer/examples.json
Original file line number Diff line number Diff line change
Expand Up @@ -18,21 +18,21 @@
"name": "3a · historical p99 batch",
"example": "planner-layering-3a",
"file": "out/example3a.json",
"story": "An ad hoc batch of five p99 reports over 1–5 years, run once. Each query gets a KLL, and sharing the input scan across all five is cheapest. Pass 2 also offers five 1-year KLL segments shared by all five reports, but the cost model charges each segment's time shift and range for every row of the 5-year scan, so it costs more."
"story": "An ad hoc batch of five p99 reports over 1–5 years, run once. Each query gets a KLL, and Pass 2 offers five 1-year KLL segments shared by all five reports. The segments win: each row of the 5-year scan lands in one segment, a time shift is free and a time range pays only for the rows it keeps, so they cost less than sharing only the input scan."
},
{
"id": "example3b",
"name": "3b · live p99 panel",
"example": "planner-layering-3b",
"file": "out/example3b.json",
"story": "p99 over the last 5 min every minute, 1M series sampled every 15 s, raw data kept by the deployment. A per-series KLL built at query time wins: keeping 1-min KLL panes from ingestion time costs 768 cost/s of memory, and rebuilding the panes at query time breaks the 200 ms bound."
"story": "p99 over the last 5 min every minute, 1M series sampled every 15 s, raw data kept by the deployment. A per-series KLL built at query time wins: keeping 1-min KLL panes from ingestion time costs 768 cost/s of memory, and rebuilding the panes at query time meets the 200 ms bound but costs slightly more."
},
{
"id": "example4a",
"name": "4a · monthly p99 reports",
"example": "planner-layering-4a",
"file": "out/example4a.json",
"story": "Example 3a repeated monthly and known in advance. The same shared-input KLL plan wins; its cost is now amortized over the month between runs. The shared 1-year segments may now be maintained at ingestion time, but keeping a million per-series KLLs per segment costs far more."
"story": "Example 3a repeated monthly and known in advance. The same shared 1-year segments win, built at query time; their cost is now amortized over the month between runs. They may now be maintained at ingestion time, but keeping a million per-series KLLs per segment costs far more."
},
{
"id": "example4b",
Expand Down
Loading