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
89 changes: 86 additions & 3 deletions crates/executor/src/capability.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ pub fn validate_summary_kernel(
grouping: &GroupingStrategy,
) -> Result<(), String> {
if grouping != &GroupingStrategy::PerSubpopulationInstance {
return Err("shared summary grouping has no registered kernel".into());
return validate_hydra_cms(family, input, grouping);
}
let keyed = match family {
SummaryFamilyType::ExactAggregate(kind, params) => {
Expand Down Expand Up @@ -111,6 +111,76 @@ pub fn validate_summary_kernel(
Ok(())
}

/// The only shared grouping with a kernel: Hydra over Count-Min. Each update
/// adds a unit or non-negative column weight for one item of one group.
fn validate_hydra_cms(
family: &SummaryFamilyType,
input: &SummaryUpdate,
grouping: &GroupingStrategy,
) -> Result<(), String> {
use planner_types::ir::schema::{SummaryInputExpr, WeightDomain};
hydra_cms_shape(family, grouping)?;
if !matches!(input.item, Some(SummaryInputExpr::Column(_))) {
return Err("HydraCms requires one item column".into());
}
if !matches!(input.weight_domain, WeightDomain::NonNegative { .. })
|| !matches!(
input.weight,
SummaryInputExpr::Constant(1.0) | SummaryInputExpr::Column(_)
)
{
return Err("HydraCms requires a unit or non-negative column weight".into());
}
Ok(())
}

/// `(shared_rows, shared_columns, width, depth)` of a HydraCms state. Its
/// family is the per-group Count-Min it emulates, with the Hydra's own width
/// and depth.
pub(crate) fn hydra_cms_shape(
family: &SummaryFamilyType,
grouping: &GroupingStrategy,
) -> Result<(usize, usize, usize, usize), String> {
use planner_types::ir::schema::{HydraKind, HydraParams};
let GroupingStrategy::SharedMultiSubpopulation {
kind: HydraKind::HydraCms,
params:
HydraParams::HydraCms {
width,
depth,
shared_rows,
shared_columns,
},
} = grouping
else {
return Err("shared summary grouping has no registered kernel".into());
};
let SummaryFamilyType::Sketch(kind, layout) = family else {
return Err("HydraCms requires a Count-Min family".into());
};
if layout != grouping {
return Err("Planner family and operator grouping disagree".into());
}
if kind.algorithm() != &SketchAlgorithm::Cms
|| kind.params()
!= &(SketchParams::Cms {
width: *width,
depth: *depth,
})
{
return Err("HydraCms family must be Count-Min with the Hydra width and depth".into());
}
if !valid_matrix(*width, *depth) || !valid_matrix(*shared_columns, *shared_rows) {
return Err("invalid HydraCms dimensions".into());
}
Ok((
*shared_rows as usize,
*shared_columns as usize,
*width as usize,
*depth as usize,
))
}

fn valid_matrix(width: u32, depth: u32) -> bool {
// Construction uses the kernel's native row hashing, so no encoded-size
// limit applies here.
Expand Down Expand Up @@ -142,9 +212,15 @@ pub fn validate_native_family(family: &SummaryFamilyType) -> Result<(), Error> {
use planner_types::ir::schema::SketchAlgorithm as A;
if let SummaryFamilyType::Sketch(kind, grouping) = family {
// Plain Count-Min is native as stored state only: it merges and reads
// its bare count, but the DAG does not build it from rows.
// its bare count, but the DAG does not build it from rows. HydraCms
// is built by the shared summary operator.
if let (A::Cms, SketchParams::Cms { width, depth }) = (kind.algorithm(), kind.params()) {
return if valid_matrix(*width, *depth) && grouping == &Default::default() {
if grouping != &GroupingStrategy::PerSubpopulationInstance {
return hydra_cms_shape(family, grouping)
.map(|_| ())
.map_err(Error::Invalid);
}
return if valid_matrix(*width, *depth) {
Ok(())
} else {
Err(Error::Invalid(
Expand Down Expand Up @@ -211,6 +287,13 @@ pub fn validate_sketch_evaluation(
| (A::UnivMon, SketchStatistic::Cardinality)
| (A::UnivMon, SketchStatistic::FrequencyL2)
| (A::UnivMon, SketchStatistic::FrequencyEntropy) => true,
// HydraCms answers a group's item frequency as well as its total.
(A::Cms, SketchStatistic::PointCount { .. })
if matches!(family, SummaryFamilyType::Sketch(_, grouping)
if grouping != &GroupingStrategy::PerSubpopulationInstance) =>
{
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
13 changes: 12 additions & 1 deletion crates/executor/src/operators/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -147,6 +147,13 @@ enum Kind {
items: Vec<usize>,
groups: Vec<usize>,
},
/// One shared state for all groups (HydraCms); `weight: None` is a unit count.
SharedSummaryBuild {
family: SummaryFamilyType,
item: usize,
weight: Option<usize>,
groups: Vec<usize>,
},
KeyedEvaluation {
state: usize,
k: usize,
Expand Down Expand Up @@ -291,6 +298,7 @@ impl PhysicalOperator<Batch, SchemaRef> for Operator {
| Kind::SemiJoin { .. }
| Kind::SummaryBuild { .. }
| Kind::KeyedSummaryBuild { .. }
| Kind::SharedSummaryBuild { .. }
| Kind::SummaryMerge { .. }
| Kind::VectorToScalar { .. }
)
Expand Down Expand Up @@ -342,7 +350,9 @@ impl PhysicalOperator<Batch, SchemaRef> for Operator {
Kind::Window { .. } => "WindowAggregate",
Kind::SemiJoin { .. } => "SemiJoin",
Kind::Join { .. } => "RelationalJoin",
Kind::SummaryBuild { .. } | Kind::KeyedSummaryBuild { .. } => "SummaryAgg",
Kind::SummaryBuild { .. }
| Kind::KeyedSummaryBuild { .. }
| Kind::SharedSummaryBuild { .. } => "SummaryAgg",
Kind::KeyedEvaluation { .. } => "SummaryEstimate",
Kind::SummaryMerge { .. } => "SummaryMerge",
Kind::Evaluation { .. } => "SummaryEvaluation",
Expand Down Expand Up @@ -397,6 +407,7 @@ impl PhysicalOperator<Batch, SchemaRef> for Operator {
Kind::SummaryBuild { .. }
| Kind::Evaluation { .. }
| Kind::KeyedSummaryBuild { .. }
| Kind::SharedSummaryBuild { .. }
| Kind::KeyedEvaluation { .. } => summary::execute(self, inputs, context),
}
}
Expand Down
130 changes: 130 additions & 0 deletions crates/executor/src/operators/summary/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,49 @@ impl Operator {
output: schema(fields),
})
}
/// One HydraCms state shared by every group, emitted as one row per group
/// that reads it. `weight: None` counts each row once.
pub fn shared_summary_build(
input: SchemaRef,
family: SummaryFamilyType,
item: usize,
weight: Option<usize>,
groups: Vec<usize>,
) -> Result<Self, Error> {
crate::values::validate_family(&family)?;
crate::factory::create_hydra_cms(&family).map_err(Error::Invalid)?;
validate_groups(&input, &groups)?;
if !matches!(
plain(&input, item)?.0,
DataType::Utf8 | DataType::Int64 | DataType::Bool
) {
return Err(invalid("HydraCms items must be Utf8, Int64 or Bool"));
}
if weight.is_some_and(|weight| !matches!(plain(&input, weight), Ok((DataType::Float64, _))))
{
return Err(invalid("HydraCms weight must be Float64"));
}
let mut fields = groups
.iter()
.map(|&i| input.fields[i].clone())
.collect::<Vec<_>>();
fields.push(SummaryField {
name: "state".into(),
dtype: family.clone(),
nullable: false,
table: None,
});
Ok(Self {
kind: Kind::SharedSummaryBuild {
family,
item,
weight,
groups,
},
inputs: vec![input],
output: schema(fields),
})
}
pub fn keyed_evaluation(
input: SchemaRef,
state: usize,
Expand Down Expand Up @@ -264,6 +307,18 @@ pub(super) fn execute<'a>(
)
})
.boxed_local()),
Kind::SharedSummaryBuild {
family,
item,
weight,
groups,
} => Ok(futures::stream::once(async move {
Batch::try_new(
output,
build_shared_summary(input, family, *item, *weight, groups, &context).await?,
)
})
.boxed_local()),
Kind::KeyedEvaluation { state, k } => Ok(input
.map(move |batch| {
let batch = batch?;
Expand Down Expand Up @@ -610,3 +665,78 @@ async fn build_keyed_summary(
})
.collect())
}

async fn build_shared_summary(
mut input: Input<'_, Batch>,
family: &SummaryFamilyType,
item: usize,
weight: Option<usize>,
groups: &[usize],
context: &RunContext,
) -> Result<Vec<Vec<Value>>, Error> {
use crate::summary_kernels::HydraCmsGroup;
let mut grid = crate::factory::create_hydra_cms(family).map_err(Error::Invalid)?;
let grid_bytes = grid.approx_memory_bytes();
let mut memory = context.reserve(grid_bytes)?;
let mut retained = grid_bytes;
let mut work = Cooperative::new(context);
// Group labels and the group's Hydra subpopulation name, an injective
// encoding of its typed key.
let mut seen = BTreeMap::<Vec<Vec<u8>>, (Vec<Value>, String)>::new();
while let Some(batch) = input.next().await {
let batch = batch?;
for row in batch.rows() {
work.checkpoint().await?;
let key = group_key(row, groups)?;
if !seen.contains_key(&key) {
let name = key
.iter()
.map(|part| part.iter().map(|b| format!("{b:02x}")).collect::<String>())
.collect::<Vec<_>>()
.join(",");
let labels = groups.iter().map(|&i| row[i].clone()).collect::<Vec<_>>();
retained +=
key_bytes(&key) + labels.iter().map(Value::bytes).sum::<usize>() + name.len();
memory.resize(retained)?;
seen.insert(key.clone(), (labels, name));
}
let count = match weight.map(|column| &row[column]) {
None => 1,
// SQL aggregates ignore NULL weights while retaining the group.
Some(Value::Null) => continue,
Some(Value::Float64(w))
if *w >= 0.0 && w.fract() == 0.0 && *w <= f64::from(i32::MAX) =>
{
*w as i32
}
Some(_) => {
return Err(Error::Operator(
"HydraCms weights must be non-negative integers".into(),
))
}
};
let item = match &row[item] {
Value::Utf8(v) => v.to_string(),
Value::Int64(v) => v.to_string(),
Value::Bool(v) => v.to_string(),
_ => return Err(Error::Operator("HydraCms item must be non-null".into())),
};
grid.update(&seen[&key].1, &item, count)
.map_err(|e| Error::Operator(e.to_string()))?;
}
}
let grid = Arc::new(grid);
Ok(seen
.into_values()
.map(|(mut labels, group)| {
labels.push(Value::Summary {
family: family.clone(),
state: Arc::new(HydraCmsGroup {
grid: grid.clone(),
group,
}),
});
labels
})
.collect())
}
6 changes: 6 additions & 0 deletions crates/executor/src/operators/unchecked.rs
Original file line number Diff line number Diff line change
Expand Up @@ -162,6 +162,12 @@ impl TryFrom<UncheckedOperator> for Operator {
items,
groups,
} => Operator::keyed_summary_build(input(0)?, family, value, items, groups)?,
Kind::SharedSummaryBuild {
family,
item,
weight,
groups,
} => Operator::shared_summary_build(input(0)?, family, item, weight, groups)?,
Kind::KeyedEvaluation { state, k } => {
Operator::keyed_evaluation(input(0)?, state, k, output.clone())?
}
Expand Down
21 changes: 21 additions & 0 deletions crates/executor/src/physical_planner/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -923,6 +923,27 @@ fn bind_operation(node: &PhysicalASAPDAGNode, inputs: &[SchemaRef]) -> Result<Op
"filtered summary update has no native implementation",
));
}
if grouping != &planner_types::ir::schema::GroupingStrategy::PerSubpopulationInstance {
crate::capability::validate_summary_kernel(family, update, grouping)
.map_err(Error::Invalid)?;
let PlannerReduction::Reduce(keys) = reduction else {
return Err(invalid("shared summary requires explicit groups"));
};
let Some(SummaryInputExpr::Column(item)) = &update.item else {
unreachable!("validated HydraCms item column")
};
let weight = match &update.weight {
SummaryInputExpr::Column(weight) => Some(named_column(input, weight)?),
_ => None,
};
return Operator::shared_summary_build(
input.clone(),
family.clone(),
named_column(input, item)?,
weight,
groups(input, keys)?,
);
}
if let Some(item) = &update.item {
let PlannerReduction::Reduce(keys) = reduction else {
return Err(invalid("keyed summary requires explicit partitions"));
Expand Down
Loading