Skip to content

feat(query-engine): execute nested PromQL aggregation DAGs - #770

Merged
milindsrivastava1997 merged 24 commits into
mainfrom
738-nested-promql-aggregations
Oct 5, 2026
Merged

milindsrivastava1997 merged 24 commits into
mainfrom
738-nested-promql-aggregations

Conversation

@milindsrivastava1997

@milindsrivastava1997 milindsrivastava1997 commented Oct 4, 2026 •

Copy link
Copy Markdown
Contributor

Closes #738

Summary

  • plan PromQL aggregation chains as explicit query-time stages
  • execute each range stage as a visible AggregateVector DAG node
  • carry the resulting label schema through the native DAG to response formatting
  • preserve plain-selector TopK metric-name labels through nested stages
  • keep grouped instant TopK buckets contiguous and ranked within each bucket
  • fall back to Prometheus when a later by-stage requires a label removed by an earlier stage

Verification

  • cargo test -p query_engine_rust
  • full pre-commit suite
  • make run DATASET=../datasets/nested-aggregations.yaml SUITE=../suites/nested-aggregations.yaml (90/90 passed)

Docker compliance matrix

Current make run-all result (2026-10-04):

Suite Result Notes
single-rate-temporal Fails Pre-existing window mismatch
single-rate-off-grid-rate Fails Pre-existing off-grid native-rate behavior
aggregations-native-dag Passes
aggregations Fails Pre-existing unsupported rate/nested temporal shapes and avg_over_time
topk-ordering Passes Confirms grouped instant TopK bucket contiguity and within-bucket ordering against Prometheus
quantiles Passes
olly-bench Fails as expected Non-CI suite

The command exits non-zero because it includes the listed known failures; the new TopK ordering suite passes.

return Ok(None);
};
let Some((anchor_labels, anchor_result)) =
self.execute_context_result(context, false, false)?

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 anchor runs with topk limiting and formatting both off (false, false). If the planned subquery is itself a topk (e.g. sum(topk(3, x)) → anchor topk(3, x) + sum), it's never truncated to k, so the outer sum adds every candidate series instead of the top 3, and it groups over unformatted labels. The same applies to the range path at L1454.

else {
return Ok(None);
};
let output = self.execute_observed_range_query_pipeline(

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.

This calls execute_observed_range_query_pipeline(...)? directly and skips map_local_execution_outcome, which the instant path and the non-nested range path both go through. When the local store has no data for the window, NoLocalData becomes an HTTP 500 instead of Ok(None), so the request is never forwarded to the fallback Prometheus.

labels: labels.labels.clone(),
}
}
_ => QueryTimeGrouping {

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.

without () with no labels falls through to QueryTimeGroupingMode::All, but in PromQL it keeps every label. For example, sum without () (sum by (job, instance) (x)) should return one series per (job, instance), but here it returns a single summed series. Map it to Without with an empty list, or mark it unsupported.


let labels = match aggregation.grouping.mode {
QueryTimeGroupingMode::All => Vec::new(),
QueryTimeGroupingMode::By => aggregation.grouping.labels.clone(),

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.

A by label that isn't in the input label set is treated as an error here, but in Prometheus it's valid and produces an empty label value. As a result, sum by (job) (sum by (instance) (x)) passes planning but fails when the plan is compiled. The client gets a 500 instead of one series with job="", and the request isn't forwarded to Prometheus.

planned_subquery = %config.planned_subquery,
"configured query-time aggregation anchor does not parse: {error}"
);
return Ok(None);

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.

Triaged: agreed. A matching configured query with an unparsable planned_subquery is a corrupt static execution plan, not a capability miss. The instant and range paths should return QueryExecutionError::Native, with regressions for both.

group.truncate(k);
results.extend(group);
}
results.sort_by(|left, right| left.labels.labels.cmp(&right.labels.labels));

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.

Triaged: agreed. The final label sort discards the descending-value order established above. Query-time TopK must return selected instant-vector elements in descending value order, using the full label set only as the tie-breaker. Range correctness remains per timestamp; a matrix cannot encode a different global rank order at each timestamp.

Comment thread asap-planner-rs/src/planner/promql.rs Outdated
else {
return None;
};
if number.val < 0.0 || number.val.fract() != 0.0 {

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.

Triaged: the loadability issue is real. The planner must omit nested plans for topk(0) and quantile values outside [0, 1], so those queries fall back to Prometheus. However, without () is valid PromQL and explicitly supported by #738: it must remain Without([]). The correct companion fix is to allow Without([]) in QueryTimeAggregation validation, not map it to All.

"max" => QueryTimeAggregationOperator::Max,
"quantile" => QueryTimeAggregationOperator::Quantile,
"topk" => QueryTimeAggregationOperator::Topk,
_ => return None,

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.

Triaged: no change for sort_desc. The #738 contract defines the supported query-time operators as sum, count, avg, min, max, quantile, and topk; sort_desc is not included and is a vector-sorting function rather than an aggregation. The separate TopK ordering defect is being fixed.

return Ok(None);
};
let Some((anchor_labels, anchor_result)) =
self.execute_context_result(context, true, false)?

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.

