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
15 changes: 8 additions & 7 deletions crates/devtools/src/bin/stage_pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -267,11 +268,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
Expand Down
19 changes: 10 additions & 9 deletions crates/devtools/tests/stage_pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -170,13 +170,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(&[
Expand All @@ -194,14 +195,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()
Expand Down
122 changes: 87 additions & 35 deletions crates/integration-tests/tests/planner_layering_example3.rs
Original file line number Diff line number Diff line change
Expand Up @@ -180,51 +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<usize> = 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<Vec<u64>> = candidates
.iter()
.flat_map(|(build, ..)| estimates_of(&c.dag, *build))
.map(|c| {
let mut widths: Vec<u64> = 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), LogicalASAPOperatorPayload::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<usize> = 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), LogicalASAPOperatorPayload::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() {
Expand Down
38 changes: 25 additions & 13 deletions crates/integration-tests/tests/planner_layering_example4.rs
Original file line number Diff line number Diff line change
Expand Up @@ -97,13 +97,17 @@ fn eh_options(run: &Run) -> BTreeMap<Materialization, (String, LogicalASAPNodeId
options_of(run, None, 5)
}

/// The physical candidates of the shared-segment candidate (Q60), keyed by
/// the materialization of its segments, with the newest segment's id.
fn segment_options(run: &Run) -> BTreeMap<Materialization, (String, LogicalASAPNodeId)> {
/// 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<Materialization, (String, LogicalASAPNodeId)> {
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| {
Expand Down Expand Up @@ -197,23 +201,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::<BTreeSet<_>>(), 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<usize> = segments
.iter()
.flat_map(|(b, ..)| readers(&p.dag, &p.query_roots, *b))
Expand Down
7 changes: 3 additions & 4 deletions crates/logical-optimizer/src/pass1/logical_candidates.rs
Original file line number Diff line number Diff line change
Expand Up @@ -711,10 +711,9 @@ fn realize(
};
build(child, coverage)?
}
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 {
Expand Down
14 changes: 8 additions & 6 deletions crates/logical-optimizer/src/pass2/identical_expressions.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -62,8 +63,9 @@ pub struct SharingVariant<Id> {
/// 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<Id: Clone>(
roots: Vec<(Id, QueryRoot)>,
metric_types: &BTreeMap<String, MetricType>,
Expand Down Expand Up @@ -93,9 +95,9 @@ pub fn stage1_logical_candidates<Id: Clone>(
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,
});
}
Expand Down
Loading