Skip to content
Merged
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
62 changes: 62 additions & 0 deletions crates/asap-physical-operators/src/operators/aggregate/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,12 @@ impl Operator {
false,
)
}
Reduction::Quantile { column, q } => {
if plain(&input, *column)?.0 != &DataType::Float64 || q.is_nan() {
return Err(invalid("quantile requires Float64 input and a numeric q"));
}
(DataType::Float64, false)
}
Reduction::Min(i) | Reduction::Max(i) => {
let (t, nullable) = plain(&input, *i)?;
if !ordered(t) {
Expand Down Expand Up @@ -120,6 +126,11 @@ pub enum Reduction {
Avg(usize),
Min(usize),
Max(usize),
/// PromQL `quantile`: linear interpolation between closest ranks.
Quantile {
column: usize,
q: f64,
},
}
pub(super) fn execute<'a>(
operator: &'a Operator,
Expand Down Expand Up @@ -198,6 +209,31 @@ async fn reduce(
Ok(output)
}

// Matches Prometheus `quantile`: NaN for no values, ±Inf outside [0, 1],
// and NaN samples ordered first.
pub(super) fn quantile(q: f64, mut values: Vec<f64>) -> f64 {
if values.is_empty() {
return f64::NAN;
}
if q < 0. {
return f64::NEG_INFINITY;
}
if q > 1. {
return f64::INFINITY;
}
values.sort_by(|a, b| match (a.is_nan(), b.is_nan()) {
(true, true) => std::cmp::Ordering::Equal,
(true, false) => std::cmp::Ordering::Less,
(false, true) => std::cmp::Ordering::Greater,
_ => a.total_cmp(b),
});
let rank = q * (values.len() - 1) as f64;
let low = rank.floor() as usize;
let high = (low + 1).min(values.len() - 1);
let weight = rank - low as f64;
values[low] * (1. - weight) + values[high] * weight
}

async fn reduce_one(
rows: &[Vec<Value>],
measure: &Reduction,
Expand All @@ -211,6 +247,18 @@ async fn reduce_one(
))
}
Reduction::Sum(i) | Reduction::Avg(i) | Reduction::Min(i) | Reduction::Max(i) => *i,
Reduction::Quantile { column, q } => {
let mut values = Vec::with_capacity(rows.len());
for row in rows {
work.checkpoint().await?;
match &row[*column] {
Value::Float64(value) => values.push(*value),
Value::Null => {}
_ => return Err(invalid("floating quantile value required")),
}
}
return Ok(Value::Float64(quantile(*q, values)));
}
};
let values = rows
.iter()
Expand Down Expand Up @@ -291,3 +339,17 @@ async fn reduce_one(
sum
}))
}

#[cfg(test)]
mod tests {
// Quantile follows Prometheus: interpolate ranks, NaN when empty, ±Inf outside [0, 1].
#[test]
fn quantile_matches_prometheus_edge_cases() {
assert!(super::quantile(0.5, vec![]).is_nan());
assert_eq!(super::quantile(0.5, vec![3.]), 3.);
assert_eq!(super::quantile(0.75, vec![4., 1., 2., 3.]), 3.25);
assert_eq!(super::quantile(-0.1, vec![1.]), f64::NEG_INFINITY);
assert_eq!(super::quantile(1.1, vec![1.]), f64::INFINITY);
assert_eq!(super::quantile(1., vec![2., f64::NAN, 1.]), 2.);
}
}
108 changes: 73 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,76 @@ 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
}
}))),
AggIntent::IRate | AggIntent::IDelta => {
let [.., (t0, v0), (t1, v1)] = points else {
return Ok(None);
};
let rate = matches!(intent, AggIntent::IRate);
// A counter reset makes the last value the increase.
let delta = if rate && v1 < v0 { *v1 } else { v1 - v0 };
match (rate, t1 - t0) {
(_, 0) => None,
(true, interval) => Some(Value::Float64(delta / (interval as f64 / 1000.))),
(false, _) => Some(Value::Float64(delta)),
}
}
AggIntent::Changes | AggIntent::Resets => {
let changed = |(a, b): (f64, f64)| match intent {
AggIntent::Changes => a != b && !(a.is_nan() && b.is_nan()),
_ => b < a,
};
let count = points
.windows(2)
.filter(|pair| changed((pair[0].1, pair[1].1)))
.count();
Some(Value::Float64(count as f64))
}
AggIntent::LastOverTime => points.last().map(|p| Value::Float64(p.1)),
AggIntent::Quantile { col: None, q, .. } => Some(Value::Float64(super::quantile(
*q,
points.iter().map(|p| p.1).collect(),
))),
_ => 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 +156,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 +167,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
30 changes: 30 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,16 @@ mod joins;
mod limit;
mod projection;
mod scope_timestamp;
mod series_labels;
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 +69,22 @@ 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,
at_ms: Option<i64>,
steps: Option<SubquerySteps>,
},
SeriesLabels {
kind: planner_types::pre_asap::VectorMatchKind,
labels: Vec<String>,
},
SeriesBinary {
operator: planner_types::post_asap::BinaryOperator,
},
Project(Vec<Expression>),
Filter(Expression),
Limit {
Expand Down Expand Up @@ -237,6 +256,9 @@ impl PhysicalOperator<Batch, Schema> for Operator {
| Kind::RangeWindow { .. }
| Kind::HistogramQuantile
| Kind::CurrentSeries { .. }
| Kind::SeriesWindow { .. }
| Kind::SeriesLabels { .. }
| Kind::SeriesBinary { .. }
| Kind::Aggregate { .. }
| Kind::Window { .. }
| Kind::Join { .. }
Expand Down Expand Up @@ -279,6 +301,9 @@ impl PhysicalOperator<Batch, Schema> for Operator {
Kind::AlignedBinary { .. } => "AlignedBinary",
Kind::RangeWindow { .. } => "RangeWindow",
Kind::HistogramQuantile => "HistogramQuantile",
Kind::SeriesWindow { .. } => "SeriesWindow",
Kind::SeriesLabels { .. } => "SeriesLabels",
Kind::SeriesBinary { .. } => "SeriesBinary",
Kind::Project(_) => "Project",
Kind::Filter(_) => "Filter",
Kind::Limit { .. } => "Limit",
Expand All @@ -295,6 +320,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 +349,10 @@ 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::SeriesLabels { .. } | Kind::SeriesBinary { .. } => {
series_labels::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
Loading
Loading