Repository navigation
refactor(query-engine): make native DAG execution self-contained #788
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: main
Are you sure you want to change the base?
Changes from 1 commit
337ae1f
b53ad94
3539309
659bff8
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -5,30 +5,78 @@ use crate::engines::simple_engine::{RangeQueryExecutionContext, StoreQueryParams | |
| use asap_types::enums::WindowType; | ||
| use asap_types::query_config::QueryTimeAggregation; | ||
| use promql_utilities::data_model::KeyByLabelNames; | ||
| use promql_utilities::query_logics::enums::Statistic; | ||
| use promql_utilities::query_logics::enums::{AggregationType, Statistic}; | ||
| use tracing::debug; | ||
|
|
||
| #[derive(Debug, Clone, Copy, PartialEq, Eq)] | ||
| pub(crate) struct NodeId(usize); | ||
|
|
||
| #[derive(Debug, Clone, Copy, PartialEq, Eq)] | ||
| #[derive(Debug, Clone, PartialEq, Eq)] | ||
| pub(crate) enum StoreReadStrategy { | ||
| WindowGrid, | ||
| SlidingExactCover, | ||
| SlidingExactCover { | ||
| output_timestamps: Vec<u64>, | ||
| lookback_ms: u64, | ||
| window_size_ms: u64, | ||
| bucket_step_ms: u64, | ||
| }, | ||
| } | ||
|
|
||
| impl StoreReadStrategy { | ||
| fn for_window( | ||
| window_type: WindowType, | ||
| output_timestamps: &[u64], | ||
| lookback_ms: u64, | ||
| window_size_ms: u64, | ||
| bucket_step_ms: u64, | ||
| ) -> Self { | ||
| match window_type { | ||
| WindowType::Tumbling => Self::WindowGrid, | ||
| WindowType::Sliding => Self::SlidingExactCover { | ||
| output_timestamps: output_timestamps.to_vec(), | ||
| lookback_ms, | ||
| window_size_ms, | ||
| bucket_step_ms, | ||
| }, | ||
| } | ||
| } | ||
| } | ||
|
|
||
| #[derive(Debug, Clone, Copy, PartialEq, Eq)] | ||
| pub(crate) enum StoreReadRole { | ||
| Values, | ||
| Keys, | ||
| } | ||
|
|
||
| #[derive(Debug, Clone)] | ||
| pub(crate) struct RangeEstimateSpec { | ||
| pub output_timestamps: Vec<u64>, | ||
| pub query_range_ms: u64, | ||
| pub buckets_per_step: usize, | ||
| pub lookback_bucket_count: usize, | ||
| pub tumbling_window_ms: u64, | ||
| pub window_type: WindowType, | ||
| pub window_size_ms: u64, | ||
| pub keys_window_type: Option<WindowType>, | ||
|
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. The keys-window settings are stored twice. The read resolves these with |
||
| pub keys_window_size_ms: Option<u64>, | ||
| pub keys_lookback_ms: Option<u64>, | ||
| pub keys_tumbling_window_ms: Option<u64>, | ||
| pub value_aggregation_type: AggregationType, | ||
| pub key_aggregation_type: AggregationType, | ||
| pub grouping_labels: KeyByLabelNames, | ||
| pub aggregated_labels: KeyByLabelNames, | ||
| pub row_label_order: KeyByLabelNames, | ||
| } | ||
|
|
||
| #[derive(Debug, Clone)] | ||
| pub(crate) enum QueryPlanNode { | ||
| StoreRead { | ||
| query: StoreQueryParams, | ||
| role: StoreReadRole, | ||
| strategy: StoreReadStrategy, | ||
| }, | ||
| ComposeWindows { | ||
| PrepareBuckets { | ||
| input: NodeId, | ||
| output_timestamps: Vec<u64>, | ||
| lookback_ms: u64, | ||
| window_size_ms: u64, | ||
| bucket_step_ms: u64, | ||
| }, | ||
| ResolveKeys { | ||
| values: NodeId, | ||
|
|
@@ -39,6 +87,7 @@ pub(crate) enum QueryPlanNode { | |
| statistic: Statistic, | ||
| query_kwargs: std::collections::HashMap<String, String>, | ||
| output_labels: KeyByLabelNames, | ||
| spec: RangeEstimateSpec, | ||
| }, | ||
| AggregateVector { | ||
| input: NodeId, | ||
|
|
@@ -49,6 +98,7 @@ pub(crate) enum QueryPlanNode { | |
| input: NodeId, | ||
| k: String, | ||
| grouping_labels: KeyByLabelNames, | ||
| row_label_order: KeyByLabelNames, | ||
| }, | ||
| Format { | ||
| input: NodeId, | ||
|
|
@@ -113,46 +163,53 @@ impl QueryPlan { | |
| query_time_aggregations: &[QueryTimeAggregation], | ||
| ) -> Result<Self, String> { | ||
| let mut nodes = Vec::new(); | ||
| let values_lookback_ms = | ||
|
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
If one copy changes (e.g. the bucket-lookback rounding this PR just fixed), the reads fetch a different cover than the estimator expects, and steps get silently skipped as "incomplete Sliding value cover". This works against #754's "each node consumes the fields that define its semantics": let the estimator take the window spec from the plan instead of recomputing it. |
||
| (context.lookback_bucket_count as u64) * context.tumbling_window_ms; | ||
| let values_read = Self::push_read( | ||
| &mut nodes, | ||
| &context.base.store_plan.values_query, | ||
| context.window_type, | ||
| ); | ||
| let values = Self::push_compose( | ||
| &mut nodes, | ||
| values_read, | ||
| &context.output_timestamps, | ||
| context.query_range_ms, | ||
| context.window_size_ms, | ||
| context.tumbling_window_ms, | ||
| StoreReadRole::Values, | ||
| StoreReadStrategy::for_window( | ||
| context.window_type, | ||
| &context.output_timestamps, | ||
| values_lookback_ms, | ||
| context.window_size_ms, | ||
| context.tumbling_window_ms, | ||
| ), | ||
| ); | ||
| let values = Self::push_prepare_buckets(&mut nodes, values_read); | ||
| let keys = context.base.store_plan.keys_query.as_ref().map(|query| { | ||
| let keys_lookback_ms = context.keys_lookback_ms.unwrap_or(context.query_range_ms); | ||
|
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Silent defaults for required key-window config. Both production builders ( Suggested fix: see the comment on
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Missing keys settings now fall back silently. The keys Example: with |
||
| let keys_window_size_ms = context | ||
| .keys_window_size_ms | ||
| .unwrap_or(context.window_size_ms); | ||
| let keys_bucket_step_ms = context | ||
| .keys_tumbling_window_ms | ||
| .unwrap_or(context.tumbling_window_ms); | ||
| let read = Self::push_read( | ||
| &mut nodes, | ||
| query, | ||
| context.keys_window_type.unwrap_or(context.window_type), | ||
| StoreReadRole::Keys, | ||
| StoreReadStrategy::for_window( | ||
| context.keys_window_type.unwrap_or(context.window_type), | ||
|
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. The DAG and the legacy oracle choose different keys read strategies. Here the strategy comes from When values are Sliding, |
||
| &context.output_timestamps, | ||
| keys_lookback_ms, | ||
| keys_window_size_ms, | ||
| keys_bucket_step_ms, | ||
| ), | ||
| ); | ||
| Self::push_compose( | ||
| &mut nodes, | ||
| read, | ||
| &context.output_timestamps, | ||
| context.keys_lookback_ms.unwrap_or(context.query_range_ms), | ||
| context | ||
| .keys_window_size_ms | ||
| .unwrap_or(context.window_size_ms), | ||
| context | ||
| .keys_tumbling_window_ms | ||
| .unwrap_or(context.tumbling_window_ms), | ||
| ) | ||
| Self::push_prepare_buckets(&mut nodes, read) | ||
| }); | ||
| let resolved = Self::push(&mut nodes, QueryPlanNode::ResolveKeys { values, keys }); | ||
| let estimate_spec = RangeEstimateSpec::from(context); | ||
| let mut root = Self::push( | ||
| &mut nodes, | ||
| QueryPlanNode::Estimate { | ||
| input: resolved, | ||
| statistic: context.base.metadata.statistic_to_compute, | ||
| query_kwargs: context.base.metadata.query_kwargs.clone(), | ||
| output_labels: context.base.metadata.query_output_labels.clone(), | ||
| spec: estimate_spec.clone(), | ||
|
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Nit:
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. The timestamp and row-order lists are copied several times. A query with 10k steps ends up holding three copies of the timestamp vector per plan. Sharing them through an |
||
| }, | ||
| ); | ||
| if options.limit_topk && context.base.metadata.statistic_to_compute == Statistic::Topk { | ||
|
|
@@ -171,6 +228,7 @@ impl QueryPlan { | |
| input: root, | ||
| k, | ||
| grouping_labels: context.base.grouping_labels.clone(), | ||
| row_label_order: estimate_spec.row_label_order.clone(), | ||
| }, | ||
| ); | ||
| } | ||
|
|
@@ -285,62 +343,52 @@ impl QueryPlan { | |
| fn push_read( | ||
| nodes: &mut Vec<QueryPlanNode>, | ||
| query: &StoreQueryParams, | ||
| window_type: WindowType, | ||
| role: StoreReadRole, | ||
| strategy: StoreReadStrategy, | ||
| ) -> NodeId { | ||
| let strategy = match window_type { | ||
| WindowType::Tumbling => StoreReadStrategy::WindowGrid, | ||
| WindowType::Sliding => StoreReadStrategy::SlidingExactCover, | ||
| }; | ||
| Self::push( | ||
| nodes, | ||
| QueryPlanNode::StoreRead { | ||
| query: query.clone(), | ||
| role, | ||
| strategy, | ||
| }, | ||
| ) | ||
| } | ||
|
|
||
| fn push_compose( | ||
| nodes: &mut Vec<QueryPlanNode>, | ||
| input: NodeId, | ||
| output_timestamps: &[u64], | ||
| lookback_ms: u64, | ||
| window_size_ms: u64, | ||
| bucket_step_ms: u64, | ||
| ) -> NodeId { | ||
| Self::push( | ||
| nodes, | ||
| QueryPlanNode::ComposeWindows { | ||
| input, | ||
| output_timestamps: output_timestamps.to_vec(), | ||
| lookback_ms, | ||
| window_size_ms, | ||
| bucket_step_ms, | ||
| }, | ||
| ) | ||
| fn push_prepare_buckets(nodes: &mut Vec<QueryPlanNode>, input: NodeId) -> NodeId { | ||
| Self::push(nodes, QueryPlanNode::PrepareBuckets { input }) | ||
| } | ||
|
|
||
| pub(crate) fn explain(&self) -> String { | ||
| let mut lines = Vec::with_capacity(self.nodes.len() + 1); | ||
| for (index, node) in self.nodes.iter().enumerate() { | ||
| let line = match node { | ||
| QueryPlanNode::StoreRead { query, strategy } => format!( | ||
| "n{index} StoreRead({strategy:?}, {}#{}, [{}, {}])", | ||
| QueryPlanNode::StoreRead { query, role, strategy } => format!( | ||
| "n{index} StoreRead({strategy:?}, role={role:?}, {}#{}, [{}, {}])", | ||
| query.metric, query.aggregation_id, query.start_timestamp, query.end_timestamp | ||
| ), | ||
| QueryPlanNode::ComposeWindows { input, output_timestamps, lookback_ms, window_size_ms, bucket_step_ms } => format!( | ||
| "n{index} ComposeWindows(n{}, outputs={:?}, lookback={lookback_ms}ms, window={window_size_ms}ms, step={bucket_step_ms}ms)", | ||
| input.0, output_timestamps | ||
| ), | ||
| QueryPlanNode::PrepareBuckets { input } => { | ||
| format!("n{index} PrepareBuckets(n{})", input.0) | ||
| } | ||
| QueryPlanNode::ResolveKeys { values, keys } => format!( | ||
| "n{index} ResolveKeys(values=n{}, keys={})", | ||
| values.0, | ||
| keys.map(|id| format!("n{}", id.0)).unwrap_or_else(|| "self".to_string()) | ||
| ), | ||
| QueryPlanNode::Estimate { input, statistic, query_kwargs, .. } => { | ||
| QueryPlanNode::Estimate { | ||
| input, | ||
| statistic, | ||
| query_kwargs, | ||
| spec, | ||
| .. | ||
| } => { | ||
| let mut kwargs: Vec<_> = query_kwargs.iter().collect(); | ||
| kwargs.sort_unstable_by_key(|(key, _)| *key); | ||
| format!("n{index} Estimate(n{}, {statistic}, {kwargs:?})", input.0) | ||
| format!( | ||
| "n{index} Estimate(n{}, {statistic}, {kwargs:?}, outputs={:?})", | ||
|
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
For a query with thousands of steps, one debug log line holds the list up to three times. Suggest printing a summary instead: the count plus the first and last timestamp. |
||
| input.0, spec.output_timestamps | ||
| ) | ||
| }, | ||
| QueryPlanNode::LimitTopK { input, k, .. } => { | ||
| format!("n{index} LimitTopK(n{}, k={k})", input.0) | ||
|
|
@@ -359,11 +407,38 @@ impl QueryPlan { | |
| } | ||
| } | ||
|
|
||
| impl From<&RangeQueryExecutionContext> for RangeEstimateSpec { | ||
| fn from(context: &RangeQueryExecutionContext) -> Self { | ||
| Self { | ||
| output_timestamps: context.output_timestamps.clone(), | ||
| query_range_ms: context.query_range_ms, | ||
| buckets_per_step: context.buckets_per_step, | ||
| lookback_bucket_count: context.lookback_bucket_count, | ||
| tumbling_window_ms: context.tumbling_window_ms, | ||
| window_type: context.window_type, | ||
| window_size_ms: context.window_size_ms, | ||
| keys_window_type: context.keys_window_type, | ||
| keys_window_size_ms: context.keys_window_size_ms, | ||
| keys_lookback_ms: context.keys_lookback_ms, | ||
| keys_tumbling_window_ms: context.keys_tumbling_window_ms, | ||
| value_aggregation_type: context.base.agg_info.aggregation_type_for_value, | ||
| key_aggregation_type: context.base.agg_info.aggregation_type_for_key, | ||
| grouping_labels: context.base.grouping_labels.clone(), | ||
| aggregated_labels: context.base.aggregated_labels.clone(), | ||
| row_label_order: crate::engines::simple_engine::SimpleEngine::topk_row_label_order( | ||
| &context.base.metadata, | ||
| &context.base.grouping_labels, | ||
| &context.base.aggregated_labels, | ||
| ), | ||
| } | ||
| } | ||
| } | ||
|
|
||
| impl QueryPlanNode { | ||
| fn kind(&self) -> &'static str { | ||
| match self { | ||
| Self::StoreRead { .. } => "StoreRead", | ||
| Self::ComposeWindows { .. } => "ComposeWindows", | ||
| Self::PrepareBuckets { .. } => "PrepareBuckets", | ||
| Self::ResolveKeys { .. } => "ResolveKeys", | ||
| Self::Estimate { .. } => "Estimate", | ||
| Self::AggregateVector { .. } => "AggregateVector", | ||
|
|
@@ -375,7 +450,7 @@ impl QueryPlanNode { | |
| fn inputs(&self) -> Vec<NodeId> { | ||
| match self { | ||
| Self::StoreRead { .. } => Vec::new(), | ||
| Self::ComposeWindows { input, .. } | ||
| Self::PrepareBuckets { input, .. } | ||
| | Self::Estimate { input, .. } | ||
| | Self::AggregateVector { input, .. } | ||
| | Self::LimitTopK { input, .. } | ||
|
|
@@ -473,7 +548,9 @@ mod tests { | |
| .explain(); | ||
|
|
||
| assert!(explanation.contains("n4 ResolveKeys(values=n1, keys=n3)")); | ||
| assert!(explanation.contains("n2 StoreRead(SlidingExactCover, requests#8")); | ||
| assert!(explanation.contains("role=Values")); | ||
| assert!(explanation.contains("role=Keys")); | ||
| assert!(explanation.contains("n2 StoreRead(SlidingExactCover {")); | ||
|
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. This assertion doesn't check the key read's parameters (lookback, window, step). The fixture (L537) sets |
||
| assert!(explanation.ends_with("root: n5")); | ||
| } | ||
|
|
||
|
|
@@ -531,6 +608,29 @@ mod tests { | |
| assert!(explanation.contains("outputs=[1000, 2000, 3000]")); | ||
| } | ||
|
|
||
| #[test] | ||
|
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. The lookback test is brittle, and the empty-read behavior is untested. No native-path test checks that an empty Keys read succeeds while an empty Values read returns |
||
| fn sliding_value_read_uses_bucket_lookback() { | ||
| let mut context = context(); | ||
| context.window_type = WindowType::Sliding; | ||
| context.query_range_ms = 1_500; | ||
| context.lookback_bucket_count = 1; | ||
|
|
||
| let explanation = QueryPlan::compile_range( | ||
| &context, | ||
| PlanOptions { | ||
| limit_topk: false, | ||
| format_output: false, | ||
| }, | ||
| &[], | ||
| ) | ||
| .unwrap() | ||
| .explain(); | ||
|
|
||
| // Store reads must match the estimator's whole-bucket lookback. | ||
| assert!(explanation.contains("lookback_ms: 1000")); | ||
| assert!(!explanation.contains("lookback_ms: 1500")); | ||
| } | ||
|
|
||
| #[test] | ||
| fn topk_formatting_is_the_plan_root() { | ||
| let mut context = context(); | ||
|
|
@@ -584,6 +684,7 @@ mod tests { | |
| statistic: Statistic::Sum, | ||
| query_kwargs: HashMap::new(), | ||
| output_labels: KeyByLabelNames::empty(), | ||
| spec: RangeEstimateSpec::from(&context()), | ||
| }], | ||
| root: NodeId(0), | ||
| }; | ||
|
|
@@ -622,15 +723,10 @@ mod tests { | |
| start_timestamp: 0, | ||
| end_timestamp: 1, | ||
| }, | ||
| role: StoreReadRole::Values, | ||
| strategy: StoreReadStrategy::WindowGrid, | ||
| }, | ||
| QueryPlanNode::ComposeWindows { | ||
| input: NodeId(0), | ||
| output_timestamps: vec![1], | ||
| lookback_ms: 1, | ||
| window_size_ms: 1, | ||
| bucket_step_ms: 1, | ||
| }, | ||
| QueryPlanNode::PrepareBuckets { input: NodeId(0) }, | ||
| ], | ||
| root: NodeId(1), | ||
| }; | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
#754: "Do not copy the existing context wholesale into a node: extract cohesive plan-owned specs with explicit invariants."
This struct is close to a field-for-field copy of
RangeQueryExecutionContext, including four parallelkeys_*Options that must be all-Some or all-None. Nothing enforces that, so the estimator.expect()s each one separately and the compiler falls back withunwrap_or.Suggest a single
keys: Option<KeyWindowSpec { window_type, window_size_ms, lookback_ms, bucket_step_ms }>(and the same for the value window). Build it once incompile_range, failing if any field is missing, and use it for both the keyStoreReadstrategy and the estimator. That makes the half-set state unrepresentable, and the read and estimate can't disagree.