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
6 changes: 5 additions & 1 deletion crates/executor/src/capability.rs
Original file line number Diff line number Diff line change
Expand Up @@ -167,7 +167,7 @@ pub fn validate_native_family(family: &SummaryFamilyType) -> Result<(), Error> {
match family {
SummaryFamilyType::ExactAggregate(..) => {}
SummaryFamilyType::Sketch(kind, _)
if matches!(kind.algorithm(), A::Kll | A::DDSketch | A::Hll) => {}
if matches!(kind.algorithm(), A::Kll | A::DDSketch | A::Hll | A::UnivMon) => {}
_ => {
return Err(Error::Invalid(
"summary family has no native DAG state implementation".into(),
Expand Down Expand Up @@ -207,6 +207,10 @@ pub fn validate_sketch_evaluation(
(A::DDSketch, _) => bare_count,
(A::Hll, SketchStatistic::Cardinality) => true,
(A::Hll, _) => bare_count,
(A::UnivMon, SketchStatistic::PointCount { value: None, .. })
| (A::UnivMon, SketchStatistic::Cardinality)
| (A::UnivMon, SketchStatistic::FrequencyL2)
| (A::UnivMon, SketchStatistic::FrequencyEntropy) => true,
// Only count intents read a Count-Min bare count, and their
// updates have unit weight; the evaluation is typed Int64 on that basis.
(A::Cms, _) => bare_count,
Expand Down
89 changes: 88 additions & 1 deletion crates/executor/src/operators/aggregate/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,17 @@ impl Operator {
for (name, reduction) in &measures {
let (t, n) = match reduction {
Reduction::Count => (DataType::Int64, false),
Reduction::Cardinality(columns) => {
if columns.is_empty() {
return Err(invalid(
"distinct aggregate requires at least one identity column",
));
}
for column in columns {
plain(&input, *column)?;
}
(DataType::Int64, false)
}
Reduction::Sum(i) | Reduction::Avg(i) => {
let (t, _) = plain(&input, *i)?;
if !matches!(t, DataType::Int64 | DataType::Float64) {
Expand All @@ -34,6 +45,15 @@ impl Operator {
}
(DataType::Float64, false)
}
Reduction::FrequencyL2(i) | Reduction::FrequencyEntropy(i) => {
if !matches!(
plain(&input, *i)?.0,
DataType::Bool | DataType::Int64 | DataType::Float64 | DataType::Utf8
) {
return Err(invalid("frequency aggregate requires a Boolean, Int64, Float64 or Utf8 identity"));
}
(DataType::Float64, false)
}
Reduction::Min(i) | Reduction::Max(i) => {
let (t, nullable) = plain(&input, *i)?;
if !ordered(t) {
Expand Down Expand Up @@ -125,10 +145,16 @@ impl Operator {
#[derive(serde::Serialize, serde::Deserialize, Clone, Debug)]
pub enum Reduction {
Count,
/// Exact distinct tuple count; a tuple with any NULL component is skipped.
Cardinality(Vec<usize>),
Sum(usize),
Avg(usize),
Min(usize),
Max(usize),
/// L2 norm of unit-update frequencies; NULL identities are skipped.
FrequencyL2(usize),
/// Shannon entropy in bits; NULL identities are skipped.
FrequencyEntropy(usize),
/// PromQL `quantile`: linear interpolation between closest ranks.
Quantile {
column: usize,
Expand Down Expand Up @@ -205,7 +231,7 @@ async fn reduce(
.map(|&i| rows[0][i].clone())
.collect::<Vec<_>>();
for measure in measures {
result.push(reduce_one(&rows, measure, input, &mut work).await?);
result.push(reduce_one(&rows, measure, input, &mut work, context).await?);
}
workspace.grow(row_bytes(&result))?;
output.push(result);
Expand Down Expand Up @@ -243,6 +269,7 @@ async fn reduce_one(
measure: &Reduction,
input: &SchemaRef,
work: &mut Cooperative,
context: &RunContext,
) -> Result<Value, Error> {
let column = match measure {
Reduction::Count => {
Expand All @@ -251,6 +278,66 @@ async fn reduce_one(
))
}
Reduction::Sum(i) | Reduction::Avg(i) | Reduction::Min(i) | Reduction::Max(i) => *i,
Reduction::Cardinality(columns) => {
let mut workspace = Workspace::new(context)?;
let mut identities = std::collections::BTreeSet::new();
for row in rows {
work.checkpoint().await?;
if columns
.iter()
.any(|column| matches!(row[*column], Value::Null))
{
continue;
}
let key = group_key(row, columns)?;
if !identities.contains(&key) {
workspace.grow(key_bytes(&key))?;
identities.insert(key);
}
}
return Ok(Value::Int64(
i64::try_from(identities.len()).map_err(|_| invalid("distinct count overflow"))?,
));
}
Reduction::FrequencyL2(column) | Reduction::FrequencyEntropy(column) => {
let mut workspace = Workspace::new(context)?;
let mut counts = BTreeMap::<Vec<u8>, u64>::new();
let mut total = 0_u64;
for row in rows {
work.checkpoint().await?;
let value = &row[*column];
if matches!(value, Value::Null) {
continue;
}
if matches!(value, Value::Float64(v) if !v.is_finite()) {
return Err(invalid(
"frequency aggregate requires a finite floating identity",
));
}
let key = value.key()?;
if !counts.contains_key(&key) {
workspace.grow(
64 + std::mem::size_of::<Vec<u8>>()
+ key.len()
+ std::mem::size_of::<u64>(),
)?;
}
*counts.entry(key).or_default() += 1;
total += 1;
}
let mut result = 0.0_f64;
for count in counts.into_values() {
work.checkpoint().await?;
let count = count as f64;
if matches!(measure, Reduction::FrequencyL2(_)) {
result = result.hypot(count);
} else {
let probability = count / total as f64;
result -= probability * probability.log2();
}
}
return Ok(Value::Float64(result));
}
Reduction::Quantile { column, q } => {
let mut values = Vec::with_capacity(rows.len());
for row in rows {
Expand Down
22 changes: 13 additions & 9 deletions crates/executor/src/operators/summary/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -99,8 +99,14 @@ impl Operator {
) -> Result<Self, Error> {
crate::values::validate_family(&family)?;
validate_groups(&input, &groups)?;
if plain(&input, value)?.0 != &DataType::Float64 {
return Err(invalid("summary numeric update requires Float64"));
let dtype = plain(&input, value)?.0;
let typed_frequency = matches!(&family,
SummaryFamilyType::Sketch(kind, _) if kind.algorithm() == &planner_types::ir::schema::SketchAlgorithm::UnivMon);
if *dtype != DataType::Float64
&& !(typed_frequency
&& matches!(dtype, DataType::Utf8 | DataType::Int64 | DataType::Bool))
{
return Err(invalid("summary update type is unsupported by its family"));
}
if let Some(time) = time {
if plain(&input, time)? != (&DataType::Timestamp, false) {
Expand Down Expand Up @@ -435,11 +441,10 @@ async fn build_summary(
states.get_mut(&key).expect("inserted group");
// SQL aggregates ignore NULL samples while retaining the group.
// A missing counter sample also contributes no observation.
let value = match row[value] {
Value::Float64(value) => value,
Value::Null => continue,
_ => return Err(invalid("summary update type")),
};
let value = &row[value];
if matches!(value, Value::Null) {
continue;
}
let timestamp = if let Some(time) = time {
let Value::Timestamp(time) = row[time] else {
return Err(invalid("summary time type"));
Expand All @@ -455,9 +460,8 @@ async fn build_summary(
));
}
updater
.validate_single_input(value)
.update_value(value, timestamp)
.map_err(Error::Operator)?;
updater.update_single(value, timestamp);
*previous = Some(timestamp);
memory.resize(updater.memory_usage_bytes() + *overhead)?;
}
Expand Down
13 changes: 13 additions & 0 deletions crates/executor/src/physical_planner/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -848,8 +848,21 @@ fn bind_operation(node: &PhysicalASAPDAGNode, inputs: &[SchemaRef]) -> Result<Op
};
let m = match m {
AggIntent::Count { .. } => Reduction::Count,
AggIntent::Cardinality { cols, .. } => {
Reduction::Cardinality(if cols.is_empty() {
vec![column(None)?]
} else {
cols.clone()
})
}
AggIntent::Sum { col } => Reduction::Sum(column(*col)?),
AggIntent::Avg { col } => Reduction::Avg(column(*col)?),
AggIntent::FrequencyL2 { col, .. } => {
Reduction::FrequencyL2(column(*col)?)
}
AggIntent::FrequencyEntropy { col, .. } => {
Reduction::FrequencyEntropy(column(*col)?)
}
AggIntent::Min { col } => Reduction::Min(column(*col)?),
AggIntent::Max { col } => Reduction::Max(column(*col)?),
_ => {
Expand Down
20 changes: 20 additions & 0 deletions crates/executor/src/summary_kernels/factory.rs
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,21 @@ pub trait AccumulatorUpdater: Send {
}
}

/// Update from the typed native row. Numeric kernels retain their existing
/// Float64 contract; frequency kernels may accept nonnumeric identities.
fn update_value(
&mut self,
value: &crate::values::Value,
timestamp_ms: i64,
) -> Result<(), String> {
let crate::values::Value::Float64(value) = value else {
return Err("summary update requires Float64".into());
};
self.validate_single_input(*value)?;
self.update_single(*value, timestamp_ms);
Ok(())
}

/// Feed a single (value, timestamp_ms) pair — for SingleSubpopulation types.
fn update_single(&mut self, value: f64, timestamp_ms: i64);

Expand Down Expand Up @@ -690,6 +705,11 @@ impl AccumulatorUpdater for HllUpdater {
}

impl AccumulatorUpdater for UnivMonUpdater {
fn update_value(&mut self, value: &crate::values::Value, _: i64) -> Result<(), String> {
self.acc
.insert_value(value)
.map_err(|error| error.to_string())
}
fn is_keyed(&self) -> bool {
false
}
Expand Down
85 changes: 76 additions & 9 deletions crates/executor/src/summary_kernels/univmon.rs
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,25 @@ impl UnivMonAccumulator {
Ok(())
}

/// SQL identities are kept in their original type: converting Int64 to
/// Float64 would collapse neighboring keys above 2^53.
pub fn insert_value(&mut self, value: &crate::values::Value) -> Result<(), Error> {
use crate::values::Value;
match value {
Value::Null => return Ok(()),
Value::Float64(number) if number.is_finite() => return self.insert_sample(*number),
Value::Bool(_) | Value::Int64(_) | Value::Utf8(_) => {}
_ => return Err("unsupported UnivMon identity type or nonfinite sample".into()),
}
let key = value.key()?;
self.inner
.bucket_size
.checked_add(1)
.ok_or("UnivMon count overflow")?;
self.inner.insert(&DataInput::Bytes(&key), 1);
Ok(())
}

fn compatible(&self, other: &Self) -> bool {
(
self.inner.heap_size,
Expand Down Expand Up @@ -113,15 +132,25 @@ impl UnivMonAccumulator {

impl AggregateCore for UnivMonAccumulator {
fn approx_memory_bytes(&self) -> usize {
std::mem::size_of::<Self>().saturating_add(
self.inner.layer_size.saturating_mul(
self.inner
.sketch_row
.saturating_mul(self.inner.sketch_col)
.saturating_mul(16)
.saturating_add(self.inner.heap_size.saturating_mul(256)),
),
)
let key_bytes = (0..self.inner.layer_size)
.flat_map(|layer| self.inner.hh_layers[layer].heap())
.map(|item| match &item.key {
asap_sketchlib::HeapItem::String(key) => key.capacity(),
asap_sketchlib::HeapItem::Bytes(key) => key.capacity(),
_ => 0,
})
.fold(0usize, usize::saturating_add);
key_bytes
.saturating_add(std::mem::size_of::<Self>())
.saturating_add(
self.inner.layer_size.saturating_mul(
self.inner
.sketch_row
.saturating_mul(self.inner.sketch_col)
.saturating_mul(16)
.saturating_add(self.inner.heap_size.saturating_mul(256)),
),
)
}
fn clone_boxed_core(&self) -> Box<dyn AggregateCore> {
Box::new(self.clone())
Expand Down Expand Up @@ -219,4 +248,42 @@ mod tests {
assert!(UnivMonAccumulator::from_sketch(terminal).is_ok());
assert!(UnivMonAccumulator::from_sketch(UnivMon::init_univmon(4, 21, 16, 2)).is_err());
}
/// Typed keys remain distinct after merging and restoring persisted native state.
#[test]
fn typed_keys_merge_and_roundtrip() {
use crate::values::Value;
let mut left = UnivMonAccumulator::new(32, 5, 1024, 4).unwrap();
let mut right = UnivMonAccumulator::new(32, 5, 1024, 4).unwrap();
for _ in 0..2 {
left.insert_value(&Value::Utf8("192.0.2.1".into())).unwrap();
right
.insert_value(&Value::Utf8("192.0.2.2".into()))
.unwrap();
}
left.merge_in_place(&right).unwrap();
let bytes = left.sketch().serialize_to_bytes().unwrap();
let restored =
UnivMonAccumulator::from_sketch(UnivMon::deserialize_from_bytes(&bytes).unwrap())
.unwrap();
for (statistic, expected) in [
(SketchStatistic::Cardinality, 2.0),
(SketchStatistic::FrequencyL2, 8.0_f64.sqrt()),
(SketchStatistic::FrequencyEntropy, 1.0),
] {
assert!((restored.estimate(&statistic).unwrap() - expected).abs() < 0.01);
}
}

/// Retained variable-length identities contribute to the runtime memory reservation.
#[test]
fn memory_accounts_for_string_identities() {
let mut state = UnivMonAccumulator::new(32, 5, 1024, 4).unwrap();
let empty = state.approx_memory_bytes();
state
.insert_value(&crate::values::Value::Utf8("x".repeat(4096).into()))
.unwrap();
assert!(state.approx_memory_bytes() >= empty + 4096);
state.clear();
assert_eq!(state.approx_memory_bytes(), empty);
}
}
16 changes: 16 additions & 0 deletions crates/executor/src/values.rs
Original file line number Diff line number Diff line change
Expand Up @@ -265,6 +265,7 @@ fn validate_state(family: &SummaryFamilyType, state: &dyn AggregateCore) -> Resu
use crate::summary_kernels::{
count_min_sketch::CountMinSketchAccumulator, datasketches_kll::DatasketchesKLLAccumulator,
dd_sketch::DDSketchAccumulator, exact::ExactAccumulator, hll_sketch::HllSketchAccumulator,
univmon::UnivMonAccumulator,
};
use planner_types::ir::schema::SketchParams;
validate_family(family)?;
Expand Down Expand Up @@ -302,6 +303,21 @@ fn validate_state(family: &SummaryFamilyType, state: &dyn AggregateCore) -> Resu
.as_any()
.downcast_ref::<HllSketchAccumulator>()
.is_some_and(|s| s.inner.precision == u32::from(*precision)),
SketchParams::UnivMon {
heap_size,
sketch_rows,
sketch_cols,
layers,
} => state
.as_any()
.downcast_ref::<UnivMonAccumulator>()
.is_some_and(|s| {
let sketch = s.sketch();
sketch.heap_size == *heap_size as usize
&& sketch.sketch_row == *sketch_rows as usize
&& sketch.sketch_col == *sketch_cols as usize
&& sketch.layer_size == *layers as usize
}),
SketchParams::Cms { width, depth } => state
.as_any()
.downcast_ref::<CountMinSketchAccumulator>()
Expand Down
Loading