Skip to content
Draft
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
77 changes: 71 additions & 6 deletions crates/executor/src/operators/aggregate/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -83,11 +83,50 @@ impl Operator {
kind: Kind::Aggregate {
groups,
measures: measures.into_iter().map(|(_, r)| r).collect(),
filters: vec![],
},
inputs: vec![input],
output: schema(fields),
})
}

/// Give each measure an optional row filter (SQL `FILTER (WHERE …)`).
/// A group is formed from all its rows, so a group with no matching row
/// still appears, with the measure's empty-input value: 0 for counts,
/// NULL for SQL SUM/AVG/MIN/MAX. No filter at all keeps the plain form.
pub fn with_measure_filters(mut self, filters: Vec<Option<Expression>>) -> Result<Self, Error> {
let Kind::Aggregate {
filters: slot,
measures,
groups,
} = &mut self.kind
else {
return Err(invalid("measure filters need an aggregate"));
};
if filters.iter().all(Option::is_none) {
return Ok(self);
}
if filters.len() != measures.len() {
return Err(invalid("one filter slot per aggregate measure required"));
}
let mut output = (*self.output).clone();
for (index, filter) in filters.iter().enumerate() {
let Some(filter) = filter else { continue };
if filter.dtype(&self.inputs[0])?.0 != DataType::Bool {
return Err(invalid("measure filter must be boolean"));
}
if matches!(
measures[index],
Reduction::Sum(_) | Reduction::Avg(_) | Reduction::Min(_) | Reduction::Max(_)
) && !self.inputs[0].has_promql_series_identity()
{
output.fields[groups.len() + index].nullable = true;
}
}
*slot = filters;
self.output = Arc::new(output);
Ok(self)
}
pub fn window(
input: SchemaRef,
intent: planner_types::ir::operator::AggIntent<ColumnRef>,
Expand Down Expand Up @@ -189,7 +228,7 @@ pub(super) fn execute<'a>(
Kind::SQLWindowSum { column } => {
let mut work = Cooperative::new(&context);
let total = reduce_one(
&rows,
&rows.iter().collect::<Vec<_>>(),
&Reduction::Sum(*column),
&operator.inputs[0],
&mut work,
Expand Down Expand Up @@ -224,8 +263,20 @@ pub(super) fn execute<'a>(
)
.await?
}
Kind::Aggregate { groups, measures } => {
reduce(rows, groups, measures, &operator.inputs[0], &context).await?
Kind::Aggregate {
groups,
measures,
filters,
} => {
reduce(
rows,
groups,
measures,
filters,
&operator.inputs[0],
&context,
)
.await?
}
_ => unreachable!(),
};
Expand All @@ -239,6 +290,7 @@ async fn reduce(
rows: Vec<Vec<Value>>,
groups: &[usize],
measures: &[Reduction],
filters: &[Option<Expression>],
input: &SchemaRef,
context: &RunContext,
) -> Result<Vec<Vec<Value>>, Error> {
Expand All @@ -265,8 +317,21 @@ async fn reduce(
.iter()
.map(|&i| rows[0][i].clone())
.collect::<Vec<_>>();
for measure in measures {
result.push(reduce_one(&rows, measure, input, &mut work, context).await?);
for (index, measure) in measures.iter().enumerate() {
let mut selected = Vec::with_capacity(rows.len());
for row in &rows {
// SQL FILTER keeps a row only when the predicate is true, not NULL.
if filters
.get(index)
.and_then(Option::as_ref)
.map_or(Ok(true), |filter| {
filter.evaluate(row).map(|v| matches!(v, Value::Bool(true)))
})?
{
selected.push(row);
}
}
result.push(reduce_one(&selected, measure, input, &mut work, context).await?);
}
workspace.grow(row_bytes(&result))?;
output.push(result);
Expand Down Expand Up @@ -300,7 +365,7 @@ pub(super) fn quantile(q: f64, mut values: Vec<f64>) -> f64 {
}

async fn reduce_one(
rows: &[Vec<Value>],
rows: &[&Vec<Value>],
measure: &Reduction,
input: &SchemaRef,
work: &mut Cooperative,
Expand Down
12 changes: 12 additions & 0 deletions crates/executor/src/operators/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -126,6 +126,9 @@ enum Kind {
Aggregate {
groups: Vec<usize>,
measures: Vec<Reduction>,
/// Per-measure row filters (SQL `FILTER (WHERE …)`); empty when none.
#[serde(default, skip_serializing_if = "Vec::is_empty")]
filters: Vec<Option<Expression>>,
},
SemiJoin {
keys: Vec<(usize, usize)>,
Expand All @@ -141,19 +144,28 @@ enum Kind {
value: Option<usize>,
time: Option<usize>,
groups: Vec<usize>,
/// Rows for which this is not true update no state; their group is kept.
#[serde(default, skip_serializing_if = "Option::is_none")]
filter: Option<Box<Expression>>,
},
KeyedSummaryBuild {
family: SummaryFamilyType,
value: usize,
items: Vec<usize>,
groups: Vec<usize>,
/// Rows for which this is not true update no state; their group is kept.
#[serde(default, skip_serializing_if = "Option::is_none")]
filter: Option<Box<Expression>>,
},
/// One shared state for all groups (HydraCms); `weight: None` is a unit count.
SharedSummaryBuild {
family: SummaryFamilyType,
item: usize,
weight: Option<usize>,
groups: Vec<usize>,
/// Rows for which this is not true update no state; their group is kept.
#[serde(default, skip_serializing_if = "Option::is_none")]
filter: Option<Box<Expression>>,
},
KeyedEvaluation {
state: usize,
Expand Down
55 changes: 52 additions & 3 deletions crates/executor/src/operators/summary/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,7 @@ impl Operator {
value,
items,
groups,
filter: None,
},
inputs: vec![input],
output: schema(fields),
Expand Down Expand Up @@ -98,6 +99,7 @@ impl Operator {
item,
weight,
groups,
filter: None,
},
inputs: vec![input],
output: schema(fields),
Expand Down Expand Up @@ -216,11 +218,27 @@ impl Operator {
value,
time,
groups,
filter: None,
},
inputs: vec![input],
output: schema(fields),
})
}
/// Update the summary only from rows where `filter` is true (a filtered
/// `SummaryAgg`). Every group still gets a state, so a group with no
/// matching row reads as an empty summary.
pub fn with_row_filter(mut self, filter: Expression) -> Result<Self, Error> {
if filter.dtype(&self.inputs[0])?.0 != DataType::Bool {
return Err(invalid("summary filter must be boolean"));
}
match &mut self.kind {
Kind::SummaryBuild { filter: slot, .. }
| Kind::KeyedSummaryBuild { filter: slot, .. }
| Kind::SharedSummaryBuild { filter: slot, .. } => *slot = Some(Box::new(filter)),
_ => return Err(invalid("a row filter needs a summary build")),
}
Ok(self)
}
pub fn summary_merge(
input: SchemaRef,
state: usize,
Expand Down Expand Up @@ -314,10 +332,11 @@ pub(super) fn execute<'a>(
value,
time,
groups,
filter,
} => Ok(futures::stream::once(async move {
Batch::try_new(
output,
build_summary(input, family, *value, *time, groups, !operator.inputs[0].has_promql_series_identity(), &context).await?,
build_summary(input, family, *value, *time, groups, filter.as_deref(), !operator.inputs[0].has_promql_series_identity(), &context).await?,
)
})
.boxed_local()),
Expand All @@ -326,10 +345,12 @@ pub(super) fn execute<'a>(
value,
items,
groups,
filter,
} => Ok(futures::stream::once(async move {
Batch::try_new(
output,
build_keyed_summary(input, family, *value, items, groups, &context).await?,
build_keyed_summary(input, family, *value, items, groups, filter.as_deref(), &context)
.await?,
)
})
.boxed_local()),
Expand All @@ -338,10 +359,12 @@ pub(super) fn execute<'a>(
item,
weight,
groups,
filter,
} => Ok(futures::stream::once(async move {
Batch::try_new(
output,
build_shared_summary(input, family, *item, *weight, groups, &context).await?,
build_shared_summary(input, family, *item, *weight, groups, filter.as_deref(), &context)
.await?,
)
})
.boxed_local()),
Expand Down Expand Up @@ -388,6 +411,12 @@ pub(super) fn execute<'a>(
return Err(invalid("summary value required"));
};
row[*state] = match query {
// A filtered or NULL-only group's quantile is SQL NULL.
SummaryEvaluation::Sketch(_)
if output.fields[*state].nullable && summary.is_empty() =>
{
Value::Null
}
SummaryEvaluation::Sketch(query) => {
let value = summary
.estimate(query)
Expand Down Expand Up @@ -462,12 +491,14 @@ pub(super) fn execute_merge<'a>(
.boxed_local())
}

#[allow(clippy::too_many_arguments)]
async fn build_summary(
mut input: Input<'_, Batch>,
family: &SummaryFamilyType,
value: Option<usize>,
time: Option<usize>,
groups: &[usize],
filter: Option<&Expression>,
emit_empty_global: bool,
context: &RunContext,
) -> Result<Vec<Vec<Value>>, Error> {
Expand Down Expand Up @@ -518,6 +549,9 @@ async fn build_summary(
)?;
states.insert(key.clone(), state);
}
if !selected(filter, row)? {
continue;
}
let (_, updater, memory, overhead, previous) =
states.get_mut(&key).expect("inserted group");
// SQL aggregates ignore NULL samples while retaining the group.
Expand Down Expand Up @@ -561,6 +595,13 @@ async fn build_summary(
})
.collect())
}
/// Whether a row passes a summary build's filter: only a true predicate
/// does, as with SQL `FILTER (WHERE …)`.
fn selected(filter: Option<&Expression>, row: &[Value]) -> Result<bool, Error> {
filter.map_or(Ok(true), |filter| {
Ok(matches!(filter.evaluate(row)?, Value::Bool(true)))
})
}
async fn merge_summary(
rows: Vec<Vec<Value>>,
state_column: usize,
Expand Down Expand Up @@ -630,6 +671,7 @@ async fn build_keyed_summary(
value: usize,
items: &[usize],
groups: &[usize],
filter: Option<&Expression>,
context: &RunContext,
) -> Result<Vec<Vec<Value>>, Error> {
use crate::{summary_kernels::weighted_frequency::WeightedFrequency, AggregateCore};
Expand Down Expand Up @@ -666,6 +708,9 @@ async fn build_keyed_summary(
),
);
}
if !selected(filter, row)? {
continue;
}
let (_, summary, reservation, overhead) = states.get_mut(&key).unwrap();
let Value::Float64(weight) = row[value] else {
return Err(invalid("weighted frequency weight type"));
Expand Down Expand Up @@ -701,6 +746,7 @@ async fn build_shared_summary(
item: usize,
weight: Option<usize>,
groups: &[usize],
filter: Option<&Expression>,
context: &RunContext,
) -> Result<Vec<Vec<Value>>, Error> {
use crate::summary_kernels::HydraCmsGroup;
Expand Down Expand Up @@ -729,6 +775,9 @@ async fn build_shared_summary(
memory.resize(retained)?;
seen.insert(key.clone(), (labels, name));
}
if !selected(filter, row)? {
continue;
}
let count = match weight.map(|column| &row[column]) {
None => 1,
// SQL aggregates ignore NULL weights while retaining the group.
Expand Down
Loading