diff --git a/crates/devtools/src/bin/stage_pipeline.rs b/crates/devtools/src/bin/stage_pipeline.rs index c135460c..2a1fa813 100644 --- a/crates/devtools/src/bin/stage_pipeline.rs +++ b/crates/devtools/src/bin/stage_pipeline.rs @@ -20,7 +20,9 @@ // the queries as written and, when Pass 2's identical-expression rule // merges something, again with identical sub-DAGs shared ("· shared // input"), and, when the summary-capability rule applies, again with -// one summary sized for its strictest consumer ("· shared summary"); in +// one summary sized for its strictest consumer ("· shared summary"), +// and, when queries read windows of one scan on a common grid, again +// with one summary per shared segment ("· shared segments"); in // enumeration order, only those with a written physical candidate. A // repeating query's mergeable alternatives also come in tumbling panes // (Pass 2's window-composition rule), e.g. "Q1 Kll · tumbling 1m panes"; @@ -267,6 +269,7 @@ fn stage_pipeline( Sharing::Independent => "", Sharing::IdenticalExpressions => " · shared input", Sharing::SummaryCapability => " · shared summary", + Sharing::WindowSegments => " · shared segments", }; let written = candidate .physical diff --git a/crates/devtools/tests/stage_pipeline.rs b/crates/devtools/tests/stage_pipeline.rs index 64d671be..9b2afcd8 100644 --- a/crates/devtools/tests/stage_pipeline.rs +++ b/crates/devtools/tests/stage_pipeline.rs @@ -183,13 +183,15 @@ fn example3b_lists_tumbling_candidates() { } } -/// Example 4, Pattern A repeated monthly: the same 486 candidates as the ad -/// hoc batch (3a), and none maintained at ingestion time. The windows (1–5 y) -/// are longer than the month between runs and no pane width fits, so -/// nothing is maintainable. The selected plan is 3a's, now amortized over -/// monthly runs instead of one run per hour of horizon. +/// Example 4, Pattern A repeated monthly: the same 487 logical candidates +/// as the ad hoc batch (3a), 486 independent or sharing input plus the one +/// sharing 1-year segments (Q60). The windows (1–5 y) are longer than the +/// month between runs and no pane width fits, so only the shared segments +/// can be maintained at ingestion time (Q61): 488 plans. The selected plan +/// is 3a's, now amortized over monthly runs instead of one run per hour of +/// horizon. #[test] -fn example4a_repeats_monthly_with_nothing_maintainable() { +fn example4a_repeats_monthly_with_only_segments_maintainable() { let once = generate(&[ "--example", "planner-layering-3a", @@ -205,10 +207,14 @@ fn example4a_repeats_monthly_with_nothing_maintainable() { let physical = monthly["stage2_physical_asap"]["candidates"] .as_array() .unwrap(); - assert_eq!(physical.len(), 486); - assert!(physical + assert_eq!(physical.len(), 488); + let maintained: Vec<_> = physical .iter() - .all(|p| !p["label"].as_str().unwrap().contains("ingestion time"))); + .map(|p| p["label"].as_str().unwrap()) + .filter(|label| label.contains("ingestion time")) + .collect(); + assert_eq!(maintained.len(), 1, "{maintained:?}"); + assert!(maintained[0].ends_with("shared segments · ingestion time: Kll ×5 panes")); let selected = |d: &Value| { d["stage3_selection"]["selected"] .as_str() diff --git a/crates/integration-tests/tests/planner_layering_common/mod.rs b/crates/integration-tests/tests/planner_layering_common/mod.rs index bc2908ef..5366cc92 100644 --- a/crates/integration-tests/tests/planner_layering_common/mod.rs +++ b/crates/integration-tests/tests/planner_layering_common/mod.rs @@ -575,6 +575,14 @@ pub fn query_option( } } +/// Whether `dag` shares a summary build across queries: Pass 2's +/// shared-segment rule (Q60), whose segments several queries merge. +pub fn shares_segments(dag: &impl ExportedDag, query_roots: &[PhysicalASAPNodeId]) -> bool { + cross_query_nodes(dag, query_roots) + .into_iter() + .any(|id| matches!(dag.payload(id), Operator::ASAP(ASAPOp::SummaryAgg { .. }))) +} + /// Nodes reachable from more than one query root. pub fn cross_query_nodes( dag: &impl ExportedDag, diff --git a/crates/integration-tests/tests/planner_layering_example3.rs b/crates/integration-tests/tests/planner_layering_example3.rs index 42aaed82..824efb26 100644 --- a/crates/integration-tests/tests/planner_layering_example3.rs +++ b/crates/integration-tests/tests/planner_layering_example3.rs @@ -158,13 +158,17 @@ fn stage1_a_keeps_five_independent_klls() { assert_eq!(sizes.len(), 1, "one ε, one KLL size: {sizes:?}"); } -/// The identical-expression rule shares only raw input (scan, range, shift), never a quantile or summary, and keeps the unshared variant. +/// The identical-expression rule shares only raw input (scan, range, shift), never a quantile or summary, and keeps the unshared variant. (The shared-segment candidate shares its segments; see below.) #[test] fn stage1_a_identical_expression_rule_shares_only_raw_input() { let run = run_a(); assert!(run.logical.iter().any(|c| c.shared_input)); assert!(run.logical.iter().any(|c| !c.shared_input)); - for c in &run.logical { + for c in run + .logical + .iter() + .filter(|c| !shares_segments(&c.dag, &c.query_roots)) + { for id in cross_query_nodes(&c.dag, &c.query_roots) { let kind = relational(c.dag.payload(id)); assert!( @@ -176,24 +180,35 @@ fn stage1_a_identical_expression_rule_shares_only_raw_input() { } } -/// One Exponential Histogram of KLLs over [T − 5y, T] serves all five queries, each through its own merge and p99 estimate. +/// The windows' boundaries lie on a 1-year grid, so one set of five 1-year +/// KLL segments over [T − 5y, T] serves all five queries (Q60, in place of +/// the spec's Exponential Histogram): each query merges the segments its +/// range covers and estimates p99 from its merge. #[test] -#[ignore = "needs Pass 2 window composition (#580): Exponential Histogram"] -fn stage1_a_window_composition_adds_one_eh_for_all_five() { +fn stage1_a_shared_segments_serve_all_five() { let run = run_a(); - let found = run.logical.iter().find(|c| { - windowed_builds(c).iter().any(|(_, a, form, readers)| { - *a == SketchAlgorithm::Kll - && matches!(form, WindowForm::ExponentialHistogram { horizon_ms } if *horizon_ms >= 5 * YEAR_MS) - && readers.len() == 5 - }) - }); - let c = found.expect("a candidate with one EH of KLLs read by all five queries"); - let (eh, ..) = windowed_builds(c) - .into_iter() - .find(|(_, _, form, _)| matches!(form, WindowForm::ExponentialHistogram { .. })) - .unwrap(); - let estimates = estimates_of(&c.dag, eh); + let c = run + .logical + .iter() + .find(|c| shares_segments(&c.dag, &c.query_roots)) + .expect("a candidate sharing window segments"); + let segments = windowed_builds(c); + assert_eq!(segments.len(), 5, "{}", c.id); + for (_, algorithm, form, _) in &segments { + assert_eq!(*algorithm, SketchAlgorithm::Kll); + assert_eq!(*form, WindowForm::Tumbling { length_ms: YEAR_MS }); + } + let read: BTreeSet = segments.iter().flat_map(|(.., r)| r.clone()).collect(); + assert_eq!(read, (0..5).collect()); + // Each query merges as many segments as its range has years. + for (q, (_, lookback, _)) in PATTERN_A.iter().enumerate() { + let covering = segments.iter().filter(|(.., r)| r.contains(&q)).count() as u64; + assert_eq!(covering, lookback / YEAR_MS, "q{}", q + 1); + } + let estimates: std::collections::BTreeMap<_, _> = segments + .iter() + .flat_map(|(build, ..)| estimates_of(&c.dag, *build)) + .collect(); assert_eq!(estimates.len(), 5, "one estimate per query"); for (estimate, statistic) in estimates { assert_eq!(statistic, SketchStatistic::Quantile { q: 0.99 }); @@ -211,16 +226,14 @@ fn stage1_a_window_composition_adds_one_eh_for_all_five() { } } -/// The shared EH candidate is added next to the independent KLL candidates, not instead of them. +/// The shared-segment candidate is added next to the independent KLL candidates, not instead of them. #[test] -#[ignore = "needs Pass 2 window composition (#580): Exponential Histogram"] fn stage1_a_keeps_independent_and_shared_window_summaries() { let run = run_a(); - let shared = run.logical.iter().any(|c| { - windowed_builds(c).iter().any(|(_, _, f, r)| { - matches!(f, WindowForm::ExponentialHistogram { .. }) && r.len() == 5 - }) - }); + let shared = run + .logical + .iter() + .any(|c| shares_segments(&c.dag, &c.query_roots)); let independent = run.logical.iter().any(|c| { let builds = windowed_builds(c); builds.len() == 5 @@ -233,7 +246,7 @@ fn stage1_a_keeps_independent_and_shared_window_summaries() { /// Pass 2 adds a candidate for each way of grouping two queries onto one shared window summary. #[test] -#[ignore = "needs Pass 2 window composition (#580): partial groupings"] +#[ignore = "partial and pairwise groupings are not generated, only all-shared segments (Q62)"] fn stage1_a_window_composition_groups_every_pair() { let run = run_a(); let groups: BTreeSet<_> = run @@ -282,7 +295,11 @@ fn stage3_a_shared_scan_is_not_costlier() { }; let cost = |c: &Logical| run.cost(&run.physical_of(c).next().unwrap().id); let mut compared = 0; - for shared in run.logical.iter().filter(|c| c.shared_input) { + for shared in run + .logical + .iter() + .filter(|c| c.shared_input && !shares_segments(&c.dag, &c.query_roots)) + { let separate = run .logical .iter() diff --git a/crates/integration-tests/tests/planner_layering_example4.rs b/crates/integration-tests/tests/planner_layering_example4.rs index fec55479..76e6aaa9 100644 --- a/crates/integration-tests/tests/planner_layering_example4.rs +++ b/crates/integration-tests/tests/planner_layering_example4.rs @@ -98,6 +98,22 @@ fn eh_options(run: &Run) -> BTreeMap BTreeMap { + let logical = run + .logical + .iter() + .find(|c| shares_segments(&c.dag, &c.query_roots)) + .expect("a candidate sharing window segments"); + run.physical_of(logical) + .map(|p| { + let (segment, ..) = sketch_builds(&p.dag)[0]; + (materialization(p, segment), (p.id.clone(), segment)) + }) + .collect() +} + fn tumbling() -> Option { Some(WindowForm::Tumbling { length_ms: PATTERN_B_INTERVAL_MS, @@ -165,7 +181,7 @@ fn stage2_materialized_output_covers_its_consumers() { /// The shared EH candidate yields A1 (query time, kept), A2 (ingestion time) and A3 (not materialized). #[test] -#[ignore = "needs Pass 2 window composition (#580): Exponential Histogram; and A3, not materialized for several consumers (Q44)"] +#[ignore = "shared segments (Q60) are not kept for one batch, and A3 is not generated (Q44); a one-off batch has no ingestion-time option (Q61)"] fn stage2_a_shared_eh_has_three_materialization_options() { let found: BTreeSet<_> = eh_options(&run_promql(&once())).into_keys().collect(); assert_eq!( @@ -176,43 +192,53 @@ fn stage2_a_shared_eh_has_three_materialization_options() { /// With data at rest, A2 is not generated; A1 and A3 remain. #[test] -#[ignore = "needs Pass 2 window composition (#580): Exponential Histogram"] +#[ignore = "shared segments (Q60) are not kept for one batch, and A3 is not generated (Q44)"] fn stage2_a_at_rest_drops_the_ingestion_time_option() { let found: BTreeSet<_> = eh_options(&run_promql(&at_rest())).into_keys().collect(); assert_eq!(found, BTreeSet::from([QueryTimeKept, NotMaterialized])); } -/// A materialized shared EH is one build node read by all five queries, charged once. +/// The shared segments (Q60, in place of the spec's EH) are five builds, +/// each built once, read by all five queries through their merges and +/// charged once, in every materialization: at query time for the one-off +/// batch, and also at ingestion time when it repeats monthly (Q61). #[test] -#[ignore = "needs Pass 2 window composition (#580): Exponential Histogram"] -fn stage2_a_materialized_eh_is_built_once_for_all_consumers() { - let run = run_promql(&once()); - for (m, (id, build)) in eh_options(&run) { - if m == NotMaterialized { - continue; +fn stage2_a_segments_are_built_once_for_all_consumers() { + for (workload, expected) in [ + (once(), BTreeSet::from([NotMaterialized])), + (monthly(), BTreeSet::from([IngestionTime, NotMaterialized])), + ] { + let run = run_promql(&workload); + let options = segment_options(&run); + assert_eq!(options.keys().copied().collect::>(), expected); + for (id, _) in options.values() { + let p = run.physical(id); + let segments = sketch_builds(&p.dag); + assert_eq!(segments.len(), 5, "{id}"); + let read: BTreeSet = segments + .iter() + .flat_map(|(b, ..)| readers(&p.dag, &p.query_roots, *b)) + .collect(); + assert_eq!(read.len(), 5, "{id}"); + let merges = p + .dag + .nodes + .iter() + .filter(|n| matches!(n.payload, Operator::ASAP(ASAPOp::SummaryMerge { .. }))) + .count(); + assert_eq!(merges, 5, "{id}"); + if let Some(cost) = run.selection.costs.get(id) { + for (b, ..) in &segments { + assert!(cost.per_node.contains_key(b), "{id}"); + } + } } - let p = run.physical(&id); - let ehs = sketch_builds(&p.dag) - .into_iter() - .filter(|(b, ..)| { - matches!( - window_form(&p.dag, *b), - WindowForm::ExponentialHistogram { .. } - ) - }) - .count(); - assert_eq!(ehs, 1, "{id}"); - assert_eq!(readers(&p.dag, &p.query_roots, build).len(), 5, "{id}"); - assert!( - run.selection.costs[&id].per_node.contains_key(&build), - "{id}" - ); } } /// A3 rebuilds the EH for each of the five queries, so it costs more than A1, which builds it once. #[test] -#[ignore = "needs Pass 2 window composition (#580): Exponential Histogram; and A3 (Q44)"] +#[ignore = "A3, not materialized for several consumers, is not generated (Q44)"] fn stage3_a_rebuilding_per_query_costs_more_than_building_once() { let run = run_promql(&once()); let options = eh_options(&run); @@ -227,7 +253,7 @@ fn stage3_a_rebuilding_per_query_costs_more_than_building_once() { /// As given (run once, ad hoc), A1 is the cheapest of the three: A2 maintains years of history for one batch. #[test] -#[ignore = "needs Pass 2 window composition (#580): Exponential Histogram"] +#[ignore = "a one-off batch has no ingestion-time option (Q61), and A3 is not generated (Q44)"] fn stage3_a_once_adhoc_prefers_the_query_time_eh() { let run = run_promql(&once()); let options = eh_options(&run); @@ -239,7 +265,7 @@ fn stage3_a_once_adhoc_prefers_the_query_time_eh() { /// Repeated monthly and predictable, A2's maintenance is shared by many batches, so A2 gains on A1. #[test] -#[ignore = "needs Pass 2 window composition (#580): Exponential Histogram, maintainable over a monthly window"] +#[ignore = "a one-off batch has no ingestion-time option to compare with (Q61)"] fn stage3_a_monthly_amortizes_ingestion_time_maintenance() { let ratio = |workload: PlanningWorkload| { let run = run_promql(&workload); diff --git a/crates/logical-optimizer/src/pass1/logical_candidates.rs b/crates/logical-optimizer/src/pass1/logical_candidates.rs index 1e57af85..90a1070a 100644 --- a/crates/logical-optimizer/src/pass1/logical_candidates.rs +++ b/crates/logical-optimizer/src/pass1/logical_candidates.rs @@ -697,7 +697,10 @@ fn realize( }; let state = match window { WindowForm::Whole => build(child)?, - WindowForm::Tumbling { pane_ms } => tumbling_state(&child, pane_ms, build)?, + WindowForm::Tumbling { pane_ms } + | WindowForm::Segments { + segment_ms: pane_ms, + } => tumbling_state(&child, pane_ms, build)?, }; let evaluation = match query { Some(query) => ASAPOp::SummaryEstimate { diff --git a/crates/logical-optimizer/src/pass2/identical_expressions.rs b/crates/logical-optimizer/src/pass2/identical_expressions.rs index c8b21f46..c6aa85b4 100644 --- a/crates/logical-optimizer/src/pass2/identical_expressions.rs +++ b/crates/logical-optimizer/src/pass2/identical_expressions.rs @@ -16,7 +16,7 @@ use asap_types::ir::{OperatorNode, QueryRoot}; use asap_types::workload::{MetricType, RootDemand}; use super::summary_capability::share_summary_capability; -use super::window_composition::add_window_forms; +use super::window_composition::{add_window_forms, share_window_segments}; use crate::pass1::logical_candidates::{ enumerate_local_logical_candidates, LocalLogicalCandidates, LogicalCandidateError, }; @@ -34,6 +34,11 @@ pub enum Sharing { /// are sized for their strictest consumer, so composition builds /// identical producers, which are merged. SummaryCapability, + /// The shared-segment rule on top of the identical-expression rule + /// ([`share_window_segments`]): queries over windows of one scan merge + /// one shared summary per segment, and the identical segments are merged + /// after composition. + WindowSegments, } impl Sharing { @@ -56,7 +61,9 @@ pub struct SharingVariant { /// last is skipped when it would repeat the identical-expression variant. /// In every variant, the window-composition rule adds tumbling forms of /// mergeable alternatives for repeating queries (`demand[i]` is the demand -/// of `roots[i]`; a root without one gets none). +/// of `roots[i]`; a root without one gets none). Last, the shared-segment +/// variant when windows of one scan can share segments: only that form, not +/// the tumbling forms, for the targets it groups. pub fn stage1_logical_candidates( roots: Vec<(Id, QueryRoot)>, metric_types: &BTreeMap, @@ -74,6 +81,7 @@ pub fn stage1_logical_candidates( }); } let base = &variants.last().expect("the independent variant").inventory; + let segments = share_window_segments(base, demand)?; if let Some(capability) = share_summary_capability(base)? { if capability.resized || variants.len() == 1 { variants.push(SharingVariant { @@ -85,6 +93,12 @@ pub fn stage1_logical_candidates( for variant in &mut variants { add_window_forms(&mut variant.inventory, demand); } + if let Some(inventory) = segments { + variants.push(SharingVariant { + sharing: Sharing::WindowSegments, + inventory, + }); + } Ok(variants) } diff --git a/crates/logical-optimizer/src/pass2/mod.rs b/crates/logical-optimizer/src/pass2/mod.rs index 301c1d44..09f37094 100644 --- a/crates/logical-optimizer/src/pass2/mod.rs +++ b/crates/logical-optimizer/src/pass2/mod.rs @@ -2,7 +2,8 @@ //! strictest accuracy requirement of its readers. The stage pipeline applies //! the identical-expression rule ([`identical_expressions`]), the //! summary-capability rule ([`summary_capability`]) and the -//! window-composition rule's tumbling windows ([`window_composition`]). +//! window-composition rule's tumbling windows and shared segments +//! ([`window_composition`]). pub mod identical_expressions; pub mod reconciliation; diff --git a/crates/logical-optimizer/src/pass2/summary_capability.rs b/crates/logical-optimizer/src/pass2/summary_capability.rs index 0aa1d923..3c0835ba 100644 --- a/crates/logical-optimizer/src/pass2/summary_capability.rs +++ b/crates/logical-optimizer/src/pass2/summary_capability.rs @@ -98,7 +98,7 @@ fn key(node: &OperatorNode) -> Option> { } /// `intent` with its accuracy requirement replaced. -fn with_accuracy(intent: &AggIntent, target: AccuracyTarget) -> AggIntent { +pub(crate) fn with_accuracy(intent: &AggIntent, target: AccuracyTarget) -> AggIntent { let mut intent = intent.clone(); match &mut intent { AggIntent::Quantile { accuracy, .. } @@ -114,7 +114,9 @@ fn with_accuracy(intent: &AggIntent, target: AccuracyTarget) -> AggIntent { /// The requirement that dominates every one in `targets`: the smallest ε /// and the smallest δ. -fn strictest<'a>(targets: impl IntoIterator) -> AccuracyTarget { +pub(crate) fn strictest<'a>( + targets: impl IntoIterator, +) -> AccuracyTarget { let (epsilon, delta) = targets .into_iter() .map(accuracy_budget) diff --git a/crates/logical-optimizer/src/pass2/window_composition.rs b/crates/logical-optimizer/src/pass2/window_composition.rs index be234fe3..2a4cc94a 100644 --- a/crates/logical-optimizer/src/pass2/window_composition.rs +++ b/crates/logical-optimizer/src/pass2/window_composition.rs @@ -16,21 +16,29 @@ //! relative time is `(-(o + (i + 1)·w), -(o + i·w)]` (W2). A `SummaryMerge` //! combines the panes (W1), which derivation accepts only for disjoint //! panes; a tumbling form is used only when the merged coverage is exactly -//! the whole-window state's. Stage 2 still -//! runs every node at query time, so until it plans materialization each -//! evaluation rebuilds all panes. +//! the whole-window state's. Stage 2 decides whether the panes are rebuilt +//! at each evaluation, maintained at ingestion time or kept. +//! +//! **Shared segments (Q60, Example 3 Pattern A).** Queries reading windows +//! of one scan at different offsets ([`share_window_segments`]) share one +//! summary per segment of a grid that every window's boundaries lie on, and +//! each query merges the segments its range covers. use std::collections::HashMap; use std::rc::Rc; use std::time::Duration; -use asap_types::ir::operator::{Source, TimeShift}; +use asap_types::ir::operator::{AggIntent, Reduction, Source, TimeShift}; use asap_types::ir::schema::{FieldDataType, GroupingStrategy}; use asap_types::ir::{ASAPOp, NonASAPOp, Operator, OperatorNode, QueryRoot, TimeRangeKind}; +use asap_types::types::AccuracyTarget; use asap_types::workload::{QueryRecurrence, RepeatedDemand, RootDemand}; -use crate::pass1::logical_candidates::{LocalLogicalCandidates, LogicalCandidateError}; -use crate::pass1::replacement::Realization; +use super::summary_capability::{strictest, with_accuracy}; +use crate::pass1::logical_candidates::{ + local_realizations_for_intent, LocalLogicalCandidates, LogicalCandidateError, +}; +use crate::pass1::replacement::{accuracy_target, Realization}; /// How a summary alternative covers its query's window. #[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] @@ -40,6 +48,9 @@ pub enum WindowForm { Whole, /// Back-to-back panes of `pane_ms`, merged at each evaluation. Tumbling { pane_ms: u64 }, + /// Segments of `segment_ms` shared by several queries' windows + /// ([`share_window_segments`]); built as tumbling panes are. + Segments { segment_ms: u64 }, } impl WindowForm { @@ -50,12 +61,20 @@ impl WindowForm { WindowForm::Tumbling { pane_ms } => { Some(format!("tumbling {} panes", duration_label(pane_ms))) } + WindowForm::Segments { segment_ms } => { + Some(format!("{} segments", duration_label(segment_ms))) + } } } } fn duration_label(ms: u64) -> String { - for (unit, size) in [("h", 3_600_000), ("m", 60_000), ("s", 1_000)] { + for (unit, size) in [ + ("d", 86_400_000), + ("h", 3_600_000), + ("m", 60_000), + ("s", 1_000), + ] { if ms.is_multiple_of(size) { return format!("{}{unit}", ms / size); } @@ -226,6 +245,158 @@ pub fn add_window_forms(inventory: &mut LocalLogicalCandidates, demand: } } +/// The shared-segment rule (Q60) over a Pass 1 inventory: targets that +/// estimate the same statistic (up to accuracy) with the same grouping and +/// no filters, over windows of one scan, share one summary per segment. The +/// segment width is the gcd of every window's lookback and offset, so each +/// window `[-(o + lookback), -o)` is a whole number of segments, and each +/// query merges the segments its range covers. Merging disjoint summaries is +/// exact (`tumbling_state` checks the derived cover), so no +/// accuracy is split: every segment is sized for the strictest consumer, as +/// the summary-capability rule sizes a shared summary. +/// +/// Only the all-shared form is offered (Q62): in the returned inventory each +/// grouped target has the segment form of its first mergeable summary as its +/// only alternative, so composition builds identical segments, which the +/// identical-expression rule merges. A group is skipped when its windows are +/// all the same (nothing to split), the union of its windows has more than +/// [`MAX_PANES`] segments, or its queries recur at different cadences. +/// `None` when no group remains. +pub fn share_window_segments( + inventory: &LocalLogicalCandidates, + demand: &[RootDemand], +) -> Result>, LogicalCandidateError> { + // The cadences of the roots reaching each node. + let mut cadences: HashMap<*const OperatorNode, Vec>> = HashMap::new(); + for (index, (_, root)) in inventory.roots.iter().enumerate() { + let cadence = demand + .get(index) + .and_then(|d| evaluation_cadence_ms(&d.recurrence)); + let operators = match root { + QueryRoot::Operator(node) => vec![node], + QueryRoot::Scalar(expr) => expr.operator_refs(), + }; + for node in operators.into_iter().flat_map(OperatorNode::reachable) { + cadences.entry(Rc::as_ptr(&node)).or_default().push(cadence); + } + } + // Per target: the estimate without its accuracy, the window, and the + // grouping, when the rule applies to it. + let keyed: Vec, &Reduction)>> = inventory + .targets + .iter() + .map(|t| match t.target.non_asap() { + Some(NonASAPOp::Aggregate { + reduction, + measures, + filters, + having: None, + child, + .. + }) if filters.iter().all(Option::is_none) => match measures.as_slice() { + [intent] + if accuracy_target(intent) + .is_some_and(|a| !matches!(a, AccuracyTarget::Exact)) + && window(child).is_some_and(|w| w.offset_ms >= 0) => + { + Some(( + with_accuracy(intent, AccuracyTarget::Exact), + window(child)?, + reduction, + )) + } + _ => None, + }, + _ => None, + }) + .collect(); + let same_scan = |a: &Rc, b: &Rc| Rc::ptr_eq(a, b) || a == b; + let mut out = inventory.clone(); + let mut grouped = vec![false; keyed.len()]; + let mut any = false; + for leader in 0..keyed.len() { + let Some((intent, first, reduction)) = &keyed[leader] else { + continue; + }; + if grouped[leader] { + continue; + } + let members: Vec = (leader..keyed.len()) + .filter(|&t| { + keyed[t].as_ref().is_some_and(|(i, w, r)| { + i == intent && r == reduction && same_scan(w.scan, first.scan) + }) + }) + .collect(); + for &t in &members { + grouped[t] = true; + } + let windows: Vec<&Window<'_>> = members + .iter() + .map(|&t| &keyed[t].as_ref().expect("keyed").1) + .collect(); + let bounds = |w: &Window<'_>| (w.offset_ms as u64, w.offset_ms as u64 + w.lookback_ms); + let reaching: Vec> = members + .iter() + .flat_map(|&t| { + cadences + .get(&Rc::as_ptr(&inventory.targets[t].target)) + .into_iter() + .flatten() + .copied() + }) + .collect(); + if members.len() < 2 + || windows.iter().all(|w| bounds(w) == bounds(windows[0])) + || reaching.iter().any(|c| *c != reaching[0]) + { + continue; + } + let width = windows + .iter() + .fold(0, |g, w| gcd(gcd(g, w.lookback_ms), w.offset_ms as u64)); + let start = windows.iter().map(|w| bounds(w).0).min().unwrap_or(0); + let end = windows.iter().map(|w| bounds(w).1).max().unwrap_or(0); + if width == 0 + || (end - start) / width > MAX_PANES + || !windows.iter().all(|w| panes_tile_window(w, width)) + { + continue; + } + let requirements: Vec = members + .iter() + .map(|&t| single_intent(inventory, t)) + .filter_map(|intent| accuracy_target(intent).cloned()) + .collect(); + let target = strictest(requirements.iter()); + let sized = with_accuracy(single_intent(inventory, leader), target); + let grouping = GroupingStrategy::default(); + let Some(summary) = local_realizations_for_intent(&sized)? + .into_iter() + .find(|r| family(r, &grouping).is_some_and(|f| f.family_merges())) + else { + continue; + }; + // The members' estimates differ only in accuracy, so the one sized + // summary serves each of them. + for &t in &members { + out.targets[t].alternatives = vec![summary.clone()]; + out.targets[t].absorbs = vec![None]; + out.targets[t].windows = vec![WindowForm::Segments { segment_ms: width }]; + out.targets[t].groupings = vec![grouping.clone()]; + } + any = true; + } + Ok(any.then_some(out)) +} + +fn single_intent(inventory: &LocalLogicalCandidates, t: usize) -> &AggIntent { + match inventory.targets[t].target.non_asap() { + Some(NonASAPOp::Aggregate { measures, .. }) => &measures[0], + _ => unreachable!("grouped targets are single-measure aggregates"), + } +} + fn family(realization: &Realization, grouping: &GroupingStrategy) -> Option { match realization { Realization::ExactAggregate { kind, params } => { @@ -571,4 +742,99 @@ mod tests { // Panes of 2m over a 5m window leave the oldest minute uncovered. assert!(tumbling_state(child, 120_000, kll_state(target)).is_err()); } + + const YEAR_MS: u64 = 365 * 86_400_000; + + /// Example 3's Pattern A: five p99 reports over [5y], [1y], [1y offset + /// 1y], [1y offset 2y] and [3y offset 2y], one ad hoc batch. + fn pattern_a() -> Vec<(usize, QueryRoot)> { + [ + "latency[5y]", + "latency[1y]", + "latency[1y] offset 1y", + "latency[1y] offset 2y", + "latency[3y] offset 2y", + ] + .iter() + .enumerate() + .map(|(i, window)| { + let root = lower_promql( + &format!("quantile_over_time(0.99, {window})"), + AccuracyTarget::EpsilonDelta { + epsilon: 0.005, + delta: 0.01, + }, + ); + let root = asap_types::ir::schema_support::with_promql_series_identity(&root).unwrap(); + (i, QueryRoot::Operator(root)) + }) + .collect() + } + + /// The windows' boundaries lie on a 1-year grid, so the five queries + /// share five 1-year KLL segments: the variant offers each query only + /// that form, and composition with the identical-expression merge builds + /// each segment once, read by the merges of every query covering it. + #[test] + fn pattern_a_shares_five_one_year_segments() { + use crate::pass1::logical_candidates::compose_logical_candidate; + use crate::pass2::identical_expressions::{ + share_identical_expressions, stage1_logical_candidates, Sharing, + }; + let variants = stage1_logical_candidates(pattern_a(), &BTreeMap::new(), &[]).unwrap(); + let segments = variants + .iter() + .find(|v| v.sharing == Sharing::WindowSegments) + .expect("a shared-segment variant"); + let form = WindowForm::Segments { + segment_ms: YEAR_MS, + }; + assert_eq!(form.label().as_deref(), Some("365d segments")); + for target in &segments.inventory.targets { + assert_eq!(target.windows, [form]); + assert!(matches!( + &target.alternatives[..], + [Realization::Sketch(kind)] if *kind.algorithm() == SketchAlgorithm::Kll + )); + } + let composed = compose_logical_candidate(&segments.inventory, &[0; 5]).unwrap(); + let composed = share_identical_expressions(&composed).unwrap(); + let roots: Vec<_> = composed + .iter() + .map(|(_, root)| match root { + QueryRoot::Operator(node) => node.clone(), + QueryRoot::Scalar(_) => unreachable!(), + }) + .collect(); + let mut builds = std::collections::HashSet::new(); + let mut merges = std::collections::HashSet::new(); + for node in roots.iter().flat_map(OperatorNode::reachable) { + match &node.operator { + Operator::ASAP(ASAPOp::SummaryAgg { .. }) => builds.insert(Rc::as_ptr(&node)), + Operator::ASAP(ASAPOp::SummaryMerge { .. }) => merges.insert(Rc::as_ptr(&node)), + _ => false, + }; + } + assert_eq!((builds.len(), merges.len()), (5, 5)); + } + + /// Queries that recur at different cadences, or read the same window, + /// share no segments. + #[test] + fn segments_need_one_cadence_and_different_windows() { + let roots = pattern_a(); + let inventory = + enumerate_local_logical_candidates(roots[..2].to_vec(), &BTreeMap::new()).unwrap(); + assert!(share_window_segments(&inventory, &[]).unwrap().is_some()); + let demand = [every(60_000), every(120_000)]; + assert!(share_window_segments(&inventory, &demand) + .unwrap() + .is_none()); + let same = enumerate_local_logical_candidates( + vec![roots[1].clone(), (1, roots[1].1.clone())], + &BTreeMap::new(), + ) + .unwrap(); + assert!(share_window_segments(&same, &[]).unwrap().is_none()); + } } diff --git a/crates/plan-selection/src/lib.rs b/crates/plan-selection/src/lib.rs index c1bc60a9..834abdeb 100644 --- a/crates/plan-selection/src/lib.rs +++ b/crates/plan-selection/src/lib.rs @@ -1131,7 +1131,8 @@ pub struct StagePipelineRun { /// The #509 stage pipeline over `roots`: Stage 1 (Pass 1 and Pass 2's /// identical-expression and summary-capability rules and the -/// window-composition rule's tumbling panes, from each root's `demand`), +/// window-composition rule's tumbling panes and shared segments, from each +/// root's `demand`), /// Stage 2 and Stage 3. The facade and the /// `stage_pipeline` devtool both run this. `display` builds and prices up to /// that many candidates for display as well (0: none). diff --git a/tools/dag-viewer/examples.json b/tools/dag-viewer/examples.json index 9ff42877..ce6b2aa7 100644 --- a/tools/dag-viewer/examples.json +++ b/tools/dag-viewer/examples.json @@ -18,7 +18,7 @@ "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. Exponential histograms are not implemented yet, so no plan shares window summaries across the reports." + "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." }, { "id": "example3b", @@ -32,7 +32,7 @@ "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." + "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." }, { "id": "example4b",