mirror of
https://github.com/GreptimeTeam/greptimedb.git
synced 2026-10-03 18:45: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,
|
subquery_expr: &SubqueryExpr,
|
||||||
) -> Result<LogicalPlan> {
|
) -> Result<LogicalPlan> {
|
||||||
let SubqueryExpr {
|
let SubqueryExpr {
|
||||||
expr, range, step, ..
|
expr,
|
||||||
|
range,
|
||||||
|
step,
|
||||||
|
offset,
|
||||||
|
..
|
||||||
} = subquery_expr;
|
} = 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;
|
let current_interval = self.ctx.interval;
|
||||||
if let Some(step) = step {
|
if let Some(step) = step {
|
||||||
self.ctx.interval = step.as_millis() as _;
|
self.ctx.interval = step.as_millis() as _;
|
||||||
}
|
}
|
||||||
let current_start = self.ctx.start;
|
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?;
|
let input = self.prom_expr_to_plan(expr, query_engine_state).await?;
|
||||||
self.ctx.interval = current_interval;
|
self.ctx.interval = current_interval;
|
||||||
self.ctx.start = current_start;
|
self.ctx.start = current_start;
|
||||||
|
self.ctx.end = current_end;
|
||||||
|
|
||||||
ensure!(!range.is_zero(), ZeroRangeSelectorSnafu);
|
ensure!(!range.is_zero(), ZeroRangeSelectorSnafu);
|
||||||
let range_ms = range.as_millis() as _;
|
let range_ms = range.as_millis() as _;
|
||||||
@@ -616,17 +653,38 @@ impl PromPlanner {
|
|||||||
.context(DataFusionPlanningSnafu)?;
|
.context(DataFusionPlanningSnafu)?;
|
||||||
let divide_plan = LogicalPlan::Extension(Extension {
|
let divide_plan = LogicalPlan::Extension(Extension {
|
||||||
node: Arc::new(SeriesDivide::new(
|
node: Arc::new(SeriesDivide::new(
|
||||||
series_key_columns,
|
series_key_columns.clone(),
|
||||||
time_index_column.clone(),
|
time_index_column.clone(),
|
||||||
sort_plan,
|
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(
|
let manipulate = RangeManipulate::new(
|
||||||
self.ctx.start,
|
self.ctx.start,
|
||||||
self.ctx.end,
|
self.ctx.end,
|
||||||
self.ctx.interval,
|
self.ctx.interval,
|
||||||
0,
|
offset_ms,
|
||||||
range_ms,
|
range_ms,
|
||||||
time_index_column,
|
time_index_column,
|
||||||
self.ctx.field_columns.clone(),
|
self.ctx.field_columns.clone(),
|
||||||
@@ -12194,6 +12252,28 @@ mod test {
|
|||||||
indie_query_plan_compare(query, expected).await;
|
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]
|
#[tokio::test]
|
||||||
async fn test_hash_join() {
|
async fn test_hash_join() {
|
||||||
let mut eval_stmt = EvalStmt {
|
let mut eval_stmt = EvalStmt {
|
||||||
|
|||||||
@@ -63,3 +63,102 @@ drop table metric_total;
|
|||||||
|
|
||||||
Affected Rows: 0
|
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]);
|
tql eval (20, 20, '1s') rate(metric_total[20s:5s]);
|
||||||
|
|
||||||
drop table metric_total;
|
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