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
4 changes: 2 additions & 2 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

7 changes: 4 additions & 3 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -35,9 +35,10 @@ sql_utilities = { path = "asap-common/dependencies/rs/sql_utilities" }
asap_types = { path = "asap-common/dependencies/rs/asap_types" }
elastic_dsl_utilities = { path = "asap-common/dependencies/rs/elastic_dsl_utilities" }
asap_planner = { path = "asap-planner-rs" }
# Pin past crates.io 0.2.2 to pick up the CMS estimate i32::MAX clamp fix
# (ProjectASAP/asap_sketchlib#76). Revert to a published version once released.
asap_sketchlib = { git = "https://github.com/ProjectASAP/asap_sketchlib", rev = "94d76f6b772a33e82991ac284355800053ad30bf" }
# Pin past the published releases to pick up the CMS estimate i32::MAX clamp fix
# and DDSketch's zero bucket and negative store. Revert to a published version
# once released.
asap_sketchlib = { git = "https://github.com/ProjectASAP/asap_sketchlib", rev = "5158bf232894e8f816f906dc535677b56bd81d5a" }
indexmap = { version = "2.0", features = ["serde"] }

[profile.release]
Expand Down
24 changes: 24 additions & 0 deletions asap-common/dependencies/rs/asap_types/src/aggregation_config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,13 @@ pub enum AggregationConfigError {
aggregation_id: u64,
aggregation_type: AggregationType,
},
#[error(
"aggregation {aggregation_id} (DDSketch) parameter 'alpha' must be a number in (0, 1), got {value:?}"
)]
InvalidAlpha {
aggregation_id: u64,
value: Option<Value>,
},
}

#[derive(Debug, Clone, Serialize, Deserialize)]
Expand Down Expand Up @@ -85,6 +92,7 @@ impl AggregationConfig {
pub fn validate(&self) -> Result<(), AggregationConfigError> {
self.mode().map(|_| ())?;
self.validate_hll_precision()?;
self.validate_ddsketch_alpha()?;
Ok(())
}

Expand All @@ -101,6 +109,22 @@ impl AggregationConfig {
)
}

/// DDSketch needs a relative-accuracy `alpha` strictly between 0 and 1;
/// `asap_sketchlib::DDSketch::new` panics on anything else.
fn validate_ddsketch_alpha(&self) -> Result<(), AggregationConfigError> {
if self.aggregation_type != AggregationType::DDSketch {
return Ok(());
}
let value = self.parameters.get("alpha");
match value.and_then(Value::as_f64) {
Some(alpha) if alpha > 0.0 && alpha < 1.0 => Ok(()),
_ => Err(AggregationConfigError::InvalidAlpha {
aggregation_id: self.aggregation_id,
value: value.cloned(),
}),
}
}

