query/promql/planner/at_modifier.rs
1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7// http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Anchored selector planning and replay for the PromQL `@` modifier.
16
17use std::sync::Arc;
18use std::time::{SystemTime, UNIX_EPOCH};
19
20use datafusion::logical_expr::{Extension, LogicalPlan, LogicalPlanBuilder};
21use datafusion::prelude::{Column, Expr as DfExpr};
22use promql::extension_plan::{InstantManipulate, Millisecond, SeriesDivide};
23use promql_parser::parser::{
24 AtModifier, Call, Expr as PromExpr, MatrixSelector, Offset, ParenExpr,
25};
26use snafu::{OptionExt, ResultExt, ensure};
27
28use crate::promql::error::{
29 AtModifierTimestampOutOfRangeSnafu, DataFusionPlanningSnafu, Result, TimeIndexNotFoundSnafu,
30};
31use crate::promql::planner::PromPlanner;
32use crate::query_engine::QueryEngineState;
33
34impl PromPlanner {
35 /// Resolve the `@` modifier of a vector or matrix selector into the timestamp its sample
36 /// window is anchored at, in milliseconds since the Unix epoch.
37 ///
38 /// Prometheus semantics:
39 /// - `@ <unix_ts>` anchors at the given timestamp,
40 /// - `@ start()` / `@ end()` anchor at the evaluation range of the whole statement,
41 /// - `offset` shifts the anchor backwards: the window ends at `anchor - offset`.
42 ///
43 /// Returns `None` when the selector has no `@` modifier.
44 fn at_ref_time(
45 &self,
46 at: &Option<AtModifier>,
47 offset: &Option<Offset>,
48 ) -> Result<Option<Millisecond>> {
49 let anchor = match at {
50 None => return Ok(None),
51 Some(AtModifier::Start) => self.ctx.stmt_start,
52 Some(AtModifier::End) => self.ctx.stmt_end,
53 Some(AtModifier::At(time)) => Self::system_time_to_millis(time)?,
54 };
55 Ok(Some(Self::anchor_sub(anchor, Self::offset_millis(offset))?))
56 }
57
58 /// Subtracts `rhs` from `lhs` on the millisecond timeline of an `@` anchor.
59 ///
60 /// A negative result is valid: `@` and `offset` accept timestamps before the Unix epoch. A
61 /// result outside the representable millisecond range is rejected like an unrepresentable
62 /// anchor ([`Self::system_time_to_millis`]) instead of clamping it, so the same class of
63 /// input always gets the same answer.
64 pub(crate) fn anchor_sub(lhs: Millisecond, rhs: Millisecond) -> Result<Millisecond> {
65 lhs.checked_sub(rhs)
66 .with_context(|| AtModifierTimestampOutOfRangeSnafu {
67 timestamp: format!("{}ms - {}ms", lhs, rhs),
68 })
69 }
70
71 /// The offset a selector with an `@` modifier is evaluated with.
72 ///
73 /// Prometheus anchors such a selector by rewriting its offset to `eval_time - anchor`
74 /// (`setOffsetForAtModifier`), so that the selector always selects its samples around `anchor`
75 /// regardless of the step being evaluated. `eval_time` is the start of the evaluation the
76 /// selector belongs to, which is `ctx.start`.
77 ///
78 /// Returns `None` when the selector has no `@` modifier.
79 pub(crate) fn at_modifier_offset(
80 &self,
81 at: &Option<AtModifier>,
82 offset: &Option<Offset>,
83 ) -> Result<Option<Millisecond>> {
84 let Some(anchor) = self.at_ref_time(at, offset)? else {
85 return Ok(None);
86 };
87 Ok(Some(Self::anchor_sub(self.ctx.start, anchor)?))
88 }
89
90 /// Whether `expr` is a call that has to be evaluated once for the whole grid, because it folds
91 /// a range selector anchored by `@` — the range argument of the call's parser signature, which
92 /// is a [`MatrixSelector`] here; only a call with such an argument can take this path, so no
93 /// function-name registry is involved.
94 ///
95 /// This is the shape that needs Prometheus' `StepInvariantExpr` the most: a range function such
96 /// as `rate` derives its result from the evaluation instant it is called at, so folding the
97 /// window once per step would let the outer evaluation grid change the result of a window that
98 /// `@` fixed. Evaluated once, at the start of the grid, the rewritten offset of
99 /// [`Self::at_modifier_offset`] places the anchor at that instant, and [`Self::replay_over_grid`]
100 /// reports the result at every step.
101 ///
102 /// Unlike Prometheus, which wraps the whole step-invariant subtree (`preprocessExprHelper`),
103 /// only the call itself is promoted here. The operators above it are not: they are still planned
104 /// at every step over the replayed result, which is safe for the row-wise ones and keeps the
105 /// promotion root narrow. The promotion root has to stay a call over one range selector, because
106 /// the replay needs one series per batch ([`Self::series_divide_plan`]) and only such a call
107 /// guarantees that the rows it emits still describe the series it was divided by. The operators
108 /// left out — an aggregation, a join, or a label rewriting call such as `label_join` — mix or
109 /// re-label the rows of different series, so they are unsafe as promotion roots even though
110 /// evaluating them after the replay is fine. A call whose input is an anchored *instant*
111 /// selector (`abs(some_metric @ 300)`) needs no promotion either: the selector anchors and
112 /// replays its sample per series on its own, and the call above it is row-wise.
113 ///
114 /// `predict_linear` is the exception among the range functions: it predicts from the evaluation
115 /// instant of each step ([`Self::create_range_eval_ts_expr`]), so it has to stay outside the
116 /// promoted subtree and follow the grid. The remaining arguments of the call have to be
117 /// literals, since the replay of the promoted result has no second vector input to divide.
118 /// Parentheses around the range argument are transparent (`rate((m[5m] @ 300))`), so they are
119 /// looked through and the call is promoted as if they were absent. Nothing else of the subtree
120 /// is unwrapped, so an outer parenthesis promotes no operator above the call.
121 fn promotes_anchored_range_call(expr: &PromExpr) -> bool {
122 let PromExpr::Call(Call { func, args }) = expr else {
123 return false;
124 };
125 // See the doc comment: the regression of `predict_linear` follows the evaluation step.
126 if func.name == "predict_linear" {
127 return false;
128 }
129 let mut anchored_range = false;
130 for arg in &args.args {
131 // Parentheses around the range argument are transparent, so the call is promoted the
132 // same way for `rate((m[5m] @ 300))` as for `rate(m[5m] @ 300)`. Only the parentheses
133 // directly around this one argument are looked through here: the promotion stays
134 // confined to a call over one anchored range selector instead of descending into an
135 // arbitrary parenthesized subtree.
136 let mut arg = arg.as_ref();
137 while let PromExpr::Paren(ParenExpr { expr }) = arg {
138 arg = expr;
139 }
140 match arg {
141 // The window is pinned by `@`, so every step folds the same samples.
142 PromExpr::MatrixSelector(MatrixSelector { vs, .. }) if vs.at.is_some() => {
143 if anchored_range {
144 return false;
145 }
146 anchored_range = true;
147 }
148 // A literal argument is the same value at every step.
149 arg if Self::try_build_literal_expr(arg).is_some() => {}
150 _ => return false,
151 }
152 }
153 anchored_range
154 }
155
156 /// Plans the anchored range call `prom_expr` ([`Self::promotes_anchored_range_call`]) as a
157 /// step-invariant subtree: the call is evaluated on a single evaluation instant (`grid_start`,
158 /// the start of the outer evaluation) and its result is then reported at every step of
159 /// `[grid_start, ctx.end]` by [`Self::replay_over_grid`].
160 ///
161 /// This is the planner's counterpart of Prometheus' `StepInvariantExpr` for the one shape it
162 /// promotes. Evaluating the call once matters for the functions that derive their result from
163 /// the step being evaluated: `rate(m[5m] @ 300)` folds its window around the anchor once, and
164 /// the extrapolation boundaries of `rate` must be derived from that same window at every step
165 /// instead of following the outer evaluation timestamp.
166 ///
167 /// Only the call itself is promoted; the operators above it are planned as usual over the
168 /// replayed result. The result of the promoted call is split into one series per batch before it
169 /// is replayed ([`Self::series_divide_plan`]), because the row-wise projection of the call does
170 /// not preserve the batch layout of the selector.
171 ///
172 /// Returns `None` when `prom_expr` is not such a call, so that the caller plans it as usual.
173 /// The selector inside the promoted call keeps its own `@` anchoring (see
174 /// [`Self::at_modifier_offset`]), and planning it with `ctx.end == ctx.start` folds its window
175 /// once for that single instant instead of expanding it over the grid, which the replay of the
176 /// call result above already does.
177 pub(crate) async fn promote_anchored_range_call(
178 &mut self,
179 prom_expr: &PromExpr,
180 timestamp_fn: bool,
181 query_engine_state: &QueryEngineState,
182 ) -> Result<Option<LogicalPlan>> {
183 let grid_start = self.ctx.start;
184 let grid_end = self.ctx.end;
185 // An instant query evaluates a single step, so there is nothing to promote.
186 if grid_start == grid_end || !Self::promotes_anchored_range_call(prom_expr) {
187 return Ok(None);
188 }
189
190 // Plan the subtree on a single evaluation instant: every selector below still anchors its
191 // window through `@`, and the functions above them derive their result from that one
192 // instant. The planner is single-use, so `ctx.end` needs no restore-on-error.
193 self.ctx.end = grid_start;
194 let anchored = self
195 .prom_expr_to_plan_inner(prom_expr, timestamp_fn, query_engine_state)
196 .await?;
197 self.ctx.end = grid_end;
198
199 let time_index_column =
200 self.ctx
201 .time_index_column
202 .clone()
203 .with_context(|| TimeIndexNotFoundSnafu {
204 table: self.ctx.table_name.clone().unwrap_or_default(),
205 })?;
206 // The replay reads one series per batch; see [`Self::series_divide_plan`] for why the
207 // layout of the selector below the promoted call does not survive it.
208 let anchored = self.series_divide_plan(anchored, &time_index_column)?;
209 Ok(Some(self.replay_over_grid(
210 anchored,
211 grid_start,
212 grid_end,
213 time_index_column,
214 )))
215 }
216
217 /// Convert the timestamp of an `@` modifier into milliseconds since the Unix epoch.
218 fn system_time_to_millis(time: &SystemTime) -> Result<Millisecond> {
219 let (millis, negative) = match time.duration_since(UNIX_EPOCH) {
220 Ok(duration) => (duration.as_millis(), false),
221 // The `@` modifier accepts timestamps before the Unix epoch, e.g. `@ -1`.
222 Err(err) => (err.duration().as_millis(), true),
223 };
224 ensure!(
225 millis <= i64::MAX as u128,
226 AtModifierTimestampOutOfRangeSnafu {
227 timestamp: time
228 .duration_since(UNIX_EPOCH)
229 .map(|duration| format!("+{}ms", duration.as_millis()))
230 .unwrap_or_else(|err| format!("-{}ms", err.duration().as_millis())),
231 }
232 );
233 let millis = millis as Millisecond;
234 Ok(if negative { -millis } else { millis })
235 }
236
237 /// Report the samples of `anchored` at every step of the evaluation grid
238 /// `[grid_start, grid_end]`.
239 ///
240 /// A selector with an `@` modifier is anchored: the sample window is selected once, around the
241 /// anchor timestamp, and every evaluation step reports that same window. Prometheus does this
242 /// by rewriting the selector's offset to `eval_time - anchor` and only fetching the samples on
243 /// the first step (`setOffsetForAtModifier` plus the `refetch` shortcut in `rangeEval`).
244 ///
245 /// The expansion reuses [`InstantManipulate`] with a lookback that spans the whole grid: every
246 /// step then picks the same sample (or the same already computed value, when `anchored` ends
247 /// with a function call such as `rate`) and stamps it with the step's timestamp.
248 ///
249 /// Every input batch of `anchored` must hold exactly one series, because [`InstantManipulate`]
250 /// takes a batch as one timeline. A leaf-level replay (`m @ 300`) consumes the [`SeriesDivide`]
251 /// of its selector directly. A promoted call is guaranteed that layout by
252 /// [`Self::series_divide_plan`], which is why [`Self::promote_anchored_range_call`] splits
253 /// its result before calling this method.
254 pub(crate) fn replay_over_grid(
255 &self,
256 anchored: LogicalPlan,
257 grid_start: Millisecond,
258 grid_end: Millisecond,
259 time_index_column: String,
260 ) -> LogicalPlan {
261 if grid_start == grid_end {
262 // A single evaluation step: `anchored` is already stamped with that timestamp.
263 return anchored;
264 }
265
266 let series_key_columns = self.series_key_columns_for_schema(anchored.schema());
267 let replayed = InstantManipulate::new(
268 grid_start,
269 grid_end,
270 // The lookback must keep the single anchored sample eligible for every step.
271 grid_end - grid_start + 1,
272 self.ctx.interval,
273 0,
274 time_index_column,
275 series_key_columns,
276 self.ctx.field_columns.first().cloned(),
277 anchored,
278 );
279 LogicalPlan::Extension(Extension {
280 node: Arc::new(replayed),
281 })
282 }
283
284 /// Sorts `input` by its series key and time index and splits it into one batch per series.
285 ///
286 /// [`InstantManipulate`] reads every input batch as one series (it takes the timeline of the
287 /// batch and reports the row selected at every step), so a batch holding several series would
288 /// lose all but one of them. A selector establishes that layout with its own [`SeriesDivide`],
289 /// but the per-series distribution requirement does not reach a promoted subtree above it:
290 /// the row-wise projection of a call sits in between, so the batch boundaries of the selector
291 /// are not preserved — in a distributed plan the promoted result can be delivered as one batch
292 /// holding every series. Sorting and dividing here restores the layout, exactly like
293 /// [`Self::prom_matrix_selector_to_plan`] does for the input of a range function.
294 ///
295 /// Series keys that are not present in `input` are dropped, since `ctx.tag_columns` may have
296 /// drifted from the actual output schema. A plan without any series key column is returned
297 /// unchanged: there is nothing to divide by.
298 fn series_divide_plan(
299 &self,
300 input: LogicalPlan,
301 time_index_column: &str,
302 ) -> Result<LogicalPlan> {
303 let series_key_columns = self.series_key_columns_for_schema(input.schema());
304 if series_key_columns.is_empty() {
305 return Ok(input);
306 }
307
308 let mut sort_exprs = series_key_columns
309 .iter()
310 .map(|name| DfExpr::Column(Column::from_name(name)).sort(true, true))
311 .collect::<Vec<_>>();
312 sort_exprs.push(DfExpr::Column(Column::from_name(time_index_column)).sort(true, true));
313 let sort_plan = LogicalPlanBuilder::from(input)
314 .sort(sort_exprs)
315 .context(DataFusionPlanningSnafu)?
316 .build()
317 .context(DataFusionPlanningSnafu)?;
318 Ok(LogicalPlan::Extension(Extension {
319 node: Arc::new(SeriesDivide::new(
320 series_key_columns,
321 time_index_column.to_string(),
322 sort_plan,
323 )),
324 }))
325 }
326}