feat(query): implement PromQL @ modifier on vector and matrix selectors (#9224)

* feat(query): implement PromQL @ modifier on vector and matrix selectors

The planner previously dropped the `at` field of selectors (`at: _`),
so `some_metric @ 300` silently returned step-following values instead
of the samples anchored at the fixed timestamp.

Anchoring follows Prometheus's `setOffsetForAtModifier` + `refetch`
semantics: resolve the anchor (`@ <ts>`, `@ start()`, `@ end()`), apply
`offset` to the anchor, rewrite the selector offset to
`eval_start - anchor`, scan only the anchored window, then report the
same window at every evaluation step via a grid-wide-replay
`InstantManipulate`. `@ start()` / `@ end()` resolve against the
statement's evaluation range (`stmt_start` / `stmt_end`), which subquery
planning does not rewrite.

Timestamps before the Unix epoch are accepted; unrepresentable ones are
rejected with `AtModifierTimestampOutOfRange` instead of wrapping.
Selectors without `@` are planned exactly as before.

Report: .e-agent/greptimedb_promql_compatibility_report_2026-09-16.md P0-2
Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>

* fix(promql): center predict_linear on the evaluation timestamp

predict_linear_impl used the window's last sample time as the evaluation
timestamp, and the UDF took only (ts_range, value_range, t) with no channel
for the current instant. After a range selector is folded for @ (or with an
offset) the same window is replayed at every step, so the prediction stayed
constant at the anchor's answer; even a plain window ended before the step
produced a stale value.

Give prom_predict_linear a 4th argument carrying the step's evaluation
instant (ms timestamp), derived in the planner from the row's time index
plus the offset the window was folded with (at_offset for @, offset_ms
otherwise, recorded on PromPlannerContext). The regression is centered on
that instant, matching Prometheus' use of enh.Ts. The mixed
float/native-histogram path forwards the same argument.

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>

* refactor(query): tighten @ modifier planner and predict_linear

- range_fold_offset is always a concrete millisecond offset, not an
  Option; the single-step replay helper no longer wraps an infallible
  plan in Result.
- predict_linear's eval timestamp is always cast to Timestamp(ms) at the
  call site, so the UDF drops its dead Int64 branch and the extra
  func_name argument.
- Drop two redundant comparison cases from the at_modifier sqlness test
  and fix the lookback window comment to the half-open (0s, 300s].

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>

* test(query): accept parse-time rejection of @ on Windows

`@ 1e16` is 10^19 milliseconds, beyond i64::MAX. A Unix SystemTime
holds it and the planner rejects the anchor it cannot represent, but a
Windows SystemTime tops out near 1.8e12 seconds, so the parser's
checked_add fails first and the same literal is rejected while parsing.
Accept either rejection path so the test passes on both platforms.

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>

* fix(query): avoid replaying label_join over rewritten series keys

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>

* fix(query): narrow anchored range call promotion

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>

* fix(query): address @ modifier review feedback

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>

* fix(query): satisfy super import format check

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>

* docs(query): address at modifier review follow-ups

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>

---------

Signed-off-by: discord9 <55937128+discord9@users.noreply.github.com>
This commit is contained in:
discord9
2026-09-28 09:38:12 +00:00
committed by GitHub
parent 03823a9a01
commit 28e415a103
9 changed files with 2529 additions and 94 deletions
+3 -1
View File
@@ -270,7 +270,7 @@ fn make_quantile_input(num_points: usize, window_size: u32) -> Vec<ColumnarValue
}
fn make_predict_linear_input(num_points: usize, window_size: u32) -> Vec<ColumnarValue> {
let (ts_range, val_range, _) = build_sliding_ranges(
let (ts_range, val_range, eval_ts) = build_sliding_ranges(
num_points,
window_size,
1,
@@ -282,6 +282,8 @@ fn make_predict_linear_input(num_points: usize, window_size: u32) -> Vec<Columna
ColumnarValue::Array(Arc::new(val_range.into_dict())),
// predict 60s into the future
ColumnarValue::Scalar(ScalarValue::Int64(Some(60))),
// evaluated at the end of each window
ColumnarValue::Array(eval_ts),
]
}
+140 -16
View File
@@ -46,6 +46,8 @@ impl PredictLinear {
RangeArray::convert_data_type(DataType::Float64),
// t
DataType::Int64,
// evaluation timestamp
DataType::Timestamp(TimeUnit::Millisecond, None),
];
create_udf(
Self::name(),
@@ -58,11 +60,12 @@ impl PredictLinear {
fn predict_linear(input: &[ColumnarValue]) -> Result<ColumnarValue, DataFusionError> {
error::ensure(
input.len() == 3,
DataFusionError::Plan("prom_predict_linear function should have 3 inputs".to_string()),
input.len() == 4,
DataFusionError::Plan("prom_predict_linear function should have 4 inputs".to_string()),
)?;
let t_col = &input[2];
let eval_ts_col = &input[3];
let ts_range = extract_range_array(&input[0])?;
let value_range = extract_range_array(&input[1])?;
@@ -128,6 +131,51 @@ impl PredictLinear {
Box::new(t_array.iter())
}
};
// The evaluation instant the window above was folded for. It is the instant the regression
// is centered on, which is *not* the last sample's timestamp: a window may end before the
// step being evaluated, and an `@`-anchored window is replayed at every step while keeping
// its own end. See `PromPlanner::create_range_eval_ts_expr`.
let eval_ts_iter: Box<dyn Iterator<Item = Option<i64>>> = match eval_ts_col {
ColumnarValue::Scalar(eval_ts_scalar) => {
let eval_ts = match eval_ts_scalar {
ScalarValue::TimestampMillisecond(Some(eval_ts), _) => *eval_ts,
// For a NULL or otherwise unusable evaluation timestamp, returns NULL array,
// which conforms to PromQL's behavior.
_ => {
let null_array = Float64Array::new_null(ts_range.len());
return Ok(ColumnarValue::Array(Arc::new(null_array)));
}
};
Box::new((0..ts_range.len()).map(move |_| Some(eval_ts)))
}
ColumnarValue::Array(eval_ts_array) => {
error::ensure(
eval_ts_array.len() == ts_range.len(),
DataFusionError::Execution(format!(
"{}: evaluation timestamp array should have the same length as other columns, found {} and {}",
Self::name(),
eval_ts_array.len(),
ts_range.len()
)),
)?;
match eval_ts_array.data_type() {
DataType::Timestamp(TimeUnit::Millisecond, _) => Box::new(
eval_ts_array
.as_any()
.downcast_ref::<TimestampMillisecondArray>()
.expect("checked by data type")
.iter(),
),
other => {
return Err(DataFusionError::Execution(format!(
"{}: expect TimestampMillisecond as evaluation timestamp array's type, found {}",
Self::name(),
other
)));
}
}
}
};
let all_timestamps = ts_range
.values()
.as_any()
@@ -140,7 +188,12 @@ impl PredictLinear {
.downcast_ref::<Float64Array>()
.unwrap();
let mut result_builder = Float64Builder::with_capacity(ts_range.len());
for (index, t) in t_iter.enumerate() {
for (index, (t, eval_ts)) in t_iter.zip(eval_ts_iter).enumerate() {
// A step without an evaluation instant has nothing to extrapolate to.
let Some(eval_ts) = eval_ts else {
result_builder.append_null();
continue;
};
match predict_linear_impl(
&ts_range,
&value_range,
@@ -148,7 +201,7 @@ impl PredictLinear {
all_values,
index,
t.unwrap(),
Self::name(),
eval_ts,
)? {
Some(value) => result_builder.append_value(value),
None => result_builder.append_null(),
@@ -167,7 +220,7 @@ fn predict_linear_impl(
all_values: &Float64Array,
index: usize,
t: i64,
func_name: &str,
eval_ts: i64,
) -> Result<Option<f64>, DataFusionError> {
let (ts_offset, ts_len) = ts_range.get_offset_length(index).unwrap();
let (value_offset, value_len) = value_range.get_offset_length(index).unwrap();
@@ -175,22 +228,28 @@ fn predict_linear_impl(
ts_len == value_len,
DataFusionError::Execution(format!(
"{}: time and value arrays in a group should have the same length, found {} and {}",
func_name, ts_len, value_len
PredictLinear::name(),
ts_len,
value_len
)),
)?;
if ts_len < 2 {
return Ok(None);
}
// last timestamp is evaluation timestamp
let evaluate_ts = all_timestamps[ts_offset + ts_len - 1];
// Like Prometheus, the regression is centered on the evaluation timestamp: the returned
// intercept is the value the window's trend predicts at the step being evaluated, and `t` is
// the horizon (in seconds) that is added on top of it. Centering on the last sample instead
// would drop the distance between the window's end and the step, which is non-zero both
// without `@` (the window can end before the step) and with it (the same anchored window is
// replayed at every step).
let (slope, intercept) = linear_regression_slices(
all_timestamps,
ts_offset,
all_values,
value_offset,
value_len,
evaluate_ts,
eval_ts,
);
if slope.is_none() || intercept.is_none() {
@@ -231,6 +290,11 @@ mod test {
(ts_range_array, value_range_array)
}
/// The evaluation timestamp at which the window built by [`build_test_range_arrays`] ends.
fn window_end_eval_ts() -> ScalarValue {
ScalarValue::TimestampMillisecond(Some(3000), None)
}
#[test]
fn calculate_predict_linear_none() {
let ts_array = Arc::new(TimestampMillisecondArray::from_iter(
@@ -244,7 +308,7 @@ mod test {
PredictLinear::scalar_udf(),
ts_array,
value_array,
vec![ScalarValue::Int64(Some(0))],
vec![ScalarValue::Int64(Some(0)), window_end_eval_ts()],
vec![None, None],
);
}
@@ -256,7 +320,7 @@ mod test {
PredictLinear::scalar_udf(),
ts_array,
value_array,
vec![ScalarValue::Int64(Some(0))],
vec![ScalarValue::Int64(Some(0)), window_end_eval_ts()],
// value at t = 0
vec![Some(38.63636363636364)],
);
@@ -269,7 +333,7 @@ mod test {
PredictLinear::scalar_udf(),
ts_array,
value_array,
vec![ScalarValue::Int64(Some(3000))],
vec![ScalarValue::Int64(Some(3000)), window_end_eval_ts()],
// value at t = 3000
vec![Some(31856.818181818187)],
);
@@ -282,7 +346,7 @@ mod test {
PredictLinear::scalar_udf(),
ts_array,
value_array,
vec![ScalarValue::Int64(Some(4200))],
vec![ScalarValue::Int64(Some(4200)), window_end_eval_ts()],
// value at t = 4200
vec![Some(44584.09090909091)],
);
@@ -295,7 +359,7 @@ mod test {
PredictLinear::scalar_udf(),
ts_array,
value_array,
vec![ScalarValue::Int64(Some(6600))],
vec![ScalarValue::Int64(Some(6600)), window_end_eval_ts()],
// value at t = 6600
vec![Some(70038.63636363638)],
);
@@ -308,7 +372,7 @@ mod test {
PredictLinear::scalar_udf(),
ts_array,
value_array,
vec![ScalarValue::Int64(Some(7800))],
vec![ScalarValue::Int64(Some(7800)), window_end_eval_ts()],
// value at t = 7800
vec![Some(82765.9090909091)],
);
@@ -327,7 +391,7 @@ mod test {
PredictLinear::scalar_udf(),
ts_array,
value_array,
vec![ScalarValue::Int64(Some(0))],
vec![ScalarValue::Int64(Some(0)), window_end_eval_ts()],
vec![Some(30.0)],
);
}
@@ -348,9 +412,69 @@ mod test {
ColumnarValue::Array(Arc::new(ts_dict)),
ColumnarValue::Array(Arc::new(value_dict)),
ColumnarValue::Scalar(ScalarValue::Int64(Some(0))),
ColumnarValue::Scalar(window_end_eval_ts()),
])
.unwrap_err();
assert!(err.to_string().contains("Empty range is not expected"));
}
/// The regression is centered on the evaluation timestamp, so the returned intercept follows
/// the step being evaluated even though the window stays the same.
#[test]
fn calculate_predict_linear_centers_on_eval_ts() {
// One sample per second with values 0, 1, 2: the trend is exactly one unit per second.
for (eval_ts, t, expected) in [
// The window ends at 2s: evaluating there returns its last value.
(2000, 0, 2.0),
// One second past the window's end the trend is at 3.
(3000, 0, 3.0),
// One second past the window's end, predicting another two seconds ahead.
(3000, 2, 5.0),
] {
let ts_values = Arc::new(TimestampMillisecondArray::from_iter(
[0i64, 1000, 2000].into_iter().map(Some),
));
let value_values = Arc::new(Float64Array::from_iter([0.0, 1.0, 2.0]));
let ts_array = RangeArray::from_ranges(ts_values, [(0, 3)]).unwrap();
let value_array = RangeArray::from_ranges(value_values, [(0, 3)]).unwrap();
simple_range_udf_runner(
PredictLinear::scalar_udf(),
ts_array,
value_array,
vec![
ScalarValue::Int64(Some(t)),
ScalarValue::TimestampMillisecond(Some(eval_ts), None),
],
vec![Some(expected)],
);
}
}
/// The evaluation timestamp arrives as a column, so the UDF accepts it as an array too.
#[test]
fn calculate_predict_linear_accepts_eval_ts_array() {
let (ts_array, value_array) = build_test_range_arrays();
let eval_ts_array = Arc::new(TimestampMillisecondArray::from_iter_values([3000]));
let result = PredictLinear::predict_linear(&[
ColumnarValue::Array(Arc::new(ts_array.into_dict())),
ColumnarValue::Array(Arc::new(value_array.into_dict())),
ColumnarValue::Scalar(ScalarValue::Int64(Some(0))),
ColumnarValue::Array(eval_ts_array),
])
.unwrap();
let result = result
.into_array(1)
.unwrap()
.as_any()
.downcast_ref::<Float64Array>()
.unwrap()
.clone();
assert_eq!(result.len(), 1);
// The same value `calculate_predict_linear_test1` gets from a scalar evaluation timestamp.
assert!((result.value(0) - 38.63636363636364).abs() < 0.0001);
}
}
+8
View File
@@ -195,6 +195,13 @@ pub enum Error {
location: Location,
},
#[snafu(display("Timestamp out of range for the `@` modifier: {}", timestamp))]
AtModifierTimestampOutOfRange {
timestamp: String,
#[snafu(implicit)]
location: Location,
},
#[snafu(display("Timestamp out of range: {} of {:?}", timestamp, unit))]
TimestampOutOfRange {
timestamp: i64,
@@ -252,6 +259,7 @@ impl ErrorExt for Error {
| SameLabelSet { .. }
| TimestampOutOfRange { .. }
| SystemTimeOutOfRange { .. }
| AtModifierTimestampOutOfRange { .. }
| InvalidRegularExpression { .. }
| InvalidDestinationLabelName { .. } => StatusCode::InvalidArguments,
+337 -71
View File
@@ -12,10 +12,9 @@
// See the License for the specific language governing permissions and
// limitations under the License.
mod at_modifier;
mod function_plans;
mod island;
mod matching_filters;
mod set_operator;
@@ -165,6 +164,12 @@ struct PromPlannerContext {
end: Millisecond,
interval: Millisecond,
lookback_delta: Millisecond,
/// Evaluation range of the whole statement, which `@ start()` and `@ end()` refer to.
///
/// Unlike [`Self::start`] and [`Self::end`], these are never rewritten while planning, so a
/// selector inside a subquery still resolves `@ start()` / `@ end()` against the statement.
stmt_start: Millisecond,
stmt_end: Millisecond,
// planner states
table_name: Option<String>,
@@ -193,6 +198,16 @@ struct PromPlannerContext {
schema_name: Option<String>,
/// The range in millisecond of range selector. None if there is no range selector.
range: Option<Millisecond>,
/// The offset in milliseconds the window of the last planned range selector is folded with,
/// or `None` when no range selector has been planned since the last read.
///
/// [`Self::start`] and the sample timestamps are compared on the shifted timeline: a range
/// payload carries `sample_timestamp + offset`, while the time index column of a folded row
/// stays the evaluation timestamp of its step. A function that reads it right after its input
/// plan is built, like `predict_linear`, consumes it instead of reading the state, so the
/// offset cannot leak from one input to another; see
/// [`PromPlanner::create_function_expr`].
range_fold_offset: Option<Millisecond>,
}
/// Result labels a vector-vector binary operation derives from its matching modifier, projected
@@ -226,11 +241,15 @@ impl BinaryResultLabels {
impl PromPlannerContext {
fn from_eval_stmt(stmt: &EvalStmt) -> Self {
let start = stmt.start.duration_since(UNIX_EPOCH).unwrap().as_millis() as Millisecond;
let end = stmt.end.duration_since(UNIX_EPOCH).unwrap().as_millis() as Millisecond;
Self {
start: stmt.start.duration_since(UNIX_EPOCH).unwrap().as_millis() as _,
end: stmt.end.duration_since(UNIX_EPOCH).unwrap().as_millis() as _,
start,
end,
interval: stmt.interval.as_millis() as _,
lookback_delta: stmt.lookback_delta.as_millis() as _,
stmt_start: start,
stmt_end: end,
..Default::default()
}
}
@@ -246,6 +265,7 @@ impl PromPlannerContext {
self.selector_matcher.clear();
self.schema_name = None;
self.range = None;
self.range_fold_offset = None;
}
/// Reset table name and schema to empty
@@ -325,6 +345,16 @@ impl PromPlanner {
timestamp_fn: bool,
query_engine_state: &QueryEngineState,
) -> Result<LogicalPlan> {
// An anchored range call is step-invariant: evaluate it once, at the start of the
// evaluation, and report its result at every step; see
// [`Self::promote_anchored_range_call`].
if let Some(plan) = self
.promote_anchored_range_call(prom_expr, timestamp_fn, query_engine_state)
.await?
{
return Ok(plan);
}
let res = match prom_expr {
PromExpr::Aggregate(expr) => {
self.prom_aggr_expr_to_plan(query_engine_state, expr)
@@ -458,6 +488,9 @@ impl PromPlanner {
divide_plan,
)
.context(DataFusionPlanningSnafu)?;
// A subquery always folds with offset 0, so its payload timestamps are already on the
// evaluation timeline a function above it reads; see [`Self::create_range_eval_ts_expr`].
self.ctx.range_fold_offset = Some(0);
Ok(LogicalPlan::Extension(Extension {
node: Arc::new(manipulate),
@@ -1579,6 +1612,44 @@ impl PromPlanner {
Ok(plan)
}
/// The offset of a selector in milliseconds. A positive offset selects samples from an earlier
/// time and moves them forward into the evaluation timeline.
fn offset_millis(offset: &Option<Offset>) -> Millisecond {
match offset {
Some(Offset::Pos(duration)) => duration.as_millis() as Millisecond,
Some(Offset::Neg(duration)) => -(duration.as_millis() as Millisecond),
None => 0,
}
}
/// The columns that identify one series, which is the series key expected by the PromQL plan
/// nodes that hold exactly one series per input batch.
fn series_key_columns(&self) -> Vec<String> {
if self.ctx.use_tsid {
vec![DATA_SCHEMA_TSID_COLUMN_NAME.to_string()]
} else {
self.ctx.tag_columns.clone()
}
}
/// Keep replay and series division on the same effective keys when a call rewrites labels.
fn series_key_columns_for_schema(&self, schema: &DFSchemaRef) -> Vec<String> {
let has_tsid = schema.fields().iter().any(|field| {
field.name() == DATA_SCHEMA_TSID_COLUMN_NAME
&& field.data_type() == &ArrowDataType::UInt64
});
if has_tsid {
vec![DATA_SCHEMA_TSID_COLUMN_NAME.to_string()]
} else {
self.ctx
.tag_columns
.iter()
.filter(|name| schema.has_column_with_unqualified_name(name))
.cloned()
.collect()
}
}
async fn prom_vector_selector_to_plan(
&mut self,
vector_selector: &VectorSelector,
@@ -1588,20 +1659,37 @@ impl PromPlanner {
name,
offset,
matchers,
at: _,
at,
} = vector_selector;
let matchers = self.preprocess_label_matchers(matchers, name)?;
if let Some(empty_plan) = self.setup_context().await? {
return Ok(empty_plan);
}
let offset_ms = match offset {
Some(Offset::Pos(duration)) => duration.as_millis() as Millisecond,
Some(Offset::Neg(duration)) => -(duration.as_millis() as Millisecond),
None => 0,
let offset_ms = Self::offset_millis(offset);
// `@` anchors the sample window at a fixed timestamp: the selector selects its samples
// around the anchor once, instead of following the outer evaluation grid. See
// [`Self::at_modifier_offset`].
let at_offset = self.at_modifier_offset(at, offset)?;
let grid_start = self.ctx.start;
let grid_end = self.ctx.end;
let normalize = match at_offset {
Some(at_offset) => {
// Select the anchored samples at the start of the evaluation, with the offset that
// re-anchors the selector.
// The planner is single-use (one `EvalStmt` produces one plan), so an error
// below aborts the whole planning and `ctx.end` needs no restore-on-error.
self.ctx.end = grid_start;
let plan = self
.selector_to_series_normalize_plan(at_offset, matchers, false)
.await?;
self.ctx.end = grid_end;
plan
}
None => {
self.selector_to_series_normalize_plan(offset_ms, matchers, false)
.await?
}
};
let normalize = self
.selector_to_series_normalize_plan(offset, matchers, false)
.await?;
let time_index_column =
self.ctx
.time_index_column
@@ -1683,24 +1771,46 @@ impl PromPlanner {
};
let field_column = self.ctx.field_columns.first().cloned();
let manipulate = InstantManipulate::new(
self.ctx.start,
self.ctx.end,
self.ctx.lookback_delta,
self.ctx.interval,
offset_ms,
time_index_column,
if self.ctx.use_tsid {
vec![DATA_SCHEMA_TSID_COLUMN_NAME.to_string()]
} else {
self.ctx.tag_columns.clone()
},
field_column,
normalize,
);
let manipulate = LogicalPlan::Extension(Extension {
node: Arc::new(manipulate),
});
let series_key_columns = self.series_key_columns();
let manipulate = match at_offset {
Some(at_offset) => {
// Select the anchored sample once, then report it at every step of the outer
// grid. The samples keep their native anchor-time timestamps, so the manipulate
// must shift them onto the evaluation timeline with the rewritten offset.
let anchored = InstantManipulate::new(
grid_start,
grid_start,
self.ctx.lookback_delta,
self.ctx.interval,
at_offset,
time_index_column.clone(),
series_key_columns,
field_column,
normalize,
);
self.replay_over_grid(
LogicalPlan::Extension(Extension {
node: Arc::new(anchored),
}),
grid_start,
grid_end,
time_index_column,
)
}
None => LogicalPlan::Extension(Extension {
node: Arc::new(InstantManipulate::new(
grid_start,
grid_end,
self.ctx.lookback_delta,
self.ctx.interval,
offset_ms,
time_index_column,
series_key_columns,
field_column,
normalize,
)),
}),
};
if let Some(timestamp_value_column) = timestamp_value_column {
self.create_timestamp_func_plan(manipulate, &timestamp_value_column)
} else {
@@ -1758,46 +1868,108 @@ impl PromPlanner {
name,
offset,
matchers,
..
at,
} = vs;
let matchers = self.preprocess_label_matchers(matchers, name)?;
ensure!(!range.is_zero(), ZeroRangeSelectorSnafu);
let range_ms = range.as_millis() as _;
self.ctx.range = Some(range_ms);
let offset_ms = match offset {
Some(Offset::Pos(duration)) => duration.as_millis() as Millisecond,
Some(Offset::Neg(duration)) => -(duration.as_millis() as Millisecond),
None => 0,
};
let offset_ms = Self::offset_millis(offset);
// `@` anchors the range selector's window at a fixed timestamp, so the same window is fed
// to the enclosing function at every evaluation step. See [`Self::at_modifier_offset`].
let at_offset = self.at_modifier_offset(at, offset)?;
let grid_start = self.ctx.start;
let grid_end = self.ctx.end;
// Some functions like rate may require special fields in the RangeManipulate plan
// so we can't skip RangeManipulate.
let normalize = match self.setup_context().await? {
Some(empty_plan) => empty_plan,
let (normalize, at_offset) = match self.setup_context().await? {
// An empty metric does not contain any sample, so anchoring cannot change the result.
// The manipulate below folds the empty input with `offset_ms`, and the recorded fold
// offset agrees with it instead of with the anchor the window cannot use.
Some(empty_plan) => {
self.ctx.range_fold_offset = Some(offset_ms);
(empty_plan, None)
}
None => {
self.selector_to_series_normalize_plan(offset, matchers, true)
.await?
let normalize = match at_offset {
Some(at_offset) => {
// Fold the anchored window once, at the start of the evaluation.
// Single-use planner: an error below aborts planning, so `ctx.end`
// needs no restore-on-error.
self.ctx.end = grid_start;
let plan = self
.selector_to_series_normalize_plan(at_offset, matchers, true)
.await?;
self.ctx.end = grid_end;
plan
}
None => {
self.selector_to_series_normalize_plan(offset_ms, matchers, true)
.await?
}
};
// Samples are shifted onto the evaluation timeline with the very same offset while
// the window is folded. Record it so that the function above the selector can
// recover the evaluation instant of the folded window; see
// [`Self::create_range_eval_ts_expr`].
self.ctx.range_fold_offset = Some(at_offset.unwrap_or(offset_ms));
(normalize, at_offset)
}
};
let manipulate = RangeManipulate::new(
self.ctx.start,
self.ctx.end,
self.ctx.interval,
offset_ms,
// TODO(ruihang): convert via Timestamp datatypes to support different time units
range_ms,
self.ctx
.time_index_column
.clone()
.expect("time index should be set in `setup_context`"),
self.ctx.field_columns.clone(),
normalize,
)
.context(DataFusionPlanningSnafu)?;
let time_index_column = self
.ctx
.time_index_column
.clone()
.expect("time index should be set in `setup_context`");
let manipulate = match at_offset {
Some(at_offset) => {
// Fold the anchored window once, then report it at every step of the outer
// grid. The samples keep their native anchor-time timestamps, so the manipulate
// must shift them onto the evaluation timeline with the rewritten offset.
let anchored = RangeManipulate::new(
grid_start,
grid_start,
self.ctx.interval,
at_offset,
// TODO(ruihang): convert via Timestamp datatypes to support different time units
range_ms,
time_index_column.clone(),
self.ctx.field_columns.clone(),
normalize,
)
.context(DataFusionPlanningSnafu)?;
self.replay_over_grid(
LogicalPlan::Extension(Extension {
node: Arc::new(anchored),
}),
grid_start,
grid_end,
time_index_column,
)
}
None => {
let manipulate = RangeManipulate::new(
grid_start,
grid_end,
self.ctx.interval,
offset_ms,
// TODO(ruihang): convert via Timestamp datatypes to support different time units
range_ms,
time_index_column,
self.ctx.field_columns.clone(),
normalize,
)
.context(DataFusionPlanningSnafu)?;
Ok(LogicalPlan::Extension(Extension {
node: Arc::new(manipulate),
}))
LogicalPlan::Extension(Extension {
node: Arc::new(manipulate),
})
}
};
Ok(manipulate)
}
async fn prom_call_expr_to_plan(
@@ -1846,11 +2018,17 @@ impl PromPlanner {
),
})
};
// The input plan records the fold offset of the range selector it is built from. Take it
// here, so that the offset of one input cannot leak into another call, and pass it to
// `create_function_expr`: the function that reads it (`predict_linear`) then depends on an
// argument instead of on planner state written by the selector below it.
let range_fold_offset = self.ctx.range_fold_offset.take();
let (mut func_exprs, new_tags) = self.create_function_expr(
func,
args.literals.clone(),
input.schema(),
query_engine_state,
range_fold_offset,
)?;
func_exprs.insert(0, self.create_time_index_column_expr()?);
func_exprs.extend_from_slice(&self.create_tag_column_exprs()?);
@@ -2034,7 +2212,7 @@ impl PromPlanner {
async fn selector_to_series_normalize_plan(
&mut self,
offset: &Option<Offset>,
offset_duration: Millisecond,
label_matchers: Matchers,
is_range_selector: bool,
) -> Result<LogicalPlan> {
@@ -2044,11 +2222,6 @@ impl PromPlanner {
let table_schema = table_scan.schema();
// make filter exprs
let offset_duration = match offset {
Some(Offset::Pos(duration)) => duration.as_millis() as Millisecond,
Some(Offset::Neg(duration)) => -(duration.as_millis() as Millisecond),
None => 0,
};
let mut scan_filters = Self::matchers_to_expr(label_matchers.clone(), table_schema)?;
if let Some(time_index_filter) =
self.build_time_index_filter(offset_duration, table_schema)?
@@ -2147,11 +2320,7 @@ impl PromPlanner {
}
// make sort plan
let series_key_columns = if self.ctx.use_tsid {
vec![DATA_SCHEMA_TSID_COLUMN_NAME.to_string()]
} else {
self.ctx.tag_columns.clone()
};
let series_key_columns = self.series_key_columns();
let sort_exprs = if self.ctx.use_tsid {
vec![
@@ -2948,6 +3117,7 @@ impl PromPlanner {
float_field: &str,
histogram_field: &str,
input_schema: &DFSchemaRef,
range_fold_offset: Option<Millisecond>,
) -> Result<Option<Vec<DfExpr>>> {
let returns_histogram = matches!(
func.name,
@@ -2987,6 +3157,14 @@ impl PromPlanner {
Box::new(other_input_exprs[0].clone()),
ArrowDataType::Int64,
));
// Same evaluation instant as in the non-mixed path; it follows the other inputs so the
// float UDF below is called as `predict_linear(ts_range, value_range, t, eval_ts)`.
// The regression is only defined for a folded window, so a missing fold offset is an
// error instead of a default.
other_input_exprs.push_back(self.create_range_eval_ts_expr(
range_fold_offset.context(ExpectRangeSelectorSnafu)?,
input_schema,
)?);
}
let timestamp_range = DfExpr::Column(Column::from_name(
@@ -3045,7 +3223,13 @@ impl PromPlanner {
.alias(histogram_field),
]
} else {
let display_name = float_expr.schema_name().to_string();
let display_name = if func.name == "predict_linear" {
// The evaluation instant is the private last argument of the mixed float UDF
// call; keep it out of the output column name like on the non-mixed path.
Self::name_without_last_arg(&float_expr)
} else {
float_expr.schema_name().to_string()
};
self.ctx.field_columns = vec![display_name.clone()];
vec![float_expr.alias(display_name)]
};
@@ -3063,6 +3247,7 @@ impl PromPlanner {
other_input_exprs: Vec<DfExpr>,
input_schema: &DFSchemaRef,
query_engine_state: &QueryEngineState,
range_fold_offset: Option<Millisecond>,
) -> Result<(Vec<DfExpr>, Vec<String>)> {
// TODO(ruihang): check function args list
let mut other_input_exprs: VecDeque<DfExpr> = other_input_exprs.into();
@@ -3075,6 +3260,7 @@ impl PromPlanner {
&float_field,
&histogram_field,
input_schema,
range_fold_offset,
)?
{
return Ok((exprs, vec![]));
@@ -3276,6 +3462,16 @@ impl PromPlanner {
Box::new(other_input_exprs[0].clone()),
ArrowDataType::Int64,
));
// The prediction starts at the evaluation instant of the step, which the
// window's fold offset recovers from the row's time index; see
// [`Self::create_range_eval_ts_expr`]. It is appended last, so the UDF is
// called as `predict_linear(ts_range, value_range, t, eval_ts)`. The
// regression is only defined for a folded window, so a missing fold offset
// is an error instead of a default.
other_input_exprs.push_back(self.create_range_eval_ts_expr(
range_fold_offset.context(ExpectRangeSelectorSnafu)?,
input_schema,
)?);
ScalarFunc::Udf(Arc::new(PredictLinear::scalar_udf()))
}
}
@@ -3648,7 +3844,17 @@ impl PromPlanner {
exprs = exprs
.into_iter()
.map(|expr| {
let display_name = expr.schema_name().to_string();
// `predict_linear` appends its private evaluation instant as the last argument
// of the UDF call; the output column is named after the call without it, so
// the injected expression stays out of the user-visible schema. The
// native-histogram drop UDF takes no such argument.
let display_name = if func.name == "predict_linear"
&& !all_field_columns_are_native_histogram_ranges
{
Self::name_without_last_arg(&expr)
} else {
expr.schema_name().to_string()
};
new_field_columns.push(display_name.clone());
Ok(expr.alias(display_name))
})
@@ -3917,6 +4123,66 @@ impl PromPlanner {
)))
}
/// Builds the evaluation instant the window of the last planned range selector is folded for,
/// as a `Timestamp(Millisecond)` expression.
///
/// The timestamp payload of a folded window is shifted onto the evaluation timeline by the
/// offset the window is folded with (`fold_offset`), while the time index column of a folded
/// row keeps the evaluation timestamp of its step. Adding the offset back yields the
/// evaluation instant on the payload timeline, which is where the regression of
/// `predict_linear` is centered: neither a plain window (which may end before the step, and
/// is additionally shifted by `offset` on the payload timeline) nor an `@`-anchored one
/// (whose end is the anchor, while the payload is shifted by `at_offset`) ends at the step it
/// is evaluated at.
///
/// The sum is computed on the millisecond representation and cast back, so that the result
/// keeps the `Timestamp(Millisecond)` type the range functions declare for it.
fn create_range_eval_ts_expr(
&self,
fold_offset: Millisecond,
input_schema: &DFSchemaRef,
) -> Result<DfExpr> {
let eval_ts = self
.create_time_index_column_expr()?
.cast_to(
&ArrowDataType::Timestamp(ArrowTimeUnit::Millisecond, None),
input_schema,
)
.context(DataFusionPlanningSnafu)?
.cast_to(&ArrowDataType::Int64, input_schema)
.context(DataFusionPlanningSnafu)?;
DfExpr::BinaryExpr(BinaryExpr {
left: Box::new(eval_ts),
op: Operator::Plus,
right: Box::new(lit(fold_offset)),
})
.cast_to(
&ArrowDataType::Timestamp(ArrowTimeUnit::Millisecond, None),
input_schema,
)
.context(DataFusionPlanningSnafu)
}
/// The name of `expr` without its last argument.
///
/// A `predict_linear` call appends the evaluation instant of its window as a private last
/// argument ([`Self::create_range_eval_ts_expr`]). Naming the output column after the call
/// the user wrote, without that argument, keeps the injected expression out of the
/// user-visible schema; the expression itself keeps every argument it needs.
fn name_without_last_arg(expr: &DfExpr) -> String {
if let DfExpr::ScalarFunction(ScalarFunction { func, args }) = expr
&& let Some((_, visible_args)) = args.split_last()
{
let visible = ScalarFunction {
func: func.clone(),
args: visible_args.to_vec(),
};
return DfExpr::ScalarFunction(visible).schema_name().to_string();
}
expr.schema_name().to_string()
}
fn create_tag_column_exprs(&self) -> Result<Vec<DfExpr>> {
let mut result = Vec::with_capacity(self.ctx.tag_columns.len());
for tag in &self.ctx.tag_columns {
+326
View File
@@ -0,0 +1,326 @@
// Copyright 2023 Greptime Team
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
//! Anchored selector planning and replay for the PromQL `@` modifier.
use std::sync::Arc;
use std::time::{SystemTime, UNIX_EPOCH};
use datafusion::logical_expr::{Extension, LogicalPlan, LogicalPlanBuilder};
use datafusion::prelude::{Column, Expr as DfExpr};
use promql::extension_plan::{InstantManipulate, Millisecond, SeriesDivide};
use promql_parser::parser::{
AtModifier, Call, Expr as PromExpr, MatrixSelector, Offset, ParenExpr,
};
use snafu::{OptionExt, ResultExt, ensure};
use crate::promql::error::{
AtModifierTimestampOutOfRangeSnafu, DataFusionPlanningSnafu, Result, TimeIndexNotFoundSnafu,
};
use crate::promql::planner::PromPlanner;
use crate::query_engine::QueryEngineState;
impl PromPlanner {
/// Resolve the `@` modifier of a vector or matrix selector into the timestamp its sample
/// window is anchored at, in milliseconds since the Unix epoch.
///
/// Prometheus semantics:
/// - `@ <unix_ts>` anchors at the given timestamp,
/// - `@ start()` / `@ end()` anchor at the evaluation range of the whole statement,
/// - `offset` shifts the anchor backwards: the window ends at `anchor - offset`.
///
/// Returns `None` when the selector has no `@` modifier.
fn at_ref_time(
&self,
at: &Option<AtModifier>,
offset: &Option<Offset>,
) -> Result<Option<Millisecond>> {
let anchor = match at {
None => return Ok(None),
Some(AtModifier::Start) => self.ctx.stmt_start,
Some(AtModifier::End) => self.ctx.stmt_end,
Some(AtModifier::At(time)) => Self::system_time_to_millis(time)?,
};
Ok(Some(Self::anchor_sub(anchor, Self::offset_millis(offset))?))
}
/// Subtracts `rhs` from `lhs` on the millisecond timeline of an `@` anchor.
///
/// A negative result is valid: `@` and `offset` accept timestamps before the Unix epoch. A
/// result outside the representable millisecond range is rejected like an unrepresentable
/// anchor ([`Self::system_time_to_millis`]) instead of clamping it, so the same class of
/// input always gets the same answer.
pub(crate) fn anchor_sub(lhs: Millisecond, rhs: Millisecond) -> Result<Millisecond> {
lhs.checked_sub(rhs)
.with_context(|| AtModifierTimestampOutOfRangeSnafu {
timestamp: format!("{}ms - {}ms", lhs, rhs),
})
}
/// The offset a selector with an `@` modifier is evaluated with.
///
/// Prometheus anchors such a selector by rewriting its offset to `eval_time - anchor`
/// (`setOffsetForAtModifier`), so that the selector always selects its samples around `anchor`
/// regardless of the step being evaluated. `eval_time` is the start of the evaluation the
/// selector belongs to, which is `ctx.start`.
///
/// Returns `None` when the selector has no `@` modifier.
pub(crate) fn at_modifier_offset(
&self,
at: &Option<AtModifier>,
offset: &Option<Offset>,
) -> Result<Option<Millisecond>> {
let Some(anchor) = self.at_ref_time(at, offset)? else {
return Ok(None);
};
Ok(Some(Self::anchor_sub(self.ctx.start, anchor)?))
}
/// Whether `expr` is a call that has to be evaluated once for the whole grid, because it folds
/// a range selector anchored by `@` — the range argument of the call's parser signature, which
/// is a [`MatrixSelector`] here; only a call with such an argument can take this path, so no
/// function-name registry is involved.
///
/// This is the shape that needs Prometheus' `StepInvariantExpr` the most: a range function such
/// as `rate` derives its result from the evaluation instant it is called at, so folding the
/// window once per step would let the outer evaluation grid change the result of a window that
/// `@` fixed. Evaluated once, at the start of the grid, the rewritten offset of
/// [`Self::at_modifier_offset`] places the anchor at that instant, and [`Self::replay_over_grid`]
/// reports the result at every step.
///
/// Unlike Prometheus, which wraps the whole step-invariant subtree (`preprocessExprHelper`),
/// only the call itself is promoted here. The operators above it are not: they are still planned
/// at every step over the replayed result, which is safe for the row-wise ones and keeps the
/// promotion root narrow. The promotion root has to stay a call over one range selector, because
/// the replay needs one series per batch ([`Self::series_divide_plan`]) and only such a call
/// guarantees that the rows it emits still describe the series it was divided by. The operators
/// left out — an aggregation, a join, or a label rewriting call such as `label_join` — mix or
/// re-label the rows of different series, so they are unsafe as promotion roots even though
/// evaluating them after the replay is fine. A call whose input is an anchored *instant*
/// selector (`abs(some_metric @ 300)`) needs no promotion either: the selector anchors and
/// replays its sample per series on its own, and the call above it is row-wise.
///
/// `predict_linear` is the exception among the range functions: it predicts from the evaluation
/// instant of each step ([`Self::create_range_eval_ts_expr`]), so it has to stay outside the
/// promoted subtree and follow the grid. The remaining arguments of the call have to be
/// literals, since the replay of the promoted result has no second vector input to divide.
/// Parentheses around the range argument are transparent (`rate((m[5m] @ 300))`), so they are
/// looked through and the call is promoted as if they were absent. Nothing else of the subtree
/// is unwrapped, so an outer parenthesis promotes no operator above the call.
fn promotes_anchored_range_call(expr: &PromExpr) -> bool {
let PromExpr::Call(Call { func, args }) = expr else {
return false;
};
// See the doc comment: the regression of `predict_linear` follows the evaluation step.
if func.name == "predict_linear" {
return false;
}
let mut anchored_range = false;
for arg in &args.args {
// Parentheses around the range argument are transparent, so the call is promoted the
// same way for `rate((m[5m] @ 300))` as for `rate(m[5m] @ 300)`. Only the parentheses
// directly around this one argument are looked through here: the promotion stays
// confined to a call over one anchored range selector instead of descending into an
// arbitrary parenthesized subtree.
let mut arg = arg.as_ref();
while let PromExpr::Paren(ParenExpr { expr }) = arg {
arg = expr;
}
match arg {
// The window is pinned by `@`, so every step folds the same samples.
PromExpr::MatrixSelector(MatrixSelector { vs, .. }) if vs.at.is_some() => {
if anchored_range {
return false;
}
anchored_range = true;
}
// A literal argument is the same value at every step.
arg if Self::try_build_literal_expr(arg).is_some() => {}
_ => return false,
}
}
anchored_range
}
/// Plans the anchored range call `prom_expr` ([`Self::promotes_anchored_range_call`]) as a
/// step-invariant subtree: the call is evaluated on a single evaluation instant (`grid_start`,
/// the start of the outer evaluation) and its result is then reported at every step of
/// `[grid_start, ctx.end]` by [`Self::replay_over_grid`].
///
/// This is the planner's counterpart of Prometheus' `StepInvariantExpr` for the one shape it
/// promotes. Evaluating the call once matters for the functions that derive their result from
/// the step being evaluated: `rate(m[5m] @ 300)` folds its window around the anchor once, and
/// the extrapolation boundaries of `rate` must be derived from that same window at every step
/// instead of following the outer evaluation timestamp.
///
/// Only the call itself is promoted; the operators above it are planned as usual over the
/// replayed result. The result of the promoted call is split into one series per batch before it
/// is replayed ([`Self::series_divide_plan`]), because the row-wise projection of the call does
/// not preserve the batch layout of the selector.
///
/// Returns `None` when `prom_expr` is not such a call, so that the caller plans it as usual.
/// The selector inside the promoted call keeps its own `@` anchoring (see
/// [`Self::at_modifier_offset`]), and planning it with `ctx.end == ctx.start` folds its window
/// once for that single instant instead of expanding it over the grid, which the replay of the
/// call result above already does.
pub(crate) async fn promote_anchored_range_call(
&mut self,
prom_expr: &PromExpr,
timestamp_fn: bool,
query_engine_state: &QueryEngineState,
) -> Result<Option<LogicalPlan>> {
let grid_start = self.ctx.start;
let grid_end = self.ctx.end;
// An instant query evaluates a single step, so there is nothing to promote.
if grid_start == grid_end || !Self::promotes_anchored_range_call(prom_expr) {
return Ok(None);
}
// Plan the subtree on a single evaluation instant: every selector below still anchors its
// window through `@`, and the functions above them derive their result from that one
// instant. The planner is single-use, so `ctx.end` needs no restore-on-error.
self.ctx.end = grid_start;
let anchored = self
.prom_expr_to_plan_inner(prom_expr, timestamp_fn, query_engine_state)
.await?;
self.ctx.end = grid_end;
let time_index_column =
self.ctx
.time_index_column
.clone()
.with_context(|| TimeIndexNotFoundSnafu {
table: self.ctx.table_name.clone().unwrap_or_default(),
})?;
// The replay reads one series per batch; see [`Self::series_divide_plan`] for why the
// layout of the selector below the promoted call does not survive it.
let anchored = self.series_divide_plan(anchored, &time_index_column)?;
Ok(Some(self.replay_over_grid(
anchored,
grid_start,
grid_end,
time_index_column,
)))
}
/// Convert the timestamp of an `@` modifier into milliseconds since the Unix epoch.
fn system_time_to_millis(time: &SystemTime) -> Result<Millisecond> {
let (millis, negative) = match time.duration_since(UNIX_EPOCH) {
Ok(duration) => (duration.as_millis(), false),
// The `@` modifier accepts timestamps before the Unix epoch, e.g. `@ -1`.
Err(err) => (err.duration().as_millis(), true),
};
ensure!(
millis <= i64::MAX as u128,
AtModifierTimestampOutOfRangeSnafu {
timestamp: time
.duration_since(UNIX_EPOCH)
.map(|duration| format!("+{}ms", duration.as_millis()))
.unwrap_or_else(|err| format!("-{}ms", err.duration().as_millis())),
}
);
let millis = millis as Millisecond;
Ok(if negative { -millis } else { millis })
}
/// Report the samples of `anchored` at every step of the evaluation grid
/// `[grid_start, grid_end]`.
///
/// A selector with an `@` modifier is anchored: the sample window is selected once, around the
/// anchor timestamp, and every evaluation step reports that same window. Prometheus does this
/// by rewriting the selector's offset to `eval_time - anchor` and only fetching the samples on
/// the first step (`setOffsetForAtModifier` plus the `refetch` shortcut in `rangeEval`).
///
/// The expansion reuses [`InstantManipulate`] with a lookback that spans the whole grid: every
/// step then picks the same sample (or the same already computed value, when `anchored` ends
/// with a function call such as `rate`) and stamps it with the step's timestamp.
///
/// Every input batch of `anchored` must hold exactly one series, because [`InstantManipulate`]
/// takes a batch as one timeline. A leaf-level replay (`m @ 300`) consumes the [`SeriesDivide`]
/// of its selector directly. A promoted call is guaranteed that layout by
/// [`Self::series_divide_plan`], which is why [`Self::promote_anchored_range_call`] splits
/// its result before calling this method.
pub(crate) fn replay_over_grid(
&self,
anchored: LogicalPlan,
grid_start: Millisecond,
grid_end: Millisecond,
time_index_column: String,
) -> LogicalPlan {
if grid_start == grid_end {
// A single evaluation step: `anchored` is already stamped with that timestamp.
return anchored;
}
let series_key_columns = self.series_key_columns_for_schema(anchored.schema());
let replayed = InstantManipulate::new(
grid_start,
grid_end,
// The lookback must keep the single anchored sample eligible for every step.
grid_end - grid_start + 1,
self.ctx.interval,
0,
time_index_column,
series_key_columns,
self.ctx.field_columns.first().cloned(),
anchored,
);
LogicalPlan::Extension(Extension {
node: Arc::new(replayed),
})
}
/// Sorts `input` by its series key and time index and splits it into one batch per series.
///
/// [`InstantManipulate`] reads every input batch as one series (it takes the timeline of the
/// batch and reports the row selected at every step), so a batch holding several series would
/// lose all but one of them. A selector establishes that layout with its own [`SeriesDivide`],
/// but the per-series distribution requirement does not reach a promoted subtree above it:
/// the row-wise projection of a call sits in between, so the batch boundaries of the selector
/// are not preserved — in a distributed plan the promoted result can be delivered as one batch
/// holding every series. Sorting and dividing here restores the layout, exactly like
/// [`Self::prom_matrix_selector_to_plan`] does for the input of a range function.
///
/// Series keys that are not present in `input` are dropped, since `ctx.tag_columns` may have
/// drifted from the actual output schema. A plan without any series key column is returned
/// unchanged: there is nothing to divide by.
fn series_divide_plan(
&self,
input: LogicalPlan,
time_index_column: &str,
) -> Result<LogicalPlan> {
let series_key_columns = self.series_key_columns_for_schema(input.schema());
if series_key_columns.is_empty() {
return Ok(input);
}
let mut sort_exprs = series_key_columns
.iter()
.map(|name| DfExpr::Column(Column::from_name(name)).sort(true, true))
.collect::<Vec<_>>();
sort_exprs.push(DfExpr::Column(Column::from_name(time_index_column)).sort(true, true));
let sort_plan = LogicalPlanBuilder::from(input)
.sort(sort_exprs)
.context(DataFusionPlanningSnafu)?
.build()
.context(DataFusionPlanningSnafu)?;
Ok(LogicalPlan::Extension(Extension {
node: Arc::new(SeriesDivide::new(
series_key_columns,
time_index_column.to_string(),
sort_plan,
)),
}))
}
}
+670 -5
View File
@@ -1929,6 +1929,600 @@ async fn tsid_is_used_for_series_divide_when_available() {
assert!(format!("{exec:?}").contains("reuse_tsid_column: true"));
}
async fn build_at_modifier_plan(query: &str, start_secs: u64, end_secs: u64) -> LogicalPlan {
let eval_stmt = build_at_modifier_eval_stmt(query, start_secs, end_secs);
let table_provider = build_test_table_provider(
&[(DEFAULT_SCHEMA_NAME.to_string(), "some_metric".to_string())],
1,
1,
)
.await;
PromPlanner::stmt_to_plan(table_provider, &eval_stmt, &build_query_engine_state())
.await
.unwrap()
}
fn build_at_modifier_eval_stmt(query: &str, start_secs: u64, end_secs: u64) -> EvalStmt {
EvalStmt {
expr: parser::parse(query).unwrap(),
start: UNIX_EPOCH
.checked_add(Duration::from_secs(start_secs))
.unwrap(),
end: UNIX_EPOCH
.checked_add(Duration::from_secs(end_secs))
.unwrap(),
interval: Duration::from_secs(5),
lookback_delta: Duration::from_secs(1),
}
}
/// Every selector of an `@` anchored query must scan around the anchor only, instead of the
/// whole evaluation range.
#[tokio::test]
async fn at_modifier_anchors_selector_scan_window() {
// `@ 100` anchors at t=100s; the lookback delta is 1s, so the scan covers (99s, 100s].
let plan = build_at_modifier_plan("some_metric @ 100", 0, 1000).await;
let plan_str = plan.display_indent_schema().to_string();
assert!(
plan_str.contains(
"some_metric.timestamp >= TimestampMillisecond(99001, None) AND some_metric.timestamp <= TimestampMillisecond(100000, None)"
),
"{plan_str}"
);
// The result is reported at the evaluation timestamps, not at the anchor.
assert!(plan_str.contains("range=[0..1000000]"), "{plan_str}");
// `@ start()` / `@ end()` resolve to the evaluation range of the statement.
let plan = build_at_modifier_plan("some_metric @ start()", 0, 1000).await;
let plan_str = plan.display_indent_schema().to_string();
assert!(
plan_str.contains(
"some_metric.timestamp >= TimestampMillisecond(-999, None) AND some_metric.timestamp <= TimestampMillisecond(0, None)"
),
"{plan_str}"
);
let plan = build_at_modifier_plan("some_metric @ end()", 0, 1000).await;
let plan_str = plan.display_indent_schema().to_string();
assert!(
plan_str.contains(
"some_metric.timestamp >= TimestampMillisecond(999001, None) AND some_metric.timestamp <= TimestampMillisecond(1000000, None)"
),
"{plan_str}"
);
// `offset` moves the anchor backwards and is not applied twice.
let plan = build_at_modifier_plan("some_metric @ 200 offset 50s", 0, 1000).await;
let plan_str = plan.display_indent_schema().to_string();
assert!(
plan_str.contains(
"some_metric.timestamp >= TimestampMillisecond(149001, None) AND some_metric.timestamp <= TimestampMillisecond(150000, None)"
),
"{plan_str}"
);
// A timestamp before the Unix epoch is accepted, as in Prometheus.
let plan = build_at_modifier_plan("some_metric @ -1", 0, 1000).await;
let plan_str = plan.display_indent_schema().to_string();
assert!(
plan_str.contains(
"some_metric.timestamp >= TimestampMillisecond(-1999, None) AND some_metric.timestamp <= TimestampMillisecond(-1000, None)"
),
"{plan_str}"
);
}
/// A call over one range selector anchored by `@` is evaluated once, at the start of the
/// evaluation, and its result is reported at every step: the window is folded around the anchor
/// instead of following the outer evaluation grid. This is the planner's counterpart of
/// Prometheus' `StepInvariantExpr` wrapper; see [`PromPlanner::promotes_anchored_range_call`].
#[tokio::test]
async fn at_modifier_promotes_anchored_range_call() {
let plan = build_at_modifier_plan("rate(some_metric[5m] @ 300)", 0, 1000).await;
let plan_str = plan.display_indent_schema().to_string();
// The scan is limited to the anchored window (offset by `eval_start - anchor`).
assert!(
plan_str.contains(
"some_metric.timestamp >= TimestampMillisecond(1, None) AND some_metric.timestamp <= TimestampMillisecond(300000, None)"
),
"{plan_str}"
);
// A single fold, at the anchor...
assert_eq!(
plan_str
.matches("PromRangeManipulate: req range=[0..0]")
.count(),
1,
"{plan_str}"
);
// ... never a fold per step of the outer grid.
assert!(
!plan_str.contains("PromRangeManipulate: req range=[0..1000000]"),
"{plan_str}"
);
// A single replay of the function result over the whole grid...
assert_eq!(
plan_str.matches("PromInstantManipulate").count(),
1,
"{plan_str}"
);
assert!(
plan_str.contains("PromInstantManipulate: range=[0..1000000], lookback=[1000001]"),
"{plan_str}"
);
// ... so `rate` itself is evaluated below that replay node, on the single evaluation
// instant of the anchored subtree, instead of once per step.
let replay = plan_str
.find("PromInstantManipulate")
.expect("instant manipulate node");
let rate = plan_str.find("prom_rate(").expect("rate projection");
assert!(
replay < rate,
"`rate` must be evaluated below the replay node:\n{plan_str}"
);
}
/// Parentheses around the range argument are transparent: `rate((some_metric[5m] @ 300))` gets
/// the same fixed-window promotion as `rate(some_metric[5m] @ 300)`, with the anchored window
/// folded once and its result replayed at every step of the grid. Only that one argument is
/// looked through, so a parenthesis above the call promotes no operator of its own and a
/// parenthesized subtree is planned exactly like the bare one.
#[tokio::test]
async fn at_modifier_promotes_parenthesized_range_argument() {
// Each form is planned exactly like its unparenthesized counterpart: parentheses below the
// call are transparent, a parenthesis around the call adds nothing, and a parenthesis above
// it does not widen the promotion.
for (query, plain) in [
(
"rate((some_metric[5m] @ 300))",
"rate(some_metric[5m] @ 300)",
),
(
"rate(((some_metric[5m] @ 300)))",
"rate(some_metric[5m] @ 300)",
),
(
"(rate(some_metric[5m] @ 300))",
"rate(some_metric[5m] @ 300)",
),
(
"abs((rate(some_metric[5m] @ 300)))",
"abs(rate(some_metric[5m] @ 300))",
),
] {
assert_eq!(
build_at_modifier_plan(query, 0, 1000)
.await
.display_indent_schema()
.to_string(),
build_at_modifier_plan(plain, 0, 1000)
.await
.display_indent_schema()
.to_string(),
"`{query}` must be planned like `{plain}`"
);
}
// The parentheses do not push the enclosing operator into the promotion either: `abs` stays
// above the replay of the promoted call, exactly as it does without them. The promotion
// itself (one anchored fold, one replay, `rate` below it) is asserted for the
// unparenthesized form by `at_modifier_promotes_anchored_range_call`, and the form above is
// planned identically to it.
let plan_str = build_at_modifier_plan("abs((rate(some_metric[5m] @ 300)))", 0, 1000)
.await
.display_indent_schema()
.to_string();
let replay = plan_str
.find("PromInstantManipulate")
.expect("instant manipulate node");
assert!(
plan_str.find("abs(").expect("`abs` projection") < replay,
"`abs` must be evaluated above the replay of the promoted call:\n{plan_str}"
);
}
/// `@ start()` and `@ end()` are fixed anchors for the whole statement, so a call using them is
/// promoted as well.
#[tokio::test]
async fn at_modifier_promotes_start_and_end_anchored_call() {
for query in [
"rate(some_metric[5m] @ start())",
"rate(some_metric[5m] @ end())",
"max_over_time(some_metric[5m] @ end())",
] {
let plan = build_at_modifier_plan(query, 0, 1000).await;
let plan_str = plan.display_indent_schema().to_string();
assert_eq!(
plan_str
.matches("PromRangeManipulate: req range=[0..0]")
.count(),
1,
"{query}:\n{plan_str}"
);
assert!(
plan_str.contains("PromInstantManipulate: range=[0..1000000], lookback=[1000001]"),
"{query}:\n{plan_str}"
);
}
// Only the call itself is promoted: a call above it (`abs`) is planned as usual and
// evaluated at every step over the replayed result of the promoted call. The window is
// still folded once, around the anchor.
let plan = build_at_modifier_plan("abs(max_over_time(some_metric[5m] @ end()))", 0, 1000).await;
let plan_str = plan.display_indent_schema().to_string();
assert_eq!(
plan_str
.matches("PromRangeManipulate: req range=[0..0]")
.count(),
1,
"{plan_str}"
);
let abs = plan_str.find("abs(").expect("`abs` projection");
let replay = plan_str
.find("PromInstantManipulate: range=[0..1000000], lookback=[1000001]")
.expect("replay node");
assert!(
abs < replay,
"`abs` must be evaluated above the replay of the promoted call:\n{plan_str}"
);
// The window is the only one folded once, and it sits below the replay: the promoted call
// feeds the grid from there.
let fold = plan_str.find("PromRangeManipulate").expect("range fold");
assert!(
replay < fold,
"the anchored window must be folded below the replay:\n{plan_str}"
);
// An aggregation above the promoted call stays above it as well: `sum` aggregates the
// replayed per-series rows at every step, instead of aggregating the single anchored
// instant and replaying the aggregation — which would also have to replay rows of several
// groups through the one-series-per-batch `InstantManipulate`.
let plan = build_at_modifier_plan("sum(rate(some_metric[5m] @ start()))", 0, 1000).await;
let plan_str = plan.display_indent_schema().to_string();
assert_eq!(
plan_str
.matches("PromRangeManipulate: req range=[0..0]")
.count(),
1,
"{plan_str}"
);
// A single replay, and it sits below the aggregation node: `sum` aggregates the replayed
// per-series rows at every step.
assert_eq!(
plan_str.matches("PromInstantManipulate").count(),
1,
"{plan_str}"
);
let aggregate = plan_str.find("Aggregate:").expect("aggregate node");
let replay = plan_str
.find("PromInstantManipulate: range=[0..1000000], lookback=[1000001]")
.expect("replay node");
assert!(
aggregate < replay,
"the aggregation must stay above the replay of the promoted call:\n{plan_str}"
);
// Without `@` nothing is promoted: the function keeps folding one window per step.
let plan = build_at_modifier_plan("rate(some_metric[5m])", 0, 1000).await;
let plan_str = plan.display_indent_schema().to_string();
assert!(
plan_str.contains("PromRangeManipulate: req range=[0..1000000]"),
"{plan_str}"
);
assert!(!plan_str.contains("lookback=[1000001]"), "{plan_str}");
}
/// A call or a unary operator above the anchored range call is planned as usual: the inner range
/// call is promoted on its own (it is the direct call over the anchored range selector), and the
/// operator above it is evaluated at every step over the replayed result.
#[tokio::test]
async fn at_modifier_promotes_inner_range_call_below_wrappers() {
for (query, wrapper) in [
("abs(rate(some_metric[5m] @ 300))", "abs(prom_rate("),
("-rate(some_metric[5m] @ 300)", "(- prom_rate("),
(
"abs(max_over_time(some_metric[5m] @ 300))",
"abs(prom_max_over_time(",
),
] {
let plan = build_at_modifier_plan(query, 0, 1000).await;
let plan_str = plan.display_indent_schema().to_string();
// The anchored window is folded once, and the wrapper sits above its replay.
assert_eq!(
plan_str
.matches("PromRangeManipulate: req range=[0..0]")
.count(),
1,
"{query}:\n{plan_str}"
);
assert_eq!(
plan_str.matches("PromInstantManipulate").count(),
1,
"{query}:\n{plan_str}"
);
let wrapper = plan_str
.find(wrapper)
.unwrap_or_else(|| panic!("no `{wrapper}` projection in:\n{plan_str}"));
let replay = plan_str
.find("PromInstantManipulate: range=[0..1000000], lookback=[1000001]")
.expect("replay node");
assert!(
wrapper < replay,
"the wrapper must be evaluated above the replay of the range call:\n{plan_str}"
);
}
}
/// `anchored + plain`: only the binary operand that is a call over the anchored range selector
/// is promoted, and the plain side keeps following the evaluation step.
#[tokio::test]
async fn at_modifier_promotes_only_anchored_binary_operand() {
let plan = build_at_modifier_plan("rate(some_metric[5m] @ 300) + some_metric", 0, 1000).await;
let plan_str = plan.display_indent_schema().to_string();
// The anchored operand is folded once and its result replayed over the whole grid.
assert_eq!(
plan_str
.matches("PromRangeManipulate: req range=[0..0]")
.count(),
1,
"{plan_str}"
);
assert!(
plan_str.contains("PromInstantManipulate: range=[0..1000000], lookback=[1000001]"),
"{plan_str}"
);
// The plain operand still selects one sample per step with the default lookback.
assert!(
plan_str.contains("PromInstantManipulate: range=[0..1000000], lookback=[1000]"),
"{plan_str}"
);
assert_eq!(
plan_str.matches("PromInstantManipulate").count(),
2,
"{plan_str}"
);
// Both operands anchored: the binary expression is not promoted as a whole, because a join
// emits the rows of several series in shared batches and the replay needs one series per
// batch. Each operand anchors and replays on its own instead, and the join runs at every
// step over those per-series results.
let plan = build_at_modifier_plan("some_metric @ 300 + some_metric @ 0", 0, 1000).await;
let plan_str = plan.display_indent_schema().to_string();
// One anchoring node and one replay per operand...
assert_eq!(
plan_str.matches("PromInstantManipulate").count(),
4,
"{plan_str}"
);
assert_eq!(
plan_str
.matches("PromInstantManipulate: range=[0..1000000], lookback=[1000001]")
.count(),
2,
"{plan_str}"
);
// ... and the two operands are anchored at different timestamps, so both select their own
// sample.
assert_eq!(
plan_str
.matches("PromInstantManipulate: range=[0..0], lookback=[1000]")
.count(),
2,
"{plan_str}"
);
}
/// Only a direct call over an anchored range selector is promoted. A value function or a unary
/// operator over an anchored *instant* selector needs no promotion: the selector anchors and
/// replays its sample per series on its own, and the operator above it is row-wise, so it can be
/// evaluated at every step over that replay.
#[tokio::test]
async fn at_modifier_does_not_promote_value_calls_over_anchored_selectors() {
for (query, value_expr) in [
("abs(some_metric @ 300)", "abs(some_metric.field_0)"),
("-some_metric @ 300", "(- some_metric.field_0)"),
] {
let plan = build_at_modifier_plan(query, 0, 1000).await;
let plan_str = plan.display_indent_schema().to_string();
// The scan is limited to the anchored sample (the lookback delta of this test is 1s).
assert!(
plan_str.contains(
"some_metric.timestamp >= TimestampMillisecond(299001, None) AND some_metric.timestamp <= TimestampMillisecond(300000, None)"
),
"{query}:\n{plan_str}"
);
// ... and the result is replayed at every step by the selector itself: the anchored
// selection, then the grid replay, with no promoted subtree on top of the operator.
assert!(
!plan_str.starts_with("PromInstantManipulate"),
"{query}:\n{plan_str}"
);
assert_eq!(
plan_str
.matches("PromInstantManipulate: range=[0..0], lookback=[1000]")
.count(),
1,
"{query}:\n{plan_str}"
);
let replay = plan_str
.find("PromInstantManipulate: range=[0..1000000], lookback=[1000001]")
.unwrap_or_else(|| panic!("replay node:\n{plan_str}"));
let value = plan_str
.find(value_expr)
.unwrap_or_else(|| panic!("no `{value_expr}` projection in:\n{plan_str}"));
assert!(
value < replay,
"`{value_expr}` must be evaluated above the per-series replay:\n{plan_str}"
);
}
}
/// A call whose argument merges series (an aggregation or a join) or whose own output reorders
/// the whole vector (`sort*`, the histogram folds) is never the promoted root either: it is
/// planned as usual over the leaf-level anchoring of its selectors, which replays every selector
/// per series and keeps the one-series-per-batch layout the replay needs.
#[tokio::test]
async fn at_modifier_keeps_multi_series_roots_out_of_promoted_subtree() {
// An aggregation below the call emits one row per group in shared batches, so the call is
// not promoted: the replay of the anchored selector stays below the aggregate.
let plan = build_at_modifier_plan("abs(sum(some_metric @ 300))", 0, 1000).await;
let plan_str = plan.display_indent_schema().to_string();
let aggregate = plan_str.find("Aggregate:").expect("aggregate node");
let replay = plan_str
.find("PromInstantManipulate: range=[0..1000000], lookback=[1000001]")
.expect("replay node");
assert!(
aggregate < replay,
"the aggregation must stay above the replay:\n{plan_str}"
);
// A join of two anchored selectors below the call: neither the join nor the call is
// promoted, so each operand is replayed on its own.
let plan = build_at_modifier_plan("abs(some_metric @ 300 + some_metric @ 0)", 0, 1000).await;
let plan_str = plan.display_indent_schema().to_string();
assert!(!plan_str.starts_with("PromInstantManipulate"), "{plan_str}");
assert_eq!(
plan_str
.matches("PromInstantManipulate: range=[0..1000000], lookback=[1000001]")
.count(),
2,
"{plan_str}"
);
// A call that reorders the whole vector keeps its sort above the per-series replay.
let plan = build_at_modifier_plan("sort_by_label(some_metric @ 300, \"tag_0\")", 0, 1000).await;
let plan_str = plan.display_indent_schema().to_string();
assert!(plan_str.starts_with("Sort:"), "{plan_str}");
assert_eq!(
plan_str
.matches("PromInstantManipulate: range=[0..1000000], lookback=[1000001]")
.count(),
1,
"{plan_str}"
);
}
/// `label_join` rewrites the labels of its input series, so it is never the promoted root: the
/// anchored selector keeps replaying one series per batch, and the join runs at every step above
/// that replay (see [`Self::promotes_anchored_range_call`]). Promoting it would replay the
/// joined rows through the labels the join just rewrote, which merges the distinct input series
/// into one timeline.
#[tokio::test]
async fn at_modifier_does_not_promote_label_join() {
for query in [
// Directly above the anchored instant selector...
"label_join(some_metric @ 300, \"tag_0\", \"-\", \"\")",
// ... and below another call, which is planned as usual over the join.
"abs(label_join(some_metric @ 300, \"tag_0\", \"-\", \"\"))",
] {
let plan = build_at_modifier_plan(query, 0, 1000).await;
let plan_str = plan.display_indent_schema().to_string();
// The join is not wrapped in a replay of its own; the only replay over the grid is the
// one of the anchored selector...
assert!(
!plan_str.starts_with("PromInstantManipulate"),
"{query}:\n{plan_str}"
);
assert_eq!(
plan_str
.matches("PromInstantManipulate: range=[0..1000000], lookback=[1000001]")
.count(),
1,
"{query}:\n{plan_str}"
);
// ... and the projected join stays above it, evaluated at every step.
let join = plan_str
.find("concat_ws(")
.unwrap_or_else(|| panic!("no `label_join` projection in:\n{plan_str}"));
let replay = plan_str
.find("PromInstantManipulate: range=[0..1000000], lookback=[1000001]")
.expect("replay node");
assert!(
join < replay,
"`label_join` must be evaluated above the per-series replay:\n{plan_str}"
);
}
// A range call below the join is still promoted on its own: the anchored window is folded
// once per series, and the join above it is evaluated at every step over that replay.
let plan = build_at_modifier_plan(
"label_join(rate(some_metric[5m] @ 300), \"tag_0\", \"-\", \"\")",
0,
1000,
)
.await;
let plan_str = plan.display_indent_schema().to_string();
assert_eq!(
plan_str
.matches("PromRangeManipulate: req range=[0..0]")
.count(),
1,
"{plan_str}"
);
assert!(!plan_str.starts_with("PromInstantManipulate"), "{plan_str}");
let join = plan_str.find("concat_ws(").expect("join projection");
let replay = plan_str
.find("PromInstantManipulate: range=[0..1000000], lookback=[1000001]")
.expect("replay node");
assert!(
join < replay,
"the join must be evaluated above the replay of the range call:\n{plan_str}"
);
}
#[test]
fn at_modifier_rejects_subtraction_overflow() {
for (anchor, offset) in [(i64::MAX, -1), (i64::MIN, 1)] {
let err = PromPlanner::anchor_sub(anchor, offset).unwrap_err();
assert_eq!(err.status_code(), StatusCode::InvalidArguments);
assert!(
err.to_string()
.contains("Timestamp out of range for the `@` modifier"),
"{err}"
);
}
assert_eq!(PromPlanner::anchor_sub(-1, 1).unwrap(), -2);
}
/// `@` beyond the representable millisecond range is rejected instead of silently wrapping.
///
/// `@ 1e16` is 10^19 milliseconds, beyond `i64::MAX`. A Unix `SystemTime` can hold it, so
/// the planner rejects the anchor it cannot represent. A Windows `SystemTime` tops out
/// below `i64::MAX` milliseconds, so the same literal is already rejected while parsing.
#[tokio::test]
async fn at_modifier_rejects_unrepresentable_timestamp() {
#[cfg(windows)]
{
let err = parser::parse("some_metric @ 1e16").unwrap_err();
assert!(
err.to_string()
.contains("timestamp out of bounds for @ modifier"),
"{err}"
);
}
#[cfg(not(windows))]
{
let eval_stmt = build_eval_stmt("some_metric @ 1e16");
let table_provider = build_test_table_provider(
&[(DEFAULT_SCHEMA_NAME.to_string(), "some_metric".to_string())],
1,
1,
)
.await;
let err =
PromPlanner::stmt_to_plan(table_provider, &eval_stmt, &build_query_engine_state())
.await
.unwrap_err();
assert!(
err.to_string()
.contains("Timestamp out of range for the `@` modifier"),
"{err}"
);
assert_eq!(err.status_code(), StatusCode::InvalidArguments);
}
}
#[tokio::test]
async fn default_binary_join_uses_tsid_when_available() {
let eval_stmt = build_eval_stmt("some_metric / some_alt_metric");
@@ -3421,7 +4015,7 @@ async fn binary_op_column_column() {
assert_eq!(plan.display_indent_schema().to_string(), expected);
}
async fn indie_query_plan_compare<T: AsRef<str>>(query: &str, expected: T) {
async fn indie_query_plan(query: &str) -> String {
let prom_expr = parser::parse(query).unwrap();
let eval_stmt = EvalStmt {
expr: prom_expr,
@@ -3449,7 +4043,12 @@ async fn indie_query_plan_compare<T: AsRef<str>>(query: &str, expected: T) {
.await
.unwrap();
assert_eq!(plan.display_indent_schema().to_string(), expected.as_ref());
plan.display_indent_schema().to_string()
}
async fn indie_query_plan_compare<T: AsRef<str>>(query: &str, expected: T) {
let plan = indie_query_plan(query).await;
assert_eq!(plan, expected.as_ref());
}
#[tokio::test]
@@ -3537,6 +4136,51 @@ async fn increase_aggr() {
indie_query_plan_compare(query, expected).await;
}
#[tokio::test]
async fn predict_linear_injects_the_eval_timestamp() {
// The fourth argument is the evaluation instant of each row, which the fold offset
// recovers from the row's time index (here: no `@` and no `offset`, so the step itself).
let query = "predict_linear(some_metric[5m], 60)";
let expected = String::from(
"Filter: prom_predict_linear(timestamp_range,field_0,Float64(60)) IS NOT NULL [timestamp:Timestamp(ms), prom_predict_linear(timestamp_range,field_0,Float64(60)):Float64;N, tag_0:Utf8]\
\n Projection: some_metric.timestamp, prom_predict_linear(timestamp_range, field_0, CAST(Float64(60) AS Int64), CAST(CAST(some_metric.timestamp AS Int64) + Int64(0) AS Timestamp(ms))) AS prom_predict_linear(timestamp_range,field_0,Float64(60)), some_metric.tag_0 [timestamp:Timestamp(ms), prom_predict_linear(timestamp_range,field_0,Float64(60)):Float64;N, tag_0:Utf8]\
\n PromRangeManipulate: req range=[0..100000000], interval=[5000], eval range=[300000], time index=[timestamp], values=[\"field_0\"] [tag_0:Utf8, timestamp:Timestamp(ms), field_0:Dictionary(Int64, Float64);N, timestamp_range:Dictionary(Int64, Timestamp(ms))]\
\n PromSeriesNormalize: offset=[0], time index=[timestamp], filter NaN: [true] [tag_0:Utf8, timestamp:Timestamp(ms), field_0:Float64;N]\
\n PromSeriesDivide: tags=[\"tag_0\"] [tag_0:Utf8, timestamp:Timestamp(ms), field_0:Float64;N]\
\n Sort: some_metric.tag_0 ASC NULLS FIRST, some_metric.timestamp ASC NULLS FIRST [tag_0:Utf8, timestamp:Timestamp(ms), field_0:Float64;N]\
\n Filter: some_metric.timestamp >= TimestampMillisecond(-299999, None) AND some_metric.timestamp <= TimestampMillisecond(100000000, None) [tag_0:Utf8, timestamp:Timestamp(ms), field_0:Float64;N]\
\n TableScan: some_metric [tag_0:Utf8, timestamp:Timestamp(ms), field_0:Float64;N]",
);
indie_query_plan_compare(query, expected).await;
}
/// The evaluation instant follows the selector's fold offset, not the projection's: an
/// `@`-anchored selector is folded with `at_offset`, which is what the window's timestamps were
/// shifted by. The evaluation starts at 0s here, so `@ 100` anchors 100s in the future and
/// yields `at_offset` = -100000ms.
#[tokio::test]
async fn predict_linear_eval_ts_follows_the_fold_offset() {
for (query, expected_offset) in [
(
"predict_linear(some_metric[5m] @ 100, 60)",
"Int64(-100000)",
),
(
"predict_linear(some_metric[5m] offset 2m, 60)",
"Int64(120000)",
),
] {
let plan = indie_query_plan(query).await;
assert!(
plan.contains(&format!(
"some_metric.timestamp AS Int64) + {expected_offset}"
)),
"{query}\n{plan}"
);
}
}
async fn native_histogram_plan(query: &str) -> String {
let table_provider = build_test_native_histogram_table_provider("some_metric").await;
let plan = PromPlanner::stmt_to_plan(
@@ -4214,6 +4858,27 @@ async fn mixed_native_histogram_ranges_use_coordinated_udfs() {
assert_eq!(plan, expected);
}
#[tokio::test]
async fn mixed_native_histogram_predict_linear_forwards_the_eval_timestamp() {
let query = "predict_linear(some_metric[5m], 60)";
let plan = PromPlanner::stmt_to_plan(
build_test_mixed_native_histogram_table_provider("some_metric").await,
&build_eval_stmt(query),
&build_query_engine_state(),
)
.await
.unwrap()
.display_indent_schema()
.to_string();
assert!(
plan.contains(
"prom_mixed_range_float(Utf8(\"predict_linear\"), timestamp_range, greptime_value, greptime_native_histogram, CAST(Float64(60) AS Int64), CAST(CAST(some_metric.timestamp AS Int64) + Int64(0) AS Timestamp(ms)))"
),
"{query}\n{plan}"
);
}
#[tokio::test]
async fn mixed_native_histogram_rate_executes_real_ranges() {
let schema = Arc::new(ArrowSchema::new(vec![
@@ -4288,7 +4953,7 @@ async fn mixed_native_histogram_rate_executes_real_ranges() {
);
let state = build_query_engine_state();
let (mut exprs, _) = planner
.create_function_expr(&call.func, vec![], input.schema(), &state)
.create_function_expr(&call.func, vec![], input.schema(), &state, None)
.unwrap();
exprs.insert(0, planner.create_time_index_column_expr().unwrap());
let plan = LogicalPlanBuilder::from(input)
@@ -7124,7 +7789,7 @@ async fn test_mixed_or_routes_float_histogram_and_label_functions() {
};
let state = build_query_engine_state();
let (mut exprs, _) = planner
.create_function_expr(&call.func, vec![], input.schema(), &state)
.create_function_expr(&call.func, vec![], input.schema(), &state, None)
.unwrap();
exprs.insert(0, planner.create_time_index_column_expr().unwrap());
exprs.extend(planner.create_tag_column_exprs().unwrap());
@@ -7177,7 +7842,7 @@ async fn test_mixed_or_routes_float_histogram_and_label_functions() {
let args = planner.create_function_args(&call.args.args).unwrap();
let state = build_query_engine_state();
let (mut exprs, _) = planner
.create_function_expr(&call.func, args.literals, input.schema(), &state)
.create_function_expr(&call.func, args.literals, input.schema(), &state, None)
.unwrap();
exprs.insert(0, planner.create_time_index_column_expr().unwrap());
exprs.extend(planner.create_tag_column_exprs().unwrap());
+1 -1
View File
@@ -475,7 +475,7 @@ async fn delta_mixed_ranges_drop_and_float_ranges_sum() {
);
let state = build_query_engine_state();
let (mut exprs, _) = planner
.create_function_expr(&call.func, vec![], input.schema(), &state)
.create_function_expr(&call.func, vec![], input.schema(), &state, None)
.unwrap();
exprs.insert(0, planner.create_time_index_column_expr().unwrap());
let plan = LogicalPlanBuilder::from(input)