mirror of
https://github.com/GreptimeTeam/greptimedb.git
synced 2026-10-03 10:35:35 +00:00
fix(promql): apply offset to subquery evaluation window
`prom_subquery_expr_to_plan` destructured `SubqueryExpr` without reading `offset`, so `<subquery>[range:step] offset <d>` planned exactly the same window as the un-offset form and silently returned data for the wrong time range. The plain vector/matrix-selector paths already threaded the offset through `selector_to_series_normalize_plan` and `RangeManipulate`. Shift the inner evaluation window back by the offset and pass the offset to the subquery's `RangeManipulate`, which maps the inner samples forward onto the evaluation timeline before bucketing them into ranges. This matches Prometheus, whose `evaluator.subqueryTimeRange` evaluates the inner expression over `(start - offset - range, end - offset]` and whose `evalSubquery` then hands the samples to the outer range-vector function as a `MatrixSelector` that still carries the subquery offset. An offset on the inner selector composes additively, as `subqueryTimes` documents. `RangeManipulate`'s protobuf message has no offset field and recovers it on decode from an immediately underlying `SeriesNormalize`. Since `RangeManipulate` is commutative in `dist_plan` and can be pushed below a `MergeScan`, insert that carrier node so the offset survives distributed planning instead of decoding as zero. Known divergence, unchanged by this commit: Prometheus anchors subquery step points on absolute epoch multiples of the step, while GreptimeDB anchors them on the evaluation start. The two agree whenever the offset is a multiple of the subquery step; the added sqlness cases stay within that range. Closes #9330 Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Signed-off-by: Aarav <aaravsjadav@gmail.com>
This commit is contained in:
@@ -548,18 +548,55 @@ impl PromPlanner {
|
||||
subquery_expr: &SubqueryExpr,
|
||||
) -> Result<LogicalPlan> {
|
||||
let SubqueryExpr {
|
||||
expr, range, step, ..
|
||||
expr,
|
||||
range,
|
||||
step,
|
||||
offset,
|
||||
..
|
||||
} = subquery_expr;
|
||||
|
||||
// Prometheus shifts a subquery's own evaluation window back by its offset: the inner
|
||||
// expression is evaluated over `(start - offset - range, end - offset]`
|
||||
// (`evaluator.subqueryTimeRange` in `promql/engine.go`), producing samples that keep
|
||||
// their real, un-shifted timestamps. `evalSubquery` then hands those samples to the
|
||||
// outer range-vector function as a `MatrixSelector` whose `VectorSelector` still
|
||||
// carries `subq.Offset`, so the function slices `(ts - offset - range, ts - offset]`.
|
||||
// The offset is therefore applied once semantically, never twice.
|
||||
//
|
||||
// Here the same shift is expressed as: evaluate the inner plan over the shifted
|
||||
// window, then let `RangeManipulate` map the resulting samples forward by `offset_ms`
|
||||
// onto the evaluation timeline before bucketing them into ranges -- exactly what the
|
||||
// plain matrix-selector path above does.
|
||||
//
|
||||
// An offset on the inner selector composes additively with this one, because the inner
|
||||
// selector subtracts its own offset from the already shifted step timestamps
|
||||
// (`subqueryTimes` documents this as "the sum of offsets ... of all subqueries in the
|
||||
// path").
|
||||
//
|
||||
// Divergence worth knowing: Prometheus anchors subquery step points on absolute epoch
|
||||
// multiples of the step, so an offset that is not a multiple of the step does not
|
||||
// rotate the grid. GreptimeDB anchors the grid on the evaluation start instead (see
|
||||
// the `ctx.start` arithmetic below, which predates this offset handling), so a
|
||||
// sub-step offset does rotate it. The two agree whenever the offset is a multiple of
|
||||
// the subquery step.
|
||||
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 current_interval = self.ctx.interval;
|
||||
if let Some(step) = step {
|
||||
self.ctx.interval = step.as_millis() as _;
|
||||
}
|
||||
let current_start = self.ctx.start;
|
||||
self.ctx.start -= range.as_millis() as i64 - self.ctx.interval;
|
||||
let current_end = self.ctx.end;
|
||||
self.ctx.start -= offset_ms + range.as_millis() as i64 - self.ctx.interval;
|
||||
self.ctx.end -= offset_ms;
|
||||
let input = self.prom_expr_to_plan(expr, query_engine_state).await?;
|
||||
self.ctx.interval = current_interval;
|
||||
self.ctx.start = current_start;
|
||||
self.ctx.end = current_end;
|
||||
|
||||
ensure!(!range.is_zero(), ZeroRangeSelectorSnafu);
|
||||
let range_ms = range.as_millis() as _;
|
||||
@@ -616,17 +653,38 @@ impl PromPlanner {
|
||||
.context(DataFusionPlanningSnafu)?;
|
||||
let divide_plan = LogicalPlan::Extension(Extension {
|
||||
node: Arc::new(SeriesDivide::new(
|
||||
series_key_columns,
|
||||
series_key_columns.clone(),
|
||||
time_index_column.clone(),
|
||||
sort_plan,
|
||||
)),
|
||||
});
|
||||
|
||||
// `RangeManipulate`'s protobuf message has no offset field: on the decode path it
|
||||
// recovers the offset from an immediately underlying `SeriesNormalize` (`local_offset`).
|
||||
// `RangeManipulate` is `Commutative` in `dist_plan`, so a subquery's node can be pushed
|
||||
// below a `MergeScan` and round-tripped through that encoding; without the carrier its
|
||||
// offset would silently decode as zero and reintroduce this bug in distributed mode.
|
||||
// Stale-marker filtering stays off: the input here is a computed inner result, not raw
|
||||
// storage samples.
|
||||
let divide_plan = if offset_ms == 0 {
|
||||
divide_plan
|
||||
} else {
|
||||
LogicalPlan::Extension(Extension {
|
||||
node: Arc::new(SeriesNormalize::new(
|
||||
offset_ms,
|
||||
time_index_column.clone(),
|
||||
false,
|
||||
series_key_columns,
|
||||
divide_plan,
|
||||
)),
|
||||
})
|
||||
};
|
||||
|
||||
let manipulate = RangeManipulate::new(
|
||||
self.ctx.start,
|
||||
self.ctx.end,
|
||||
self.ctx.interval,
|
||||
0,
|
||||
offset_ms,
|
||||
range_ms,
|
||||
time_index_column,
|
||||
self.ctx.field_columns.clone(),
|
||||
@@ -12194,6 +12252,28 @@ mod test {
|
||||
indie_query_plan_compare(query, expected).await;
|
||||
}
|
||||
|
||||
/// `offset` on a subquery must shift the inner evaluation window back and be
|
||||
/// carried into the outer range manipulation. See
|
||||
/// <https://github.com/GreptimeTeam/greptimedb/issues/9330>.
|
||||
#[tokio::test]
|
||||
async fn count_over_time_subquery_with_offset() {
|
||||
let query = "count_over_time(some_metric[10m:1m] offset 5m)";
|
||||
let expected = String::from(
|
||||
"Filter: prom_count_over_time(timestamp_range,field_0) IS NOT NULL [timestamp:Timestamp(ms), prom_count_over_time(timestamp_range,field_0):Float64;N, tag_0:Utf8]\
|
||||
\n Projection: some_metric.timestamp, prom_count_over_time(timestamp_range, field_0) AS prom_count_over_time(timestamp_range,field_0), some_metric.tag_0 [timestamp:Timestamp(ms), prom_count_over_time(timestamp_range,field_0):Float64;N, tag_0:Utf8]\
|
||||
\n PromRangeManipulate: req range=[0..100000000], interval=[5000], eval range=[600000], 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=[300000], time index=[timestamp], filter NaN: [false] [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 PromInstantManipulate: range=[-840000..99700000], lookback=[1000], interval=[60000], time index=[timestamp] [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(-840999, None) AND some_metric.timestamp <= TimestampMillisecond(99700000, 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;
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_hash_join() {
|
||||
let mut eval_stmt = EvalStmt {
|
||||
|
||||
@@ -63,3 +63,102 @@ drop table metric_total;
|
||||
|
||||
Affected Rows: 0
|
||||
|
||||
-- Offset on a subquery shifts the subquery's own evaluation window back by the offset.
|
||||
-- Reference: Prometheus `evaluator.subqueryTimeRange` (promql/engine.go) evaluates the inner
|
||||
-- expression over `(start - offset - range, end - offset]`; `evalSubquery` then passes the
|
||||
-- resulting samples, which keep their real timestamps, to the outer range-vector function as a
|
||||
-- `MatrixSelector` that still carries the subquery offset. See also
|
||||
-- promql/promqltest/testdata/subquery.test.
|
||||
--
|
||||
-- Every offset below is a whole multiple of the subquery step. Prometheus anchors subquery step
|
||||
-- points on absolute epoch multiples of the step, while GreptimeDB anchors them on the
|
||||
-- evaluation start, so the two only agree for step-multiple offsets; sub-step offsets are
|
||||
-- deliberately not pinned here.
|
||||
create table subquery_offset_total (
|
||||
ts timestamp time index,
|
||||
host string primary key,
|
||||
val double,
|
||||
);
|
||||
|
||||
Affected Rows: 0
|
||||
|
||||
insert into subquery_offset_total values
|
||||
(0, 'a', 1),
|
||||
(10000, 'a', 2),
|
||||
(20000, 'a', 3),
|
||||
(30000, 'a', 4),
|
||||
(40000, 'a', 5),
|
||||
(50000, 'a', 6),
|
||||
(60000, 'a', 7);
|
||||
|
||||
Affected Rows: 7
|
||||
|
||||
-- baseline: no offset at t=60 covers the 10s subquery points in (40s, 60s] -> 6 + 7
|
||||
tql eval (60, 60, '1s') sum_over_time(subquery_offset_total[20s:10s]);
|
||||
|
||||
+---------------------+----------------------------------+------+
|
||||
| ts | prom_sum_over_time(ts_range,val) | host |
|
||||
+---------------------+----------------------------------+------+
|
||||
| 1970-01-01T00:01:00 | 13.0 | a |
|
||||
+---------------------+----------------------------------+------+
|
||||
|
||||
-- the same subquery evaluated at t=30 -> 3 + 4
|
||||
tql eval (30, 30, '1s') sum_over_time(subquery_offset_total[20s:10s]);
|
||||
|
||||
+---------------------+----------------------------------+------+
|
||||
| ts | prom_sum_over_time(ts_range,val) | host |
|
||||
+---------------------+----------------------------------+------+
|
||||
| 1970-01-01T00:00:30 | 7.0 | a |
|
||||
+---------------------+----------------------------------+------+
|
||||
|
||||
-- `offset 30s` at t=60 must equal the un-offset subquery at t=30
|
||||
tql eval (60, 60, '1s') sum_over_time(subquery_offset_total[20s:10s] offset 30s);
|
||||
|
||||
+---------------------+----------------------------------+------+
|
||||
| ts | prom_sum_over_time(ts_range,val) | host |
|
||||
+---------------------+----------------------------------+------+
|
||||
| 1970-01-01T00:01:00 | 7.0 | a |
|
||||
+---------------------+----------------------------------+------+
|
||||
|
||||
-- a zero subquery offset is not expressible: the shared duration check in `promql-parser`
|
||||
-- rejects any zero duration literal, matching Prometheus's own `parseDuration`
|
||||
-- ("duration must be greater than 0", cf. its `foo[0m]` parser test). The no-op case is
|
||||
-- therefore covered by the offset-free baseline above.
|
||||
tql eval (60, 60, '1s') sum_over_time(subquery_offset_total[20s:10s] offset 0s);
|
||||
|
||||
Error: 2000(InvalidSyntax), duration must be greater than 0
|
||||
|
||||
-- a negative offset looks ahead of the evaluation time
|
||||
tql eval (30, 30, '1s') sum_over_time(subquery_offset_total[20s:10s] offset -30s);
|
||||
|
||||
+---------------------+----------------------------------+------+
|
||||
| ts | prom_sum_over_time(ts_range,val) | host |
|
||||
+---------------------+----------------------------------+------+
|
||||
| 1970-01-01T00:00:30 | 13.0 | a |
|
||||
+---------------------+----------------------------------+------+
|
||||
|
||||
-- an offset on the inner selector composes additively with the subquery offset: Prometheus
|
||||
-- `subqueryTimes` accumulates "the sum of offsets and ranges of all subqueries in the path",
|
||||
-- and the inner selector subtracts its own offset from the already shifted step timestamps.
|
||||
-- 20s + 10s therefore behaves like the un-offset subquery at t=30.
|
||||
tql eval (60, 60, '1s') sum_over_time((subquery_offset_total offset 10s)[20s:10s] offset 20s);
|
||||
|
||||
+---------------------+----------------------------------+------+
|
||||
| ts | prom_sum_over_time(ts_range,val) | host |
|
||||
+---------------------+----------------------------------+------+
|
||||
| 1970-01-01T00:01:00 | 7.0 | a |
|
||||
+---------------------+----------------------------------+------+
|
||||
|
||||
-- ... and the inner offset alone accounts for the same total shift
|
||||
tql eval (60, 60, '1s') sum_over_time((subquery_offset_total offset 30s)[20s:10s]);
|
||||
|
||||
+---------------------+----------------------------------+------+
|
||||
| ts | prom_sum_over_time(ts_range,val) | host |
|
||||
+---------------------+----------------------------------+------+
|
||||
| 1970-01-01T00:01:00 | 7.0 | a |
|
||||
+---------------------+----------------------------------+------+
|
||||
|
||||
drop table subquery_offset_total;
|
||||
|
||||
Affected Rows: 0
|
||||
|
||||
|
||||
@@ -20,3 +20,58 @@ tql eval (10, 10, '1s') rate(metric_total[20s:10s]);
|
||||
tql eval (20, 20, '1s') rate(metric_total[20s:5s]);
|
||||
|
||||
drop table metric_total;
|
||||
|
||||
-- Offset on a subquery shifts the subquery's own evaluation window back by the offset.
|
||||
-- Reference: Prometheus `evaluator.subqueryTimeRange` (promql/engine.go) evaluates the inner
|
||||
-- expression over `(start - offset - range, end - offset]`; `evalSubquery` then passes the
|
||||
-- resulting samples, which keep their real timestamps, to the outer range-vector function as a
|
||||
-- `MatrixSelector` that still carries the subquery offset. See also
|
||||
-- promql/promqltest/testdata/subquery.test.
|
||||
--
|
||||
-- Every offset below is a whole multiple of the subquery step. Prometheus anchors subquery step
|
||||
-- points on absolute epoch multiples of the step, while GreptimeDB anchors them on the
|
||||
-- evaluation start, so the two only agree for step-multiple offsets; sub-step offsets are
|
||||
-- deliberately not pinned here.
|
||||
create table subquery_offset_total (
|
||||
ts timestamp time index,
|
||||
host string primary key,
|
||||
val double,
|
||||
);
|
||||
|
||||
insert into subquery_offset_total values
|
||||
(0, 'a', 1),
|
||||
(10000, 'a', 2),
|
||||
(20000, 'a', 3),
|
||||
(30000, 'a', 4),
|
||||
(40000, 'a', 5),
|
||||
(50000, 'a', 6),
|
||||
(60000, 'a', 7);
|
||||
|
||||
-- baseline: no offset at t=60 covers the 10s subquery points in (40s, 60s] -> 6 + 7
|
||||
tql eval (60, 60, '1s') sum_over_time(subquery_offset_total[20s:10s]);
|
||||
|
||||
-- the same subquery evaluated at t=30 -> 3 + 4
|
||||
tql eval (30, 30, '1s') sum_over_time(subquery_offset_total[20s:10s]);
|
||||
|
||||
-- `offset 30s` at t=60 must equal the un-offset subquery at t=30
|
||||
tql eval (60, 60, '1s') sum_over_time(subquery_offset_total[20s:10s] offset 30s);
|
||||
|
||||
-- a zero subquery offset is not expressible: the shared duration check in `promql-parser`
|
||||
-- rejects any zero duration literal, matching Prometheus's own `parseDuration`
|
||||
-- ("duration must be greater than 0", cf. its `foo[0m]` parser test). The no-op case is
|
||||
-- therefore covered by the offset-free baseline above.
|
||||
tql eval (60, 60, '1s') sum_over_time(subquery_offset_total[20s:10s] offset 0s);
|
||||
|
||||
-- a negative offset looks ahead of the evaluation time
|
||||
tql eval (30, 30, '1s') sum_over_time(subquery_offset_total[20s:10s] offset -30s);
|
||||
|
||||
-- an offset on the inner selector composes additively with the subquery offset: Prometheus
|
||||
-- `subqueryTimes` accumulates "the sum of offsets and ranges of all subqueries in the path",
|
||||
-- and the inner selector subtracts its own offset from the already shifted step timestamps.
|
||||
-- 20s + 10s therefore behaves like the un-offset subquery at t=30.
|
||||
tql eval (60, 60, '1s') sum_over_time((subquery_offset_total offset 10s)[20s:10s] offset 20s);
|
||||
|
||||
-- ... and the inner offset alone accounts for the same total shift
|
||||
tql eval (60, 60, '1s') sum_over_time((subquery_offset_total offset 30s)[20s:10s]);
|
||||
|
||||
drop table subquery_offset_total;
|
||||
|
||||
Reference in New Issue
Block a user