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
2 changes: 1 addition & 1 deletion crates/asap-aware-mapping/src/query_physical_lowering.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1029,7 +1029,7 @@ fn promql_binary_operation(
BinaryOpKind::Set(PromQLVectorSetOpKind::And) => PromqlBinaryOperation::And,
BinaryOpKind::Set(PromQLVectorSetOpKind::Or) => PromqlBinaryOperation::Or,
BinaryOpKind::Set(PromQLVectorSetOpKind::Unless) => PromqlBinaryOperation::Unless,
BinaryOpKind::Arithmetic(_) | BinaryOpKind::Compare(_) => {
BinaryOpKind::Arithmetic(_) | BinaryOpKind::Compare(_) | BinaryOpKind::CompareBool(_) => {
PromqlBinaryOperation::ArithmeticOrComparison
}
}
Expand Down
34 changes: 18 additions & 16 deletions crates/asap-aware-mapping/src/replacement.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2804,13 +2804,12 @@ fn construct_summary_agg(
} else {
reduction.clone()
};
let per_series = matches!(reduction, Reduction::PerEntity);
let by: Vec<usize> = reduction
.group_keys()
.map(|g| g.to_vec())
.unwrap_or_default();
let out_schema = node.output_schema()?;
let state_idx = summary_col_index(&out_schema, &by, per_series);
let measures = match node {
QueryExpr::Aggregate { measures, .. } => measures.len(),
_ => 1,
};
let state_idx = summary_col_index(&out_schema, reduction, measures);

let readout_schema = if keyed_heap
&& matches!(node, QueryExpr::Aggregate { child, .. } if is_snapshot_weighted_topk(intent, child))
Expand Down Expand Up @@ -3558,16 +3557,19 @@ fn compose_guarantee(
/// cross-series output is `by ++ [agg]` (the column after the keys);
/// a per-series reduction keeps every label and replaces the sample value
/// (named `value` — mirror `per_series_reduction_schema`'s fallback).
/// `per_series` is the caller's already-read `Reduction` (issue #165) —
/// this never re-derives it, so it can't disagree with the caller.
fn summary_col_index(out_schema: &Schema, by: &[usize], per_series: bool) -> usize {
if per_series {
out_schema
/// `reduction` is the caller's already-read `Reduction` (issue #165).
/// `without` output is `kept labels ++ measures`, and its keys are the
/// *excluded* labels, so the state column follows the kept labels instead.
fn summary_col_index(out_schema: &Schema, reduction: &Reduction, measures: usize) -> usize {
match reduction {
Reduction::PerEntity => out_schema
.column_id("value")
.or_else(|| (0..out_schema.columns.len()).find(|&i| Some(i) != out_schema.time_index))
.unwrap_or(0)
} else {
by.len()
.unwrap_or(0),
Reduction::Reduce(keys) if keys.is_without() => {
out_schema.columns.len().saturating_sub(measures)
}
Reduction::Reduce(keys) => keys.len(),
}
}

Expand Down Expand Up @@ -7387,7 +7389,7 @@ mod tests {
Pass,
),
// classic-bucket histogram_quantile is not re-sketchable (#79)
(A::HistogramQuantile { q: 0.99 }, Pass),
(A::HistogramQuantile { q: 0.99, le: 0 }, Pass),
// counter-derivative / range-vector functions (#44)
(A::Changes, Pass),
(A::Delta, Pass),
Expand Down Expand Up @@ -10120,7 +10122,7 @@ mod tests {
// decree. All three stay whole logical subtrees.
for intent in [
AggIntent::Avg { col: None },
AggIntent::HistogramQuantile { q: 0.99 },
AggIntent::HistogramQuantile { q: 0.99, le: 0 },
AggIntent::Quantile {
col: None,
q: 0.99,
Expand Down
33 changes: 23 additions & 10 deletions crates/asap-physical-operators/src/expressions/arithmetic.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@ pub fn evaluate_binary(
right: f64,
) -> Result<crate::values::Value, crate::Error> {
use crate::{values::Value, Error};
use planner_types::pre_asap::{ArithmeticOpKind, BinaryOpKind, CompareOpKind};
use planner_types::pre_asap::{ArithmeticOpKind, BinaryOpKind};
let invalid =
|| Error::Invalid("unsupported binary operation or invalid checked-division domain".into());
if operator.vector_match.is_some() {
Expand All @@ -49,15 +49,28 @@ pub fn evaluate_binary(
BinaryOpKind::Arithmetic(ref op) => {
Value::Float64(evaluate_float64_arithmetic(op, left, right))
}
BinaryOpKind::Compare(ref op) => Value::Bool(match op {
CompareOpKind::Eq => left == right,
CompareOpKind::Ne => left != right,
CompareOpKind::Lt => left < right,
CompareOpKind::Le => left <= right,
CompareOpKind::Gt => left > right,
CompareOpKind::Ge => left >= right,
_ => return Err(invalid()),
}),
BinaryOpKind::Compare(ref op) => Value::Bool(compare(op, left, right).ok_or_else(invalid)?),
BinaryOpKind::CompareBool(ref op) => {
Value::Float64(if compare(op, left, right).ok_or_else(invalid)? {
1.
} else {
0.
})
}
_ => return Err(invalid()),
})
}

/// IEEE comparison, as Go's: NaN is unequal to everything, itself included.
fn compare(op: &planner_types::pre_asap::CompareOpKind, left: f64, right: f64) -> Option<bool> {
use planner_types::pre_asap::CompareOpKind;
Some(match op {
CompareOpKind::Eq => left == right,
CompareOpKind::Ne => left != right,
CompareOpKind::Lt => left < right,
CompareOpKind::Le => left <= right,
CompareOpKind::Gt => left > right,
CompareOpKind::Ge => left >= right,
_ => return None,
})
}
8 changes: 8 additions & 0 deletions crates/asap-physical-operators/src/expressions/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,14 @@ impl Expression {
| CompareOpKind::Gt
| CompareOpKind::Ge,
) => DataType::Bool,
BinaryOpKind::CompareBool(
CompareOpKind::Eq
| CompareOpKind::Ne
| CompareOpKind::Lt
| CompareOpKind::Le
| CompareOpKind::Gt
| CompareOpKind::Ge,
) => DataType::Float64,
_ => return Err(invalid("unsupported binary operation")),
};
Ok((dtype, n || m))
Expand Down
70 changes: 64 additions & 6 deletions crates/asap-physical-operators/src/operators/aggregate/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -302,9 +302,9 @@ async fn reduce_one(
}
return Ok(best.cloned().unwrap_or(Value::Null));
}
let mut count = 0usize;
let dtype = plain(input, column)?.0;
if dtype == &DataType::Int64 {
let mut count = 0usize;
let mut sum = 0i128;
for v in values {
work.checkpoint().await?;
Expand All @@ -324,22 +324,80 @@ async fn reduce_one(
))
};
}
let mut sum = -0.0;
let mut floats = Vec::with_capacity(rows.len());
for v in values {
work.checkpoint().await?;
let Value::Float64(v) = v else {
return Err(invalid("floating aggregate value required"));
};
sum += v;
count += 1;
floats.push(*v);
}
Ok(Value::Float64(if matches!(measure, Reduction::Avg(_)) {
sum / count as f64
promql_avg(&floats)
} else {
sum
// Prometheus starts `sum` from the first value; adding it to 0 gives
// the same sum and compensation.
promql_sum(0., &floats)
}))
}

/// Prometheus `kahansum.Inc`: Kahan-Neumaier compensated addition, with the
/// compensation cleared once the sum is infinite.
fn kahan_inc(inc: f64, sum: f64, c: f64) -> (f64, f64) {
let t = sum + inc;
let c = if t.is_infinite() {
0.
} else if sum.abs() >= inc.abs() {
c + ((sum - t) + inc)
} else {
c + ((inc - t) + sum)
};
(t, c)
}

/// Prometheus `sum` and `sum_over_time` from `start`.
pub(in crate::operators) fn promql_sum(start: f64, values: &[f64]) -> f64 {
let (sum, c) = values
.iter()
.fold((start, 0.), |(sum, c), &v| kahan_inc(v, sum, c));
if sum.is_infinite() {
sum
} else {
sum + c
}
}

/// Prometheus `avg` and `avg_over_time`: a compensated sum divided by the
/// count until the running sum would overflow, then an incremental mean.
/// NaN for no values.
pub(in crate::operators) fn promql_avg(values: &[f64]) -> f64 {
let Some((&first, rest)) = values.split_first() else {
return f64::NAN;
};
let (mut sum, mut c, mut mean, mut incremental) = (first, 0., 0., false);
for (i, &v) in rest.iter().enumerate() {
let count = (i + 2) as f64;
if !incremental {
let (next, next_c) = kahan_inc(v, sum, c);
if !next.is_infinite() {
(sum, c) = (next, next_c);
continue;
}
incremental = true;
mean = sum / (count - 1.);
c /= count - 1.;
}
let q = (count - 1.) / count;
(mean, c) = kahan_inc(v / count, q * mean, q * c);
}
let count = values.len() as f64;
if incremental {
mean + c
} else {
sum / count + c / count
}
}

#[cfg(test)]
mod tests {
// Quantile follows Prometheus: interpolate ranks, NaN when empty, ±Inf outside [0, 1].
Expand Down
81 changes: 72 additions & 9 deletions crates/asap-physical-operators/src/operators/aggregate/temporal.rs
Original file line number Diff line number Diff line change
Expand Up @@ -52,7 +52,7 @@ pub(in crate::operators) async fn reduce(
let mut output = Vec::new();
for (_, (mut keys, buckets, points)) in grouped {
work.checkpoint().await?;
let result = if let AggIntent::HistogramQuantile { q } = intent {
let result = if let AggIntent::HistogramQuantile { q, .. } = intent {
Some(Value::Float64(bucket_quantile(*q, buckets, context).await?))
} else {
let points = cooperative_sort(points, |a, b| a.0.cmp(&b.0), context).await?;
Expand Down Expand Up @@ -92,10 +92,14 @@ pub(in crate::operators) fn window_value(
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::Sum { .. } | AggIntent::Avg { .. } => {
let values = points.iter().map(|p| p.1).collect::<Vec<_>>();
Some(Value::Float64(if matches!(intent, AggIntent::Sum { .. }) {
super::promql_sum(0., &values)
} else {
super::promql_avg(&values)
}))
}
AggIntent::Min { .. } => Some(Value::Float64(points.iter().fold(f64::NAN, |a, p| {
if a.is_nan() || p.1 < a {
p.1
Expand Down Expand Up @@ -176,7 +180,7 @@ fn rate(points: &[(i64, f64)], start: i64, end: i64, counter: bool) -> Option<f6
Some(delta * (span + to_start + to_end) / span / ((end as f64 - start as f64) / 1000.))
}

async fn bucket_quantile(
pub(in crate::operators) async fn bucket_quantile(
q: f64,
mut b: Vec<(f64, f64)>,
context: &RunContext,
Expand Down Expand Up @@ -211,7 +215,7 @@ async fn bucket_quantile(
let mut prev = buckets[0].1;
for p in buckets.iter_mut().skip(1) {
work.checkpoint().await?;
if p.1 < prev || (p.1 - prev).abs() <= 1e-12 * (p.1.abs() + prev.abs()) {
if p.1 < prev || almost_equal(prev, p.1) {
p.1 = prev;
}
prev = p.1;
Expand All @@ -221,7 +225,13 @@ async fn bucket_quantile(
return Ok(f64::NAN);
}
let rank = q * count;
let idx = buckets[..buckets.len() - 1].partition_point(|p| p.1 < rank);
let idx = buckets[..buckets.len() - 1].partition_point(|p| {
// Go searches for `count >= rank`; a NaN comparison is never a match.
!matches!(
p.1.partial_cmp(&rank),
Some(std::cmp::Ordering::Greater | std::cmp::Ordering::Equal)
)
});
if idx == buckets.len() - 1 {
return Ok(buckets[idx - 1].0);
}
Expand All @@ -230,7 +240,21 @@ async fn bucket_quantile(
}
let (start, base) = if idx == 0 { (0., 0.) } else { buckets[idx - 1] };
let (end, upper) = buckets[idx];
Ok(start + (end - start) * (rank - base) / (upper - base))
Ok(start + (end - start) * ((rank - base) / (upper - base)))
}

/// Prometheus `almost.Equal` with its bucket tolerance of 1e-12.
fn almost_equal(a: f64, b: f64) -> bool {
const EPSILON: f64 = 1e-12;
if a == b || (a.is_nan() && b.is_nan()) {
return true;
}
let sum = a.abs() + b.abs();
let diff = (a - b).abs();
if a == 0. || b == 0. || sum < f64::MIN_POSITIVE {
return diff < EPSILON * f64::MIN_POSITIVE;
}
diff / sum.min(f64::MAX) < EPSILON
}

#[cfg(test)]
Expand Down Expand Up @@ -335,4 +359,43 @@ mod tests {
assert!(bucket_quantile(0.5, vec![(1., 2.), (2., 4.)]).is_nan());
assert_eq!(bucket_quantile(-0.1, vec![]), f64::NEG_INFINITY);
}

// Matches Prometheus bit for bit: interpolation divides before scaling,
// and an infinite count is not "almost equal" to a finite one.
#[test]
fn histogram_matches_prometheus_arithmetic() {
let context = RunContext::new(
Scope::Query {
evaluation_time_ms: 0,
revision: 0,
},
Limits::default(),
)
.unwrap();
let bucket_quantile = |q, buckets| {
futures::executor::block_on(super::bucket_quantile(q, buckets, &context)).unwrap()
};
// 0.5 + (0.8 - 0.5) * ((19.95 - 10) / 11) in Go.
assert_eq!(
bucket_quantile(0.95, vec![(0.5, 10.), (0.8, 21.), (f64::INFINITY, 21.)]),
0.7713636363636364
);
// Rank ∞ lies past every finite bucket.
assert_eq!(
bucket_quantile(
0.5,
vec![(1., 1.), (2., 2.), (f64::INFINITY, f64::INFINITY)]
),
2.
);
// A NaN rank (0 · ∞, or a NaN total) finds no bucket, as Go's sort.Search.
assert_eq!(
bucket_quantile(0., vec![(1., 1.), (2., 2.), (f64::INFINITY, f64::INFINITY)]),
2.
);
assert_eq!(
bucket_quantile(0.5, vec![(1., 1.), (2., 2.), (f64::INFINITY, f64::NAN)]),
2.
);
}
}
Loading
Loading