Skip to content
Open
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
61 changes: 56 additions & 5 deletions asap-common/dependencies/rs/asap_types/src/streaming_config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -58,12 +58,17 @@ impl StreamingConfig {
for aggregation in &query_config.aggregations {
let aggregation_id = aggregation.aggregation_id;
if let Some(num_aggregates) = aggregation.num_aggregates_to_retain {
// OLD: Keep last value only (for backwards compatibility)
retention_map.insert(aggregation_id, num_aggregates);
retention_map
.entry(aggregation_id)
.and_modify(|retain| *retain = (*retain).max(num_aggregates))
.or_insert(num_aggregates);
}

// NEW: Sum up num_aggregates_to_retain across all queries
*read_count_threshold_map.entry(aggregation_id).or_insert(0) +=
num_aggregates;
if let Some(read_count_threshold) = aggregation.read_count_threshold {
let threshold = read_count_threshold_map.entry(aggregation_id).or_insert(0);
*threshold = threshold
.checked_add(read_count_threshold)
.expect("read-count threshold overflowed");
}
}
}
Expand Down Expand Up @@ -136,6 +141,10 @@ impl Default for StreamingConfig {
mod tests {
use super::*;
use crate::aggregation_config::AggregationConfigError;
use crate::aggregation_reference::AggregationReference;
use crate::enums::{CleanupPolicy, QueryLanguage};
use crate::inference_config::InferenceConfig;
use crate::query_config::QueryConfig;

#[test]
fn rejects_heap_config_with_invalid_sub_type() {
Expand Down Expand Up @@ -365,4 +374,46 @@ aggregations:
StreamingConfig::from_yaml_data(&yaml, None)
.expect("MinMax config with 'MAX' subtype must be accepted");
}

#[test]
fn shared_references_keep_the_largest_circular_buffer_retention() {
// A later short-range query must not reduce a shared aggregation's retention.
let yaml = minmax_yaml("MinMax", "max");
for references in [[7, 2], [2, 7]] {
let mut inference =
InferenceConfig::new(QueryLanguage::promql, CleanupPolicy::CircularBuffer);
inference.query_configs = references
.into_iter()
.map(|retain| {
QueryConfig::new(format!("max(metric[{retain}m])"))
.add_aggregation(AggregationReference::new(1, Some(retain)))
})
.collect();

let config = StreamingConfig::from_yaml_data(&yaml, Some(&inference))
.expect("shared aggregation references should load");

assert_eq!(config[1].num_aggregates_to_retain, Some(7));
assert_eq!(config[1].read_count_threshold, None);
}
}

#[test]
fn shared_references_sum_read_based_thresholds() {
// Each shared query consumes a read before the aggregate can be cleaned up.
let yaml = minmax_yaml("MinMax", "max");
let mut inference = InferenceConfig::new(QueryLanguage::promql, CleanupPolicy::ReadBased);
inference.query_configs = vec![
QueryConfig::new("max(metric[10m])".into())
.add_aggregation(AggregationReference::with_read_count_threshold(1, Some(7))),
QueryConfig::new("max(metric[1m])".into())
.add_aggregation(AggregationReference::with_read_count_threshold(1, Some(2))),
];

let config = StreamingConfig::from_yaml_data(&yaml, Some(&inference))
.expect("shared aggregation references should load");

assert_eq!(config[1].num_aggregates_to_retain, None);
assert_eq!(config[1].read_count_threshold, Some(9));
}
}
165 changes: 153 additions & 12 deletions asap-planner-rs/src/optimizer/translator.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,35 +2,46 @@ use asap_types::aggregation_reference::AggregationReference;
use asap_types::inference_config::InferenceConfig;
use asap_types::query_config::QueryConfig;
use asap_types::streaming_config::StreamingConfig;
use std::collections::HashMap;

use super::solution::{OptimizerSolution, QueryMethod};

/// Translate an `OptimizerSolution` into the deployment artifacts consumed by
/// Arroyo and the query engine.
///
pub fn translate(solution: &OptimizerSolution) -> (StreamingConfig, InferenceConfig) {
let streaming_config = build_streaming_config(solution);
let inference_config = build_inference_config(solution);
let streaming_config = build_streaming_config(solution, &inference_config);
(streaming_config, inference_config)
}

fn build_streaming_config(solution: &OptimizerSolution) -> StreamingConfig {
// Deployed configs map directly to AggregationConfigs — the types are the same.
StreamingConfig::new(solution.deployed_configs().clone())
fn build_streaming_config(
solution: &OptimizerSolution,
inference_config: &InferenceConfig,
) -> StreamingConfig {
let mut configs = solution.deployed_configs().clone();
for (aggregation_id, threshold) in read_count_thresholds(inference_config) {
let config = configs
.get_mut(&aggregation_id)
.expect("every assigned aggregation must be deployed");
config.num_aggregates_to_retain = None;
config.read_count_threshold = Some(threshold);
}
StreamingConfig::new(configs)
}

fn build_inference_config(solution: &OptimizerSolution) -> InferenceConfig {
use asap_types::enums::{CleanupPolicy, QueryLanguage};

let mut inference = InferenceConfig::new(QueryLanguage::promql, CleanupPolicy::NoCleanup);
let mut inference = InferenceConfig::new(QueryLanguage::promql, CleanupPolicy::ReadBased);

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Subtract windows are deleted while queries still need them (this comment is about Subtract => 2 at line 89, which takes effect because of the switch to ReadBased here). Subtract => 2 assumes the engine reads only two prefix-sum checkpoints. The engine has no subtract path (IncrementalMerger is listed as future work in window_merger), so Subtract runs as a normal range query that reads all n windows on every run.

Example: sum_over_time(m[5m]) with t_repeat = 60s and scrape = 60s. This picks tumbling W = 60s with n = 5, and Sum is subtractable, so the method is Subtract with threshold 2. Each window is needed by 5 consecutive runs but is deleted after the 2nd read. Runs 3–5 merge only about 2 of the 5 windows, so the sum comes back too low and no error is raised.


for assignment in &solution.assignments {
let aggregation_id = assignment.aggregation_id;
let retain = retention_count_for_assignment(&assignment.query_method);
let agg_ref = AggregationReference::new(aggregation_id, Some(retain));
let cleanup_count = cleanup_count_for_assignment(&assignment.query_method);
let agg_ref =
AggregationReference::with_read_count_threshold(aggregation_id, Some(cleanup_count));
let key_ref = assignment.key_aggregation_id.map(|key_id| {
let key_retain = solution.deployed_configs()[&key_id].num_aggregates_to_retain;
AggregationReference::new(key_id, key_retain)
AggregationReference::with_read_count_threshold(key_id, Some(cleanup_count))

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The key tracker gets the value aggregation's threshold, which doesn't match how often its panes are read. build_key_config makes the DeltaSet tracker tumble at the value's slide S and keep ceil(range / S) panes. Each pane is read about range/S times when t_repeat = S. The threshold passed here is the value's count instead: range/W for sliding Merge, or 2 for Subtract.

Example: W = 120s, S = 60s, range = 240s gives a key threshold of 2. Each key pane is deleted after 2 of the 4 runs that need it. Later runs then build the key set from missing deltas, and keys drop out of the query output.

});

for query_string in &assignment.item.query_strings {
Expand All @@ -47,9 +58,26 @@ fn build_inference_config(solution: &OptimizerSolution) -> InferenceConfig {
inference
}

/// For a Merge assignment, the number of retained windows to configure in the
/// inference config (num_aggregates_to_retain on the AggregationReference).
pub fn retention_count_for_assignment(query_method: &QueryMethod) -> u64 {
fn read_count_thresholds(inference_config: &InferenceConfig) -> HashMap<u64, u64> {
let mut thresholds = HashMap::new();
for query_config in &inference_config.query_configs {
for aggregation in &query_config.aggregations {
let Some(cleanup_count) = aggregation.read_count_threshold else {
continue;
};
let threshold = thresholds
.entry(aggregation.aggregation_id)
.or_insert(0_u64);
*threshold = threshold
.checked_add(cleanup_count)
.expect("read-count threshold overflowed");
}
}
thresholds
}

/// Number of aggregate reads an assignment needs before its windows may be cleaned up.
pub fn cleanup_count_for_assignment(query_method: &QueryMethod) -> u64 {
match query_method {
QueryMethod::Direct => 1,
QueryMethod::Merge { num_windows } => *num_windows,

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Memory grows without limit when queries run less often than the window slides. Merge { num_windows } => n assumes each window is read n times. That only holds when t_repeat equals the slide. window_candidates allows W and S smaller than t_repeat.

  • range = 300s, t_repeat = 300s, W = 60s gives Merge{5}. Each window is read once per run, so its count stays at 1 and never reaches 5.
  • sliding W = 120s, S = 60s, t_repeat = 120s: every other window is never read.

In both cases those windows are never deleted. num_aggregates_to_retain used to cap retention, and it is now None, so store memory grows forever. The rule-based planner avoids the first case by using ceil(lookback / effective_repeat).

Expand Down Expand Up @@ -77,3 +105,116 @@ impl TranslationSummary {
}
}
}

#[cfg(test)]
mod tests {
use std::collections::HashMap;

use asap_types::aggregation_config::AggregationConfig;
use asap_types::enums::{CleanupPolicy, WindowType};
use asap_types::query_requirements::QueryRequirements;
use promql_utilities::data_model::KeyByLabelNames;
use promql_utilities::query_logics::enums::{AggregationType, Statistic};

use super::super::solution::{AQEAssignment, OptimizerItem};
use super::*;

fn assignment(aggregation_id: u64, query_method: QueryMethod) -> AQEAssignment {
AQEAssignment {
item: OptimizerItem {
requirements: QueryRequirements {
metric: "metric".into(),
statistics: vec![Statistic::Sum],
data_range_ms: 60_000,
grouping_labels: KeyByLabelNames::empty(),
spatial_filter_normalized: String::new(),
topk_count_events: None,
topk_by_labels: None,
},
query_strings: vec!["sum(metric)".into()],
query_frequency_hz: 1.0,
t_repeat_ms: 60_000,
accuracy_sla: 0.01,
latency_sla: 1.0,
},
aggregation_id,
key_aggregation_id: None,
query_method,
estimated_query_cost_per_sec: 0.0,
}
}

#[test]
fn translation_uses_read_based_cleanup_and_sums_shared_thresholds() {
let mut solution = OptimizerSolution::empty();
let aggregation_id = solution.register_config(AggregationConfig::new(
0,
AggregationType::Sum,
"sum".into(),
HashMap::new(),
KeyByLabelNames::empty(),
KeyByLabelNames::empty(),
KeyByLabelNames::empty(),
String::new(),
60_000,
60_000,
WindowType::Tumbling,
String::new(),
"metric".into(),
None,
None,
None,
None,
));
solution.assignments = vec![
assignment(aggregation_id, QueryMethod::Merge { num_windows: 3 }),
assignment(aggregation_id, QueryMethod::Direct),
];

let (streaming, inference) = translate(&solution);

assert_eq!(inference.cleanup_policy, CleanupPolicy::ReadBased);
assert_eq!(streaming[aggregation_id].read_count_threshold, Some(4));
assert!(inference.query_configs.iter().all(|query_config| {
query_config.aggregations.iter().all(|aggregation| {
aggregation.num_aggregates_to_retain.is_none()
&& aggregation.read_count_threshold.is_some()
})
}));
}

#[test]
fn translation_counts_each_query_string_in_a_shared_threshold() {
let mut solution = OptimizerSolution::empty();
let aggregation_id = solution.register_config(AggregationConfig::new(
0,
AggregationType::Sum,
"sum".into(),
HashMap::new(),
KeyByLabelNames::empty(),
KeyByLabelNames::empty(),
KeyByLabelNames::empty(),
String::new(),
60_000,
60_000,
WindowType::Tumbling,
String::new(),
"metric".into(),
None,
None,
None,
None,
));
let mut shared_assignment =
assignment(aggregation_id, QueryMethod::Merge { num_windows: 3 });
shared_assignment
.item
.query_strings
.push("sum by (job) (metric)".into());
solution.assignments = vec![shared_assignment];

let (streaming, _) = translate(&solution);

assert_eq!(streaming[aggregation_id].read_count_threshold, Some(6));
}
}
Loading