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
53 changes: 51 additions & 2 deletions crates/devtools/src/bin/stage_pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,10 @@
// The deployment inputs are the built-in cost and accuracy models and the
// reference executor's capabilities (`asap_executor::capabilities`).
//
// - deployment: the deployment inputs Stage 3 used (the executor's
// capabilities, summarized; the cost model and its Stage3Calibration;
// the accuracy model's name).
//
// 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
// displayed candidate; the facade's dynamic program selects the same winner
Expand All @@ -46,7 +50,9 @@ use asap_logical_optimizer::pass1::logical_candidates::{
};
use asap_logical_optimizer::Realization;
use asap_plan_selection::PlanningModels;
use asap_plan_selection::{plan_stages, Selection, Sharing, MAX_ENUMERATED_CANDIDATES};
use asap_plan_selection::{
plan_stages, Selection, Sharing, COST_MODEL, COST_PER_SECOND, MAX_ENUMERATED_CANDIDATES,
};
use asap_types::ir::export::{
compile_logical_asap_workload, LogicalASAPDAG, LogicalASAPDAGDocument,
};
Expand Down Expand Up @@ -138,11 +144,12 @@ fn stage_pipeline(workload: &PlanningWorkload, max_candidates: usize) -> Result<
.collect();
let data = workload.data_workload.clone().unwrap_or_default();
let capabilities = asap_executor::capabilities();
let models = PlanningModels::builtin().with_capabilities(&capabilities);
let run = plan_stages(
roots.into_iter().enumerate().collect(),
&demand,
&data,
PlanningModels::builtin().with_capabilities(&capabilities),
models,
max_candidates.max(1),
)
.map_err(|e| format!("planning: {e}"))?;
Expand Down Expand Up @@ -190,6 +197,7 @@ fn stage_pipeline(workload: &PlanningWorkload, max_candidates: usize) -> Result<
Ok(json!({
"format": "asap-stage-pipeline/v1",
"workload": { "queries": workload_queries(workload) },
"deployment": deployment_json(&models),
"stage0_logical": { "dag": stage0 },
"stage1_logical_asap": {
"combinations": combinations,
Expand All @@ -201,6 +209,47 @@ fn stage_pipeline(workload: &PlanningWorkload, max_candidates: usize) -> Result<
}))
}

/// The deployment inputs Stage 3 used: the executor's capabilities, one
/// line per summary, and the cost and accuracy models.
fn deployment_json(models: &PlanningModels<'_>) -> Value {
let caps = models.capabilities;
let summaries = caps.summaries.as_ref().map(|summaries| {
summaries
.iter()
.map(|s| {
let readouts: Vec<_> = s.readouts.iter().map(|r| format!("{r:?}")).collect();
json!({ "summary": s.name(), "readouts": readouts })
})
.collect::<Vec<_>>()
});
let c = &models.calibration;
json!({
"capabilities": {
"source": "asap_executor::capabilities",
"summaries": summaries,
"ingestion_time": caps.ingestion_time,
"query_time_retention": caps.query_time_retention,
"memory_budget_bytes": caps.memory_budget_bytes,
"raw_data_retained": caps.raw_data_retained,
"raw_bytes_per_sample": caps.raw_bytes_per_sample,
},
"cost_model": {
"name": COST_MODEL,
"unit": COST_PER_SECOND,
"calibration": {
"version": c.version,
"cost_per_cpu_op": c.cost_per_cpu_op,
"cost_per_scan_byte": c.cost_per_scan_byte,
"cost_per_retained_byte_second": c.cost_per_retained_byte_second,
"horizon_s": c.horizon_s,
"latency_ms_per_cost_unit": c.latency_ms_per_cost_unit,
},
},
// `PlanningModels::builtin()`'s accuracy model and evidence.
"accuracy_model": { "name": "DefaultAccuracyModel", "evidence": "none" },
})
}

