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
2 changes: 1 addition & 1 deletion crates/devtools/src/bin/stage_pipeline.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
// cargo run -p asap-devtools --bin stage_pipeline -- \
// --example planner-layering-1 --max-candidates 128 --out planner-layering-example1.json
// --example planner-layering-1 --max-candidates 160 --out planner-layering-example1.json
// (also planner-layering-2: #509 Example 2, three SQL statistics over
// `flows`; planner-layering-3a and planner-layering-3b: Example 3, Patterns A
// and B; planner-layering-4a: Example 4, Pattern A repeated monthly;
Expand Down
12 changes: 7 additions & 5 deletions crates/devtools/tests/stage_pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,8 +10,9 @@ use serde_json::Value;

const COMMITTED: &str = "../../tools/dag-viewer/examples/planner-layering-example1.json";

/// Every Example 1 logical candidate (88) is displayed; the default cap is 64.
const EXAMPLE1: [&str; 4] = ["--example", "planner-layering-1", "--max-candidates", "128"];
/// Every Example 1 plan (136, from 88 logical candidates) is displayed; the
/// default cap is 64.
const EXAMPLE1: [&str; 4] = ["--example", "planner-layering-1", "--max-candidates", "160"];

fn generate(args: &[&str]) -> Value {
// Tests run in parallel and may generate the same document.
Expand Down Expand Up @@ -75,9 +76,10 @@ fn example1_document_is_valid_and_committed_fixture_is_current() {
.iter()
.map(|p| p["from_logical"].as_str().unwrap())
.collect();
// One all-query-time candidate per logical one, plus one with Q2's sum
// panes at ingestion time for each of the 24 with panes.
assert_eq!(physical.len(), candidates.len() + 24);
// One all-query-time candidate per logical one, plus, for each of the
// 24 with panes, one with Q2's sum panes at ingestion time and one with
// them kept at query time.
assert_eq!(physical.len(), candidates.len() + 48);
assert_eq!(sources, ids);
for candidate in physical {
assert!(candidate.get("cost").is_none(), "Stage 2 has no cost");
Expand Down
1 change: 1 addition & 0 deletions crates/executor/tests/hydra_cms_execution.rs
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,7 @@ fn node(id: u32, payload: Payload, output: Schema) -> PhysicalASAPDAGNode {
id: LogicalASAPNodeId(id),
payload,
output_state: ExecutionDataState::QUERY_ROWS,
kept: false,
output_schema: output,
guarantee: None,
coverage: None,
Expand Down
6 changes: 6 additions & 0 deletions crates/executor/tests/physical_dag.rs
Original file line number Diff line number Diff line change
Expand Up @@ -499,6 +499,7 @@ fn bind_post_asap_before_execution() {
id: planner_types::ir::export::LogicalASAPNodeId(id),
payload,
output_state: ExecutionDataState::QUERY_ROWS,
kept: false,
output_schema: (*schema).clone(),
guarantee: None,
};
Expand Down Expand Up @@ -654,6 +655,7 @@ fn source_batches_must_match_the_bound_schema() {
},
},
output_state: ExecutionDataState::QUERY_ROWS,
kept: false,
output_schema: (*expected).clone(),
guarantee: None,
}],
Expand Down Expand Up @@ -734,6 +736,7 @@ fn planner_semijoin_sort_limit_contract_at_both_phases() {
payload,
output_schema: (**schema).clone(),
output_state: ExecutionDataState::QUERY_ROWS,
kept: false,
guarantee: None,
};
let edge =
Expand Down Expand Up @@ -1211,6 +1214,7 @@ fn grouped_temporal_schema_compiles_and_executes_topk() {
operator: operation,
},
output_state: ExecutionDataState::QUERY_ROWS,
kept: false,
output_schema: (*input).clone(),
guarantee: None,
};
Expand Down Expand Up @@ -1297,6 +1301,7 @@ fn certified_pruning_rejects_missing_authoritative_values_after_recovery() {
id: planner_types::ir::export::LogicalASAPNodeId(2),
output_schema: (*schema).clone(),
output_state: ExecutionDataState::QUERY_ROWS,
kept: false,
guarantee: None,
payload: PhysicalASAPOperatorPayload::Relational {
operator: planner_types::ir::export::NonASAPOpKind::Join {
Expand Down Expand Up @@ -1415,6 +1420,7 @@ fn compiled_ingestion_binary_preserves_alignment_and_rejects_missing_updates() {
id: planner_types::ir::export::LogicalASAPNodeId(2),
output_schema: (*input).clone(),
output_state: ExecutionDataState::INGESTION_ROWS,
kept: false,
guarantee: None,
payload: PhysicalASAPOperatorPayload::Relational {
operator: planner_types::ir::export::NonASAPOpKind::BinaryOp {
Expand Down
3 changes: 3 additions & 0 deletions crates/executor/tests/physical_semantics.rs
Original file line number Diff line number Diff line change
Expand Up @@ -387,6 +387,7 @@ fn global_extrema_bind_with_planner_derived_schema() {
},
},
output_state: ExecutionDataState::QUERY_ROWS,
kept: false,
output_schema: (*output).clone(),
guarantee: None,
};
Expand Down Expand Up @@ -771,6 +772,7 @@ fn exact_frequency_intents_execute_typed_keys_and_empty_input() {
},
},
output_state: ExecutionDataState::QUERY_ROWS,
kept: false,
output_schema: (*schema(&[(name, DataType::Float64, false)])).clone(),
guarantee: None,
};
Expand Down Expand Up @@ -858,6 +860,7 @@ fn exact_cardinality_binds_and_executes_typed_tuples() {
},
},
output_state: ExecutionDataState::QUERY_ROWS,
kept: false,
output_schema: (*schema(&[("distinct", DataType::Int64, false)])).clone(),
guarantee: None,
};
Expand Down
7 changes: 7 additions & 0 deletions crates/executor/tests/precompute_population.rs
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,7 @@ fn finalized_shared_panes_rebuild_one_global_summary_after_recovery() {
id: planner_types::ir::export::LogicalASAPNodeId(0),
payload: PhysicalASAPOperatorPayload::SummaryMerge,
output_state: ExecutionDataState::INGESTION_SUMMARY,
kept: false,
output_schema: state_schema.clone(),
guarantee: None,
},
Expand All @@ -67,6 +68,7 @@ fn finalized_shared_panes_rebuild_one_global_summary_after_recovery() {
id: planner_types::ir::export::LogicalASAPNodeId(1),
payload: PhysicalASAPOperatorPayload::FinalizeExactAccumulator,
output_state: ExecutionDataState::INGESTION_ROWS,
kept: false,
output_schema: value_schema.clone(),
guarantee: None,
},
Expand All @@ -85,6 +87,7 @@ fn finalized_shared_panes_rebuild_one_global_summary_after_recovery() {
},
},
output_state: ExecutionDataState::INGESTION_ROWS,
kept: false,
output_schema: value_schema.clone(),
guarantee: None,
},
Expand All @@ -102,6 +105,7 @@ fn finalized_shared_panes_rebuild_one_global_summary_after_recovery() {
filter: None,
},
output_state: ExecutionDataState::INGESTION_SUMMARY,
kept: false,
output_schema: state_schema.clone(),
guarantee: None,
},
Expand Down Expand Up @@ -252,6 +256,7 @@ fn state_dag(
id: planner_types::ir::export::LogicalASAPNodeId(0),
payload: PhysicalASAPOperatorPayload::SummaryMerge,
output_state: ExecutionDataState::INGESTION_SUMMARY,
kept: false,
output_schema: logical_schema(family.clone()),
guarantee: None,
}];
Expand All @@ -268,6 +273,7 @@ fn state_dag(
id: planner_types::ir::export::LogicalASAPNodeId(read_id),
payload: PhysicalASAPOperatorPayload::FinalizeExactAccumulator,
output_state: ExecutionDataState::INGESTION_ROWS,
kept: false,
output_schema: logical_schema(FieldDataType::Plain(DataType::Float64)),
guarantee: None,
});
Expand All @@ -283,6 +289,7 @@ fn state_dag(
filter: None,
},
output_state: ExecutionDataState::INGESTION_SUMMARY,
kept: false,
output_schema: logical_schema(target),
guarantee: None,
});
Expand Down
2 changes: 2 additions & 0 deletions crates/executor/tests/promql_binary.rs
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,7 @@ fn program_for_bool(operator: BinaryOperator, return_bool: bool) -> CompiledPhys
},
},
output_state: ExecutionDataState::QUERY_ROWS,
kept: false,
output_schema: (*schema).clone(),
guarantee: None,
};
Expand Down Expand Up @@ -390,6 +391,7 @@ fn stored_series_evaluations_support_filters_and_sets() {
} else {
ExecutionDataState::QUERY_ROWS
},
kept: false,
output_schema: if id < 2 {
(*state_schema).clone()
} else {
Expand Down
1 change: 1 addition & 0 deletions crates/executor/tests/raw_scan.rs
Original file line number Diff line number Diff line change
Expand Up @@ -74,6 +74,7 @@ fn plan(
id: planner_types::ir::export::LogicalASAPNodeId(id),
payload,
output_state: state,
kept: false,
output_schema: (**schema).clone(),
guarantee: None,
};
Expand Down
2 changes: 2 additions & 0 deletions crates/executor/tests/summary_projection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -62,6 +62,7 @@ fn post_asap_summary_projection_survives_recovery() {
payload: PhysicalASAPOperatorPayload::SummaryMerge,
output_schema: (*schema).clone(),
output_state: ExecutionDataState::INGESTION_SUMMARY,
kept: false,
guarantee: None,
},
PhysicalASAPDAGNode {
Expand All @@ -81,6 +82,7 @@ fn post_asap_summary_projection_survives_recovery() {
},
output_schema: output.clone(),
output_state: ExecutionDataState::INGESTION_SUMMARY,
kept: false,
guarantee: None,
},
],
Expand Down
33 changes: 21 additions & 12 deletions crates/integration-tests/tests/planner_layering_common/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -208,10 +208,9 @@ pub fn run_promql(workload: &PlanningWorkload) -> Run {
}

