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
16 changes: 11 additions & 5 deletions crates/devtools/src/bin/stage_pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,11 +17,13 @@
// enumeration order and capped by `--max-candidates` (default 64). A
// repeating query's mergeable alternatives also come in tumbling panes
// (Pass 2's window-composition rule), e.g. "Q1 Kll · tumbling 1m panes";
// - stage2_physical_asap: one physical candidate per logical candidate
// (operator implementation only, everything at query time), no cost;
// - stage2_physical_asap: per logical candidate, its physical candidates
// (operator implementation; everything at query time, then one per
// down-closed set of summaries maintained at ingestion time, labeled
// e.g. "· ingestion time: Kll ×5 panes"), no cost;
// - stage3_selection: per-candidate costs, the selected candidate, and
// every other candidate as rejected (`valid: false`, including one that
// could not be built) or costlier.
// every other candidate as rejected (`valid: false`: inaccurate, over a
// latency bound, or could not be built) or costlier.
//
// Everything is the library's `plan_selection::plan_stages`, the function the
// facade runs; this tool only serializes it. Stage 3 here is over every
Expand Down Expand Up @@ -166,7 +168,11 @@ fn stage_pipeline(workload: &PlanningWorkload, max_candidates: usize) -> Result<
candidates
.push(json!({ "id": format!("L{index}"), "label": label, "dag": export(&roots)? }));
}
if let Some(p) = &candidate.physical {
for p in &candidate.physical {
let label = match p.materialization.as_str() {
"" => label.clone(),
m => format!("{label} · {m}"),
};
stage2.push(
json!({ "id": p.id, "from_logical": p.from_logical, "label": label, "dag": p.dag }),
);
Expand Down
8 changes: 5 additions & 3 deletions crates/devtools/tests/stage_pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ use serde_json::Value;

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

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

fn generate(args: &[&str]) -> Value {
Expand Down Expand Up @@ -51,7 +51,7 @@ fn validated_dag(value: &Value, queries: usize) {

/// Every logical DAG is a well-formed flat DAG and every physical DAG validates;
/// each has one root per query; ids are unique;
/// Stage 2 maps Stage 1 one-to-one with no cost; Stage 3 accounts for every
/// Stage 2 covers every Stage 1 candidate, with no cost; Stage 3 accounts for every
/// Stage 2 candidate once and costs exactly the valid ones; the committed
/// fixture is current.
#[test]
Expand Down Expand Up @@ -86,7 +86,9 @@ fn example1_document_is_valid_and_committed_fixture_is_current() {
.iter()
.map(|p| p["from_logical"].as_str().unwrap())
.collect();
assert_eq!(physical.len(), candidates.len());
// 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);
assert_eq!(sources, ids);
for candidate in physical {
assert!(candidate.get("cost").is_none(), "Stage 2 has no cost");
Expand Down
20 changes: 11 additions & 9 deletions crates/integration-tests/tests/planner_layering_common/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -357,7 +357,7 @@ pub fn run_stages(workload: &PlanningWorkload, roots: Vec<QueryRoot>) -> Run {
let mut logical = Vec::new();
let mut physical = Vec::new();
for c in enumeration.candidates {
let p = c.physical.expect("every candidate builds");
assert!(!c.physical.is_empty(), "every candidate builds");
let roots: Vec<_> = c
.logical
.expect("composes")
Expand All @@ -366,18 +366,20 @@ pub fn run_stages(workload: &PlanningWorkload, roots: Vec<QueryRoot>) -> Run {
.collect();
let (dag, query_roots) = export(&roots);
logical.push(Logical {
id: p.from_logical.clone(),
id: c.physical[0].from_logical.clone(),
shared_input: c.sharing.merges_after_composition(),
dag,
query_roots,
});
physical.push(Physical {
id: p.id.clone(),
from_logical: p.from_logical.clone(),
dag: p.dag.clone(),
query_roots: p.dag.roots.clone(),
stage2: p,
});
for p in c.physical {
physical.push(Physical {
id: p.id.clone(),
from_logical: p.from_logical.clone(),
dag: p.dag.clone(),
query_roots: p.dag.roots.clone(),
stage2: p,
});
}
}
Run {
stage0,
Expand Down
Loading