fn validate_hll_precision(&self) -> Result<(), AggregationConfigError> {
if self.aggregation_type == AggregationType::HLL {
match self.parameters.get("precision") {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,11 @@ pub fn compatible_agg_types(stat: Statistic) -> &'static [AggregationType] {
Statistic::Min | Statistic::Max => {
&[AggregationType::MinMax, AggregationType::MultipleMinMax]
}
Statistic::Quantile => &[AggregationType::DatasketchesKLL, AggregationType::HydraKLL],
Statistic::Quantile => &[
AggregationType::DatasketchesKLL,
AggregationType::DDSketch,
AggregationType::HydraKLL,
],
Statistic::Rate | Statistic::Increase => {
&[AggregationType::Increase, AggregationType::MultipleIncrease]
}
Expand Down Expand Up @@ -474,6 +478,11 @@ mod tests {
.parameters
.insert("precision".to_string(), serde_json::Value::from(14));
}
if config.aggregation_type == AggregationType::DDSketch {
config
.parameters
.insert("alpha".to_string(), serde_json::Value::from(0.01));
}
config
}

Expand Down Expand Up @@ -678,6 +687,25 @@ mod tests {
assert_eq!(result.unwrap().aggregation_id_for_value, 3);
}

#[test]
fn quantile_matches_ddsketch() {
let configs = single_config(make_config(
4,
"lat",
"DDSketch",
"",
300_000,
"tumbling",
&[],
"",
));
let result = find_compatible_aggregation(
&configs,
&req("lat", &[Statistic::Quantile], 300_000, &[], ""),
);
assert_eq!(result.unwrap().aggregation_id_for_value, 4);
}

#[test]
fn no_match_wrong_metric() {
let configs = single_config(make_config(
Expand Down
57 changes: 57 additions & 0 deletions asap-common/dependencies/rs/asap_types/src/streaming_config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -136,6 +136,7 @@ impl Default for StreamingConfig {
mod tests {
use super::*;
use crate::aggregation_config::AggregationConfigError;
use promql_utilities::query_logics::enums::AggregationType;

#[test]
fn rejects_heap_config_with_invalid_sub_type() {
Expand Down Expand Up @@ -365,4 +366,60 @@ aggregations:
StreamingConfig::from_yaml_data(&yaml, None)
.expect("MinMax config with 'MAX' subtype must be accepted");
}

fn ddsketch_yaml(parameters_yaml: &str) -> Value {
serde_yaml::from_str(&format!(
r#"
aggregations:
- aggregationId: 1
aggregationType: DDSketch
aggregationSubType: ''
parameters:
{parameters_yaml}
labels:
grouping: [label_0]
aggregated: []
rollup: [instance]
metric: data
windowSizeMs: 60000
slideIntervalMs: 60000
windowType: tumbling
spatialFilter: ''
"#
))
.unwrap()
}

#[test]
fn accepts_ddsketch_config_with_alpha() {
let config = StreamingConfig::from_yaml_data(&ddsketch_yaml("alpha: 0.01"), None)
.expect("DDSketch config with alpha in (0, 1) must load");
let agg = config.get_aggregation_config(1).unwrap();
assert_eq!(agg.aggregation_type, AggregationType::DDSketch);
assert_eq!(agg.parameters["alpha"], serde_json::json!(0.01));
}

#[test]
fn rejects_ddsketch_config_with_missing_or_out_of_range_alpha() {
for parameters in [
"{}",
"alpha: 0",
"alpha: 1",
"alpha: -0.1",
r#"alpha: "0.01""#,
] {
let error = StreamingConfig::from_yaml_data(&ddsketch_yaml(parameters), None)
.expect_err("DDSketch config without a valid alpha must be rejected");
assert!(
matches!(
error.downcast_ref::<AggregationConfigError>(),
Some(AggregationConfigError::InvalidAlpha {
aggregation_id: 1,
..
})
),
"parameters {parameters}: unexpected error {error}"
);
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -278,6 +278,7 @@ pub enum AggregationType {
Increase,
MinMax,
DatasketchesKLL,
DDSketch,
// ---------- multi-population (keyed) ----------
MultipleSum,
MultipleIncrease,
Expand All @@ -298,6 +299,7 @@ impl AggregationType {
AggregationType::Increase => "Increase",
AggregationType::MinMax => "MinMax",
AggregationType::DatasketchesKLL => "DatasketchesKLL",
AggregationType::DDSketch => "DDSketch",
AggregationType::MultipleSum => "MultipleSum",
AggregationType::MultipleIncrease => "MultipleIncrease",
AggregationType::MultipleMinMax => "MultipleMinMax",
Expand Down Expand Up @@ -362,6 +364,7 @@ impl FromStr for AggregationType {
"Increase" => Ok(AggregationType::Increase),
"MinMax" => Ok(AggregationType::MinMax),
"DatasketchesKLL" => Ok(AggregationType::DatasketchesKLL),
"DDSketch" => Ok(AggregationType::DDSketch),
"MultipleSum" => Ok(AggregationType::MultipleSum),
"MultipleIncrease" => Ok(AggregationType::MultipleIncrease),
"MultipleMinMax" => Ok(AggregationType::MultipleMinMax),
Expand All @@ -380,6 +383,7 @@ impl FromStr for AggregationType {
"DatasketchesKLLAccumulator" | "KLL" | "kll" | "datasketches_kll" => {
Ok(AggregationType::DatasketchesKLL)
}
"DDSketchAccumulator" | "ddsketch" | "dd" => Ok(AggregationType::DDSketch),
"MultipleSumAccumulator" | "multiple_sum" => Ok(AggregationType::MultipleSum),
"MultipleIncreaseAccumulator" | "multiple_increase" => {
Ok(AggregationType::MultipleIncrease)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -68,7 +68,8 @@ pub fn does_precompute_operator_support_subpopulations(
AggregationType::Increase
| AggregationType::MinMax
| AggregationType::Sum
| AggregationType::DatasketchesKLL => false,
| AggregationType::DatasketchesKLL
| AggregationType::DDSketch => false,

// Multi-key operators
AggregationType::MultipleIncrease
Expand Down
7 changes: 7 additions & 0 deletions asap-planner-rs/src/config/input.rs
Original file line number Diff line number Diff line change
Expand Up @@ -173,6 +173,8 @@ pub struct SketchParameterOverrides {
pub count_min_sketch_with_heap: Option<CmsHeapParams>,
#[serde(rename = "DatasketchesKLL")]
pub datasketches_kll: Option<KllParams>,
#[serde(rename = "DDSketch")]
pub ddsketch: Option<DDSketchParams>,
#[serde(rename = "HydraKLL")]
pub hydra_kll: Option<HydraParams>,
#[serde(rename = "HLL")]
Expand All @@ -198,6 +200,11 @@ pub struct KllParams {
pub k: u64,
}

#[derive(Debug, Clone, Deserialize)]
pub struct DDSketchParams {
pub alpha: f64,
}

#[derive(Debug, Clone, Deserialize)]
pub struct HydraParams {
pub row_num: u64,
Expand Down
2 changes: 2 additions & 0 deletions asap-planner-rs/src/optimizer/candidate_gen.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,8 @@ const OPTIMIZER_SKIPPED_AGG_TYPES: &[AggregationType] = &[
AggregationType::Sum,
AggregationType::MinMax,
AggregationType::Increase,
// No parameter grid or cost rows yet.
AggregationType::DDSketch,
];

/// A candidate streaming config for one optimizer item, ready for cost evaluation.
Expand Down
2 changes: 2 additions & 0 deletions asap-planner-rs/src/optimizer/sketch_properties.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,8 @@ pub fn sketch_properties(t: AggregationType) -> SketchProperties {
AggregationType::Increase => p(true, false, false),
AggregationType::MinMax => p(true, false, false),
AggregationType::DatasketchesKLL => p(true, false, false),
// Bucket counts could be subtracted, but the accumulator only merges.
AggregationType::DDSketch => p(true, false, false),
AggregationType::MultipleSum => p(true, true, true),
AggregationType::MultipleIncrease => p(true, false, true),
AggregationType::MultipleMinMax => p(true, false, true),
Expand Down
37 changes: 37 additions & 0 deletions asap-planner-rs/src/planner/sketch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ const DEFAULT_CMS_DEPTH: u64 = 3;
const DEFAULT_CMS_WIDTH: u64 = 1024;
const DEFAULT_CMS_HEAP_MULT: u64 = 4;
const DEFAULT_KLL_K: u64 = 500;
const DEFAULT_DDSKETCH_ALPHA: f64 = 0.01;
const DEFAULT_HYDRA_ROW: u64 = 3;
const DEFAULT_HYDRA_COL: u64 = 1024;
const DEFAULT_HYDRA_K: u64 = 20;
Expand Down Expand Up @@ -98,6 +99,16 @@ pub fn build_sketch_parameters(
Ok(m)
}

AggregationType::DDSketch => {
let alpha = sketch_params
.and_then(|p| p.ddsketch.as_ref())
.map(|p| p.alpha)
.unwrap_or(DEFAULT_DDSKETCH_ALPHA);
let mut m = HashMap::new();
m.insert("alpha".to_string(), serde_json::json!(alpha));
Ok(m)
}

AggregationType::HLL => {
let precision = sketch_params
.and_then(|p| p.hll.as_ref())
Expand Down Expand Up @@ -175,3 +186,29 @@ pub fn build_sketch_parameters_from_promql(
sketch_params,
)
}

#[cfg(test)]
mod tests {
use super::*;
use crate::config::input::DDSketchParams;

#[test]
fn ddsketch_uses_default_alpha_without_override() {
let params =
build_sketch_parameters(AggregationType::DDSketch, "", None, None, None).unwrap();
assert_eq!(params.len(), 1);
assert_eq!(params["alpha"], serde_json::json!(DEFAULT_DDSKETCH_ALPHA));
}

#[test]
fn ddsketch_alpha_override_is_applied() {
let overrides = SketchParameterOverrides {
ddsketch: Some(DDSketchParams { alpha: 0.02 }),
..Default::default()
};
let params =
build_sketch_parameters(AggregationType::DDSketch, "", None, None, Some(&overrides))
.unwrap();
assert_eq!(params["alpha"], serde_json::json!(0.02));
}
}
11 changes: 10 additions & 1 deletion asap-query-engine/src/engines/merge_utils.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,14 +3,15 @@
//! (`NaiveMerger::merge_all`), so the two stay behaviorally identical.
//!
//! Tries a batch merge for accumulator types that support one (currently
//! `DatasketchesKLL` and `CountMinSketch`), falling back to a sequential
//! `DatasketchesKLL`, `DDSketch` and `CountMinSketch`), falling back to a sequential
//! pairwise fold otherwise or if the batch merge itself fails. The fold
//! aborts on the first `merge_with` error instead of skipping it, so a
//! caller can't get a silently-partial merge back as `Ok`.

use crate::data_model::{AggregateCore, AggregationType};
use crate::precompute_operators::count_min_sketch_accumulator::CountMinSketchAccumulator;
use crate::precompute_operators::datasketches_kll_accumulator::DatasketchesKLLAccumulator;
use crate::precompute_operators::ddsketch_accumulator::DDSketchAccumulator;
use tracing::warn;

/// Precondition: `accumulators` is non-empty. Callers already special-case
Expand Down Expand Up @@ -38,6 +39,14 @@ pub(crate) fn merge_accumulators_batch(
e
),
}
} else if accumulator_type == AggregationType::DDSketch {
match DDSketchAccumulator::merge_multiple(accumulators) {
Ok(merged) => return Ok(Box::new(merged)),
Err(e) => warn!(
"Batch merge failed: {}. Falling back to sequential merge.",
e
),
}
} else if accumulator_type == AggregationType::CountMinSketch {
match CountMinSketchAccumulator::merge_multiple(accumulators) {
Ok(merged) => return Ok(Box::new(merged)),
Expand Down
Loading
Loading