fn stage3_json(selection: &Selection) -> Value {
let costs: serde_json::Map<_, _> = selection
.costs
Expand Down
35 changes: 34 additions & 1 deletion crates/devtools/tests/stage_pipeline.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,9 +14,12 @@ const COMMITTED: &str = "../../tools/dag-viewer/examples/planner-layering-exampl
const EXAMPLE1: [&str; 4] = ["--example", "planner-layering-1", "--max-candidates", "128"];

fn generate(args: &[&str]) -> Value {
// Tests run in parallel and may generate the same document.
static NEXT: std::sync::atomic::AtomicUsize = std::sync::atomic::AtomicUsize::new(0);
let out = std::env::temp_dir().join(format!(
"stage_pipeline_{}_{}.json",
"stage_pipeline_{}_{}_{}.json",
std::process::id(),
NEXT.fetch_add(1, std::sync::atomic::Ordering::Relaxed),
args.join("_")
));
let status = Command::new(env!("CARGO_BIN_EXE_stage_pipeline"))
Expand Down Expand Up @@ -191,3 +194,33 @@ fn example4a_repeats_monthly_with_nothing_maintainable() {
};
assert!(cost(&monthly) < cost(&once));
}

/// The document records the deployment inputs Stage 3 used: the executor's
/// capabilities (Hydra summaries named by their kind), the cost model with
/// its calibration, including the memory weight and raw-retention settings,
/// and the accuracy model's name.
#[test]
fn document_records_the_deployment_inputs() {
let document = generate(&EXAMPLE1);
let deployment = &document["deployment"];
let capabilities = &deployment["capabilities"];
assert_eq!(capabilities["source"], "asap_executor::capabilities");
assert_eq!(capabilities["ingestion_time"], true);
assert_eq!(capabilities["raw_data_retained"], true);
assert_eq!(capabilities["raw_bytes_per_sample"], 16);
assert!(capabilities["memory_budget_bytes"].is_null());
let summaries = capabilities["summaries"].as_array().unwrap();
let hydra = summaries
.iter()
.find(|s| s["summary"] == "HydraCms")
.expect("HydraCms is listed");
assert_eq!(
hydra["readouts"],
serde_json::json!(["TotalCount", "ItemCount"])
);
let calibration = &deployment["cost_model"]["calibration"];
assert_eq!(deployment["cost_model"]["name"], "analytical-cost-v2");
assert_eq!(calibration["cost_per_retained_byte_second"], 1.25e-7);
assert_eq!(calibration["version"], "illustrative-v2");
assert_eq!(deployment["accuracy_model"]["name"], "DefaultAccuracyModel");
}
5 changes: 4 additions & 1 deletion crates/plan-selection/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -100,6 +100,9 @@ use asap_physical_optimizer::materialization::MAX_PHYSICAL_PER_LOGICAL;
/// [`Stage3Calibration::ILLUSTRATIVE`]. See
/// `docs/design_docs/proposals/stage3-cost-model.md`.
pub const COST_PER_SECOND: &str = "cost_per_second";
/// The analytical cost model Stage 3 prices with, as reported in each
/// candidate's cost `source`.
pub const COST_MODEL: &str = "analytical-cost-v2";

/// Groups assumed for every `by (...)` reduction, absent group-count evidence.
const DEFAULT_GROUP_COUNT: u64 = 100;
Expand Down Expand Up @@ -1877,7 +1880,7 @@ fn price_nodes(
total: per_node.values().map(|n| n.cost).sum(),
unit: COST_PER_SECOND,
source: format!(
"analytical-cost-v2 (illustrative statistics, calibration {})",
"{COST_MODEL} (illustrative statistics, calibration {})",
calibration.version
),
per_node,
Expand Down
7 changes: 7 additions & 0 deletions crates/types/src/deployment.rs
Original file line number Diff line number Diff line change
Expand Up @@ -101,6 +101,13 @@ pub struct SummarySupport {
pub readouts: BTreeSet<Readout>,
}

impl SummarySupport {
/// E.g. "Kll", "HydraCms" or "exact Sum".
pub fn name(&self) -> String {
name(&self.family, &self.layout)
}
}

/// A summary family, without its sizing parameters.
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
pub enum SummaryFamily {
Expand Down
100 changes: 100 additions & 0 deletions tools/dag-viewer/examples/planner-layering-example1.json
Original file line number Diff line number Diff line change
@@ -1,4 +1,104 @@
{
"deployment": {
"accuracy_model": {
"evidence": "none",
"name": "DefaultAccuracyModel"
},
"capabilities": {
"ingestion_time": true,
"memory_budget_bytes": null,
"query_time_retention": false,
"raw_bytes_per_sample": 16,
"raw_data_retained": true,
"source": "asap_executor::capabilities",
"summaries": [
{
"readouts": [],
"summary": "exact Sum"
},
{
"readouts": [],
"summary": "exact Count"
},
{
"readouts": [],
"summary": "exact Min"
},
{
"readouts": [],
"summary": "exact Max"
},
{
"readouts": [],
"summary": "exact Increase"
},
{
"readouts": [],
"summary": "exact Rate"
},
{
"readouts": [
"Quantile"
],
"summary": "Kll"
},
{
"readouts": [
"Quantile",
"TotalCount"
],
"summary": "DDSketch"
},
{
"readouts": [
"TotalCount",
"Cardinality"
],
"summary": "Hll"
},
{
"readouts": [
"TotalCount",
"ItemCount"
],
"summary": "HydraCms"
},
{
"readouts": [
"TopK"
],
"summary": "CmsWithHeap"
},
{
"readouts": [
"TopK"
],
"summary": "CountSketchWithHeap"
},
{
"readouts": [
"TotalCount",
"Cardinality",
"FrequencyL2",
"FrequencyEntropy"
],
"summary": "UnivMon"
}
]
},
"cost_model": {
"calibration": {
"cost_per_cpu_op": 1e-6,
"cost_per_retained_byte_second": 1.25e-7,
"cost_per_scan_byte": 1e-7,
"horizon_s": 3600.0,
"latency_ms_per_cost_unit": 1.0,
"version": "illustrative-v2"
},
"name": "analytical-cost-v2",
"unit": "cost_per_second"
}
},
"format": "asap-stage-pipeline/v1",
"stage0_logical": {
"dag": {
Expand Down