When the anchor is a topk over a bare metric (e.g. topk(3, transfer_events)), anchor_labels starts with __name__, but running it with formatting off yields rows without a __name__ value, so labels and row values are misaligned by one.

  • sum by (srcip) (topk(3, transfer_events)) / sum without (...): index lookup for srcip hits position 1 on a 1-value row → "result labels do not match the configured label schema" error instead of an answer or Prometheus fallback.
  • topk(1, topk(3, transfer_events)): convert_query_result_to_prometheus zips them, giving __name__="10.0.0.x" and no srcip.

Same issue on the range path (query_plan.rs:155). The new test only covers plain sum(...), which drops all labels and so masks this — worth adding a by (...) case.

.iter()
.map(|index| {
index.map_or_else(
|| Ok(String::new()),

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.

A grouping label missing from the inner result becomes "" and is emitted as a real label. E.g. max by (instance) (sum by (job) (x)) returns {instance=""}, whereas Prometheus returns {}. This will diverge in Prometheus comparisons (e.g. the compliance suite). Missing labels should be dropped from the output series rather than set to empty.

@milindsrivastava1997 milindsrivastava1997 left a comment

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.

Code review: 3 findings inline (1 high, 1 medium, 1 low).

results.sort_by(|a, b| {
Self::cmp_topk_value_desc(a.value, &a.labels.labels, b.value, &b.labels.labels)
});
Self::sort_instant_topk_results(

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.

High: grouped bare-selector topk fails (or sorts the wrong column) when enable_topk_formatting=false.

For a bare-selector topk (keep_metric_name), query_output_labels starts with __name__, but with enable_topk_formatting=false the metric name is never added to the result values, so every grouping-label position is off by one.

E.g. topk by (job) (3, x): labels [__name__, instance, job] put job at index 2, but the values are [instance, job], so get(2) is None → "Topk result labels do not match the configured output schema". If the grouping label isn't last, there's no error, but the sort uses the wrong column.

Callers that pass false:

  • the binary-arm instant path (promql.rs:510, e.g. topk by (job) (3, x) + 0), which worked before this PR
  • the new nested instant path (execute_context_result(context, true, false), e.g. sum(topk by (job) (3, x)))

Fix: strip the leading __name__ before computing positions, as topk_row_label_order already does.

let planned_subquery = query_data
.get("planned_subquery")
.and_then(|v| v.as_str())
.ok_or_else(|| anyhow::anyhow!("Missing planned_subquery field"))?

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.

Medium: making planned_subquery / query_time_aggregations required breaks the checked-in inference configs.

These configs have neither field and now fail with "Missing planned_subquery field":

  • asap-query-engine/examples/promql/inference_config.yaml
  • asap-query-engine/examples/sql/inference_config.yaml
  • asap-tools/execution-utilities/elastic-asap-benchmarking/inference_config.yaml

src/bin/test_e2e_precompute.rs loads the PromQL example directly, so that binary breaks right away. Older planner outputs fail the same way.

Fix: add the fields to the configs, or default them (planned_subquery = query, empty pipeline).

let indices = label_indices(input_labels, &grouping_labels);
let mut groups: BTreeMap<Vec<String>, Vec<InstantVectorElement>> = BTreeMap::new();
for element in input {
if !element.value.is_finite() {

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.

Low: a NaN/Inf from the inner query fails the whole request.

Here and in aggregate_values (line 108), any non-finite input returns an error. That reaches the caller as QueryExecutionError::Native, not as a "no local data" result that falls back to Prometheus.

E.g. if one series of sum by (job) (...) or a topk sketch yields NaN or ±Inf, max(...) or topk(3, ...) over it fails. Prometheus would return a result (NaN goes through sum/max, and topk ranks it).

Fix: pass NaN/Inf through the way PromQL does, or return Ok(None) so the query falls back.

@milindsrivastava1997
milindsrivastava1997 merged commit cbabfd6 into main Oct 5, 2026
10 checks passed
@milindsrivastava1997
milindsrivastava1997 deleted the 738-nested-promql-aggregations branch October 5, 2026 14:12
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Tracking: asap-planner query coverage for olly-bench

1 participant