Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
18 commits
Select commit Hold shift + click to select a range
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
77 changes: 38 additions & 39 deletions crates/asap-aware-mapping/src/analytical_cost.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@ use serde::{Deserialize, Serialize};

use crate::physical_operator_statistics::{
validate_comparison_scopes, ComparisonScope, EdgeStatistics, OperatorStatistics,
OperatorStatisticsProvider, PromqlEdgeStatistics, PromqlValueKind, SourceCoverage,
OperatorStatisticsProvider, PromqlEdgeStatistics, PromqlValueKind, ScanSelection,
};

/// Version of the analytical formulas applied to evidenced physical plans.
Expand Down Expand Up @@ -383,9 +383,9 @@ pub struct PhysicalDAGNode {
pub operator: PhysicalOperator,
pub children: Vec<String>,
/// Exact comparison-scope coverage consumed by a scan. Non-scan nodes
/// leave this empty. Reusing `SourceCoverage` prevents a physical plan
/// leave this empty. Reusing `ScanSelection` prevents a physical plan
/// from naming a source independently of its snapshot and predicates.
pub source_coverage: Option<SourceCoverage>,
pub scan_selection: Option<ScanSelection>,
/// Maximum transient edge buffer, distinct from logical `output_bytes`.
pub output_buffer_bytes: u64,
/// State that remains live after this node finishes (zero for ordinary
Expand Down Expand Up @@ -573,9 +573,10 @@ pub fn estimate_physical_dag_with_cache(
let node_statistics = &resolved_statistics[id];
match node.operator {
PhysicalOperator::Scan => {
let coverage = node.source_coverage.as_ref().ok_or_else(|| {
AnalyticalCostError::MissingScanSourceCoverage(node.id.clone())
})?;
let coverage = node
.scan_selection
.as_ref()
.ok_or_else(|| AnalyticalCostError::MissingScanSelection(node.id.clone()))?;
if !scope.sources.contains(coverage) {
return Err(AnalyticalCostError::ScanOutsideComparisonScope(
node.id.clone(),
Expand All @@ -585,9 +586,9 @@ pub fn estimate_physical_dag_with_cache(
consumed_sources.push(coverage);
}
}
_ if node.source_coverage.is_some() => {
_ if node.scan_selection.is_some() => {
return Err(AnalyticalCostError::InvalidPhysicalDAG(
"only scan nodes may declare source coverage",
"only scan nodes may declare scan selection",
));
}
_ => {}
Expand Down Expand Up @@ -2199,9 +2200,9 @@ pub enum AnalyticalCostError {
UnsupportedSummaryOperation(&'static str),
#[error("required comparison-scope field {0} is missing")]
MissingComparisonScope(&'static str),
#[error("scan node {0} does not declare source coverage")]
MissingScanSourceCoverage(String),
#[error("scan node {0} reads source coverage outside the comparison scope")]
#[error("scan node {0} does not declare scan selection")]
MissingScanSelection(String),
#[error("scan node {0} reads scan selection outside the comparison scope")]
ScanOutsideComparisonScope(String),
#[error("raw and candidate comparison scopes differ in {0}")]
ComparisonScopeMismatch(&'static str),
Expand Down Expand Up @@ -2234,7 +2235,7 @@ mod tests {
use crate::physical_operator_statistics::{
validate_comparison_scopes, BinaryEdgeStatistics, ComparisonScope, EdgeStatistics,
OperatorStatistics, PartitionStatistics, PromqlBinaryEdgeStatistics, PromqlEdgeStatistics,
PromqlUnaryEdgeStatistics, PromqlValueKind, SourceCoverage, UnaryEdgeStatistics,
PromqlUnaryEdgeStatistics, PromqlValueKind, ScanSelection, UnaryEdgeStatistics,
};

/// Analytical estimates reuse the shared dimensions while preserving exact
Expand Down Expand Up @@ -2475,7 +2476,7 @@ mod tests {
id: "scan".into(),
operator: PhysicalOperator::Scan,
children: vec![],
source_coverage: Some(comparison_scope().sources[0].clone()),
scan_selection: Some(comparison_scope().sources[0].clone()),
output_buffer_bytes: 8,
retained_bytes: 0,
execution: ExecutionMultiplicity::PerEvaluation,
Expand All @@ -2484,7 +2485,7 @@ mod tests {
id: "filter".into(),
operator: filter_operator(),
children: vec!["scan".into()],
source_coverage: None,
scan_selection: None,
output_buffer_bytes: 8,
retained_bytes: 0,
execution: ExecutionMultiplicity::PerEvaluation,
Expand Down Expand Up @@ -2814,7 +2815,7 @@ mod tests {
id: "scan".into(),
operator: PhysicalOperator::Scan,
children: vec![],
source_coverage: Some(coverage),
scan_selection: Some(coverage),
output_buffer_bytes: 10,
retained_bytes: 0,
execution: ExecutionMultiplicity::PerEvaluation,
Expand All @@ -2823,7 +2824,7 @@ mod tests {
id: "left".into(),
operator: filter_operator(),
children: vec!["scan".into()],
source_coverage: None,
scan_selection: None,
output_buffer_bytes: 4,
retained_bytes: 0,
execution: ExecutionMultiplicity::PerEvaluation,
Expand All @@ -2832,7 +2833,7 @@ mod tests {
id: "right".into(),
operator: filter_operator(),
children: vec!["scan".into()],
source_coverage: None,
scan_selection: None,
output_buffer_bytes: 4,
retained_bytes: 0,
execution: ExecutionMultiplicity::PerEvaluation,
Expand All @@ -2841,7 +2842,7 @@ mod tests {
id: "root".into(),
operator: PhysicalOperator::Concat,
children: vec!["left".into(), "right".into()],
source_coverage: None,
scan_selection: None,
output_buffer_bytes: 8,
retained_bytes: 0,
execution: ExecutionMultiplicity::PerEvaluation,
Expand Down Expand Up @@ -2905,7 +2906,7 @@ mod tests {
id: "scan".into(),
operator: PhysicalOperator::Scan,
children: vec![],
source_coverage: Some(coverage),
scan_selection: Some(coverage),
output_buffer_bytes: 10,
retained_bytes: 0,
execution: ExecutionMultiplicity::Once,
Expand All @@ -2914,7 +2915,7 @@ mod tests {
id: "state".into(),
operator: aggregate_operator(),
children: vec!["scan".into()],
source_coverage: None,
scan_selection: None,
output_buffer_bytes: 16,
retained_bytes: 32,
execution: ExecutionMultiplicity::Once,
Expand All @@ -2926,7 +2927,7 @@ mod tests {
offset: 0,
},
children: vec!["state".into()],
source_coverage: None,
scan_selection: None,
output_buffer_bytes: 16,
retained_bytes: 0,
execution: ExecutionMultiplicity::PerEvaluation,
Expand Down Expand Up @@ -2987,7 +2988,7 @@ mod tests {
lookback: Some(DurationMs(300_000)),
as_of: Some(TimestampMs(1_000)),
},
sources: vec![SourceCoverage {
sources: vec![ScanSelection {
source: Source::Table {
table_ref: "metrics".into(),
},
Expand Down Expand Up @@ -3059,7 +3060,7 @@ mod tests {
id: "scan".into(),
operator: PhysicalOperator::Scan,
children: vec![],
source_coverage: Some(coverage),
scan_selection: Some(coverage),
output_buffer_bytes: 10,
retained_bytes: 0,
execution: ExecutionMultiplicity::PerEvaluation,
Expand All @@ -3068,7 +3069,7 @@ mod tests {
id: "filter".into(),
operator: filter_operator(),
children: vec!["scan".into()],
source_coverage: None,
scan_selection: None,
output_buffer_bytes: 4,
retained_bytes: 0,
execution: ExecutionMultiplicity::PerEvaluation,
Expand Down Expand Up @@ -3134,7 +3135,7 @@ mod tests {
id: "scan".into(),
operator: PhysicalOperator::Scan,
children: vec![],
source_coverage: Some(coverage),
scan_selection: Some(coverage),
output_buffer_bytes: 10,
retained_bytes: 0,
execution: ExecutionMultiplicity::PerEvaluation,
Expand All @@ -3143,7 +3144,7 @@ mod tests {
id: "filter".into(),
operator: filter_operator(),
children: vec!["scan".into()],
source_coverage: None,
scan_selection: None,
output_buffer_bytes: 4,
retained_bytes: 0,
execution: ExecutionMultiplicity::PerEvaluation,
Expand Down Expand Up @@ -3197,7 +3198,7 @@ mod tests {
id: "scan".into(),
operator: PhysicalOperator::Scan,
children: vec![],
source_coverage: Some(coverage),
scan_selection: Some(coverage),
output_buffer_bytes: 10,
retained_bytes: 0,
execution: ExecutionMultiplicity::PerEvaluation,
Expand All @@ -3206,7 +3207,7 @@ mod tests {
id: "filter".into(),
operator: filter_operator(),
children: vec!["scan".into()],
source_coverage: None,
scan_selection: None,
output_buffer_bytes: 0,
retained_bytes: 0,
execution: ExecutionMultiplicity::PerEvaluation,
Expand Down Expand Up @@ -3245,7 +3246,7 @@ mod tests {
id: "scan".into(),
operator: PhysicalOperator::Scan,
children: vec![],
source_coverage: Some(SourceCoverage {
scan_selection: Some(ScanSelection {
source: asap_types::pre_asap::query_expr::Source::Table {
table_ref: "other_metrics".into(),
},
Expand Down Expand Up @@ -3283,7 +3284,7 @@ mod tests {
id: "scan".into(),
operator: PhysicalOperator::Scan,
children: vec![],
source_coverage: None,
scan_selection: None,
output_buffer_bytes: 10,
retained_bytes: 0,
execution: ExecutionMultiplicity::PerEvaluation,
Expand All @@ -3302,17 +3303,15 @@ mod tests {

assert_eq!(
estimate_physical_dag(&nodes, "scan", &comparison_scope(), &provided),
Err(AnalyticalCostError::MissingScanSourceCoverage(
"scan".into()
))
Err(AnalyticalCostError::MissingScanSelection("scan".into()))
);
}

#[test]
fn physical_dag_rejects_an_unconsumed_scope_source() {
let mut scope = comparison_scope();
let coverage = scope.sources[0].clone();
scope.sources.push(SourceCoverage {
scope.sources.push(ScanSelection {
source: asap_types::pre_asap::query_expr::Source::Table {
table_ref: "auxiliary".into(),
},
Expand All @@ -3324,7 +3323,7 @@ mod tests {
id: "scan".into(),
operator: PhysicalOperator::Scan,
children: vec![],
source_coverage: Some(coverage),
scan_selection: Some(coverage),
output_buffer_bytes: 10,
retained_bytes: 0,
execution: ExecutionMultiplicity::PerEvaluation,
Expand Down Expand Up @@ -3357,7 +3356,7 @@ mod tests {
id: "scan".into(),
operator: PhysicalOperator::Scan,
children: vec![],
source_coverage: Some(coverage),
scan_selection: Some(coverage),
output_buffer_bytes: 10,
retained_bytes: 0,
execution: ExecutionMultiplicity::PerEvaluation,
Expand All @@ -3366,7 +3365,7 @@ mod tests {
id: "aggregate".into(),
operator: aggregate_operator(),
children: vec!["scan".into()],
source_coverage: None,
scan_selection: None,
output_buffer_bytes: 16,
retained_bytes: 32,
execution: ExecutionMultiplicity::Once,
Expand Down Expand Up @@ -3411,7 +3410,7 @@ mod tests {
id: "scan".into(),
operator: PhysicalOperator::Scan,
children: vec![],
source_coverage: Some(coverage),
scan_selection: Some(coverage),
output_buffer_bytes: 10,
retained_bytes: 0,
execution: ExecutionMultiplicity::PerEvaluation,
Expand Down Expand Up @@ -3642,7 +3641,7 @@ mod tests {
id: "scan".into(),
operator: PhysicalOperator::Scan,
children: vec![],
source_coverage: Some(comparison_scope().sources[0].clone()),
scan_selection: Some(comparison_scope().sources[0].clone()),
output_buffer_bytes: 0,
retained_bytes: 0,
execution: ExecutionMultiplicity::PerEvaluation,
Expand Down
8 changes: 4 additions & 4 deletions crates/asap-aware-mapping/src/physical_operator_statistics.rs
Original file line number Diff line number Diff line change
Expand Up @@ -26,12 +26,12 @@ pub struct ComparisonScope {
pub horizon: DurationMs,
pub recurrence: QueryRecurrence,
pub time_selection: TimeSelection,
pub sources: Vec<SourceCoverage>,
pub sources: Vec<ScanSelection>,
}

/// Exact source selection covered by a physical plan.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct SourceCoverage {
pub struct ScanSelection {
pub source: Source,
/// Provider-owned stable identifier for the physical source contents,
/// such as a catalog snapshot, table version, or object generation.
Expand All @@ -53,7 +53,7 @@ impl ComparisonScope {
query: &QueryWorkloadEntry,
planning_time: TimestampMs,
horizon: DurationMs,
sources: Vec<SourceCoverage>,
sources: Vec<ScanSelection>,
) -> Result<Self, AnalyticalCostError> {
let scope = Self {
data_arrival: data.arrival,
Expand Down Expand Up @@ -87,7 +87,7 @@ impl ComparisonScope {
.any(|(index, source)| self.sources[..index].contains(source))
{
return Err(AnalyticalCostError::MissingComparisonScope(
"duplicate source coverage",
"duplicate scan selection",
));
}
if self
Expand Down
12 changes: 6 additions & 6 deletions crates/asap-aware-mapping/src/physical_plan_cost_model.rs
Original file line number Diff line number Diff line change
Expand Up @@ -367,7 +367,7 @@ mod tests {

use crate::analytical_cost::{ExecutionMultiplicity, PhysicalDAGNode, PhysicalOperator};
use crate::physical_operator_statistics::{
EdgeStatistics, OperatorStatistics, SourceCoverage, UnaryEdgeStatistics,
EdgeStatistics, OperatorStatistics, ScanSelection, UnaryEdgeStatistics,
};
use crate::replacement::ReplacementStrategy;

Expand Down Expand Up @@ -438,7 +438,7 @@ mod tests {
lookback: Some(DurationMs(10_000)),
as_of: Some(TimestampMs(1_000)),
},
sources: vec![SourceCoverage {
sources: vec![ScanSelection {
source: Source::Table {
table_ref: "events".into(),
},
Expand Down Expand Up @@ -506,7 +506,7 @@ mod tests {
id: "candidate-scan".into(),
operator: PhysicalOperator::Scan,
children: vec![],
source_coverage: Some(scope.sources[0].clone()),
scan_selection: Some(scope.sources[0].clone()),
output_buffer_bytes: 8,
retained_bytes: 0,
execution: ExecutionMultiplicity::Once,
Expand All @@ -518,7 +518,7 @@ mod tests {
accumulator_count: 1,
},
children: vec!["candidate-scan".into()],
source_coverage: None,
scan_selection: None,
output_buffer_bytes: 8,
retained_bytes: 8,
execution: ExecutionMultiplicity::Once,
Expand All @@ -527,7 +527,7 @@ mod tests {
id: "candidate-read".into(),
operator: PhysicalOperator::PassThrough,
children: vec!["candidate-state".into()],
source_coverage: None,
scan_selection: None,
output_buffer_bytes: 8,
retained_bytes: 0,
execution: ExecutionMultiplicity::PerEvaluation,
Expand Down Expand Up @@ -1006,7 +1006,7 @@ mod tests {
) -> Result<PhysicalDAG, AnalyticalCostError> {
let mut dag = self.0.summary_physical_dag(snapshot, summary, target)?;
dag.nodes[0]
.source_coverage
.scan_selection
.as_mut()
.unwrap()
.source_snapshot_id = "other".into();
Expand Down
Loading
Loading