/// Whether `form` answers a query window by merging several summaries.
pub fn needs_merge(form: WindowForm, query_window_ms: u64) -> bool {
pub fn needs_merge(form: WindowForm) -> bool {
match form {
WindowForm::None => false,
WindowForm::Sliding { length_ms, .. } => length_ms < query_window_ms,
WindowForm::Tumbling { .. } | WindowForm::ExponentialHistogram { .. } => true,
}
}
Expand Down Expand Up @@ -667,8 +666,6 @@ pub fn lookback_ms(workload: &PlanningWorkload, query: usize) -> u64 {
pub enum WindowForm {
/// Rebuilt from the query window's raw samples (no window summary).
None,
/// Windows of `length_ms` starting every `slide_ms`.
Sliding { length_ms: u64, slide_ms: u64 },
/// Back-to-back windows of `length_ms`.
Tumbling { length_ms: u64 },
/// EH buckets covering `horizon_ms` of history.
Expand All @@ -678,8 +675,8 @@ pub enum WindowForm {
/// A build merged by a `SummaryMerge` is a tumbling pane, of the width its
/// `TimeRange` input reads (#580: pane `i` is `TimeRange(w)` over
/// `TimeShift(i·w)` over the scan); any other build is rebuilt from its
/// query window. Sliding windows and Exponential Histograms are not planned
/// yet.
/// query window. Exponential Histograms are not planned yet; a sliding form
/// is superseded by tumbling panes (Q63).
pub fn window_form(dag: &impl ExportedDag, build: LogicalASAPNodeId) -> WindowForm {
let merged = dag
.consumers(build)
Expand Down Expand Up @@ -718,13 +715,14 @@ pub enum Materialization {
NotMaterialized,
}

/// Stage 2 runs a node at ingestion time or at query time, recomputed at
/// each evaluation; it does not keep query-time output yet (Example 4 B3),
/// so `QueryTimeKept` does not occur.
/// Stage 2 runs a node at ingestion time or at query time; a query-time
/// node is recomputed at each evaluation unless the export marks it kept
/// across evaluations (Example 4 B3).
pub fn materialization(p: &Physical, node: LogicalASAPNodeId) -> Materialization {
let n = p.dag.nodes.iter().find(|n| n.id == node).expect("node");
match n.output_state.timing {
ExecutionTiming::IngestionTime => Materialization::IngestionTime,
ExecutionTiming::QueryTime if n.kept => Materialization::QueryTimeKept,
ExecutionTiming::QueryTime => Materialization::NotMaterialized,
}
}
Expand All @@ -734,10 +732,12 @@ pub fn materialization(p: &Physical, node: LogicalASAPNodeId) -> Materialization
/// (`stage3-cost-model.md`): ingestion-time work read at query time through
/// a merge of `N` panes of width `w` keeps `(N + 1) · w`; read directly, the
/// window being built and the completed one, `2 · window`. Taken over every
/// query-time reader the node's ingestion-time work feeds. `None` for a
/// query-time node.
/// query-time reader the node's ingestion-time work feeds. A kept pane of a
/// merge of `N` panes of width `w` stays until the window no longer covers
/// it: `N · w`, the newest pane built from raw data and `N − 1` kept. `None`
/// for a node recomputed at each evaluation.
pub fn retention_ms(p: &Physical, node: LogicalASAPNodeId) -> Option<u64> {
if !runs_at_ingestion(p, node) {
if materialization(p, node) == Materialization::NotMaterialized {
return None;
}
// The longest raw range an ingestion-time node reads.
Expand All @@ -753,6 +753,15 @@ pub fn retention_ms(p: &Physical, node: LogicalASAPNodeId) -> Option<u64> {
.max()
.unwrap_or(0)
};
if materialization(p, node) == Materialization::QueryTimeKept {
return p
.dag
.consumers(node)
.into_iter()
.filter(|&c| matches!(p.dag.payload(c), LogicalASAPOperatorPayload::SummaryMerge))
.map(|merge| p.dag.producers(merge).len() as u64 * window(node))
.max();
}
let mut kept = 0;
let mut stack = vec![node];
let mut seen = HashSet::new();
Expand Down
Loading