From e2d15c4b6747408cc40694a5a2cf6409ba03a5b2 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Sun, 4 Oct 2026 17:31:32 +0000 Subject: [PATCH 1/2] feat(planner): Pass 2 summary-capability rule as a Stage 1 sharing variant Targets with the same summary input data and window (same input, grouping and quantile column) are re-sized for their strictest consumer, all or nothing per key (#580 W5), reusing the legacy reconciliation's accuracy_budget/dominates argument. Stage 1 adds this inventory as a third variant (Sharing::SummaryCapability) next to the independent and identical-expression ones; composition then builds identical summary producers, which are merged, and Stage 3 chooses. Also: - the tree DP now checks coupling between targets that read a common (or equal) input; it skipped those pairs, so it missed merged producers; - a composed summary evaluation keeps the measure's explicit output name (SQL), instead of the derived one. Co-Authored-By: Claude Opus 5.5 --- crates/devtools/src/bin/stage_pipeline.rs | 19 +- .../tests/planner_layering_example1.rs | 10 +- .../src/pass1/logical_candidates.rs | 29 +- .../src/pass2/identical_expressions.rs | 53 +++- crates/logical-optimizer/src/pass2/mod.rs | 4 +- .../src/pass2/reconciliation.rs | 2 +- .../src/pass2/summary_capability.rs | 286 ++++++++++++++++++ crates/plan-selection/src/lib.rs | 80 ++--- crates/planner/src/pass/stage_pipeline.rs | 8 +- .../planner/tests/stage_pipeline_selection.rs | 132 +++++++- crates/planner/tests/summary_sharing.rs | 54 +++- 11 files changed, 609 insertions(+), 68 deletions(-) create mode 100644 crates/logical-optimizer/src/pass2/summary_capability.rs diff --git a/crates/devtools/src/bin/stage_pipeline.rs b/crates/devtools/src/bin/stage_pipeline.rs index bf3518a38..177e10d73 100644 --- a/crates/devtools/src/bin/stage_pipeline.rs +++ b/crates/devtools/src/bin/stage_pipeline.rs @@ -10,8 +10,9 @@ // local alternative for every target (Pass 1, the Cartesian product), for // the queries as written and, when Pass 2's identical-expression rule // merges something, again with identical sub-DAGs shared ("· shared -// input"); in enumeration order and capped by `--max-candidates` -// (default 64); +// input"), and, when the summary-capability rule applies, again with +// one summary sized for its strictest consumer ("· shared summary"); in +// enumeration order and capped by `--max-candidates` (default 64); // - stage2_physical_asap: one physical candidate per logical candidate // (operator implementation only, everything at query time), no cost; // - stage3_selection: per-candidate costs, the selected candidate, and @@ -34,7 +35,7 @@ use asap_logical_optimizer::pass1::logical_candidates::{ }; use asap_logical_optimizer::Realization; use asap_plan_selection::PlanningModels; -use asap_plan_selection::{plan_stages, Selection, MAX_ENUMERATED_CANDIDATES}; +use asap_plan_selection::{plan_stages, Selection, Sharing, MAX_ENUMERATED_CANDIDATES}; use asap_types::ir::flat::{flatten, FlatDag}; use asap_types::ir::schema::SketchAlgorithm; use asap_types::ir::schema_support::with_promql_series_identity; @@ -132,13 +133,13 @@ fn stage_pipeline(workload: &PlanningWorkload, max_candidates: usize) -> Result< let mut candidates = Vec::new(); let mut stage2 = Vec::new(); for candidate in &enumeration.candidates { - // Candidates of the shared variant are numbered after the independent ones. + // Candidates of each variant are numbered after the previous variants'. let mut offset = 0; let variant = run .stage1 .iter() .find(|v| { - let found = v.shared == candidate.shared; + let found = v.sharing == candidate.sharing; if !found { offset += combination_count(&v.inventory); } @@ -148,9 +149,11 @@ fn stage_pipeline(workload: &PlanningWorkload, max_candidates: usize) -> Result< let inventory = &variant.inventory; let index = offset + choice_index(inventory, &candidate.choice) + 1; let mut label = label(inventory, &target_owners(inventory), &candidate.choice); - if candidate.shared { - label += " · shared input"; - } + label += match candidate.sharing { + Sharing::Independent => "", + Sharing::IdenticalExpressions => " · shared input", + Sharing::SummaryCapability => " · shared summary", + }; if let Some(logical) = &candidate.logical { let roots: Vec<_> = logical.iter().map(|(_, root)| root.clone()).collect(); candidates diff --git a/crates/integration-tests/tests/planner_layering_example1.rs b/crates/integration-tests/tests/planner_layering_example1.rs index 92bd1806d..3adcd7df0 100644 --- a/crates/integration-tests/tests/planner_layering_example1.rs +++ b/crates/integration-tests/tests/planner_layering_example1.rs @@ -170,7 +170,15 @@ mod stages { .into_iter() .map(|c| { let physical = c.physical.expect("every Example 1 candidate builds"); - let label = format!("{:?}{}", c.choice, if c.shared { " shared" } else { "" }); + let label = format!( + "{:?}{}", + c.choice, + if c.sharing.merges_after_composition() { + " shared" + } else { + "" + } + ); let roots = c.logical.expect("composes"); candidate( physical.from_logical, diff --git a/crates/logical-optimizer/src/pass1/logical_candidates.rs b/crates/logical-optimizer/src/pass1/logical_candidates.rs index 741b42b27..e5c5aa85d 100644 --- a/crates/logical-optimizer/src/pass1/logical_candidates.rs +++ b/crates/logical-optimizer/src/pass1/logical_candidates.rs @@ -564,7 +564,34 @@ fn realize( }, None => ASAPOp::FinalizeExactAccumulator { child: state }, }; - Ok(OperatorNode::new_shared(Operator::ASAP(evaluation))?) + let mut evaluation = OperatorNode::new(Operator::ASAP(evaluation))?; + keep_output_name(target, &mut evaluation); + Ok(Rc::new(evaluation)) +} + +/// The evaluation answers `target`, so a measure the query named explicitly +/// (SQL `approx_percentile_cont(...)`) keeps its name; a synthetic name is +/// left as derived. +fn keep_output_name(target: &OperatorNode, evaluation: &mut OperatorNode) { + let Some(NonASAPOp::Aggregate { output_names, .. }) = target.non_asap() else { + return; + }; + let [name] = output_names.as_slice() else { + return; + }; + if name.is_empty() || evaluation.schema.fields.len() != target.schema.fields.len() { + return; + } + for (field, named) in evaluation + .schema + .fields + .iter_mut() + .zip(&target.schema.fields) + { + if named.name == *name { + field.name = name.clone(); + } + } } /// Whether every value `node` outputs is a sample of a metric declared a diff --git a/crates/logical-optimizer/src/pass2/identical_expressions.rs b/crates/logical-optimizer/src/pass2/identical_expressions.rs index 56361cdc2..bbc59c444 100644 --- a/crates/logical-optimizer/src/pass2/identical_expressions.rs +++ b/crates/logical-optimizer/src/pass2/identical_expressions.rs @@ -15,35 +15,68 @@ use asap_types::ir::cse::share_common_sub_dags; use asap_types::ir::{OperatorNode, QueryRoot}; use asap_types::workload::MetricType; +use super::summary_capability::share_summary_capability; use crate::pass1::logical_candidates::{ enumerate_local_logical_candidates, LocalLogicalCandidates, LogicalCandidateError, }; -/// One input-sharing form of the workload, with its Pass 1 alternatives. +/// Which Pass 2 sharing a Stage 1 variant applies. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum Sharing { + /// Pass 1 over the queries as written. + Independent, + /// The identical-expression rule: identical sub-DAGs across queries are + /// merged, before Pass 1 and again after composition. + IdenticalExpressions, + /// The summary-capability rule on top of the identical-expression rule + /// ([`super::summary_capability`]): targets that can share one summary + /// are sized for their strictest consumer, so composition builds + /// identical producers, which are merged. + SummaryCapability, +} + +impl Sharing { + /// Whether identical sub-DAGs are merged after composition. + pub fn merges_after_composition(self) -> bool { + self != Sharing::Independent + } +} + +/// One sharing form of the workload, with its Pass 1 alternatives. #[derive(Debug, Clone)] pub struct SharingVariant { - /// Whether identical sub-DAGs across queries are merged. - pub shared: bool, + pub sharing: Sharing, pub inventory: LocalLogicalCandidates, } -/// Stage 1 = Pass 1 + Pass 2's identical-expression rule: the independent -/// variant first, then the shared one when sharing merges at least one node. +/// Stage 1 = Pass 1 + Pass 2: the independent variant first, then the +/// identical-expression variant when sharing merges at least one node, then +/// the summary-capability variant when two targets can share a summary. The +/// last is skipped when it would repeat the identical-expression variant. pub fn stage1_logical_candidates( roots: Vec<(Id, QueryRoot)>, metric_types: &BTreeMap, ) -> Result>, LogicalCandidateError> { let shared = share_identical_expressions(&roots); let mut variants = vec![SharingVariant { - shared: false, + sharing: Sharing::Independent, inventory: enumerate_local_logical_candidates(roots, metric_types)?, }]; if let Some(roots) = shared { variants.push(SharingVariant { - shared: true, + sharing: Sharing::IdenticalExpressions, inventory: enumerate_local_logical_candidates(roots, metric_types)?, }); } + let base = &variants.last().expect("the independent variant").inventory; + if let Some(capability) = share_summary_capability(base)? { + if capability.resized || variants.len() == 1 { + variants.push(SharingVariant { + sharing: Sharing::SummaryCapability, + inventory: capability.inventory, + }); + } + } Ok(variants) } @@ -114,8 +147,8 @@ mod tests { ) .unwrap(); assert_eq!( - variants.iter().map(|v| v.shared).collect::>(), - [false, true] + variants.iter().map(|v| v.sharing).collect::>(), + [Sharing::Independent, Sharing::IdenticalExpressions] ); assert_eq!( variants[0].inventory.targets.len(), @@ -145,6 +178,6 @@ mod tests { ) .unwrap(); assert_eq!(variants.len(), 1); - assert!(!variants[0].shared); + assert_eq!(variants[0].sharing, Sharing::Independent); } } diff --git a/crates/logical-optimizer/src/pass2/mod.rs b/crates/logical-optimizer/src/pass2/mod.rs index 520b774c0..2be12deda 100644 --- a/crates/logical-optimizer/src/pass2/mod.rs +++ b/crates/logical-optimizer/src/pass2/mod.rs @@ -1,7 +1,9 @@ //! Pass 2: ASAP-aware sharing across targets. A shared summary must meet the //! strictest accuracy requirement of its readers. The stage pipeline applies -//! the identical-expression rule ([`identical_expressions`]). +//! the identical-expression rule ([`identical_expressions`]) and the +//! summary-capability rule ([`summary_capability`]). pub mod identical_expressions; pub mod reconciliation; +pub mod summary_capability; pub mod topk_reuse; diff --git a/crates/logical-optimizer/src/pass2/reconciliation.rs b/crates/logical-optimizer/src/pass2/reconciliation.rs index 051bb6425..983999b63 100644 --- a/crates/logical-optimizer/src/pass2/reconciliation.rs +++ b/crates/logical-optimizer/src/pass2/reconciliation.rs @@ -239,7 +239,7 @@ fn same_intent_except_accuracy(a: &AggIntent, b: &AggIntent) -> bool { /// — a different `Realization` family, not a point on the same sizing /// curve — so the numeric comparison alone does not mean what it means for /// two approximate targets. See the module docs for the full reasoning. -fn dominates(tighter: &AccuracyTarget, looser: &AccuracyTarget) -> bool { +pub(crate) fn dominates(tighter: &AccuracyTarget, looser: &AccuracyTarget) -> bool { if matches!(tighter, AccuracyTarget::Exact) || matches!(looser, AccuracyTarget::Exact) { return false; } diff --git a/crates/logical-optimizer/src/pass2/summary_capability.rs b/crates/logical-optimizer/src/pass2/summary_capability.rs new file mode 100644 index 000000000..79213746f --- /dev/null +++ b/crates/logical-optimizer/src/pass2/summary_capability.rs @@ -0,0 +1,286 @@ +//! Pass 2's summary-capability rule (#509): computations with the same +//! summary input data (source, filters, update expression, grouping) and the +//! same window can share one summary that supports every requested estimate, +//! sized for the strictest consumer. +//! +//! The rule is all-or-nothing per key (#580 decision W5): every approximate +//! target of a key is re-sized for the key's strictest requirement, so their +//! summary producers are identical and the shared variant's merge after +//! composition reaches one state. Sizing reuses the legacy reconciliation's +//! argument ([`super::reconciliation`]): requirements resolve through +//! [`accuracy_budget`] and every shipped sizing formula is monotonic in +//! `(ε, δ)`, so a summary sized for a requirement that dominates every +//! consumer's meets each of them. Stage 3 still checks each query against its +//! own target. +//! +//! A shared state changes cost non-additively, so this is a Stage 1 variant +//! next to the independent one, and Stage 3 chooses. + +use std::rc::Rc; + +use asap_types::ir::operator::{AggIntent, Reduction}; +use asap_types::ir::schema::ColumnId; +use asap_types::ir::{NonASAPOp, OperatorNode}; +use asap_types::types::AccuracyTarget; + +use super::reconciliation::dominates; +use crate::pass1::logical_candidates::{ + local_realizations_for_intent, LocalLogicalCandidates, LogicalCandidateError, +}; +use crate::pass1::replacement::{accuracy_budget, accuracy_target}; + +/// The estimates one summary serves: any quantile of one column from one +/// quantile summary (KLL, DDSketch). +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum Capability { + Quantile, +} + +/// What a summary for a target would ingest and which estimates it would +/// serve. Targets with equal keys can share one summary. +struct Key<'a> { + capability: Capability, + column: Option, + reduction: &'a Reduction, + /// The input, window included (`TimeRange`, `WHERE`, ...). + child: &'a Rc, +} + +impl Key<'_> { + fn matches(&self, other: &Key<'_>) -> bool { + self.capability == other.capability + && self.column == other.column + && self.reduction == other.reduction + && (Rc::ptr_eq(self.child, other.child) || self.child == other.child) + } +} + +/// The key of a single-measure, unfiltered aggregate with an approximate +/// requirement, or `None` when the rule does not apply to it. +fn key(node: &OperatorNode) -> Option> { + let Some(NonASAPOp::Aggregate { + reduction, + measures, + filters, + having: None, + child, + .. + }) = node.non_asap() + else { + return None; + }; + let ([intent], true) = (measures.as_slice(), filters.iter().all(Option::is_none)) else { + return None; + }; + if matches!(accuracy_target(intent)?, AccuracyTarget::Exact) { + return None; + } + let (capability, column) = match intent { + AggIntent::Quantile { col, .. } => (Capability::Quantile, *col), + _ => return None, + }; + Some(Key { + capability, + column, + reduction, + child, + }) +} + +/// `intent` with its accuracy requirement replaced. +fn with_accuracy(intent: &AggIntent, target: AccuracyTarget) -> AggIntent { + let mut intent = intent.clone(); + match &mut intent { + AggIntent::Quantile { accuracy, .. } + | AggIntent::Cardinality { accuracy, .. } + | AggIntent::FrequencyL2 { accuracy, .. } + | AggIntent::FrequencyEntropy { accuracy, .. } + | AggIntent::Count { accuracy } + | AggIntent::TopK { accuracy, .. } => *accuracy = target, + _ => {} + } + intent +} + +/// The requirement that dominates every one in `targets`: the smallest ε +/// and the smallest δ. +fn strictest<'a>(targets: impl IntoIterator) -> AccuracyTarget { + let (epsilon, delta) = targets + .into_iter() + .map(accuracy_budget) + .fold((f64::INFINITY, f64::INFINITY), |(e, d), (e2, d2)| { + (e.min(e2), d.min(d2)) + }); + AccuracyTarget::EpsilonDelta { epsilon, delta } +} + +/// The summary-capability rule over a Pass 1 inventory. +#[derive(Debug, Clone)] +pub struct CapabilitySharing { + /// `inventory` with every target of a shared key re-sized for the key's + /// strictest requirement. + pub inventory: LocalLogicalCandidates, + /// Whether any target was re-sized. When none was, the rule adds only + /// the merge of identical producers after composition. + pub resized: bool, +} + +/// Apply the rule to `inventory`, or `None` when no two targets share a key. +pub fn share_summary_capability( + inventory: &LocalLogicalCandidates, +) -> Result>, LogicalCandidateError> { + let keys: Vec<_> = inventory.targets.iter().map(|t| key(&t.target)).collect(); + let mut group = vec![None::; keys.len()]; + for t in 0..keys.len() { + let Some(key) = &keys[t] else { continue }; + group[t] = Some( + (0..t) + .find(|&u| keys[u].as_ref().is_some_and(|other| key.matches(other))) + .map_or(t, |u| group[u].expect("keyed")), + ); + } + let mut out = inventory.clone(); + let mut shared = false; + let mut resized = false; + for leader in 0..keys.len() { + let members: Vec = (0..keys.len()) + .filter(|&t| group[t] == Some(leader)) + .collect(); + if members.len() < 2 { + continue; + } + shared = true; + let intents: Vec<&AggIntent> = members + .iter() + .map(|&t| single_intent(inventory, t)) + .collect(); + let requirements: Vec<&AccuracyTarget> = intents + .iter() + .map(|intent| accuracy_target(intent).expect("keyed targets are approximate")) + .collect(); + let target = strictest(requirements.iter().copied()); + debug_assert!(requirements.iter().all(|r| dominates(&target, r))); + for ((&t, intent), requirement) in members.iter().zip(intents).zip(requirements) { + if accuracy_budget(requirement) == accuracy_budget(&target) { + continue; + } + let alternatives = + local_realizations_for_intent(&with_accuracy(intent, target.clone()))?; + // Keyed targets never absorb the target beneath them (only a + // whole-expression top-k does), so only the sizes change. + debug_assert!(out.targets[t].absorbs.iter().all(Option::is_none)); + out.targets[t].absorbs = vec![None; alternatives.len()]; + out.targets[t].alternatives = alternatives; + resized = true; + } + } + Ok(shared.then_some(CapabilitySharing { + inventory: out, + resized, + })) +} + +fn single_intent(inventory: &LocalLogicalCandidates, t: usize) -> &AggIntent { + match inventory.targets[t].target.non_asap() { + Some(NonASAPOp::Aggregate { measures, .. }) => &measures[0], + _ => unreachable!("keyed targets are single-measure aggregates"), + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::pass1::logical_candidates::enumerate_local_logical_candidates; + use crate::pass1::replacement::Realization; + use crate::test_support::lower_promql; + use asap_types::ir::schema::{SketchAlgorithm, SketchParams}; + use asap_types::ir::QueryRoot; + + fn inventory(queries: &[(&str, f64)]) -> LocalLogicalCandidates { + let roots = queries + .iter() + .enumerate() + .map(|(i, (q, epsilon))| { + let root = lower_promql(q, AccuracyTarget::Epsilon(*epsilon)); + let root = + asap_types::ir::schema_support::with_promql_series_identity(&root).unwrap(); + (i, QueryRoot::Operator(root)) + }) + .collect(); + enumerate_local_logical_candidates(roots).unwrap() + } + + fn kll_k(inventory: &LocalLogicalCandidates, t: usize) -> u32 { + inventory.targets[t] + .alternatives + .iter() + .find_map(|a| match a { + Realization::Sketch(kind) => match kind.params() { + SketchParams::Kll { k } => Some(*k), + _ => None, + }, + _ => None, + }) + .expect("a KLL alternative") + } + + /// p50 at ε = 0.01 and p99 at ε = 0.001 over one input: both targets' + /// summaries are sized for ε = 0.001, so their KLLs are identical. + #[test] + fn quantiles_are_sized_for_the_strictest_consumer() { + let base = inventory(&[ + ("quantile_over_time(0.5, lat[5m])", 0.01), + ("quantile_over_time(0.99, lat[5m])", 0.001), + ]); + assert!(kll_k(&base, 0) < kll_k(&base, 1)); + let shared = share_summary_capability(&base).unwrap().expect("one key"); + assert!(shared.resized); + assert_eq!(kll_k(&shared.inventory, 0), kll_k(&base, 1)); + assert_eq!(kll_k(&shared.inventory, 1), kll_k(&base, 1)); + let ddsketch = |inv: &LocalLogicalCandidates, t: usize| { + inv.targets[t].alternatives.iter().find_map(|a| match a { + Realization::Sketch(kind) if *kind.algorithm() == SketchAlgorithm::DDSketch => { + Some(kind.clone()) + } + _ => None, + }) + }; + assert_eq!(ddsketch(&shared.inventory, 0), ddsketch(&base, 1)); + } + + /// Equal requirements need no re-sizing; the key still groups them. + #[test] + fn equal_requirements_are_grouped_without_resizing() { + let base = inventory(&[ + ("quantile_over_time(0.5, lat[5m])", 0.01), + ("quantile_over_time(0.99, lat[5m])", 0.01), + ]); + let shared = share_summary_capability(&base).unwrap().expect("one key"); + assert!(!shared.resized); + } + + /// A different window, selector or estimate family is a different key. + #[test] + fn different_summary_input_or_window_is_not_shared() { + for queries in [ + [ + ("quantile_over_time(0.5, lat[5m])", 0.01), + ("quantile_over_time(0.99, lat[10m])", 0.001), + ], + [ + ("quantile_over_time(0.5, lat{job=\"a\"}[5m])", 0.01), + ("quantile_over_time(0.99, lat{job=\"b\"}[5m])", 0.001), + ], + [ + ("quantile_over_time(0.5, lat[5m])", 0.01), + ("count_over_time(lat[5m])", 0.001), + ], + ] { + let base = inventory(&queries); + assert!( + share_summary_capability(&base).unwrap().is_none(), + "{queries:?}" + ); + } + } +} diff --git a/crates/plan-selection/src/lib.rs b/crates/plan-selection/src/lib.rs index dd8d63d35..e7badb9b5 100644 --- a/crates/plan-selection/src/lib.rs +++ b/crates/plan-selection/src/lib.rs @@ -80,10 +80,10 @@ use asap_logical_optimizer::pass1::logical_candidates::{ choice_index, combination_count, compose_logical_candidate, enumerate_choices, nested_targets, read_targets, LocalLogicalCandidates, LogicalCandidateError, }; -pub use asap_logical_optimizer::pass2::identical_expressions::SharingVariant; use asap_logical_optimizer::pass2::identical_expressions::{ share_identical_expressions, stage1_logical_candidates, }; +pub use asap_logical_optimizer::pass2::identical_expressions::{Sharing, SharingVariant}; use asap_physical_optimizer::implementation::physical_candidates::{ stage2_physical, PhysicalCandidate, }; @@ -358,8 +358,7 @@ fn assess( /// stay unique across variants. struct Variant<'a, Id> { inventory: &'a LocalLogicalCandidates, - /// Identical sub-DAGs are merged, after composition too. - shared: bool, + sharing: Sharing, /// Candidates numbered before this variant's. offset: usize, } @@ -379,7 +378,7 @@ fn variants(stage1: &[SharingVariant]) -> Vec> { .map(|v| { let variant = Variant { inventory: &v.inventory, - shared: v.shared, + sharing: v.sharing, offset, }; offset = offset.saturating_add(combination_count(&v.inventory)); @@ -398,7 +397,7 @@ impl Variant<'_, Id> { /// The independent variant alone: Pass 1 without Pass 2. pub fn independent(inventory: LocalLogicalCandidates) -> Vec> { vec![SharingVariant { - shared: false, + sharing: Sharing::Independent, inventory, }] } @@ -413,14 +412,14 @@ pub fn realize_choice( realize( Variant { inventory, - shared: false, + sharing: Sharing::Independent, offset: 0, }, choice, ) } -/// Stage 1 → Stage 2 for `choice` in `variant`. A shared variant also +/// Stage 1 → Stage 2 for `choice` in `variant`. A sharing variant also /// merges identical sub-DAGs after composition, so queries that chose the /// same summary producer reach one node; the returned Stage 1 candidate is /// the merged one. @@ -431,7 +430,7 @@ fn realize( let index = variant.number(choice); let mut logical = compose_logical_candidate(variant.inventory, choice) .map_err(|e| format!("Stage 1: {e}"))?; - if variant.shared { + if variant.sharing.merges_after_composition() { if let Some(merged) = share_identical_expressions(&logical) { logical = merged; } @@ -453,8 +452,8 @@ fn realize( /// One built combination; `physical` is `None` when it could not be built. #[derive(Debug, Clone)] pub struct EnumeratedCandidate { - /// From the shared variant (Pass 2's identical-expression rule). - pub shared: bool, + /// The Pass 2 variant it comes from. + pub sharing: Sharing, pub choice: Vec, pub logical: Option>, pub physical: Option, @@ -510,7 +509,7 @@ fn exhaustive( } }; candidates.push(EnumeratedCandidate { - shared: variant.shared, + sharing: variant.sharing, choice, logical, physical, @@ -544,10 +543,11 @@ fn exhaustive( /// The plan [`select_plan`] chose. #[derive(Debug, Clone)] pub struct SelectedPlan { - /// From the shared variant (Pass 2's identical-expression rule). - pub shared: bool, + /// The Pass 2 variant it comes from. + pub sharing: Sharing, pub choice: Vec, - /// The chosen Stage 1 candidate (identical sub-DAGs merged when `shared`). + /// The chosen Stage 1 candidate (identical sub-DAGs merged unless + /// `sharing` is independent). pub logical: Vec<(Id, QueryRoot)>, /// Stage 2 of `logical`; `roots` follow `logical`'s order. pub physical: PhysicalCandidate, @@ -646,10 +646,11 @@ fn costlier(id: &str, total: f64, best: f64) -> Rejection { /// not depend on choices. Stage 3 prices per node and sizes every /// realization of a target alike, so these hold unless a target's choice /// changes what a target reading its output costs or whether it can be -/// built, or, in a shared variant, two targets reading one input build +/// built, or, in a sharing variant, two targets reading one input build /// identical producers that are then merged. That coupling is checked: -/// every pair of choices for a target and a target beneath it (and, when -/// shared, for two targets reading a common input) is built, and must cost +/// every pair of choices for a target and a target beneath it (and, in a +/// sharing variant, for two targets reading a common or equal input) is +/// built, and must cost /// the sum of their single changes and be admissible exactly when both are. /// On coupling, or when the winner fails the full check, every combination /// of the variant is built instead if there are at most @@ -699,21 +700,27 @@ fn select_variant( .collect() }) .collect(); - let mut pairs: Vec<(usize, usize)> = Vec::new(); + // `(t, u, nested)`: `u` is read by `t`, or (not nested) both read a + // common input. + let mut pairs: Vec<(usize, usize, bool)> = Vec::new(); for (t, by_choice) in reads.iter().enumerate() { for &u in by_choice.iter().flatten() { - if !pairs.contains(&(t, u)) { - pairs.push((t, u)); + if !pairs.contains(&(t, u, true)) { + pairs.push((t, u, true)); } } } - if variant.shared { - pairs.extend(common_input_pairs(inventory, &beneath)); + if variant.sharing.merges_after_composition() { + pairs.extend( + common_input_pairs(inventory, &beneath) + .into_iter() + .map(|(t, u)| (t, u, false)), + ); } let mut coupling = None; - 'pairs: for &(t, u) in &pairs { + 'pairs: for &(t, u, nested) in &pairs { for c in 1..local[t].len() { - if !reads[t][c].contains(&u) { + if nested && !reads[t][c].contains(&u) { // `c` does not read `u`'s output (it absorbs it, or reads a // target beneath it only through another). continue; @@ -766,21 +773,26 @@ fn select_variant( } /// Pairs of targets, neither beneath the other, that read a common input -/// node: in a shared variant their producers may be merged. +/// node, or equal ones: in a sharing variant their producers may be merged. fn common_input_pairs( inventory: &LocalLogicalCandidates, beneath: &[Vec], ) -> Vec<(usize, usize)> { - let inputs: Vec> = inventory + let inputs: Vec>> = inventory .targets .iter() - .map(|t| t.target.children().into_iter().map(Rc::as_ptr).collect()) + .map(|t| t.target.children()) .collect(); + let same = |a: &Rc, b: &Rc| Rc::ptr_eq(a, b) || a == b; let mut pairs = Vec::new(); for t in 0..inputs.len() { for u in t + 1..inputs.len() { let nested = beneath[t].contains(&u) || beneath[u].contains(&t); - if !nested && inputs[t].iter().any(|p| inputs[u].contains(p)) { + if !nested + && inputs[t] + .iter() + .any(|a| inputs[u].iter().any(|b| same(a, b))) + { pairs.push((t, u)); } } @@ -833,7 +845,7 @@ fn finish( let (logical, physical) = realize(variant, &choice)?; let cost = assess(&physical, demand, data, models)?; Ok(SelectedPlan { - shared: variant.shared, + sharing: variant.sharing, choice, logical, selection: Selection { @@ -914,7 +926,7 @@ pub struct StagePipelineRun { } /// The #509 stage pipeline over `roots`: Stage 1 (Pass 1 and Pass 2's -/// identical-expression rule), Stage 2 and Stage 3. The facade and the +/// identical-expression and summary-capability rules), 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). pub fn plan_stages( @@ -1748,7 +1760,7 @@ mod tests { let plan = select_variant( Variant { inventory: &inventory, - shared: false, + sharing: Sharing::Independent, offset: 0, }, &targets, @@ -1784,7 +1796,7 @@ mod tests { let plan = select_variant( Variant { inventory: &inventory, - shared: false, + sharing: Sharing::Independent, offset: 0, }, &targets, @@ -1860,7 +1872,7 @@ mod tests { let plan = select_variant( Variant { inventory: &inventory, - shared: false, + sharing: Sharing::Independent, offset: 0, }, &targets, @@ -1903,7 +1915,7 @@ mod tests { let plan = finish( Variant { inventory: &inventory, - shared: true, + sharing: Sharing::IdenticalExpressions, offset: 0, }, &no_targets(&inventory), diff --git a/crates/planner/src/pass/stage_pipeline.rs b/crates/planner/src/pass/stage_pipeline.rs index dd6cbf29e..5d397e71c 100644 --- a/crates/planner/src/pass/stage_pipeline.rs +++ b/crates/planner/src/pass/stage_pipeline.rs @@ -3,9 +3,11 @@ //! //! [`plan_stages`] runs them: Stage 1 lists each target's local alternatives //! (Pass 1) with and without identical sub-DAGs shared across queries (Pass -//! 2's identical-expression rule), Stage 2 implements a candidate physically -//! (everything at query time), and Stage 3 checks accuracy, prices it per -//! second from each entry's recurrence and chooses. Pass 2's other rules are not planned yet. +//! 2's identical-expression rule) and with one summary sized for its +//! strictest consumer (Pass 2's summary-capability rule), Stage 2 implements +//! a candidate physically (everything at query time), and Stage 3 checks +//! accuracy, prices it per second from each entry's recurrence and chooses. +//! Pass 2's window-composition rule is not planned yet. use asap_types::ir::schema_support::with_promql_series_identity; use asap_types::ir::QueryRoot; diff --git a/crates/planner/tests/stage_pipeline_selection.rs b/crates/planner/tests/stage_pipeline_selection.rs index ac4fbeded..d2d85c1c0 100644 --- a/crates/planner/tests/stage_pipeline_selection.rs +++ b/crates/planner/tests/stage_pipeline_selection.rs @@ -5,7 +5,7 @@ use asap_frontend_sql::{lower_sql_dialect, SqlCatalog}; use asap_logical_optimizer::pass2::identical_expressions::{ - stage1_logical_candidates, SharingVariant, + stage1_logical_candidates, Sharing, SharingVariant, }; use asap_plan_selection::PlanningModels; use asap_plan_selection::{ @@ -148,6 +148,19 @@ fn assert_dp_matches_exhaustive( workload: &PlanningWorkload, combinations: usize, ) -> String { + let (selected, method, _) = selects_exhaustive_minimum(inventory, workload, combinations); + assert_eq!(method, SelectionMethod::TreeDp); + selected +} + +/// [`select_plan`] chooses what building and pricing every combination of +/// every variant chooses, however it gets there; returns the winner's id, +/// how `select_plan` found it, and its variant. +fn selects_exhaustive_minimum( + inventory: &Inventory, + workload: &PlanningWorkload, + combinations: usize, +) -> (String, SelectionMethod, Sharing) { let targets = targets(workload); let data = workload.data_workload.clone().unwrap_or_default(); let models = PlanningModels::builtin(); @@ -172,10 +185,17 @@ fn assert_dp_matches_exhaustive( .expect("winner was built"); let plan = select_plan(inventory, &targets, &data, models).expect("selects"); - assert_eq!(plan.selection.method, SelectionMethod::TreeDp); - assert_eq!((plan.shared, &plan.choice), (winner.shared, &winner.choice)); + assert!(plan.selection.guaranteed_optimal()); + assert_eq!( + (plan.sharing, &plan.choice), + (winner.sharing, &winner.choice) + ); assert_eq!(plan.selection.selected, exhaustive.selection.selected); - exhaustive.selection.selected + ( + exhaustive.selection.selected, + plan.selection.method, + plan.sharing, + ) } /// #509 Example 1: the dynamic program picks the cheapest of its 64 @@ -275,6 +295,110 @@ async fn sql_dp_equals_exhaustive() { assert_dp_matches_exhaustive(&inventory, &workload, 15); } +/// PromQL queries, each with its own ε (δ = 0.001). +fn promql_with(queries: &[(&str, f64)]) -> PlanningWorkload { + let mut workload = promql(&[], 1_000); + workload.query_workload.query_batch = Some( + queries + .iter() + .map(|(q, epsilon)| { + batch( + q, + AccuracyTarget::EpsilonDelta { + epsilon: *epsilon, + delta: 0.001, + }, + ) + }) + .collect(), + ); + workload +} + +/// p50 at ε=0.01 and p99 at ε=0.001 over one input: Stage 1 has an +/// independent, an identical-expression and a summary-capability variant +/// (pass-through, KLL, DDSketch per target: 9 combinations each). Sharing +/// one KLL couples the two targets, so selection builds every combination +/// of that variant, and picks the exhaustive minimum: the shared KLL. +#[test] +fn summary_capability_dp_equals_exhaustive() { + let workload = promql_with(&[ + ("quantile_over_time(0.5, lat[5m])", 0.01), + ("quantile_over_time(0.99, lat[5m])", 0.001), + ]); + let stage1 = promql_inventory(&workload); + assert_eq!( + stage1.iter().map(|v| v.sharing).collect::>(), + [ + Sharing::Independent, + Sharing::IdenticalExpressions, + Sharing::SummaryCapability + ] + ); + let (_, method, sharing) = selects_exhaustive_minimum(&stage1, &workload, 27); + assert_eq!(method, SelectionMethod::Exhaustive); + assert_eq!(sharing, Sharing::SummaryCapability); +} + +/// The same with an unrelated third query (pass-through or exact max): the +/// shared KLL is still the minimum over all 3 × 18 combinations. +#[test] +fn summary_capability_with_an_unrelated_query_dp_equals_exhaustive() { + let workload = promql_with(&[ + ("quantile_over_time(0.5, lat[5m])", 0.01), + ("quantile_over_time(0.99, lat[5m])", 0.001), + ("max_over_time(other[5m])", 0.01), + ]); + let stage1 = promql_inventory(&workload); + assert_eq!(stage1.len(), 3); + let (_, _, sharing) = selects_exhaustive_minimum(&stage1, &workload, 54); + assert_eq!(sharing, Sharing::SummaryCapability); +} + +/// SQL p50 and p99 over one filtered column: pre-ASAP CSE merges nothing +/// (the scan has no unique key), so the summary-capability variant is the +/// only sharing one, and the shared KLL is the exhaustive minimum. +#[tokio::test] +async fn sql_summary_capability_dp_equals_exhaustive() { + let accuracy = AccuracyTarget::Epsilon(0.01); + let queries = [ + "SELECT approx_percentile_cont(l_extendedprice, 0.5) FROM lineitem WHERE l_orderkey > 10", + "SELECT approx_percentile_cont(l_extendedprice, 0.99) FROM lineitem WHERE l_orderkey > 10", + ]; + let workload = PlanningWorkload { + query_workload: QueryWorkload { + language: QueryLanguage::SQL(SqlDialect::DataFusionSQL), + query_batch: Some(queries.iter().map(|q| batch(q, accuracy.clone())).collect()), + repeating_queries: None, + }, + data_workload: Some(DataWorkload { + arrival: DataArrival::AtRest, + ..Default::default() + }), + }; + let catalog = SqlCatalog::new().with_table( + "lineitem", + Schema::new(vec![ + Field::plain("l_orderkey", DataType::Int64, false), + Field::plain("l_extendedprice", DataType::Float64, false), + ]), + ); + let mut roots = Vec::new(); + for (index, query) in queries.iter().enumerate() { + let root = lower_sql_dialect(query, &catalog, SqlDialect::DataFusionSQL, accuracy.clone()) + .await + .expect("lowers"); + roots.push((index, QueryRoot::Operator(root))); + } + let stage1 = stage1_logical_candidates(roots).expect("Stage 1"); + assert_eq!( + stage1.iter().map(|v| v.sharing).collect::>(), + [Sharing::Independent, Sharing::SummaryCapability] + ); + let (_, _, sharing) = selects_exhaustive_minimum(&stage1, &workload, 18); + assert_eq!(sharing, Sharing::SummaryCapability); +} + /// Through the facade, Example 1 selects the exhaustive winner, P58: both /// queries exact, Q1's rate and sum and Q2's sum as exact accumulators, over /// one shared range selector. diff --git a/crates/planner/tests/summary_sharing.rs b/crates/planner/tests/summary_sharing.rs index 5aee0c3db..85ffdb5c4 100644 --- a/crates/planner/tests/summary_sharing.rs +++ b/crates/planner/tests/summary_sharing.rs @@ -235,7 +235,6 @@ fn unique_deployments(output: &PlanOutput) -> usize { /// equal-params subset of summary capability. Both plans hold the same `Rc`, /// so a consumer maintains it once. #[tokio::test] -#[ignore = "Stage 3 selects the raw plan; query-time summaries never cost less until Stage 2 plans materialization: #580"] async fn quantiles_with_equal_params_share_one_producer() { let output = plan_promql(&[ ("quantile_over_time(0.5, lat[5m])", 0.01), @@ -315,7 +314,9 @@ fn kll_k_for(epsilon: f64) -> u32 { /// for the strictest consumer when the cost model prefers that candidate; each /// reader's guarantee meets its own target. #[tokio::test] -#[ignore = "Pass 2 cross-query sharing is not planned by the stage pipeline: #580"] +#[ignore = "the stage pipeline shares the KLL sized for the strictest consumer, but attaches no \ + guarantee to plan roots, and Stage 3 ignores PlanningModels.cost, so the looser query \ + alone selects the raw plan: #580"] async fn quantiles_share_one_producer_sized_for_the_strictest_consumer() { let p50 = ("quantile_over_time(0.5, lat[5m])", 0.01); let p99 = ("quantile_over_time(0.99, lat[5m])", 0.001); @@ -338,10 +339,26 @@ async fn quantiles_share_one_producer_sized_for_the_strictest_consumer() { assert_eq!(kll_k(&alone.plans[0]), kll_k_for(0.01)); } +/// The stage pipeline's summary-capability rule: p50 at ε=0.01 and p99 at +/// ε=0.001 over one input read one KLL sized for ε=0.001, and Stage 3 accepts +/// each query against its own target. +#[tokio::test] +async fn stage_pipeline_shares_one_kll_sized_for_the_strictest_consumer() { + let p50 = ("quantile_over_time(0.5, lat[5m])", 0.01); + let p99 = ("quantile_over_time(0.99, lat[5m])", 0.001); + let output = plan_promql(&[p50, p99]).await; + assert!(same_states(&states(&output))); + assert_eq!(unique_deployments(&output), 1); + for plan in &output.plans { + assert_eq!(kll_k(plan), kll_k_for(0.001)); + } + let selection = output.selection.expect("Stage 3 ran"); + assert!(selection.costs.contains_key(&selection.selected)); +} + /// Cross-series quantiles name their KLL state after the input column, not the /// quantile, so p50 and p99 over one selector share it. #[tokio::test] -#[ignore = "Pass 2 cross-query sharing is not planned by the stage pipeline: #580"] async fn cross_series_p50_and_p99_share_one_producer() { let output = plan_promql(&[("quantile(0.5, lat)", 0.01), ("quantile(0.99, lat)", 0.01)]).await; assert!(same_states(&states(&output))); @@ -362,7 +379,6 @@ async fn identical_ungrouped_queries_share_their_producers() { /// The SQL frontend reaches the same sharing for two copies of one filtered /// percentile. #[tokio::test] -#[ignore = "Stage 3 selects the raw plan; query-time summaries never cost less until Stage 2 plans materialization: #580"] async fn identical_sql_percentiles_share_one_producer() { let query = "SELECT approx_percentile_cont(l_extendedprice, 0.5) FROM lineitem WHERE l_orderkey > 10"; @@ -371,11 +387,39 @@ async fn identical_sql_percentiles_share_one_producer() { assert_eq!(unique_deployments(&output), 1); } +/// SQL p50 and p99 over one filtered column share one KLL through the +/// summary-capability variant, although pre-ASAP CSE cannot merge their +/// unkeyed scans; each query keeps its own output column. +#[tokio::test] +async fn stage_pipeline_shares_sql_p50_and_p99() { + let output = plan_sql(&[ + "SELECT approx_percentile_cont(l_extendedprice, 0.5) FROM lineitem WHERE l_orderkey > 10", + "SELECT approx_percentile_cont(l_extendedprice, 0.99) FROM lineitem WHERE l_orderkey > 10", + ]) + .await; + assert!(same_states(&states(&output))); + assert_eq!(unique_deployments(&output), 1); + let names: Vec<_> = output + .plans + .iter() + .map(|plan| plan.root.schema.fields[0].name.clone()) + .collect(); + assert_eq!( + names, + [ + "approx_percentile_cont(lineitem.l_extendedprice,Float64(0.5))", + "approx_percentile_cont(lineitem.l_extendedprice,Float64(0.99))", + ] + ); +} + /// The quantile is a evaluation parameter: SQL p50 and p99 over one filtered /// column build one KLL, named after its input, while each query keeps its /// own output column. #[tokio::test] -#[ignore = "Pass 2 cross-query sharing is not planned by the stage pipeline: #580"] +#[ignore = "p50 and p99 now share one KLL with their own output names; the unshared cases \ + (different filter or column) select the raw plan, since an unshared query-time \ + summary never costs less until Stage 2 plans materialization: #580"] async fn sql_p50_and_p99_share_one_producer() { let p50 = "SELECT approx_percentile_cont(l_extendedprice, 0.5) FROM lineitem WHERE l_orderkey > 10"; From 4f7e5571284d22d484e315ec8c55857a1c7632f9 Mon Sep 17 00:00:00 2001 From: zzylol <50204836+zzylol@users.noreply.github.com> Date: Sun, 4 Oct 2026 18:19:53 +0000 Subject: [PATCH 2/2] fix: integrate with #593 and #594 Pass 1 takes the declared metric types since #593, and #594's year-long scan test selects its variant by the Sharing enum instead of a shared flag. Co-Authored-By: Claude Opus 5.5 --- crates/logical-optimizer/src/pass2/summary_capability.rs | 2 +- crates/plan-selection/src/lib.rs | 8 ++++---- crates/planner/tests/stage_pipeline_selection.rs | 2 +- 3 files changed, 6 insertions(+), 6 deletions(-) diff --git a/crates/logical-optimizer/src/pass2/summary_capability.rs b/crates/logical-optimizer/src/pass2/summary_capability.rs index 79213746f..586f2d714 100644 --- a/crates/logical-optimizer/src/pass2/summary_capability.rs +++ b/crates/logical-optimizer/src/pass2/summary_capability.rs @@ -207,7 +207,7 @@ mod tests { (i, QueryRoot::Operator(root)) }) .collect(); - enumerate_local_logical_candidates(roots).unwrap() + enumerate_local_logical_candidates(roots, &Default::default()).unwrap() } fn kll_k(inventory: &LocalLogicalCandidates, t: usize) -> u32 { diff --git a/crates/plan-selection/src/lib.rs b/crates/plan-selection/src/lib.rs index e7badb9b5..3711d1f8d 100644 --- a/crates/plan-selection/src/lib.rs +++ b/crates/plan-selection/src/lib.rs @@ -2254,10 +2254,10 @@ mod tests { ..every_10s(None) }) .collect(); - let total = |shared: bool| { + let total = |sharing: Sharing| { let variant = *variants(&stage1) .iter() - .find(|v| v.shared == shared) + .find(|v| v.sharing == sharing) .unwrap(); let raw = vec![0; variant.inventory.targets.len()]; let (_, candidate) = realize(variant, &raw).unwrap(); @@ -2271,8 +2271,8 @@ mod tests { .collect(); (cost.total, scan_rows) }; - let (shared, shared_scans) = total(true); - let (separate, separate_scans) = total(false); + let (shared, shared_scans) = total(Sharing::IdenticalExpressions); + let (separate, separate_scans) = total(Sharing::Independent); let year = 365.0 * 24.0 * 3_600.0 * 1_000_000.0 / 15.0; assert_eq!(shared_scans.len(), 1); assert!( diff --git a/crates/planner/tests/stage_pipeline_selection.rs b/crates/planner/tests/stage_pipeline_selection.rs index d2d85c1c0..ff48c51fc 100644 --- a/crates/planner/tests/stage_pipeline_selection.rs +++ b/crates/planner/tests/stage_pipeline_selection.rs @@ -390,7 +390,7 @@ async fn sql_summary_capability_dp_equals_exhaustive() { .expect("lowers"); roots.push((index, QueryRoot::Operator(root))); } - let stage1 = stage1_logical_candidates(roots).expect("Stage 1"); + let stage1 = stage1_logical_candidates(roots, &Default::default()).expect("Stage 1"); assert_eq!( stage1.iter().map(|v| v.sharing).collect::>(), [Sharing::Independent, Sharing::SummaryCapability]