Skip to content
Closed
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
24 changes: 24 additions & 0 deletions control_plane/src/backend_client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -375,6 +375,28 @@ impl BackendClient {
query_plan: &crate::query_plan::QueryPlan,
storage_routing: Option<serde_json::Value>,
adaptation_evidence: &[crate::physical::compiler::RuntimeAdaptationEvidence],
) -> std::result::Result<(), BackendPostError> {
self.post_physical_plan_with_sidecars(
precompute_plan,
transmission_plan,
query_plan,
storage_routing,
adaptation_evidence,
None,
None,
)
.await
}

pub async fn post_physical_plan_with_sidecars(
&self,
precompute_plan: &crate::physical::compiler::PrecomputePlan,
transmission_plan: &crate::physical::compiler::TransmissionPlan,
query_plan: &crate::query_plan::QueryPlan,
storage_routing: Option<serde_json::Value>,
adaptation_evidence: &[crate::physical::compiler::RuntimeAdaptationEvidence],
metricsql_plan: Option<&crate::query_plan::MetricsQlPlanCatalog>,
clickhouse_sql: Option<serde_json::Value>,
) -> std::result::Result<(), BackendPostError> {
// Compatibility replanner has no Planner-selected query/collector DAG.
// Still publish the actual catalog and bind every provided projection.
Expand Down Expand Up @@ -407,6 +429,8 @@ impl BackendClient {
"precompute_plan": precompute_plan,
"transmission_plan": transmission_plan,
"query_plan": query_plan,
"metricsql_plan": metricsql_plan,
"clickhouse_sql": clickhouse_sql,
"storage_routing": storage_routing,
"adaptation_evidence": adaptation_evidence,
}))
Expand Down
60 changes: 20 additions & 40 deletions control_plane/src/clickhouse.rs
Original file line number Diff line number Diff line change
Expand Up @@ -11,12 +11,11 @@ use planner_types::workload::SqlDialect;
use crate::physical::compiler::{PrecomputePlan, TransmissionPlan};
use crate::physical::post_asap::{cost_model::ControlPlaneCostModel, PhysicalExpr};
use crate::query_plan::{
ClickHousePlanningContext, FallbackPolicy, FixedEvaluationRange, InstantExecution,
MaterializationBinding, PhysicalGrouping, QueryLanguage, QueryPlan, QueryPlanEntry,
FallbackPolicy, InstantExecution, MaterializationBinding, PhysicalGrouping, QueryPlanEntry,
};
use asap_types::summary_catalog::SummaryCatalog;
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
use std::collections::{BTreeSet, HashMap};

