From f3671b56a6c8157de8f673f3217d2b6ad9873717 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Mon, 5 Oct 2026 03:28:08 +0000 Subject: [PATCH 1/2] feat(pass2): the planner chooses the segment grid (Q67) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The shared-segment rule (Q60) offered one fixed grid, the gcd of every window's lookback and offset. It now offers two segmentations, each as its own Stage 1 variant, and Stage 3's cost decides: - split only at the windows' own boundaries (fewest segments); for Example 3a four segments, the oldest two years wide; - the even gcd grid (3a: five 1-year segments), when it differs. WindowForm::Segments carries the grid unit and a bitmask of where each segment ends, so a target's segments may differ in width; tumbling_state and the coverage proof take the widths. Sharing::WindowSegments carries the number of distinct segments, labeled "· shared segments ×4" / "×5". Each segment is still sized for the group's strictest reader. Stage 2 does not maintain a merge of panes of different widths at ingestion time or keep it: it is no stream of panes. So in 4a only the even grid has an ingestion-time option. Example 3a and 4a select the four boundary segments: they build the same rows as the five 1-year ones and merge fewer states (3a: 18,104.0039 vs 18,104.0044 per second). Co-Authored-By: Claude Opus 5.5 --- crates/devtools/src/bin/stage_pipeline.rs | 15 +- crates/devtools/tests/stage_pipeline.rs | 19 +- .../tests/planner_layering_example3.rs | 123 ++++-- .../tests/planner_layering_example4.rs | 38 +- .../src/pass1/logical_candidates.rs | 7 +- .../src/pass2/identical_expressions.rs | 14 +- .../src/pass2/window_composition.rs | 356 ++++++++++++------ .../physical-optimizer/src/materialization.rs | 9 +- tools/dag-viewer/examples.json | 4 +- 9 files changed, 398 insertions(+), 187 deletions(-) diff --git a/crates/devtools/src/bin/stage_pipeline.rs b/crates/devtools/src/bin/stage_pipeline.rs index 2a1fa813..76459821 100644 --- a/crates/devtools/src/bin/stage_pipeline.rs +++ b/crates/devtools/src/bin/stage_pipeline.rs @@ -22,8 +22,9 @@ // input"), and, when the summary-capability rule applies, again with // 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 +// with one summary per shared segment, once per segmentation, by its +// number of segments ("· shared segments ×4"); 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"; // - stage2_physical_asap: per logical candidate, its physical candidates @@ -265,11 +266,11 @@ fn stage_pipeline( let inventory = &variant.inventory; let index = offset + choice_index(inventory, &candidate.choice) + 1; let mut label = label(inventory, &target_owners(inventory), &candidate.choice); - label += match candidate.sharing { - Sharing::Independent => "", - Sharing::IdenticalExpressions => " · shared input", - Sharing::SummaryCapability => " · shared summary", - Sharing::WindowSegments => " · shared segments", + label += &match candidate.sharing { + Sharing::Independent => String::new(), + Sharing::IdenticalExpressions => " · shared input".into(), + Sharing::SummaryCapability => " · shared summary".into(), + Sharing::WindowSegments { segments } => format!(" · shared segments ×{segments}"), }; let written = candidate .physical diff --git a/crates/devtools/tests/stage_pipeline.rs b/crates/devtools/tests/stage_pipeline.rs index 9b2afcd8..5c4bd367 100644 --- a/crates/devtools/tests/stage_pipeline.rs +++ b/crates/devtools/tests/stage_pipeline.rs @@ -183,13 +183,14 @@ fn example3b_lists_tumbling_candidates() { } } -/// 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. +/// Example 4, Pattern A repeated monthly: the same 488 logical candidates +/// as the ad hoc batch (3a), 486 independent or sharing input plus the two +/// sharing segments (Q60, Q67): four split at the windows' boundaries and +/// five 1-year ones. The windows (1–5 y) are longer than the month between +/// runs and no pane width fits, so only the even 1-year segments can be +/// maintained at ingestion time (Q61); the uneven ones are no stream of +/// panes: 489 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_only_segments_maintainable() { let once = generate(&[ @@ -207,14 +208,14 @@ fn example4a_repeats_monthly_with_only_segments_maintainable() { let physical = monthly["stage2_physical_asap"]["candidates"] .as_array() .unwrap(); - assert_eq!(physical.len(), 488); + assert_eq!(physical.len(), 489); let maintained: Vec<_> = physical .iter() .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")); + assert!(maintained[0].ends_with("shared segments ×5 · ingestion time: Kll ×5 panes")); let selected = |d: &Value| { d["stage3_selection"]["selected"] .as_str() diff --git a/crates/integration-tests/tests/planner_layering_example3.rs b/crates/integration-tests/tests/planner_layering_example3.rs index 2eec9dfa..36f87143 100644 --- a/crates/integration-tests/tests/planner_layering_example3.rs +++ b/crates/integration-tests/tests/planner_layering_example3.rs @@ -180,52 +180,103 @@ fn stage1_a_identical_expression_rule_shares_only_raw_input() { } } -/// 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. +/// The planner chooses the segment grid (Q67): one candidate splits the +/// windows at their boundaries into four KLL segments over [T − 5y, T], the +/// oldest two years wide; another uses the even grid of five 1-year +/// segments. Either 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] fn stage1_a_shared_segments_serve_all_five() { let run = run_a(); - let c = run + let candidates: Vec<_> = 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 + .filter(|c| shares_segments(&c.dag, &c.query_roots)) + .collect(); + let mut widths: Vec> = candidates .iter() - .flat_map(|(build, ..)| estimates_of(&c.dag, *build)) + .map(|c| { + let mut widths: Vec = windowed_builds(c) + .iter() + .map(|(_, _, form, _)| match form { + WindowForm::Tumbling { length_ms } => length_ms / YEAR_MS, + other => panic!("{}: {other:?}", c.id), + }) + .collect(); + widths.sort(); + widths + }) .collect(); - assert_eq!(estimates.len(), 5, "one estimate per query"); - for (estimate, statistic) in estimates { - assert_eq!(statistic, SketchStatistic::Quantile { q: 0.99 }); - let merged = c.dag.producers(estimate).into_iter().any(|p| { - matches!( - c.dag.payload(p), - Operator::ASAP(ASAPOp::SummaryMerge { .. }) - ) - }); - assert!( - merged, - "{}: estimate {estimate:?} does not read a merge", - c.id - ); + widths.sort(); + assert_eq!(widths, [vec![1, 1, 1, 1, 1], vec![1, 1, 1, 2]]); + for c in candidates { + let segments = windowed_builds(c); + for (_, algorithm, ..) in &segments { + assert_eq!(*algorithm, SketchAlgorithm::Kll); + } + let read: BTreeSet = segments.iter().flat_map(|(.., r)| r.clone()).collect(); + assert_eq!(read, (0..5).collect()); + // The segments each query merges add up to its range. + for (q, (_, lookback, _)) in PATTERN_A.iter().enumerate() { + let covered: u64 = segments + .iter() + .filter(|(.., r)| r.contains(&q)) + .map(|(_, _, form, _)| match form { + WindowForm::Tumbling { length_ms } => *length_ms, + _ => unreachable!(), + }) + .sum(); + assert_eq!(covered, *lookback, "{} q{}", c.id, 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 }); + let merged = c + .dag + .producers(estimate) + .into_iter() + .any(|p| matches!(c.dag.payload(p), Operator::ASAP(ASAPOp::SummaryMerge { .. }))); + assert!( + merged, + "{}: estimate {estimate:?} does not read a merge", + c.id + ); + } } } +/// Stage 3 decides the grid: the four boundary segments cost as much to +/// build as the five 1-year ones (each row lands in one segment) but merge +/// fewer states, so they are cheaper, and cheapest overall. +#[test] +fn stage3_a_selects_the_cheaper_segment_grid() { + let run = run_a(); + let mut costs: Vec<(usize, f64, &str)> = run + .logical + .iter() + .filter(|c| shares_segments(&c.dag, &c.query_roots)) + .map(|c| { + let p = run.physical_of(c).next().unwrap(); + ( + windowed_builds(c).len(), + run.cost(&p.id).unwrap(), + p.id.as_str(), + ) + }) + .collect(); + costs.sort_by_key(|(segments, ..)| *segments); + let [(4, four, selected), (5, five, _)] = costs[..] else { + panic!("{costs:?}"); + }; + assert!(four < five, "×4 {four} vs ×5 {five}"); + assert_eq!(run.selection.selected, selected); +} + /// The shared-segment candidate is added next to the independent KLL candidates, not instead of them. #[test] fn stage1_a_keeps_independent_and_shared_window_summaries() { diff --git a/crates/integration-tests/tests/planner_layering_example4.rs b/crates/integration-tests/tests/planner_layering_example4.rs index 541eafd6..55d0cbd4 100644 --- a/crates/integration-tests/tests/planner_layering_example4.rs +++ b/crates/integration-tests/tests/planner_layering_example4.rs @@ -98,13 +98,17 @@ fn eh_options(run: &Run) -> BTreeMap BTreeMap { +/// The physical candidates of the shared-segment candidate (Q60) with +/// `count` segments (Q67), keyed by the materialization of its segments, +/// with the newest segment's id. +fn segment_options( + run: &Run, + count: usize, +) -> BTreeMap { let logical = run .logical .iter() - .find(|c| shares_segments(&c.dag, &c.query_roots)) + .find(|c| shares_segments(&c.dag, &c.query_roots) && sketch_builds(&c.dag).len() == count) .expect("a candidate sharing window segments"); run.physical_of(logical) .map(|p| { @@ -198,23 +202,31 @@ fn stage2_a_at_rest_drops_the_ingestion_time_option() { assert_eq!(found, BTreeSet::from([QueryTimeKept, NotMaterialized])); } -/// 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). +/// The shared segments (Q60, in place of the spec's EH), four split at the +/// windows' boundaries or five on the 1-year grid (Q67), are 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, for +/// the even grid only, also at ingestion time when it repeats monthly +/// (Q61). Uneven segments are no stream of panes. #[test] fn stage2_a_segments_are_built_once_for_all_consumers() { - for (workload, expected) in [ - (once(), BTreeSet::from([NotMaterialized])), - (monthly(), BTreeSet::from([IngestionTime, NotMaterialized])), + for (workload, count, expected) in [ + (once(), 4, BTreeSet::from([NotMaterialized])), + (once(), 5, BTreeSet::from([NotMaterialized])), + (monthly(), 4, BTreeSet::from([NotMaterialized])), + ( + monthly(), + 5, + BTreeSet::from([IngestionTime, NotMaterialized]), + ), ] { let run = run_promql(&workload); - let options = segment_options(&run); + let options = segment_options(&run, count); 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}"); + assert_eq!(segments.len(), count, "{id}"); let read: BTreeSet = segments .iter() .flat_map(|(b, ..)| readers(&p.dag, &p.query_roots, *b)) diff --git a/crates/logical-optimizer/src/pass1/logical_candidates.rs b/crates/logical-optimizer/src/pass1/logical_candidates.rs index 90a1070a..e499fd6c 100644 --- a/crates/logical-optimizer/src/pass1/logical_candidates.rs +++ b/crates/logical-optimizer/src/pass1/logical_candidates.rs @@ -697,10 +697,9 @@ fn realize( }; let state = match window { WindowForm::Whole => build(child)?, - WindowForm::Tumbling { pane_ms } - | WindowForm::Segments { - segment_ms: pane_ms, - } => tumbling_state(&child, pane_ms, build)?, + WindowForm::Tumbling { .. } | WindowForm::Segments { .. } => { + tumbling_state(&child, window, 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 c6aa85b4..cc46a718 100644 --- a/crates/logical-optimizer/src/pass2/identical_expressions.rs +++ b/crates/logical-optimizer/src/pass2/identical_expressions.rs @@ -37,8 +37,9 @@ pub enum Sharing { /// 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, + /// after composition. One variant per segmentation, by its number of + /// distinct segments (Q67). + WindowSegments { segments: usize }, } impl Sharing { @@ -62,8 +63,9 @@ pub struct SharingVariant { /// 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). 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. +/// variants when windows of one scan can share segments, one per +/// segmentation (Q67): only that form, not the tumbling forms, for the +/// targets it groups. pub fn stage1_logical_candidates( roots: Vec<(Id, QueryRoot)>, metric_types: &BTreeMap, @@ -93,9 +95,9 @@ pub fn stage1_logical_candidates( for variant in &mut variants { add_window_forms(&mut variant.inventory, demand); } - if let Some(inventory) = segments { + for (segments, inventory) in segments { variants.push(SharingVariant { - sharing: Sharing::WindowSegments, + sharing: Sharing::WindowSegments { segments }, inventory, }); } diff --git a/crates/logical-optimizer/src/pass2/window_composition.rs b/crates/logical-optimizer/src/pass2/window_composition.rs index 2a4cc94a..d1057155 100644 --- a/crates/logical-optimizer/src/pass2/window_composition.rs +++ b/crates/logical-optimizer/src/pass2/window_composition.rs @@ -16,15 +16,16 @@ //! 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 decides whether the panes are rebuilt -//! at each evaluation, maintained at ingestion time or kept. +//! 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::collections::{BTreeMap, BTreeSet, HashMap}; use std::rc::Rc; use std::time::Duration; @@ -48,21 +49,51 @@ 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 }, + /// Segments shared by several queries' windows + /// ([`share_window_segments`]); built as tumbling panes are. Counting + /// units of `unit_ms` back from the window's newest end, a segment ends + /// after unit `j + 1` when bit `j` of `ends` is set; the highest set bit + /// is the window's oldest end. + Segments { unit_ms: u64, ends: u64 }, } impl WindowForm { - /// E.g. "tumbling 1m panes"; `None` for the whole window. + /// E.g. "tumbling 1m panes", "365d segments" (even) or "365d+730d + /// segments"; + /// `None` for the whole window. pub fn label(self) -> Option { match self { WindowForm::Whole => None, WindowForm::Tumbling { pane_ms } => { Some(format!("tumbling {} panes", duration_label(pane_ms))) } - WindowForm::Segments { segment_ms } => { - Some(format!("{} segments", duration_label(segment_ms))) + WindowForm::Segments { .. } => { + let mut widths = self.widths(0); + if widths.iter().all(|w| *w == widths[0]) { + widths.truncate(1); + } + let widths: Vec<_> = widths.into_iter().map(duration_label).collect(); + Some(format!("{} segments", widths.join("+"))) + } + } + } + + /// The widths of the panes or segments a window of `lookback_ms` is + /// split into, newest first. Segments know their own window. + pub fn widths(self, lookback_ms: u64) -> Vec { + match self { + WindowForm::Whole => vec![lookback_ms], + WindowForm::Tumbling { pane_ms } => vec![pane_ms; (lookback_ms / pane_ms) as usize], + WindowForm::Segments { unit_ms, ends } => { + let mut widths = Vec::new(); + let mut start = 0; + for j in 0..64 { + if ends & (1 << j) != 0 { + widths.push((j + 1 - start) * unit_ms); + start = j + 1; + } + } + widths } } } @@ -177,11 +208,25 @@ pub fn pane_source(target: &OperatorNode) -> Option<&Rc> { } } -/// Whether panes of `pane_ms` tile the window exactly. That the panes' -/// derived coverage is the whole window is checked again when a tumbling -/// form is realized ([`tumbling_state`]). -fn panes_tile_window(window: &Window<'_>, pane_ms: u64) -> bool { - pane_ms > 0 && window.lookback_ms.is_multiple_of(pane_ms) +/// The panes of `widths` (newest first) over `window`: each one's offset +/// back from the evaluation, and its width. +fn panes(window: &Window<'_>, widths: &[u64]) -> Option> { + let mut offset_ms = window.offset_ms; + widths + .iter() + .map(|&width| { + let pane = (offset_ms, width); + offset_ms = offset_ms.checked_add(i64::try_from(width).ok()?)?; + Some(pane) + }) + .collect() +} + +/// Whether panes of `widths` tile the window exactly. That the panes' +/// derived coverage is the whole window is checked again when a form is +/// realized ([`tumbling_state`]). +fn panes_tile_window(window: &Window<'_>, widths: &[u64]) -> bool { + widths.iter().all(|&w| w > 0) && widths.iter().sum::() == window.lookback_ms } /// Add the tumbling form of every mergeable alternative whose target reads @@ -222,7 +267,8 @@ pub fn add_window_forms(inventory: &mut LocalLogicalCandidates, demand: else { continue; }; - if !panes_tile_window(&window, pane_ms) { + let form = WindowForm::Tumbling { pane_ms }; + if !panes_tile_window(&window, &form.widths(window.lookback_ms)) { continue; } let tumbling: Vec<(Realization, GroupingStrategy)> = target @@ -239,7 +285,7 @@ pub fn add_window_forms(inventory: &mut LocalLogicalCandidates, demand: for (alternative, grouping) in tumbling { target.alternatives.push(alternative); target.absorbs.push(None); - target.windows.push(WindowForm::Tumbling { pane_ms }); + target.windows.push(form); target.groupings.push(grouping); } } @@ -247,25 +293,31 @@ 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. +/// no filters, over windows of one scan, share one summary per segment, 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. +/// +/// The planner chooses the segment grid (Q67): two segmentations are +/// offered, each as its own inventory, and Stage 3's cost decides. One +/// splits the windows only at their own boundaries (the fewest segments); +/// the other is the even grid whose width is the gcd of every window's +/// lookback and offset, offered when it differs. Each window +/// `[-(o + lookback), -o)` is a whole number of segments in both. /// -/// 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. +/// Only the all-shared form is offered (Q62): in each 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`] grid units, or its queries recur at different +/// cadences. Returns each inventory with its number of distinct segments, +/// boundary split first; empty when no group remains. pub fn share_window_segments( inventory: &LocalLogicalCandidates, demand: &[RootDemand], -) -> Result>, LogicalCandidateError> { +) -> 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() { @@ -311,7 +363,12 @@ pub fn share_window_segments( }) .collect(); let same_scan = |a: &Rc, b: &Rc| Rc::ptr_eq(a, b) || a == b; - let mut out = inventory.clone(); + // Boundary split, then even grid: the inventory and its distinct + // segments, as `[offset, offset + width)`. + let mut out = [ + (inventory.clone(), BTreeSet::new()), + (inventory.clone(), BTreeSet::new()), + ]; let mut grouped = vec![false; keyed.len()]; let mut any = false; for leader in 0..keyed.len() { @@ -357,9 +414,29 @@ pub fn share_window_segments( .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)) + if width == 0 || (end - start) / width > MAX_PANES { + continue; + } + // Each window's segments end at the units where some window's + // boundary lies, or at every unit. + let boundaries: BTreeSet = windows + .iter() + .flat_map(|w| [bounds(w).0, bounds(w).1]) + .collect(); + let forms = |w: &Window<'_>| { + let units = w.lookback_ms / width; + let ends = (1..=units) + .filter(|j| boundaries.contains(&(bounds(w).0 + j * width))) + .fold(0u64, |ends, j| ends | 1 << (j - 1)); + let even = u64::MAX.checked_shr(64 - units as u32).unwrap_or(0); + [ends, even].map(|ends| WindowForm::Segments { + unit_ms: width, + ends, + }) + }; + if !windows + .iter() + .all(|w| forms(w).iter().all(|f| panes_tile_window(w, &f.widths(0)))) { continue; } @@ -379,15 +456,29 @@ pub fn share_window_segments( }; // 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()]; + for (&t, w) in members.iter().zip(&windows) { + for ((out, segments), form) in out.iter_mut().zip(forms(w)) { + out.targets[t].alternatives = vec![summary.clone()]; + out.targets[t].absorbs = vec![None]; + out.targets[t].windows = vec![form]; + out.targets[t].groupings = vec![grouping.clone()]; + segments.extend(panes(w, &form.widths(0)).expect("covered above")); + } } any = true; } - Ok(any.then_some(out)) + if !any { + return Ok(Vec::new()); + } + let [split, even] = out; + let mut variants = vec![split]; + if even.1 != variants[0].1 { + variants.push(even); + } + Ok(variants + .into_iter() + .map(|(inventory, segments)| (segments.len(), inventory)) + .collect()) } fn single_intent(inventory: &LocalLogicalCandidates, t: usize) -> &AggIntent { @@ -407,13 +498,14 @@ fn family(realization: &Realization, grouping: &GroupingStrategy) -> Option, - pane_ms: u64, + form: WindowForm, pane: impl Fn(Rc) -> Result, LogicalCandidateError>, ) -> Result, LogicalCandidateError> { let unsupported = || LogicalCandidateError::Unsupported("tumbling panes over this window"); @@ -421,13 +513,10 @@ pub(crate) fn tumbling_state( let Some(NonASAPOp::TimeRange { kind, .. }) = input.non_asap() else { return Err(unsupported()); }; - let width = i64::try_from(pane_ms).map_err(|_| unsupported())?; let mut panes = Vec::new(); - for i in 0..(window.lookback_ms / pane_ms) as i64 { - let offset_ms = window - .offset_ms - .checked_add(i * width) - .ok_or_else(unsupported)?; + for (offset_ms, width) in + self::panes(&window, &form.widths(window.lookback_ms)).ok_or_else(unsupported)? + { let shifted = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::TimeShift { child: Rc::clone(window.scan), shift: TimeShift { @@ -436,7 +525,7 @@ pub(crate) fn tumbling_state( }, }))?; let range = OperatorNode::new_shared(Operator::NonASAP(NonASAPOp::TimeRange { - range: Duration::from_millis(pane_ms), + range: Duration::from_millis(width), kind: *kind, child: shifted, }))?; @@ -457,11 +546,10 @@ mod tests { use super::*; use crate::pass1::logical_candidates::enumerate_local_logical_candidates; use crate::test_support::lower_promql; + use std::ops::Bound; use asap_types::ir::schema::SketchAlgorithm; use asap_types::types::AccuracyTarget; use asap_types::workload::{Predictability, RepetitionInterval}; - use std::collections::BTreeMap; - use std::ops::Bound; fn every(ms: u32) -> RootDemand { RootDemand { @@ -587,7 +675,9 @@ mod tests { let Some(NonASAPOp::Aggregate { child, .. }) = target.target.non_asap() else { panic!("aggregate") }; - let merged = tumbling_state(child, 60_000, kll_state(target)).unwrap(); + let merged = + tumbling_state(child, WindowForm::Tumbling { pane_ms: 60_000 }, kll_state(target)) + .unwrap(); let Operator::ASAP(ASAPOp::SummaryMerge { children }) = &merged.operator else { panic!("merge") }; @@ -676,18 +766,16 @@ mod tests { target: &crate::pass1::logical_candidates::LocalLogicalTarget, ) -> impl Fn(Rc) -> Result, LogicalCandidateError> + '_ { move |input| { - Ok(OperatorNode::new_shared(Operator::ASAP( - ASAPOp::SummaryAgg { - child: input, - family: family(&target.alternatives[1], &GroupingStrategy::default()).unwrap(), - input: asap_types::ir::schema::SummaryUpdate::column( - asap_types::ir::scalar::ColumnRef::SampleValue, - ), - reduction: asap_types::ir::operator::Reduction::PerEntity, - grouping: GroupingStrategy::default(), - filter: None, - }, - ))?) + Ok(OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryAgg { + child: input, + family: family(&target.alternatives[1], &GroupingStrategy::default()).unwrap(), + input: asap_types::ir::schema::SummaryUpdate::column( + asap_types::ir::scalar::ColumnRef::SampleValue, + ), + reduction: asap_types::ir::operator::Reduction::PerEntity, + grouping: GroupingStrategy::default(), + filter: None, + }))?) } } @@ -704,8 +792,9 @@ mod tests { offset_ms: 0, scan: &scan, }; - assert!(panes_tile_window(&window, 60_000)); - assert!(!panes_tile_window(&window, 70_000)); + assert!(panes_tile_window(&window, &[60_000; 5])); + assert!(panes_tile_window(&window, &[60_000, 120_000, 120_000])); + assert!(!panes_tile_window(&window, &[70_000; 4])); let inventory = inventory("quantile_over_time(0.99, m[5m])", &[every(60_000)]); let target = &inventory.targets[0]; @@ -740,7 +829,10 @@ mod tests { merge(vec![pane(0, 60_000), pane(60_000, 60_000)]).unwrap(); assert!(merge(vec![pane(0, 120_000), pane(60_000, 60_000)]).is_err()); // Panes of 2m over a 5m window leave the oldest minute uncovered. - assert!(tumbling_state(child, 120_000, kll_state(target)).is_err()); + assert!( + tumbling_state(child, WindowForm::Tumbling { pane_ms: 120_000 }, kll_state(target)) + .is_err() + ); } const YEAR_MS: u64 = 365 * 86_400_000; @@ -771,51 +863,97 @@ mod tests { .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. + /// The planner chooses the grid (Q67): split at the windows' boundaries + /// (0, 1, 2, 3 and 5 years back), the five queries share four KLL + /// segments, the oldest two years wide; on the even 1-year grid (the gcd + /// of every lookback and offset) they share five. Each is its own + /// variant offering 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() { + fn pattern_a_shares_four_boundary_or_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 + let segmented: Vec<_> = 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!(), - }) + .filter(|v| matches!(v.sharing, Sharing::WindowSegments { .. })) .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, + let sharing: Vec<_> = segmented.iter().map(|v| v.sharing).collect(); + assert_eq!( + sharing, + [ + Sharing::WindowSegments { segments: 4 }, + Sharing::WindowSegments { segments: 5 } + ] + ); + // Newest first: [5y] and [3y offset 2y] end at 3 years back. + let split = |ends| WindowForm::Segments { + unit_ms: YEAR_MS, + ends, + }; + let boundary = [split(0b10111), split(1), split(1), split(1), split(0b101)]; + assert_eq!( + boundary[0].label().as_deref(), + Some("365d+365d+365d+730d segments") + ); + assert_eq!( + boundary[0].widths(0), + [YEAR_MS, YEAR_MS, YEAR_MS, 2 * YEAR_MS] + ); + let even = [split(0b11111), split(1), split(1), split(1), split(0b111)]; + assert_eq!(even[0].label().as_deref(), Some("365d segments")); + for (variant, forms) in segmented.iter().zip([boundary, even]) { + for (target, form) in variant.inventory.targets.iter().zip(forms) { + assert_eq!(target.windows, [form]); + assert!(matches!( + &target.alternatives[..], + [Realization::Sketch(kind)] if *kind.algorithm() == SketchAlgorithm::Kll + )); + } + let composed = compose_logical_candidate(&variant.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, + }; + } + let Sharing::WindowSegments { segments } = variant.sharing else { + unreachable!() }; + assert_eq!((builds.len(), merges.len()), (segments, 5)); } - assert_eq!((builds.len(), merges.len()), (5, 5)); + } + + /// When the windows' boundaries already lie on the even grid, the two + /// segmentations are the same and only one variant is offered. + #[test] + fn identical_segmentations_are_offered_once() { + let roots = pattern_a(); + let inventory = enumerate_local_logical_candidates( + vec![roots[1].clone(), roots[2].clone()], + &BTreeMap::new(), + ) + .unwrap(); + let offered: Vec<_> = share_window_segments(&inventory, &[]) + .unwrap() + .into_iter() + .map(|(segments, _)| segments) + .collect(); + assert_eq!(offered, [2]); } /// Queries that recur at different cadences, or read the same window, @@ -825,16 +963,16 @@ mod tests { 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()); + assert!(!share_window_segments(&inventory, &[]).unwrap().is_empty()); let demand = [every(60_000), every(120_000)]; assert!(share_window_segments(&inventory, &demand) .unwrap() - .is_none()); + .is_empty()); 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()); + assert!(share_window_segments(&same, &[]).unwrap().is_empty()); } } diff --git a/crates/physical-optimizer/src/materialization.rs b/crates/physical-optimizer/src/materialization.rs index 13333672..e46c828b 100644 --- a/crates/physical-optimizer/src/materialization.rs +++ b/crates/physical-optimizer/src/materialization.rs @@ -28,6 +28,8 @@ //! **Units.** The panes merged by one `SummaryMerge` are decided together: //! a pane chain is maintained as one stream of panes, and a mixed chain is //! never cheaper. Panes shared by two merges join both chains into one unit. +//! A merge of panes of different widths (uneven shared segments, Q67) is no +//! stream of panes, so its chain stays at query time. //! //! **Down-closed sets.** Everything upstream of an ingestion-time node also //! runs at ingestion time (#509), so a unit is at ingestion time only if @@ -178,7 +180,12 @@ impl MaterializationSpace { .into_iter() .filter_map(|c| index.get(&Rc::as_ptr(c)).copied()) .collect(); - if members.iter().any(|&m| !eligible[m]) { + // Panes of different widths (uneven shared segments, Q67) are + // not one stream of panes: no pane is a later one shifted back. + let uneven = members + .windows(2) + .any(|pair| window_of(&summaries[pair[0]]) != window_of(&summaries[pair[1]])); + if uneven || members.iter().any(|&m| !eligible[m]) { for &m in &members { usable[m] = false; } diff --git a/tools/dag-viewer/examples.json b/tools/dag-viewer/examples.json index 75857229..60a44bc9 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 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." + "story": "An ad hoc batch of five p99 reports over 1–5 years, run once. Each query gets a KLL, and Pass 2 offers KLL segments shared by all five reports, cut two ways: at the windows' boundaries (four segments, the oldest two years wide) and on the even 1-year grid (five). 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. The four boundary segments are cheapest, by the merges they save." }, { "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 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." + "story": "Example 3a repeated monthly and known in advance. The same four shared boundary segments win, built at query time; their cost is now amortized over the month between runs. The even 1-year segments may now be maintained at ingestion time, but keeping a million per-series KLLs per segment costs far more." }, { "id": "example4b", From b2dea0d59207e8807bcaa61f7ab113b515f1529a Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Thu, 8 Oct 2026 15:16:46 +0000 Subject: [PATCH 2/2] refactor(pass2): segment grids derive their coverage; tiling check replaces the declared proof Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01W7qG9aFyPij5uWsyAJCxDW --- .../tests/planner_layering_example3.rs | 11 +++-- .../src/pass2/window_composition.rs | 46 +++++++++++-------- 2 files changed, 33 insertions(+), 24 deletions(-) diff --git a/crates/integration-tests/tests/planner_layering_example3.rs b/crates/integration-tests/tests/planner_layering_example3.rs index 36f87143..901f00b3 100644 --- a/crates/integration-tests/tests/planner_layering_example3.rs +++ b/crates/integration-tests/tests/planner_layering_example3.rs @@ -236,11 +236,12 @@ fn stage1_a_shared_segments_serve_all_five() { assert_eq!(estimates.len(), 5, "one estimate per query"); for (estimate, statistic) in estimates { assert_eq!(statistic, SketchStatistic::Quantile { q: 0.99 }); - let merged = c - .dag - .producers(estimate) - .into_iter() - .any(|p| matches!(c.dag.payload(p), Operator::ASAP(ASAPOp::SummaryMerge { .. }))); + let merged = c.dag.producers(estimate).into_iter().any(|p| { + matches!( + c.dag.payload(p), + Operator::ASAP(ASAPOp::SummaryMerge { .. }) + ) + }); assert!( merged, "{}: estimate {estimate:?} does not read a merge", diff --git a/crates/logical-optimizer/src/pass2/window_composition.rs b/crates/logical-optimizer/src/pass2/window_composition.rs index d1057155..aa6a825b 100644 --- a/crates/logical-optimizer/src/pass2/window_composition.rs +++ b/crates/logical-optimizer/src/pass2/window_composition.rs @@ -25,7 +25,7 @@ //! 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::{BTreeMap, BTreeSet, HashMap}; +use std::collections::{BTreeSet, HashMap}; use std::rc::Rc; use std::time::Duration; @@ -546,10 +546,11 @@ mod tests { use super::*; use crate::pass1::logical_candidates::enumerate_local_logical_candidates; use crate::test_support::lower_promql; - use std::ops::Bound; use asap_types::ir::schema::SketchAlgorithm; use asap_types::types::AccuracyTarget; use asap_types::workload::{Predictability, RepetitionInterval}; + use std::collections::BTreeMap; + use std::ops::Bound; fn every(ms: u32) -> RootDemand { RootDemand { @@ -675,9 +676,12 @@ mod tests { let Some(NonASAPOp::Aggregate { child, .. }) = target.target.non_asap() else { panic!("aggregate") }; - let merged = - tumbling_state(child, WindowForm::Tumbling { pane_ms: 60_000 }, kll_state(target)) - .unwrap(); + let merged = tumbling_state( + child, + WindowForm::Tumbling { pane_ms: 60_000 }, + kll_state(target), + ) + .unwrap(); let Operator::ASAP(ASAPOp::SummaryMerge { children }) = &merged.operator else { panic!("merge") }; @@ -766,16 +770,18 @@ mod tests { target: &crate::pass1::logical_candidates::LocalLogicalTarget, ) -> impl Fn(Rc) -> Result, LogicalCandidateError> + '_ { move |input| { - Ok(OperatorNode::new_shared(Operator::ASAP(ASAPOp::SummaryAgg { - child: input, - family: family(&target.alternatives[1], &GroupingStrategy::default()).unwrap(), - input: asap_types::ir::schema::SummaryUpdate::column( - asap_types::ir::scalar::ColumnRef::SampleValue, - ), - reduction: asap_types::ir::operator::Reduction::PerEntity, - grouping: GroupingStrategy::default(), - filter: None, - }))?) + Ok(OperatorNode::new_shared(Operator::ASAP( + ASAPOp::SummaryAgg { + child: input, + family: family(&target.alternatives[1], &GroupingStrategy::default()).unwrap(), + input: asap_types::ir::schema::SummaryUpdate::column( + asap_types::ir::scalar::ColumnRef::SampleValue, + ), + reduction: asap_types::ir::operator::Reduction::PerEntity, + grouping: GroupingStrategy::default(), + filter: None, + }, + ))?) } } @@ -829,10 +835,12 @@ mod tests { merge(vec![pane(0, 60_000), pane(60_000, 60_000)]).unwrap(); assert!(merge(vec![pane(0, 120_000), pane(60_000, 60_000)]).is_err()); // Panes of 2m over a 5m window leave the oldest minute uncovered. - assert!( - tumbling_state(child, WindowForm::Tumbling { pane_ms: 120_000 }, kll_state(target)) - .is_err() - ); + assert!(tumbling_state( + child, + WindowForm::Tumbling { pane_ms: 120_000 }, + kll_state(target) + ) + .is_err()); } const YEAR_MS: u64 = 365 * 86_400_000;