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
19 changes: 11 additions & 8 deletions crates/devtools/src/bin/stage_pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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::export::{
compile_logical_asap_workload, LogicalASAPDAG, LogicalASAPDAGDocument,
};
Expand Down Expand Up @@ -134,13 +135,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);
}
Expand All @@ -150,9 +151,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
Expand Down
10 changes: 9 additions & 1 deletion crates/integration-tests/tests/planner_layering_example1.rs
Original file line number Diff line number Diff line change
Expand Up @@ -164,7 +164,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,
Expand Down
29 changes: 28 additions & 1 deletion crates/logical-optimizer/src/pass1/logical_candidates.rs
Original file line number Diff line number Diff line change
Expand Up @@ -574,7 +574,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
Expand Down
53 changes: 43 additions & 10 deletions crates/logical-optimizer/src/pass2/identical_expressions.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Id> {
/// Whether identical sub-DAGs across queries are merged.
pub shared: bool,
pub sharing: Sharing,
pub inventory: LocalLogicalCandidates<Id>,
}

/// 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<Id: Clone>(
roots: Vec<(Id, QueryRoot)>,
metric_types: &BTreeMap<String, MetricType>,
) -> Result<Vec<SharingVariant<Id>>, 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)
}

Expand Down Expand Up @@ -114,8 +147,8 @@ mod tests {
)
.unwrap();
assert_eq!(
variants.iter().map(|v| v.shared).collect::<Vec<_>>(),
[false, true]
variants.iter().map(|v| v.sharing).collect::<Vec<_>>(),
[Sharing::Independent, Sharing::IdenticalExpressions]
);
assert_eq!(
variants[0].inventory.targets.len(),
Expand Down Expand Up @@ -145,6 +178,6 @@ mod tests {
)
.unwrap();
assert_eq!(variants.len(), 1);
assert!(!variants[0].shared);
assert_eq!(variants[0].sharing, Sharing::Independent);
}
}
4 changes: 3 additions & 1 deletion crates/logical-optimizer/src/pass2/mod.rs
Original file line number Diff line number Diff line change
@@ -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;
2 changes: 1 addition & 1 deletion crates/logical-optimizer/src/pass2/reconciliation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
Expand Down
Loading