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
5 changes: 4 additions & 1 deletion crates/devtools/src/bin/stage_pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -269,6 +271,7 @@ fn stage_pipeline(
Sharing::Independent => "",
Sharing::IdenticalExpressions => " · shared input",
Sharing::SummaryCapability => " · shared summary",
Sharing::WindowSegments => " · shared segments",
};
let written = candidate
.physical
Expand Down
24 changes: 15 additions & 9 deletions crates/devtools/tests/stage_pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -170,13 +170,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",
Expand All @@ -192,10 +194,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()
Expand Down
11 changes: 11 additions & 0 deletions crates/integration-tests/tests/planner_layering_common/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -550,6 +550,17 @@ 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: &[LogicalASAPNodeId]) -> bool {
cross_query_nodes(dag, query_roots).into_iter().any(|id| {
matches!(
dag.payload(id),
LogicalASAPOperatorPayload::SummaryAgg { .. }
)
})
}

/// Nodes reachable from more than one query root.
pub fn cross_query_nodes(
dag: &impl ExportedDag,
Expand Down
71 changes: 44 additions & 27 deletions crates/integration-tests/tests/planner_layering_example3.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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!(
Expand All @@ -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<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
.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 });
Expand All @@ -210,16 +225,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
Expand All @@ -232,7 +245,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
Expand Down Expand Up @@ -281,7 +294,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()
Expand Down
82 changes: 54 additions & 28 deletions crates/integration-tests/tests/planner_layering_example4.rs
Original file line number Diff line number Diff line change
Expand Up @@ -97,6 +97,22 @@ 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)> {
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<WindowForm> {
Some(WindowForm::Tumbling {
length_ms: PATTERN_B_INTERVAL_MS,
Expand Down Expand Up @@ -164,7 +180,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!(
Expand All @@ -175,43 +191,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::<BTreeSet<_>>(), 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<usize> = 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, LogicalASAPOperatorPayload::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);
Expand All @@ -226,7 +252,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);
Expand All @@ -238,7 +264,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);
Expand Down
5 changes: 4 additions & 1 deletion crates/logical-optimizer/src/pass1/logical_candidates.rs
Original file line number Diff line number Diff line change
Expand Up @@ -711,7 +711,10 @@ fn realize(
};
build(child, coverage)?
}
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 {
Expand Down
18 changes: 16 additions & 2 deletions crates/logical-optimizer/src/pass2/identical_expressions.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
};
Expand All @@ -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 {
Expand All @@ -56,7 +61,9 @@ pub struct SharingVariant<Id> {
/// 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<Id: Clone>(
roots: Vec<(Id, QueryRoot)>,
metric_types: &BTreeMap<String, MetricType>,
Expand All @@ -74,6 +81,7 @@ pub fn stage1_logical_candidates<Id: Clone>(
});
}
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 {
Expand All @@ -85,6 +93,12 @@ pub fn stage1_logical_candidates<Id: Clone>(
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)
}

Expand Down
3 changes: 2 additions & 1 deletion crates/logical-optimizer/src/pass2/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Loading