#[derive(Debug, thiserror::Error)]
pub enum ClickHousePlanningError {
Expand Down Expand Up @@ -90,13 +89,13 @@ pub struct ClickHouseSqlWorkloadEntry {
pub cumulative: bool,
}

/// Physical-plan components produced for the normal atomic install path.
/// Wire bundle consumed by the independent backend SQL catalog.
#[derive(Debug, Serialize)]
pub struct ClickHouseCompiledBundle {
pub sds: SummaryCatalog,
pub tables: HashMap<String, planner_types::pre_asap::Schema>,
pub accuracy: AccuracyTarget,
pub query_plan: QueryPlan,
pub plans: Vec<serde_json::Value>,
pub precompute_plan: PrecomputePlan,
pub transmission_plan: TransmissionPlan,
}
Expand All @@ -115,7 +114,7 @@ pub async fn compile_clickhouse_workload(
let catalog = SqlCatalog {
tables: request.tables.clone(),
};
let mut entries = std::collections::BTreeMap::new();
let mut plans = Vec::with_capacity(request.queries.len());
for query in &request.queries {
let planned = plan_clickhouse_sql(&query.sql, &catalog, request.accuracy.clone()).await?;
let PhysicalExpr::Committed(crate::physical::post_asap::PostAsapPlan::Summary(root)) =
Expand All @@ -126,14 +125,7 @@ pub async fn compile_clickhouse_workload(
));
};
let executable = QueryPlanEntry::compile_bound_relational(
query.sql.clone(),
planned.canonical_sql.clone(),
&root,
FixedEvaluationRange {
start_ms: query.start_ms,
end_ms: query.end_ms,
cumulative: query.cumulative,
},
InstantExecution {
lookback_ms: query.end_ms.saturating_sub(query.start_ms),
full_history: query.start_ms == 0,
Expand All @@ -153,6 +145,10 @@ pub async fn compile_clickhouse_workload(
));
}
let bindings = executable.materialization_bindings();
let materializations = bindings
.iter()
.map(|binding| binding.materialization.fingerprint())
.collect::<BTreeSet<_>>();
let identities = bindings
.iter()
.map(|binding| {
Expand All @@ -167,37 +163,21 @@ pub async fn compile_clickhouse_workload(
})
})
.collect::<Result<Vec<_>, _>>()?;
// Descriptor references are already represented by each DAG's
// MaterializationBinding and validated through SummaryCatalog.
let _descriptor_ids = identities
.iter()
.map(|identity| {
(
&identity.summary_descriptor_id,
&identity.data_descriptor_id,
)
})
.collect::<Vec<_>>();
let identity = QueryPlan::catalog_key(QueryLanguage::ClickHouseSql, &planned.canonical_sql);
if entries.insert(identity.clone(), executable).is_some() {
return Err(ClickHousePlanningError::Lower(format!(
"duplicate canonical SQL query identity `{identity}`"
)));
}
plans.push(serde_json::json!({
"sql": planned.canonical_sql,
"runtime": { "start_ms": query.start_ms, "end_ms": query.end_ms,
"cumulative": query.cumulative,
"materializations": materializations,
"executable": executable },
"descriptors": { "summaries": identities.iter().map(|identity| identity.summary_descriptor_id.clone()).collect::<BTreeSet<_>>(),
"data": identities.iter().map(|identity| identity.data_descriptor_id.clone()).collect::<BTreeSet<_>>() }
}));
}
Ok(ClickHouseCompiledBundle {
sds: request.sds.clone(),
tables: request.tables.clone(),
accuracy: request.accuracy.clone(),
query_plan: QueryPlan {
plan_id: request.sds.plan_id,
plan_version: request.sds.plan_version,
clickhouse_context: Some(ClickHousePlanningContext {
tables: request.tables.clone(),
accuracy: request.accuracy.clone(),
}),
entries,
},
plans,
precompute_plan: request.precompute_plan.clone(),
transmission_plan: request.transmission_plan.clone(),
})
Expand Down Expand Up @@ -226,7 +206,7 @@ fn bind_selected_node(
window_ms: selected.slide_interval.saturating_mul(1000),
pane_origin_ms: selected.pane_origin_ms,
readout_lookback_ms: source_window.map(|seconds| seconds.saturating_mul(1000)),
item_labels: selected.aggregated_labels.labels.clone(),
item_labels: Default::default(),
})
}

