Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
79 changes: 44 additions & 35 deletions crates/asap-physical-operators/src/operators/aggregate/temporal.rs
Original file line number Diff line number Diff line change
Expand Up @@ -65,38 +65,7 @@ pub(in crate::operators) async fn reduce(
"duplicate or out-of-window timestamp".into(),
));
}
match intent {
AggIntent::Rate => rate(&points, start, end).map(Value::Float64),
AggIntent::Increase => rate(&points, start, end)
.map(|v| Value::Float64(v * (end as f64 - start as f64) / 1000.)),
AggIntent::Count { .. } => Some(Value::Int64(
i64::try_from(points.len())
.map_err(|_| Error::Invalid("count overflow".into()))?,
)),
AggIntent::Sum { .. } => Some(Value::Float64(points.iter().map(|p| p.1).sum())),
AggIntent::Avg { .. } => Some(Value::Float64(
points.iter().map(|p| p.1).sum::<f64>() / points.len() as f64,
)),
AggIntent::Min { .. } => {
Some(Value::Float64(points.iter().fold(f64::NAN, |a, p| {
if a.is_nan() || p.1 < a {
p.1
} else {
a
}
})))
}
AggIntent::Max { .. } => {
Some(Value::Float64(points.iter().fold(f64::NAN, |a, p| {
if a.is_nan() || p.1 > a {
p.1
} else {
a
}
})))
}
_ => return Err(Error::Invalid("unsupported temporal intent".into())),
}
window_value(intent, &points, start, end)?
};
if let Some(result) = result {
keys.push(result);
Expand All @@ -106,7 +75,47 @@ pub(in crate::operators) async fn reduce(
Ok(output)
}

fn rate(points: &[(i64, f64)], start: i64, end: i64) -> Option<f64> {
/// One series' value over its sorted samples in the window `(start, end]`.
/// `None` means PromQL emits no sample for this series.
pub(in crate::operators) fn window_value(
intent: &AggIntent<ColumnRef>,
points: &[(i64, f64)],
start: i64,
end: i64,
) -> Result<Option<Value>, Error> {
Ok(match intent {
AggIntent::Rate => rate(points, start, end, true).map(Value::Float64),
AggIntent::Delta => rate(points, start, end, false)
.map(|v| Value::Float64(v * (end as f64 - start as f64) / 1000.)),
AggIntent::Increase => rate(points, start, end, true)
.map(|v| Value::Float64(v * (end as f64 - start as f64) / 1000.)),
AggIntent::Count { .. } => Some(Value::Int64(
i64::try_from(points.len()).map_err(|_| Error::Invalid("count overflow".into()))?,
)),
AggIntent::Sum { .. } => Some(Value::Float64(points.iter().map(|p| p.1).sum())),
AggIntent::Avg { .. } => Some(Value::Float64(
points.iter().map(|p| p.1).sum::<f64>() / points.len() as f64,
)),
AggIntent::Min { .. } => Some(Value::Float64(points.iter().fold(f64::NAN, |a, p| {
if a.is_nan() || p.1 < a {
p.1
} else {
a
}
}))),
AggIntent::Max { .. } => Some(Value::Float64(points.iter().fold(f64::NAN, |a, p| {
if a.is_nan() || p.1 > a {
p.1
} else {
a
}
}))),
_ => return Err(Error::Invalid("unsupported temporal intent".into())),
})
}

/// Prometheus `extrapolatedRate`; `counter` enables reset correction and the zero bound.
fn rate(points: &[(i64, f64)], start: i64, end: i64, counter: bool) -> Option<f64> {
if points.len() < 2 {
return None;
}
Expand All @@ -118,7 +127,7 @@ fn rate(points: &[(i64, f64)], start: i64, end: i64) -> Option<f64> {
}
let mut delta = last - first;
for pair in points.windows(2) {
if pair[1].1 < pair[0].1 {
if counter && pair[1].1 < pair[0].1 {
delta += pair[0].1;
}
}
Expand All @@ -129,7 +138,7 @@ fn rate(points: &[(i64, f64)], start: i64, end: i64) -> Option<f64> {
to_start = average / 2.;
}
// Apply the zero bound after the sparse-window half-interval cap.
if delta > 0. && first >= 0. {
if counter && delta > 0. && first >= 0. {
to_start = to_start.min(span * first / delta);
}
if to_end >= average * 1.1 {
Expand Down
14 changes: 14 additions & 0 deletions crates/asap-physical-operators/src/operators/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,13 +23,15 @@ mod joins;
mod limit;
mod projection;
mod scope_timestamp;
mod series_window;
mod sort;
mod source;
mod summary;
mod unchecked;
pub(crate) mod vector_binary;
pub(crate) mod vector_window;
pub use aggregate::Reduction;
pub use series_window::SubquerySteps;
pub use sort::SortKey;
pub use summary::ReadoutQuery;
#[derive(Clone, serde::Serialize, serde::Deserialize)]
Expand Down Expand Up @@ -66,6 +68,14 @@ enum Kind {
intent: Box<planner_types::pre_asap::AggIntent<ColumnRef>>,
},
HistogramQuantile,
SeriesWindow {
function: Option<Box<planner_types::pre_asap::AggIntent<ColumnRef>>>,
coordinate: usize,
value: usize,
range_ms: i64,
offset_ms: i64,
steps: Option<SubquerySteps>,
},
Project(Vec<Expression>),
Filter(Expression),
Limit {
Expand Down Expand Up @@ -237,6 +247,7 @@ impl PhysicalOperator<Batch, Schema> for Operator {
| Kind::RangeWindow { .. }
| Kind::HistogramQuantile
| Kind::CurrentSeries { .. }
| Kind::SeriesWindow { .. }
| Kind::Aggregate { .. }
| Kind::Window { .. }
| Kind::Join { .. }
Expand Down Expand Up @@ -279,6 +290,7 @@ impl PhysicalOperator<Batch, Schema> for Operator {
Kind::AlignedBinary { .. } => "AlignedBinary",
Kind::RangeWindow { .. } => "RangeWindow",
Kind::HistogramQuantile => "HistogramQuantile",
Kind::SeriesWindow { .. } => "SeriesWindow",
Kind::Project(_) => "Project",
Kind::Filter(_) => "Filter",
Kind::Limit { .. } => "Limit",
Expand All @@ -295,6 +307,7 @@ impl PhysicalOperator<Batch, Schema> for Operator {
}
fn validate_context(&self, context: &RunContext) -> Result<(), Error> {
current_series::validate_context(self, context)?;
series_window::validate_context(self, context)?;
self.readout_range(context).map(|_| ())
}
fn input_schemas(&self) -> Vec<Schema> {
Expand Down Expand Up @@ -323,6 +336,7 @@ impl PhysicalOperator<Batch, Schema> for Operator {
Kind::Project(_) => projection::execute(self, inputs, context),
Kind::CurrentSeries { .. } => current_series::execute(self, inputs, context),
Kind::ScopeTimestamp { .. } => scope_timestamp::execute(self, inputs, context),
Kind::SeriesWindow { .. } => series_window::execute(self, inputs, context),
Kind::Filter(_) => filter::execute(self, inputs, context),
Kind::Limit { .. } => limit::execute(self, inputs, context),
Kind::Sort { .. } => sort::execute(self, inputs, context),
Expand Down
226 changes: 226 additions & 0 deletions crates/asap-physical-operators/src/operators/series_window.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,226 @@
//! PromQL per-series evaluation over the samples before an evaluation instant.
use super::*;
use planner_types::pre_asap::AggIntent;

/// A PromQL subquery grid: every multiple of `step_ms` in
/// `(T - offset_ms - range_ms, T - offset_ms]`, where `T` is the query time.
#[derive(Clone, Copy, Debug, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct SubquerySteps {
pub range_ms: i64,
pub step_ms: i64,
pub offset_ms: i64,
}

const STALE_MARKER: u64 = 0x7ff0_0000_0000_0002;
const MAX_SUBQUERY_STEPS: i64 = 100_000;

impl Operator {
/// Evaluate each series at instant `t` over its samples in
/// `(t - offset_ms - range_ms, t - offset_ms]`. `t` is the query time, or
/// each step of `steps`. `function: None` is instant selection: the latest
/// sample, absent if it is a stale marker. Range functions ignore stale
/// markers. A series is every column except the time and `value` columns.
/// Output rows keep the input schema, with time `t` and the result value.
pub fn series_window(
input: Schema,
function: Option<AggIntent<ColumnRef>>,
range_ms: i64,
offset_ms: i64,
steps: Option<SubquerySteps>,
) -> Result<Self, Error> {
let coordinate = input
.time_index
.ok_or_else(|| invalid("series window requires a time column"))?;
let values = input
.fields
.iter()
.enumerate()
.filter(|(_, f)| f.name == "value")
.map(|(i, _)| i)
.collect::<Vec<_>>();
let [value] = values.as_slice() else {
return Err(invalid("series window requires one value column"));
};
if plain(&input, coordinate)? != (&DataType::Timestamp, false)
|| plain(&input, *value)? != (&DataType::Float64, false)
{
return Err(invalid(
"series window requires non-null time and Float64 value",
));
}
if range_ms <= 0 || steps.is_some_and(|s| s.range_ms <= 0 || s.step_ms <= 0) {
return Err(invalid("series window ranges and steps must be positive"));
}
// Bounds per-run work independently of the data, as the backend's grid does.
if steps.is_some_and(|s| s.range_ms / s.step_ms > MAX_SUBQUERY_STEPS) {
return Err(invalid("subquery exceeds 100000 steps"));
}
if !matches!(
function,
None | Some(
AggIntent::Rate
| AggIntent::Increase
| AggIntent::Delta
| AggIntent::Count { .. }
| AggIntent::Sum { col: None }
| AggIntent::Avg { col: None }
| AggIntent::Min { col: None }
| AggIntent::Max { col: None }
)
) {
return Err(invalid("unsupported PromQL range function"));
}
Ok(Self {
kind: Kind::SeriesWindow {
function: function.map(Box::new),
coordinate,
value: *value,
range_ms,
offset_ms,
steps,
},
inputs: vec![input.clone()],
output: input,
})
}
}

/// The first and last evaluation instants and the step between them. The grid
/// is iterated, not allocated: its size depends only on the query.
fn evaluation_times(
context: &RunContext,
steps: Option<SubquerySteps>,
) -> Result<(i64, i64, i64), Error> {
let crate::runtime::Scope::Query {
evaluation_time_ms, ..
} = context.scope
else {
return Err(invalid("series window requires a query evaluation time"));
};
let Some(steps) = steps else {
return Ok((evaluation_time_ms, evaluation_time_ms, 1));
};
let overflow = || invalid("subquery grid overflows");
let end = evaluation_time_ms
.checked_sub(steps.offset_ms)
.ok_or_else(overflow)?;
let start = end.checked_sub(steps.range_ms).ok_or_else(overflow)?;
let first = (start.div_euclid(steps.step_ms) + 1)
.checked_mul(steps.step_ms)
.ok_or_else(overflow)?;
Ok((first, end, steps.step_ms))
}

pub(super) fn execute<'a>(
operator: &'a Operator,
mut inputs: Vec<Input<'a, Batch>>,
context: RunContext,
) -> Result<OutputStream<'a, Batch>, Error> {
let Kind::SeriesWindow {
function,
coordinate,
value,
range_ms,
offset_ms,
steps,
} = &operator.kind
else {
unreachable!()
};
let (coordinate, value) = (*coordinate, *value);
let input = inputs
.pop()
.ok_or_else(|| invalid("series window input missing"))?;
let times = evaluation_times(&context, *steps)?;
Ok(futures::stream::once(async move {
let (rows, _memory) = collect_rows(input, &context).await?;
let mut work = Cooperative::new(&context);
let mut workspace = Workspace::new(&context)?;
let identity = (0..operator.output.fields.len())
.filter(|&i| i != coordinate && i != value)
.collect::<Vec<_>>();
let mut series = BTreeMap::<Vec<Vec<u8>>, (usize, Vec<(i64, f64)>)>::new();
for (index, row) in rows.iter().enumerate() {
work.checkpoint().await?;
let key = group_key(row, &identity)?;
let (Value::Timestamp(time), Value::Float64(sample)) = (&row[coordinate], &row[value])
else {
return Err(invalid("series window requires time and value samples"));
};
workspace.grow(16)?;
if !series.contains_key(&key) {
workspace.grow(key_bytes(&key) + 64)?;
}
series
.entry(key)
.or_insert_with(|| (index, Vec::new()))
.1
.push((*time, *sample));
}
for (_, points) in series.values_mut() {
work.checkpoint().await?;
points.sort_by_key(|p| p.0);
if points.windows(2).any(|p| p[0].0 == p[1].0) {
return Err(invalid("duplicate sample timestamp for one series"));
}
}
let mut output = Vec::new();
let (mut time, last, step) = times;
while time <= last && !series.is_empty() {
work.checkpoint().await?;
let overflow = || invalid("series window overflows");
let end = time.checked_sub(*offset_ms).ok_or_else(overflow)?;
let start = end.checked_sub(*range_ms).ok_or_else(overflow)?;
for (template, points) in series.values() {
work.checkpoint().await?;
// PromQL ranges are left-open: a sample at `start` is outside.
let first = points.partition_point(|p| p.0 <= start);
let last = points.partition_point(|p| p.0 <= end);
let points = &points[first..last];
let result = match function {
None => points
.last()
.filter(|p| p.1.to_bits() != STALE_MARKER)
.map(|p| p.1),
Some(intent) => {
let fresh = points
.iter()
.copied()
.filter(|p| p.1.to_bits() != STALE_MARKER)
.collect::<Vec<_>>();
if fresh.is_empty() {
None
} else {
match aggregate::temporal::window_value(intent, &fresh, start, end)? {
Some(Value::Float64(v)) => Some(v),
Some(Value::Int64(v)) => Some(v as f64),
Some(_) => return Err(invalid("invalid range function result")),
None => None,
}
}
}
};
if let Some(result) = result {
let mut row = rows[*template].clone();
row[coordinate] = Value::Timestamp(time);
row[value] = Value::Float64(result);
workspace.grow(row_bytes(&row))?;
output.push(row);
}
}
let Some(next) = time.checked_add(step) else {
break;
};
time = next;
}
Batch::try_new(operator.output.clone(), output)
})
.boxed_local())
}

pub(super) fn validate_context(operator: &Operator, context: &RunContext) -> Result<(), Error> {
if let Kind::SeriesWindow { steps, .. } = operator.kind {
evaluation_times(context, steps)?;
}
Ok(())
}
Loading
Loading