Skip to content
Open
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.
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.);
}
}
113 changes: 106 additions & 7 deletions crates/asap-physical-operators/src/physical_planner/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,8 @@ use planner_types::{
SketchQuery, SummaryFamilyType, SummaryInputExpr, ValueOperation,
},
pre_asap::{
AggIntent, ColumnRef, CompareOpKind, GroupKeys, QueryExpr, Reduction as PlannerReduction,
AggIntent, ColumnRef, CompareOpKind, DataType, GroupKeys, QueryExpr,
Reduction as PlannerReduction,
},
};
use std::{
Expand Down Expand Up @@ -43,6 +44,8 @@ pub use candidates::{
mod compiled;
pub use compiled::{CompiledPhysicalDag, InputContract};

mod row_values;

/// Compile computation without opening or retaining deployment readers.
/// Input contracts identify explicit boundaries selected by maintenance planning.
pub fn compile(
Expand Down Expand Up @@ -143,7 +146,28 @@ fn compile_internal(
},
)
});
// Scalar literal operands of query-time arithmetic are folded into the consumer.
let mut literals = BTreeMap::<NodeId, (f64, bool)>::new();
for edge in edges {
let consumer = u64::from(edge.consumer.0);
if let (
Payload::Fallback { expression },
Some(PostAsapDagNode {
payload: Payload::Binary { .. },
..
}),
) = (
&nodes[&u64::from(edge.producer.0)].payload,
nodes.get(&consumer),
) {
if let Some(value) = row_values::scalar_literal(expression) {
let left = edge.role == planner_types::post_asap::EdgeRole::Left;
if literals.insert(consumer, (value, left)).is_some() {
return Err(invalid("binary with two scalar literals is not folded"));
}
continue;
}
}
dependencies
.entry(u64::from(edge.consumer.0))
.or_default()
Expand Down Expand Up @@ -176,7 +200,7 @@ fn compile_internal(
let mut graph = CompiledPhysicalDag::new(roots.to_vec());
for id in ordered {
let node = nodes[&id];
let auxiliary = helper_id(id, 0);
let mut auxiliary = helper_id(id, 0);
let output = Arc::new(node.output_schema.clone());
crate::values::validate_schema(&output)?;
if let Some(source) = sources.remove(&id) {
Expand Down Expand Up @@ -247,11 +271,6 @@ fn compile_internal(
use planner_types::post_asap::maintained_population::{
PopulationInput, PopulationReadout,
};
let PopulationReadout::TopK { k } = readout else {
return Err(invalid(
"native population readout does not support this operation",
));
};
let [producer] = inputs.as_slice() else {
return Err(invalid("population readout requires one input"));
};
Expand All @@ -272,6 +291,19 @@ fn compile_internal(
));
}
let input = schemas[0].clone();
let PopulationReadout::TopK { k } = readout else {
let mut chain =
row_values::population_aggregate(&input, &spec.grouping, readout)?;
let last = chain.pop().expect("nonempty chain");
let mut inputs = inputs;
for operator in chain {
graph.add(auxiliary, inputs, operator)?;
inputs = vec![auxiliary];
auxiliary -= 1;
}
graph.add(id, inputs, last.with_output_schema(output)?)?;
continue;
};
let groups = spec
.grouping
.iter()
Expand Down Expand Up @@ -357,6 +389,73 @@ fn compile_internal(
)?;
continue;
}
if let Payload::Binary { operator } = &node.payload {
let query_time = node.output_state.timing
== planner_types::post_asap::ExecutionTiming::QueryTime;
if let Some(&(value, left)) = literals.get(&id) {
let [input] = schemas.as_slice() else {
return Err(invalid("scalar binary requires one row input"));
};
if !query_time {
return Err(invalid("scalar literal binary must run at query time"));
}
let project = row_values::scalar_binary(input, operator, value, left)
.map_err(|error| invalid(format!("node {id}: {error}")))?;
graph.add(id, inputs, project.with_output_schema(output)?)?;
continue;
}
let label_map = |schema: &Schema| {
schema
.fields
.iter()
.any(|f| matches!(f.dtype, SummaryFamilyType::Plain(DataType::Map { .. })))
};
if let (true, [left, right]) = (query_time, schemas.as_slice()) {
if !label_map(left) && !label_map(right) {
let (join, project) = row_values::grouped_binary(left, right, operator)
.map_err(|error| invalid(format!("node {id}: {error}")))?;
graph.add(auxiliary, inputs, join)?;
graph.add(id, vec![auxiliary], project.with_output_schema(output)?)?;
auxiliary -= 1;
continue;
}
}
}
if let Payload::Value {
operation: ValueOperation::FinalizeExactAccumulator,
} = &node.payload
{
// Exact counts read out as Int64; PromQL declares a Float64 sample.
let readout = bind_operation(node, &schemas)
.map_err(|error| invalid(format!("node {id}: {error}")))?;
let actual = readout.schema();
let converted = actual.fields.iter().zip(&output.fields).position(|(a, d)| {
a.dtype == SummaryFamilyType::Plain(DataType::Int64)
&& d.dtype == SummaryFamilyType::Plain(DataType::Float64)
});
if let Some(column) = converted {
let columns = actual
.fields
.iter()
.enumerate()
.map(|(i, field)| {
(
field.name.clone(),
if i == column {
Expression::ExactFloat64(i)
} else {
Expression::Column(i)
},
)
})
.collect();
let project = Operator::project(actual, columns)?.with_output_schema(output)?;
graph.add(auxiliary, inputs, readout)?;
graph.add(id, vec![auxiliary], project)?;
auxiliary -= 1;
continue;
}
}
let mut operator = compile_node(node, &schemas)
.map_err(|error| invalid(format!("node {id}: {error}")))?;
if operator.is_counter_readout() {
Expand Down
Loading
Loading