Expand Down
1 change: 0 additions & 1 deletion control_plane/src/emit/backend_push.rs
Original file line number Diff line number Diff line change
Expand Up @@ -241,7 +241,6 @@ async fn push_documents_coupled(
let query_plan = crate::query_plan::QueryPlan {
plan_id: precompute_plan.envelope.plan_id,
plan_version: precompute_plan.envelope.plan_version,
clickhouse_context: None,
entries: Default::default(),
};
let transmission_plan = match crate::physical::compiler::TransmissionPlan::build(
Expand Down
11 changes: 9 additions & 2 deletions control_plane/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -822,13 +822,20 @@ async fn handle_compile_and_publish_clickhouse_plan(
)
.into_response();
};
let empty_query_plan = control_plane::query_plan::QueryPlan {
plan_id,
plan_version,
entries: Default::default(),
};
if let Err(error) = client
.post_physical_plan_typed(
.post_physical_plan_with_sidecars(
&bundle.precompute_plan,
&bundle.transmission_plan,
&bundle.query_plan,
&empty_query_plan,
None,
&[],
None,
Some(serde_json::to_value(&bundle).expect("SQL bundle serializes")),
)
.await
{
Expand Down
119 changes: 108 additions & 11 deletions control_plane/src/physical/compiler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@ use planner_types::post_asap::{
SketchQuery, SummaryExpr, SummaryFamilyType, SummaryMaintenanceLifecycle,
SummaryMaintenanceMode, SummaryNode, SummaryWindowFramework,
};
use planner_types::pre_asap::QueryExpr;
use planner_types::pre_asap::{AggIntent, QueryExpr, Reduction};
use planner_types::workload::{
AccuracyRequirement, DataArrival, DataWorkload, DurationMs, Evidence, EvidenceSource,
Predictability, Query, QueryLanguage, QueryRecurrence, QueryRequirements, QueryTimeScope,
Expand Down Expand Up @@ -737,6 +737,7 @@ pub struct PhysicalPlan {
pub precompute_plan: PrecomputePlan,
pub transmission_plan: TransmissionPlan,
pub query_plan: QueryPlan,
pub metricsql_plan: Option<crate::query_plan::MetricsQlPlanCatalog>,
/// Lifecycle component only, not a complete physical-plan comparison.
pub lifecycle_estimates: Vec<MaterializationLifecycleEstimate>,
pub cost_comparison: Option<super::workload_cost::WorkloadCostComparison>,
Expand Down Expand Up @@ -2047,6 +2048,41 @@ fn preserve_invalid_exact_fallback_roots(
Ok(())
}

/// Reject MetricsQL shapes whose VictoriaMetrics series-reduction semantics
/// are not represented faithfully by the current canonical executor.
pub fn validate_metricsql_acceleration_shape(expr: &QueryExpr) -> Result<(), &'static str> {
let QueryExpr::Aggregate {
reduction: Reduction::Reduce(_),
measures: outer_measures,
child,
..
} = expr
else {
return Ok(());
};
let QueryExpr::Aggregate {
reduction: Reduction::PerEntity,
measures: inner_measures,
child: inner_child,
..
} = child.as_ref()
else {
return Ok(());
};
if outer_measures.len() == 1
&& matches!(outer_measures[0], AggIntent::Sum { .. })
&& inner_measures.len() == 1
&& matches!(inner_child.as_ref(), QueryExpr::TimeRange { .. })
{
return match inner_measures[0] {
AggIntent::Rate => Err("nested sum(rate(...)) has VictoriaMetrics extrapolation semantics that the shared DAG cannot yet preserve"),
AggIntent::Increase => Err("nested sum(increase(...)) has VictoriaMetrics extrapolation semantics that the shared DAG cannot yet preserve"),
_ => Ok(()),
};
}
Ok(())
}

impl PhysicalCompiler {
pub fn compile(
&self,
Expand All @@ -2061,6 +2097,20 @@ impl PhysicalCompiler {
request: PlanningRequest,
environment: DeploymentEnvironment,
) -> Result<PhysicalPlan, CompileError> {
for query in &request.queries {
let expr = asap_frontend_metricsql::lower_metricsql(
&query.query_string,
query.accuracy.clone(),
)
.map_err(|error| CompileError::Query {
query_id: query.query_id.clone(),
reason: format!("frontend.metricsql.lowering: {error}"),
})?;
validate_metricsql_acceleration_shape(&expr).map_err(|reason| CompileError::Query {
query_id: query.query_id.clone(),
reason: format!("frontend.metricsql.unsupported: {reason}"),
})?;
}
self.compile_language(request, environment, true)
}

Expand Down Expand Up @@ -2506,6 +2556,7 @@ impl PhysicalCompiler {
.map(asap_types::PrecomputeMaterialization::policy_fingerprint)
.collect();
let mut query_entries = BTreeMap::new();
let mut metricsql_entries = BTreeMap::new();
for query in &request.queries {
let canonical = if metricsql {
asap_frontend_metricsql::canonical_metricsql(&query.query_string).map_err(
Expand Down Expand Up @@ -2609,11 +2660,19 @@ impl PhysicalCompiler {
// backend-local range index leaf.
crate::query_plan::logical::finalize_residuals(&mut entry)?;
}
if metricsql {
entry.language = crate::query_plan::QueryLanguage::MetricsQl;
}
let catalog_key = QueryPlan::catalog_key(entry.language, &canonical);
if query_entries.insert(catalog_key, entry).is_some() {
let duplicate = if metricsql {
let sidecar = crate::query_plan::MetricsQlPlanEntry {
query_id: entry.query_id.clone(),
canonical_metricsql: canonical.clone(),
executable: entry.executable(),
};
metricsql_entries
.insert(canonical.clone(), sidecar)
.is_some()
} else {
query_entries.insert(canonical.clone(), entry).is_some()
};
if duplicate {
return Err(CompileError::Query {
query_id: query.query_id.clone(),
reason: format!("duplicate canonical query identity `{canonical}`"),
Expand All @@ -2623,15 +2682,25 @@ impl PhysicalCompiler {
let query_plan = QueryPlan {
plan_id,
plan_version: envelope.plan_version,
clickhouse_context: None,
entries: query_entries,
};
let metricsql_plan = metricsql.then_some(crate::query_plan::MetricsQlPlanCatalog {
plan_id,
plan_version: envelope.plan_version,
entries: metricsql_entries,
});
for materialization in &mut precompute_plan.materializations {
let fingerprint = materialization.policy_fingerprint();
let max_lookback_ms = query_plan
.entries
.values()
.flat_map(QueryPlanEntry::materialization_bindings)
.chain(
metricsql_plan
.iter()
.flat_map(|catalog| catalog.entries.values())
.flat_map(|entry| entry.executable.materialization_bindings()),
)
.filter(|binding| binding.materialization.fingerprint() == fingerprint)
.filter_map(|binding| binding.readout_lookback_ms)
.max();
Expand Down Expand Up @@ -2701,6 +2770,7 @@ impl PhysicalCompiler {
precompute_plan,
transmission_plan,
query_plan,
metricsql_plan,
lifecycle_estimates: lifecycle_estimates.into_values().collect(),
cost_comparison: None,
})
Expand Down Expand Up @@ -4382,7 +4452,7 @@ mod tests {
}

#[test]
fn metricsql_compilation_publishes_a_language_tagged_query_entry() {
fn metricsql_compilation_publishes_an_independent_sidecar_entry() {
let query = "default_rollup(m[1m])";
let mut workload = request("vm-q", "last_over_time(m[1m])");
let accuracy = workload.queries[0].accuracy.clone();
Expand All @@ -4395,11 +4465,38 @@ mod tests {
.unwrap();
let identity = asap_frontend_metricsql::canonical_metricsql(query).unwrap();
let entry = plan
.query_plan
.lookup_canonical(crate::query_plan::QueryLanguage::MetricsQl, &identity)
.metricsql_plan
.as_ref()
.unwrap()
.lookup(&identity)
.unwrap();
assert_eq!(entry.query_id, "vm-q");
assert_eq!(entry.language, crate::query_plan::QueryLanguage::MetricsQl);
assert!(plan.query_plan.entries.is_empty());
}

#[test]
fn metricsql_non_equivalent_counter_rollups_fail_closed_before_publication() {
for (query, stage) in [
("sum(rate(m[5s]))", "rate"),
("sum(increase(m[5s]))", "increase"),
] {
let mut workload = request("vm-q", query);
let accuracy = workload.queries[0].accuracy.clone();
let canonical =
asap_frontend_metricsql::lower_metricsql(query, accuracy.clone()).unwrap();
workload.queries[0].post_asap = crate::planner_selection::select_summary(
&canonical,
&crate::physical::post_asap::cost_model::ForcedFamilyCostModel::new(
accuracy,
planner_types::post_asap::SketchAlgorithm::Kll,
),
)
.unwrap();
let error = PhysicalCompiler
.compile_metricsql(workload, environment(10_000))
.unwrap_err();
assert!(error.to_string().contains(stage), "{error}");
}
}

#[test]
Expand Down
Loading
Loading