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:
Aarav
2026-09-25 12:59:31 +05:30
co-authored by Claude Opus 5
parent 194bc2fb3c
commit e132c26733
3 changed files with 238 additions and 4 deletions
+84 -4
View File
@@ -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;