Skip to content

Commit 337ae1f

Browse files
refactor(query-engine): make native DAG execution self-contained
1 parent cbabfd6 commit 337ae1f

2 files changed

Lines changed: 270 additions & 159 deletions

File tree

‎asap-query-engine/src/engines/query_plan.rs‎

Lines changed: 167 additions & 71 deletions
Original file line numberDiff line numberDiff line change
@@ -5,30 +5,78 @@ use crate::engines::simple_engine::{RangeQueryExecutionContext, StoreQueryParams
55
use asap_types::enums::WindowType;
66
use asap_types::query_config::QueryTimeAggregation;
77
use promql_utilities::data_model::KeyByLabelNames;
8-
use promql_utilities::query_logics::enums::Statistic;
8+
use promql_utilities::query_logics::enums::{AggregationType, Statistic};
99
use tracing::debug;
1010

1111
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1212
pub(crate) struct NodeId(usize);
1313

14-
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
14+
#[derive(Debug, Clone, PartialEq, Eq)]
1515
pub(crate) enum StoreReadStrategy {
1616
WindowGrid,
17-
SlidingExactCover,
17+
SlidingExactCover {
18+
output_timestamps: Vec<u64>,
19+
lookback_ms: u64,
20+
window_size_ms: u64,
21+
bucket_step_ms: u64,
22+
},
23+
}
24+
25+
impl StoreReadStrategy {
26+
fn for_window(
27+
window_type: WindowType,
28+
output_timestamps: &[u64],
29+
lookback_ms: u64,
30+
window_size_ms: u64,
31+
bucket_step_ms: u64,
32+
) -> Self {
33+
match window_type {
34+
WindowType::Tumbling => Self::WindowGrid,
35+
WindowType::Sliding => Self::SlidingExactCover {
36+
output_timestamps: output_timestamps.to_vec(),
37+
lookback_ms,
38+
window_size_ms,
39+
bucket_step_ms,
40+
},
41+
}
42+
}
43+
}
44+
45+
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
46+
pub(crate) enum StoreReadRole {
47+
Values,
48+
Keys,
49+
}
50+
51+
#[derive(Debug, Clone)]
52+
pub(crate) struct RangeEstimateSpec {
53+
pub output_timestamps: Vec<u64>,
54+
pub query_range_ms: u64,
55+
pub buckets_per_step: usize,
56+
pub lookback_bucket_count: usize,
57+
pub tumbling_window_ms: u64,
58+
pub window_type: WindowType,
59+
pub window_size_ms: u64,
60+
pub keys_window_type: Option<WindowType>,
61+
pub keys_window_size_ms: Option<u64>,
62+
pub keys_lookback_ms: Option<u64>,
63+
pub keys_tumbling_window_ms: Option<u64>,
64+
pub value_aggregation_type: AggregationType,
65+
pub key_aggregation_type: AggregationType,
66+
pub grouping_labels: KeyByLabelNames,
67+
pub aggregated_labels: KeyByLabelNames,
68+
pub row_label_order: KeyByLabelNames,
1869
}
1970

2071
#[derive(Debug, Clone)]
2172
pub(crate) enum QueryPlanNode {
2273
StoreRead {
2374
query: StoreQueryParams,
75+
role: StoreReadRole,
2476
strategy: StoreReadStrategy,
2577
},
26-
ComposeWindows {
78+
PrepareBuckets {
2779
input: NodeId,
28-
output_timestamps: Vec<u64>,
29-
lookback_ms: u64,
30-
window_size_ms: u64,
31-
bucket_step_ms: u64,
3280
},
3381
ResolveKeys {
3482
values: NodeId,
@@ -39,6 +87,7 @@ pub(crate) enum QueryPlanNode {
3987
statistic: Statistic,
4088
query_kwargs: std::collections::HashMap<String, String>,
4189
output_labels: KeyByLabelNames,
90+
spec: RangeEstimateSpec,
4291
},
4392
AggregateVector {
4493
input: NodeId,
@@ -49,6 +98,7 @@ pub(crate) enum QueryPlanNode {
4998
input: NodeId,
5099
k: String,
51100
grouping_labels: KeyByLabelNames,
101+
row_label_order: KeyByLabelNames,
52102
},
53103
Format {
54104
input: NodeId,
@@ -113,46 +163,53 @@ impl QueryPlan {
113163
query_time_aggregations: &[QueryTimeAggregation],
114164
) -> Result<Self, String> {
115165
let mut nodes = Vec::new();
166+
let values_lookback_ms =
167+
(context.lookback_bucket_count as u64) * context.tumbling_window_ms;
116168
let values_read = Self::push_read(
117169
&mut nodes,
118170
&context.base.store_plan.values_query,
119-
context.window_type,
120-
);
121-
let values = Self::push_compose(
122-
&mut nodes,
123-
values_read,
124-
&context.output_timestamps,
125-
context.query_range_ms,
126-
context.window_size_ms,
127-
context.tumbling_window_ms,
171+
StoreReadRole::Values,
172+
StoreReadStrategy::for_window(
173+
context.window_type,
174+
&context.output_timestamps,
175+
values_lookback_ms,
176+
context.window_size_ms,
177+
context.tumbling_window_ms,
178+
),
128179
);
180+
let values = Self::push_prepare_buckets(&mut nodes, values_read);
129181
let keys = context.base.store_plan.keys_query.as_ref().map(|query| {
182+
let keys_lookback_ms = context.keys_lookback_ms.unwrap_or(context.query_range_ms);
183+
let keys_window_size_ms = context
184+
.keys_window_size_ms
185+
.unwrap_or(context.window_size_ms);
186+
let keys_bucket_step_ms = context
187+
.keys_tumbling_window_ms
188+
.unwrap_or(context.tumbling_window_ms);
130189
let read = Self::push_read(
131190
&mut nodes,
132191
query,
133-
context.keys_window_type.unwrap_or(context.window_type),
192+
StoreReadRole::Keys,
193+
StoreReadStrategy::for_window(
194+
context.keys_window_type.unwrap_or(context.window_type),
195+
&context.output_timestamps,
196+
keys_lookback_ms,
197+
keys_window_size_ms,
198+
keys_bucket_step_ms,
199+
),
134200
);
135-
Self::push_compose(
136-
&mut nodes,
137-
read,
138-
&context.output_timestamps,
139-
context.keys_lookback_ms.unwrap_or(context.query_range_ms),
140-
context
141-
.keys_window_size_ms
142-
.unwrap_or(context.window_size_ms),
143-
context
144-
.keys_tumbling_window_ms
145-
.unwrap_or(context.tumbling_window_ms),
146-
)
201+
Self::push_prepare_buckets(&mut nodes, read)
147202
});
148203
let resolved = Self::push(&mut nodes, QueryPlanNode::ResolveKeys { values, keys });
204+
let estimate_spec = RangeEstimateSpec::from(context);
149205
let mut root = Self::push(
150206
&mut nodes,
151207
QueryPlanNode::Estimate {
152208
input: resolved,
153209
statistic: context.base.metadata.statistic_to_compute,
154210
query_kwargs: context.base.metadata.query_kwargs.clone(),
155211
output_labels: context.base.metadata.query_output_labels.clone(),
212+
spec: estimate_spec.clone(),
156213
},
157214
);
158215
if options.limit_topk && context.base.metadata.statistic_to_compute == Statistic::Topk {
@@ -171,6 +228,7 @@ impl QueryPlan {
171228
input: root,
172229
k,
173230
grouping_labels: context.base.grouping_labels.clone(),
231+
row_label_order: estimate_spec.row_label_order.clone(),
174232
},
175233
);
176234
}
@@ -285,62 +343,52 @@ impl QueryPlan {
285343
fn push_read(
286344
nodes: &mut Vec<QueryPlanNode>,
287345
query: &StoreQueryParams,
288-
window_type: WindowType,
346+
role: StoreReadRole,
347+
strategy: StoreReadStrategy,
289348
) -> NodeId {
290-
let strategy = match window_type {
291-
WindowType::Tumbling => StoreReadStrategy::WindowGrid,
292-
WindowType::Sliding => StoreReadStrategy::SlidingExactCover,
293-
};
294349
Self::push(
295350
nodes,
296351
QueryPlanNode::StoreRead {
297352
query: query.clone(),
353+
role,
298354
strategy,
299355
},
300356
)
301357
}
302358

303-
fn push_compose(
304-
nodes: &mut Vec<QueryPlanNode>,
305-
input: NodeId,
306-
output_timestamps: &[u64],
307-
lookback_ms: u64,
308-
window_size_ms: u64,
309-
bucket_step_ms: u64,
310-
) -> NodeId {
311-
Self::push(
312-
nodes,
313-
QueryPlanNode::ComposeWindows {
314-
input,
315-
output_timestamps: output_timestamps.to_vec(),
316-
lookback_ms,
317-
window_size_ms,
318-
bucket_step_ms,
319-
},
320-
)
359+
fn push_prepare_buckets(nodes: &mut Vec<QueryPlanNode>, input: NodeId) -> NodeId {
360+
Self::push(nodes, QueryPlanNode::PrepareBuckets { input })
321361
}
322362

323363
pub(crate) fn explain(&self) -> String {
324364
let mut lines = Vec::with_capacity(self.nodes.len() + 1);
325365
for (index, node) in self.nodes.iter().enumerate() {
326366
let line = match node {
327-
QueryPlanNode::StoreRead { query, strategy } => format!(
328-
"n{index} StoreRead({strategy:?}, {}#{}, [{}, {}])",
367+
QueryPlanNode::StoreRead { query, role, strategy } => format!(
368+
"n{index} StoreRead({strategy:?}, role={role:?}, {}#{}, [{}, {}])",
329369
query.metric, query.aggregation_id, query.start_timestamp, query.end_timestamp
330370
),
331-
QueryPlanNode::ComposeWindows { input, output_timestamps, lookback_ms, window_size_ms, bucket_step_ms } => format!(
332-
"n{index} ComposeWindows(n{}, outputs={:?}, lookback={lookback_ms}ms, window={window_size_ms}ms, step={bucket_step_ms}ms)",
333-
input.0, output_timestamps
334-
),
371+
QueryPlanNode::PrepareBuckets { input } => {
372+
format!("n{index} PrepareBuckets(n{})", input.0)
373+
}
335374
QueryPlanNode::ResolveKeys { values, keys } => format!(
336375
"n{index} ResolveKeys(values=n{}, keys={})",
337376
values.0,
338377
keys.map(|id| format!("n{}", id.0)).unwrap_or_else(|| "self".to_string())
339378
),
340-
QueryPlanNode::Estimate { input, statistic, query_kwargs, .. } => {
379+
QueryPlanNode::Estimate {
380+
input,
381+
statistic,
382+
query_kwargs,
383+
spec,
384+
..
385+
} => {
341386
let mut kwargs: Vec<_> = query_kwargs.iter().collect();
342387
kwargs.sort_unstable_by_key(|(key, _)| *key);
343-
format!("n{index} Estimate(n{}, {statistic}, {kwargs:?})", input.0)
388+
format!(
389+
"n{index} Estimate(n{}, {statistic}, {kwargs:?}, outputs={:?})",
390+
input.0, spec.output_timestamps
391+
)
344392
},
345393
QueryPlanNode::LimitTopK { input, k, .. } => {
346394
format!("n{index} LimitTopK(n{}, k={k})", input.0)
@@ -359,11 +407,38 @@ impl QueryPlan {
359407
}
360408
}
361409

410+
impl From<&RangeQueryExecutionContext> for RangeEstimateSpec {
411+
fn from(context: &RangeQueryExecutionContext) -> Self {
412+
Self {
413+
output_timestamps: context.output_timestamps.clone(),
414+
query_range_ms: context.query_range_ms,
415+
buckets_per_step: context.buckets_per_step,
416+
lookback_bucket_count: context.lookback_bucket_count,
417+
tumbling_window_ms: context.tumbling_window_ms,
418+
window_type: context.window_type,
419+
window_size_ms: context.window_size_ms,
420+
keys_window_type: context.keys_window_type,
421+
keys_window_size_ms: context.keys_window_size_ms,
422+
keys_lookback_ms: context.keys_lookback_ms,
423+
keys_tumbling_window_ms: context.keys_tumbling_window_ms,
424+
value_aggregation_type: context.base.agg_info.aggregation_type_for_value,
425+
key_aggregation_type: context.base.agg_info.aggregation_type_for_key,
426+
grouping_labels: context.base.grouping_labels.clone(),
427+
aggregated_labels: context.base.aggregated_labels.clone(),
428+
row_label_order: crate::engines::simple_engine::SimpleEngine::topk_row_label_order(
429+
&context.base.metadata,
430+
&context.base.grouping_labels,
431+
&context.base.aggregated_labels,
432+
),
433+
}
434+
}
435+
}
436+
362437
impl QueryPlanNode {
363438
fn kind(&self) -> &'static str {
364439
match self {
365440
Self::StoreRead { .. } => "StoreRead",
366-
Self::ComposeWindows { .. } => "ComposeWindows",
441+
Self::PrepareBuckets { .. } => "PrepareBuckets",
367442
Self::ResolveKeys { .. } => "ResolveKeys",
368443
Self::Estimate { .. } => "Estimate",
369444
Self::AggregateVector { .. } => "AggregateVector",
@@ -375,7 +450,7 @@ impl QueryPlanNode {
375450
fn inputs(&self) -> Vec<NodeId> {
376451
match self {
377452
Self::StoreRead { .. } => Vec::new(),
378-
Self::ComposeWindows { input, .. }
453+
Self::PrepareBuckets { input, .. }
379454
| Self::Estimate { input, .. }
380455
| Self::AggregateVector { input, .. }
381456
| Self::LimitTopK { input, .. }
@@ -473,7 +548,9 @@ mod tests {
473548
.explain();
474549

475550
assert!(explanation.contains("n4 ResolveKeys(values=n1, keys=n3)"));
476-
assert!(explanation.contains("n2 StoreRead(SlidingExactCover, requests#8"));
551+
assert!(explanation.contains("role=Values"));
552+
assert!(explanation.contains("role=Keys"));
553+
assert!(explanation.contains("n2 StoreRead(SlidingExactCover {"));
477554
assert!(explanation.ends_with("root: n5"));
478555
}
479556

@@ -531,6 +608,29 @@ mod tests {
531608
assert!(explanation.contains("outputs=[1000, 2000, 3000]"));
532609
}
533610

611+
#[test]
612+
fn sliding_value_read_uses_bucket_lookback() {
613+
let mut context = context();
614+
context.window_type = WindowType::Sliding;
615+
context.query_range_ms = 1_500;
616+
context.lookback_bucket_count = 1;
617+
618+
let explanation = QueryPlan::compile_range(
619+
&context,
620+
PlanOptions {
621+
limit_topk: false,
622+
format_output: false,
623+
},
624+
&[],
625+
)
626+
.unwrap()
627+
.explain();
628+
629+
// Store reads must match the estimator's whole-bucket lookback.
630+
assert!(explanation.contains("lookback_ms: 1000"));
631+
assert!(!explanation.contains("lookback_ms: 1500"));
632+
}
633+
534634
#[test]
535635
fn topk_formatting_is_the_plan_root() {
536636
let mut context = context();
@@ -584,6 +684,7 @@ mod tests {
584684
statistic: Statistic::Sum,
585685
query_kwargs: HashMap::new(),
586686
output_labels: KeyByLabelNames::empty(),
687+
spec: RangeEstimateSpec::from(&context()),
587688
}],
588689
root: NodeId(0),
589690
};
@@ -622,15 +723,10 @@ mod tests {
622723
start_timestamp: 0,
623724
end_timestamp: 1,
624725
},
726+
role: StoreReadRole::Values,
625727
strategy: StoreReadStrategy::WindowGrid,
626728
},
627-
QueryPlanNode::ComposeWindows {
628-
input: NodeId(0),
629-
output_timestamps: vec![1],
630-
lookback_ms: 1,
631-
window_size_ms: 1,
632-
bucket_step_ms: 1,
633-
},
729+
QueryPlanNode::PrepareBuckets { input: NodeId(0) },
634730
],
635731
root: NodeId(1),
636732
};

0 commit comments

Comments
 (0)