diff --git a/Cargo.lock b/Cargo.lock index dc602750..c9863683 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -154,8 +154,8 @@ dependencies = [ [[package]] name = "asap_sketchlib" -version = "0.2.2" -source = "git+https://github.com/ProjectASAP/asap_sketchlib?rev=94d76f6b772a33e82991ac284355800053ad30bf#94d76f6b772a33e82991ac284355800053ad30bf" +version = "0.3.0" +source = "git+https://github.com/ProjectASAP/asap_sketchlib?rev=5158bf232894e8f816f906dc535677b56bd81d5a#5158bf232894e8f816f906dc535677b56bd81d5a" dependencies = [ "bytes", "prost", diff --git a/Cargo.toml b/Cargo.toml index f58970c5..06e9fca7 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -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] diff --git a/asap-common/dependencies/rs/asap_types/src/aggregation_config.rs b/asap-common/dependencies/rs/asap_types/src/aggregation_config.rs index 62baac29..cd8553dc 100644 --- a/asap-common/dependencies/rs/asap_types/src/aggregation_config.rs +++ b/asap-common/dependencies/rs/asap_types/src/aggregation_config.rs @@ -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, + }, } #[derive(Debug, Clone, Serialize, Deserialize)] @@ -85,6 +92,7 @@ impl AggregationConfig { pub fn validate(&self) -> Result<(), AggregationConfigError> { self.mode().map(|_| ())?; self.validate_hll_precision()?; + self.validate_ddsketch_alpha()?; Ok(()) } @@ -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") { diff --git a/asap-common/dependencies/rs/asap_types/src/capability_matching.rs b/asap-common/dependencies/rs/asap_types/src/capability_matching.rs index a4f72bee..c6f754ff 100644 --- a/asap-common/dependencies/rs/asap_types/src/capability_matching.rs +++ b/asap-common/dependencies/rs/asap_types/src/capability_matching.rs @@ -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] } @@ -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 } @@ -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( diff --git a/asap-common/dependencies/rs/asap_types/src/streaming_config.rs b/asap-common/dependencies/rs/asap_types/src/streaming_config.rs index e73dbe37..88d4a8f9 100644 --- a/asap-common/dependencies/rs/asap_types/src/streaming_config.rs +++ b/asap-common/dependencies/rs/asap_types/src/streaming_config.rs @@ -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() { @@ -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::(), + Some(AggregationConfigError::InvalidAlpha { + aggregation_id: 1, + .. + }) + ), + "parameters {parameters}: unexpected error {error}" + ); + } + } } diff --git a/asap-common/dependencies/rs/promql_utilities/src/query_logics/enums.rs b/asap-common/dependencies/rs/promql_utilities/src/query_logics/enums.rs index 4ceb7793..a1fa5144 100644 --- a/asap-common/dependencies/rs/promql_utilities/src/query_logics/enums.rs +++ b/asap-common/dependencies/rs/promql_utilities/src/query_logics/enums.rs @@ -278,6 +278,7 @@ pub enum AggregationType { Increase, MinMax, DatasketchesKLL, + DDSketch, // ---------- multi-population (keyed) ---------- MultipleSum, MultipleIncrease, @@ -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", @@ -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), @@ -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) diff --git a/asap-common/dependencies/rs/promql_utilities/src/query_logics/logics.rs b/asap-common/dependencies/rs/promql_utilities/src/query_logics/logics.rs index 0f94fd9d..595ddc27 100644 --- a/asap-common/dependencies/rs/promql_utilities/src/query_logics/logics.rs +++ b/asap-common/dependencies/rs/promql_utilities/src/query_logics/logics.rs @@ -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 diff --git a/asap-planner-rs/src/config/input.rs b/asap-planner-rs/src/config/input.rs index 8e260323..ccaf088a 100644 --- a/asap-planner-rs/src/config/input.rs +++ b/asap-planner-rs/src/config/input.rs @@ -173,6 +173,8 @@ pub struct SketchParameterOverrides { pub count_min_sketch_with_heap: Option, #[serde(rename = "DatasketchesKLL")] pub datasketches_kll: Option, + #[serde(rename = "DDSketch")] + pub ddsketch: Option, #[serde(rename = "HydraKLL")] pub hydra_kll: Option, #[serde(rename = "HLL")] @@ -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, diff --git a/asap-planner-rs/src/optimizer/candidate_gen.rs b/asap-planner-rs/src/optimizer/candidate_gen.rs index 53e46223..42edfdec 100644 --- a/asap-planner-rs/src/optimizer/candidate_gen.rs +++ b/asap-planner-rs/src/optimizer/candidate_gen.rs @@ -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. diff --git a/asap-planner-rs/src/optimizer/sketch_properties.rs b/asap-planner-rs/src/optimizer/sketch_properties.rs index 435d93a2..f914e3a0 100644 --- a/asap-planner-rs/src/optimizer/sketch_properties.rs +++ b/asap-planner-rs/src/optimizer/sketch_properties.rs @@ -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), diff --git a/asap-planner-rs/src/planner/sketch.rs b/asap-planner-rs/src/planner/sketch.rs index c41e8b9a..2391c0a6 100644 --- a/asap-planner-rs/src/planner/sketch.rs +++ b/asap-planner-rs/src/planner/sketch.rs @@ -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; @@ -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()) @@ -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)); + } +} diff --git a/asap-query-engine/src/engines/merge_utils.rs b/asap-query-engine/src/engines/merge_utils.rs index 6740cc1f..34be9d4a 100644 --- a/asap-query-engine/src/engines/merge_utils.rs +++ b/asap-query-engine/src/engines/merge_utils.rs @@ -3,7 +3,7 @@ //! (`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`. @@ -11,6 +11,7 @@ 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 @@ -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)), diff --git a/asap-query-engine/src/engines/simple_engine/mod.rs b/asap-query-engine/src/engines/simple_engine/mod.rs index 347d5d1a..c0bb98bd 100644 --- a/asap-query-engine/src/engines/simple_engine/mod.rs +++ b/asap-query-engine/src/engines/simple_engine/mod.rs @@ -4067,7 +4067,8 @@ mod merge_accumulators_regression_tests_596 { use crate::engines::simple_engine::SimpleEngine; use crate::engines::window_merger::{NaiveMerger, WindowMerger}; use crate::precompute_operators::{ - AccumulatorError, CountMinSketchAccumulator, DatasketchesKLLAccumulator, SumAccumulator, + AccumulatorError, CountMinSketchAccumulator, DDSketchAccumulator, + DatasketchesKLLAccumulator, SumAccumulator, }; use crate::stores::{Store, TimestampedBucketsMap}; use crate::tests::test_utilities::{ @@ -4276,6 +4277,44 @@ mod merge_accumulators_regression_tests_596 { ); } + #[test] + fn merge_accumulators_naive_merger_and_oracle_agree_on_ddsketch_batch() { + let boxes: Vec> = (0..3) + .map(|chunk| { + let mut dd = DDSketchAccumulator::new(0.01); + for i in 1..=100 { + dd.update((chunk * 100 + i) as f64); + } + Box::new(dd) as Box + }) + .collect(); + + let oracle = oracle_sequential_fold(&boxes); + let engine_result = test_engine() + .merge_accumulators(boxes.to_vec()) + .expect("SimpleEngine::merge_accumulators should merge a same-typed DDSketch batch"); + let naive_result = naive_merger_result(boxes.to_vec()) + .expect("NaiveMerger should merge a same-typed DDSketch batch"); + + let as_dd = |acc: &dyn AggregateCore| { + acc.as_any() + .downcast_ref::() + .unwrap() + .clone() + }; + let (oracle, engine, naive) = ( + as_dd(oracle.as_ref()), + as_dd(engine_result.as_ref()), + as_dd(naive_result.as_ref()), + ); + assert_eq!(engine.count(), oracle.count()); + assert_eq!(naive.count(), oracle.count()); + for q in [0.0, 0.5, 0.99, 1.0] { + assert_eq!(engine.get_quantile(q), oracle.get_quantile(q), "q={q}"); + assert_eq!(naive.get_quantile(q), oracle.get_quantile(q), "q={q}"); + } + } + // ---- Property 2: fold order / non-commutativity ---- /// Mock accumulator whose `merge_with` concatenates logs and adopts the diff --git a/asap-query-engine/src/precompute_engine/accumulator_factory.rs b/asap-query-engine/src/precompute_engine/accumulator_factory.rs index 9a1f1008..e55149b5 100644 --- a/asap-query-engine/src/precompute_engine/accumulator_factory.rs +++ b/asap-query-engine/src/precompute_engine/accumulator_factory.rs @@ -1,9 +1,9 @@ use crate::data_model::{AggregateCore, AggregationType, KeyByLabelValues, Measurement}; use crate::precompute_operators::{ - CountMinSketchAccumulator, CountMinSketchWithHeapAccumulator, DatasketchesKLLAccumulator, - DeltaSetAggregatorAccumulator, HllAccumulator, HydraKllSketchAccumulator, IncreaseAccumulator, - MinMaxAccumulator, MultipleIncreaseAccumulator, MultipleMinMaxAccumulator, - MultipleSumAccumulator, SetAggregatorAccumulator, SumAccumulator, + CountMinSketchAccumulator, CountMinSketchWithHeapAccumulator, DDSketchAccumulator, + DatasketchesKLLAccumulator, DeltaSetAggregatorAccumulator, HllAccumulator, + HydraKllSketchAccumulator, IncreaseAccumulator, MinMaxAccumulator, MultipleIncreaseAccumulator, + MultipleMinMaxAccumulator, MultipleSumAccumulator, SetAggregatorAccumulator, SumAccumulator, }; use asap_types::aggregation_config::AggregationConfig; use asap_types::aggregation_mode::{AggregationMode, CountMode, MinMaxMode}; @@ -299,6 +299,51 @@ impl AccumulatorUpdater for KllAccumulatorUpdater { } } +// --------------------------------------------------------------------------- +// DDSketchAccumulatorUpdater +// --------------------------------------------------------------------------- + +pub struct DDSketchAccumulatorUpdater { + acc: DDSketchAccumulator, + alpha: f64, +} + +impl DDSketchAccumulatorUpdater { + pub fn new(alpha: f64) -> Self { + Self { + acc: DDSketchAccumulator::new(alpha), + alpha, + } + } +} + +impl AccumulatorUpdater for DDSketchAccumulatorUpdater { + fn update_single(&mut self, value: f64, _timestamp_ms: i64) { + self.acc.update(value); + } + + fn update_keyed(&mut self, _key: &KeyByLabelValues, value: f64, timestamp_ms: i64) { + self.update_single(value, timestamp_ms); + } + + impl_accumulator_methods!(acc); + + fn reset(&mut self) { + self.acc = DDSketchAccumulator::new(self.alpha); + } + + fn is_keyed(&self) -> bool { + false + } + + fn memory_usage_bytes(&self) -> usize { + // One u64 counter per bucket in the positive and negative dense stores. + std::mem::size_of::() + + std::mem::size_of_val(self.acc.inner.store_counts()) + + std::mem::size_of_val(self.acc.inner.negative_store_counts()) + } +} + // --------------------------------------------------------------------------- // HllAccumulatorUpdater // --------------------------------------------------------------------------- @@ -820,6 +865,13 @@ fn kll_k_param(config: &AggregationConfig) -> Result { .ok_or_else(|| "KLL config missing required parameter (tried: K, k)".to_string()) } +/// Extract the DDSketch relative-accuracy `alpha` parameter. +fn ddsketch_alpha_param(config: &AggregationConfig) -> f64 { + config.parameters["alpha"] + .as_f64() + .expect("validation guarantees an alpha in (0, 1)") +} + /// Extract `(row_num, col_num)` for CMS / HydraKLL configs. /// /// Accepts the planner-canonical `depth`/`width` names first, then falls back @@ -918,6 +970,9 @@ pub fn create_accumulator_updater( AggregationType::DatasketchesKLL => { Ok(Box::new(KllAccumulatorUpdater::new(kll_k_param(config)?))) } + AggregationType::DDSketch => Ok(Box::new(DDSketchAccumulatorUpdater::new( + ddsketch_alpha_param(config), + ))), AggregationType::MultipleSum => Ok(Box::new(MultipleSumAccumulatorUpdater::new( count_events(config), ))), @@ -1129,6 +1184,7 @@ mod tests { }; for (agg_type, sub_type, params) in [ (AggregationType::DatasketchesKLL, "", kll_params_required()), + (AggregationType::DDSketch, "", ddsketch_params_required()), ( AggregationType::CountMinSketch, "sum", @@ -1445,6 +1501,83 @@ mod tests { p } + fn ddsketch_params_required() -> std::collections::HashMap { + std::collections::HashMap::from([("alpha".to_string(), serde_json::json!(0.01))]) + } + + fn ddsketch_config( + params: std::collections::HashMap, + ) -> AggregationConfig { + AggregationConfig::new( + 7, + AggregationType::DDSketch, + String::new(), + params, + promql_utilities::data_model::key_by_label_names::KeyByLabelNames::new(vec![]), + promql_utilities::data_model::key_by_label_names::KeyByLabelNames::new(vec![]), + promql_utilities::data_model::key_by_label_names::KeyByLabelNames::new(vec![]), + String::new(), + 60_000, + 0, + WindowType::Tumbling, + "m".to_string(), + "m".to_string(), + None, + None, + None, + None, + ) + } + + #[test] + fn test_ddsketch_updater_via_factory_uses_configured_alpha() { + let config = ddsketch_config(std::collections::HashMap::from([( + "alpha".to_string(), + serde_json::json!(0.02), + )])); + let mut updater = create_accumulator_updater(&config).unwrap(); + assert!(!updater.is_keyed()); + for i in 1..=100 { + updater.update_single(i as f64, i * 1000); + } + let acc = updater.take_accumulator(); + let dd = acc + .as_any() + .downcast_ref::() + .expect("AggregationType::DDSketch → DDSketchAccumulator"); + assert_eq!(dd.alpha(), 0.02); + assert_eq!(dd.count(), 100); + } + + #[test] + fn test_ddsketch_updater_reset_clears_state() { + let mut updater = DDSketchAccumulatorUpdater::new(0.01); + for i in 1..=50 { + updater.update_single(i as f64, 0); + } + updater.reset(); + let acc = updater.take_accumulator(); + let dd = acc.as_any().downcast_ref::().unwrap(); + assert_eq!(dd.count(), 0); + assert_eq!(dd.alpha(), 0.01); + } + + #[test] + fn test_ddsketch_missing_or_invalid_alpha_returns_err() { + for params in [ + std::collections::HashMap::new(), + std::collections::HashMap::from([("alpha".to_string(), serde_json::json!(0.0))]), + std::collections::HashMap::from([("alpha".to_string(), serde_json::json!(1.0))]), + std::collections::HashMap::from([("alpha".to_string(), serde_json::json!("0.01"))]), + ] { + let config = ddsketch_config(params.clone()); + assert!( + create_accumulator_updater(&config).is_err(), + "params {params:?} must be rejected" + ); + } + } + fn cms_params_required() -> std::collections::HashMap { let mut p = std::collections::HashMap::new(); p.insert("row_num".to_string(), serde_json::json!(4_u64)); diff --git a/asap-query-engine/src/precompute_operators/ddsketch_accumulator.rs b/asap-query-engine/src/precompute_operators/ddsketch_accumulator.rs new file mode 100644 index 00000000..b8041eff --- /dev/null +++ b/asap-query-engine/src/precompute_operators/ddsketch_accumulator.rs @@ -0,0 +1,347 @@ +use crate::data_model::{ + AggregateCore, AggregationType, MergeableAccumulator, SerializableToSink, + SingleSubpopulationAggregate, +}; +use asap_sketchlib::DDSketch; +use asap_types::traits::SerializationError; +use base64::{engine::general_purpose, Engine as _}; +use serde_json::Value; +use std::collections::HashMap; + +use promql_utilities::query_logics::enums::Statistic; + +/// DDSketch accumulator — wraps `asap_sketchlib::DDSketch`, the same +/// implementation sketch-bench measures for the planner's cost table. +/// +/// `alpha` is the relative-accuracy bound: a returned quantile `x̂` of a true +/// value `x` satisfies `|x̂ - x| <= alpha * |x|`. Negative values go to a +/// mirrored store, and zeros and magnitudes too small to index go to a zero +/// bucket. Non-finite values and magnitudes too large to index are dropped. +#[derive(Clone, Debug)] +pub struct DDSketchAccumulator { + pub inner: DDSketch, +} + +impl DDSketchAccumulator { + /// Panics unless `alpha` is in `(0, 1)`; configs are validated before this + /// is reached (`AggregationConfig::validate`). + pub fn new(alpha: f64) -> Self { + Self { + inner: DDSketch::new(alpha), + } + } + + pub fn update(&mut self, value: f64) { + self.inner.add(&value); + } + + pub fn count(&self) -> u64 { + self.inner.get_count() + } + + pub fn alpha(&self) -> f64 { + self.inner.alpha() + } + + /// `NaN` for an empty sketch, matching an empty window having no quantile. + pub fn get_quantile(&self, quantile: f64) -> f64 { + self.inner + .get_value_at_quantile(quantile) + .unwrap_or(f64::NAN) + } + + fn merge_inner(&mut self, other: &DDSketchAccumulator) -> Result<(), String> { + self.inner.merge(&other.inner) + } + + /// Merges a batch in place into one copy of the first sketch, instead of + /// cloning the running result at every step as `merge_with` does. + pub fn merge_multiple( + accumulators: &[Box], + ) -> Result> { + let mut dds = accumulators.iter().map(|acc| { + acc.as_any() + .downcast_ref::() + .ok_or_else(|| { + format!( + "Cannot merge DDSketchAccumulator with {}", + acc.get_accumulator_type() + ) + }) + }); + let mut merged = dds.next().ok_or("No accumulators to merge")??.clone(); + for dd in dds { + merged.merge_inner(dd?)?; + } + Ok(merged) + } +} + +impl SerializableToSink for DDSketchAccumulator { + fn serialize_to_json(&self) -> Result { + let sketch_bytes = self.serialize_to_bytes()?; + let sketch_b64 = general_purpose::STANDARD.encode(&sketch_bytes); + Ok(serde_json::json!({ "sketch": sketch_b64 })) + } + + fn serialize_to_bytes(&self) -> Result, SerializationError> { + self.inner + .serialize_to_bytes() + .map_err(|e| SerializationError::Bytes { + type_name: "DDSketchAccumulator", + source: e.to_string().into(), + }) + } +} + +impl AggregateCore for DDSketchAccumulator { + fn clone_boxed_core(&self) -> Box { + Box::new(self.clone()) + } + + fn type_name(&self) -> &'static str { + "DDSketchAccumulator" + } + + fn as_any(&self) -> &dyn std::any::Any { + self + } + + fn merge_with( + &self, + other: &dyn AggregateCore, + ) -> Result, Box> { + let other_dd = other + .as_any() + .downcast_ref::() + .ok_or_else(|| { + format!( + "Cannot merge DDSketchAccumulator with {}", + other.get_accumulator_type() + ) + })?; + let mut merged = self.clone(); + merged.merge_inner(other_dd)?; + Ok(Box::new(merged)) + } + + fn get_accumulator_type(&self) -> AggregationType { + AggregationType::DDSketch + } + + fn get_keys(&self) -> Option> { + None + } + + fn query_statistic( + &self, + statistic: Statistic, + _key: &Option, + query_kwargs: &HashMap, + ) -> Result> { + self.query(statistic, Some(query_kwargs)) + } +} + +impl SingleSubpopulationAggregate for DDSketchAccumulator { + fn query( + &self, + statistic: Statistic, + query_kwargs: Option<&HashMap>, + ) -> Result> { + match statistic { + Statistic::Quantile => { + let quantile = query_kwargs + .and_then(|kwargs| kwargs.get("quantile")) + .ok_or("Missing quantile parameter for quantile query")? + .parse::() + .map_err(|_| "Invalid quantile parameter format")?; + + if !(0.0..=1.0).contains(&quantile) { + return Err("Quantile must be between 0.0 and 1.0".into()); + } + + Ok(self.get_quantile(quantile)) + } + _ => Err(format!("Unsupported statistic in DDSketchAccumulator: {statistic:?}").into()), + } + } + + fn clone_boxed(&self) -> Box { + Box::new(self.clone()) + } +} + +impl MergeableAccumulator for DDSketchAccumulator { + fn merge_accumulators( + accumulators: Vec, + ) -> Result> { + let mut iter = accumulators.into_iter(); + let mut merged = iter.next().ok_or("No accumulators to merge")?; + for acc in iter { + merged.merge_inner(&acc)?; + } + Ok(merged) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + const ALPHA: f64 = 0.01; + + /// Relative error of `estimate` against `truth`, the bound DDSketch guarantees. + fn rel_err(estimate: f64, truth: f64) -> f64 { + (estimate - truth).abs() / truth + } + + #[test] + fn quantiles_are_within_relative_accuracy() { + let mut dd = DDSketchAccumulator::new(ALPHA); + for i in 1..=10_000 { + dd.update(i as f64); + } + assert_eq!(dd.count(), 10_000); + // Rank q·n of 1..=n is q·n itself, so the true quantile is known exactly. + for q in [0.5, 0.75, 0.9, 0.95, 0.99] { + let truth = (q * 10_000.0_f64).ceil(); + let estimate = dd.get_quantile(q); + assert!( + rel_err(estimate, truth) <= ALPHA, + "q={q}: estimate {estimate} vs truth {truth}" + ); + } + } + + #[test] + fn zero_and_negative_values_are_counted() { + let mut dd = DDSketchAccumulator::new(ALPHA); + dd.update(-5.0); + dd.update(0.0); + dd.update(f64::NAN); + dd.update(42.0); + assert_eq!(dd.count(), 3); + assert_eq!(dd.get_quantile(0.0), -5.0); + assert_eq!(dd.get_quantile(0.5), 0.0); + assert_eq!(dd.get_quantile(1.0), 42.0); + } + + #[test] + fn empty_sketch_quantile_is_nan() { + assert!(DDSketchAccumulator::new(ALPHA).get_quantile(0.5).is_nan()); + } + + #[test] + fn query_reads_quantile_kwarg_and_rejects_other_statistics() { + let mut dd = DDSketchAccumulator::new(ALPHA); + for i in 1..=100 { + dd.update(i as f64); + } + let mut kwargs = HashMap::new(); + kwargs.insert("quantile".to_string(), "0.5".to_string()); + let median = dd.query(Statistic::Quantile, Some(&kwargs)).unwrap(); + assert!(rel_err(median, 50.0) <= ALPHA, "median {median}"); + + kwargs.insert("quantile".to_string(), "1.5".to_string()); + assert!(dd.query(Statistic::Quantile, Some(&kwargs)).is_err()); + assert!(dd.query(Statistic::Quantile, None).is_err()); + assert!(dd.query(Statistic::Sum, Some(&kwargs)).is_err()); + } + + #[test] + fn merge_equals_single_sketch_over_all_values() { + let mut whole = DDSketchAccumulator::new(ALPHA); + let mut parts = vec![ + DDSketchAccumulator::new(ALPHA), + DDSketchAccumulator::new(ALPHA), + DDSketchAccumulator::new(ALPHA), + ]; + for i in 1..=3_000 { + whole.update(i as f64); + parts[i % 3].update(i as f64); + } + + let merged = DDSketchAccumulator::merge_accumulators(parts.clone()).unwrap(); + assert_eq!(merged.count(), whole.count()); + for q in [0.0, 0.5, 0.99, 1.0] { + assert_eq!(merged.get_quantile(q), whole.get_quantile(q), "q={q}"); + } + + // The trait-object path the window merger uses agrees with the direct merge. + let boxed = parts[0].merge_with(&parts[1]).unwrap(); + let boxed = boxed.merge_with(&parts[2]).unwrap(); + let via_trait = boxed + .as_any() + .downcast_ref::() + .unwrap(); + assert_eq!(via_trait.get_quantile(0.5), whole.get_quantile(0.5)); + } + + #[test] + fn merge_multiple_equals_single_sketch_and_rejects_bad_batches() { + let mut whole = DDSketchAccumulator::new(ALPHA); + let mut parts: Vec> = Vec::new(); + for chunk in 0..4 { + let mut part = DDSketchAccumulator::new(ALPHA); + for i in 1..=250 { + let value = (chunk * 250 + i) as f64; + whole.update(value); + part.update(value); + } + parts.push(Box::new(part)); + } + let merged = DDSketchAccumulator::merge_multiple(&parts).unwrap(); + assert_eq!(merged.count(), whole.count()); + for q in [0.0, 0.5, 0.99, 1.0] { + assert_eq!(merged.get_quantile(q), whole.get_quantile(q), "q={q}"); + } + // The batch merge leaves its inputs untouched. + let first = parts[0] + .as_any() + .downcast_ref::() + .unwrap(); + assert_eq!(first.count(), 250); + + assert!(DDSketchAccumulator::merge_multiple(&[]).is_err()); + parts.push(Box::new(DDSketchAccumulator::new(0.02))); + assert!(DDSketchAccumulator::merge_multiple(&parts).is_err()); + parts.pop(); + parts.push(Box::new(crate::precompute_operators::SumAccumulator::new())); + assert!(DDSketchAccumulator::merge_multiple(&parts).is_err()); + } + + #[test] + fn merge_rejects_mismatched_alpha_and_type() { + let mut a = DDSketchAccumulator::new(0.01); + let mut b = DDSketchAccumulator::new(0.02); + a.update(1.0); + b.update(2.0); + assert!(DDSketchAccumulator::merge_accumulators(vec![a.clone(), b.clone()]).is_err()); + assert!(a.merge_with(&b).is_err()); + + let other = crate::precompute_operators::SumAccumulator::new(); + assert!(a.merge_with(&other).is_err()); + } + + #[test] + fn serializes_to_bytes_and_json() { + let mut dd = DDSketchAccumulator::new(ALPHA); + dd.update(3.0); + let bytes = dd.serialize_to_bytes().unwrap(); + let restored = DDSketch::deserialize_from_bytes(&bytes).unwrap(); + assert_eq!(restored.get_count(), 1); + + let json = dd.serialize_to_json().unwrap(); + let b64 = json["sketch"].as_str().unwrap(); + assert_eq!(general_purpose::STANDARD.decode(b64).unwrap(), bytes); + } + + #[test] + fn reports_type() { + let dd = DDSketchAccumulator::new(ALPHA); + assert_eq!(dd.type_name(), "DDSketchAccumulator"); + assert_eq!(dd.get_accumulator_type(), AggregationType::DDSketch); + assert!(dd.get_keys().is_none()); + } +} diff --git a/asap-query-engine/src/precompute_operators/mod.rs b/asap-query-engine/src/precompute_operators/mod.rs index 36ab865b..b7e9424c 100644 --- a/asap-query-engine/src/precompute_operators/mod.rs +++ b/asap-query-engine/src/precompute_operators/mod.rs @@ -1,6 +1,7 @@ pub mod count_min_sketch_accumulator; pub mod count_min_sketch_with_heap_accumulator; pub mod datasketches_kll_accumulator; +pub mod ddsketch_accumulator; pub mod delta_set_aggregator_accumulator; pub mod error; pub mod hll_accumulator; @@ -16,6 +17,7 @@ pub mod sum_accumulator; pub use count_min_sketch_accumulator::*; pub use count_min_sketch_with_heap_accumulator::*; pub use datasketches_kll_accumulator::*; +pub use ddsketch_accumulator::*; pub use delta_set_aggregator_accumulator::*; pub use error::AccumulatorError; pub use hll_accumulator::*; diff --git a/asap-query-engine/tests/e2e_precompute_equivalence.rs b/asap-query-engine/tests/e2e_precompute_equivalence.rs index 5d6592a5..7ff11222 100644 --- a/asap-query-engine/tests/e2e_precompute_equivalence.rs +++ b/asap-query-engine/tests/e2e_precompute_equivalence.rs @@ -1701,6 +1701,89 @@ async fn e2e_grouped_quantile_preserves_output_label_shape() { ); } +/// A grouped quantile served by a DDSketch aggregation: samples from two +/// instances per job flow through remote write and precompute into one +/// DDSketch per group, and each returned quantile stays within the sketch's +/// relative-accuracy bound. Zeros are part of the population, so a sketch +/// that dropped them would answer outside the bound. +#[tokio::test] +async fn e2e_grouped_quantile_over_ddsketch_is_within_alpha() { + let port = 19428u16; + let metric = "dd_latency"; + let query = "quantile by (job) (0.9, dd_latency)"; + let alpha = 0.01; + let mut config = make_agg_config( + 17, + metric, + AggregationType::DDSketch, + "", + 1_000, + 0, + vec!["job"], + ); + config.parameters.insert("alpha".to_string(), json!(alpha)); + + // Inside the (1s, 2s] window, instance "a" sends scale * 1..=50 and "b" + // sends scale * 51..=100, and each sends ten zeros. Samples go in time + // order across instances, since the watermark is per group, and one later + // sample per series closes the window. + let groups = [("frontend", 1.0), ("backend", 10.0)]; + let instances = [("a", 1), ("b", 51)]; + let mut samples = Vec::new(); + for (job, scale) in groups { + for i in 0..50 { + for (instance, first) in instances { + let value = scale * (first + i) as f64; + let labels = vec![("job", job), ("instance", instance)]; + samples.push(make_timeseries(metric, labels, 1_005 + 5 * i, value)); + } + } + for j in 0..10 { + for (instance, _) in instances { + let labels = vec![("job", job), ("instance", instance)]; + samples.push(make_timeseries(metric, labels, 1_600 + j, 0.0)); + } + } + for (instance, _) in instances { + let labels = vec![("job", job), ("instance", instance)]; + samples.push(make_timeseries(metric, labels, 3_500, scale)); + } + } + let (engine, query) = NativeDagScenario { + port, + metric, + query, + aggregation_configs: vec![config], + schema_labels: vec!["instance".to_string(), "job".to_string()], + samples, + evaluation_time_seconds: 2.0, + base_interval_ms: 1_000, + } + .build_engine() + .await; + + let (_, result) = engine + .handle_query_promql(query, 2.0) + .expect("DDSketch quantile should execute") + .expect("DDSketch quantile should match the configured aggregation"); + let QueryResult::Vector(vector) = result else { + panic!("expected instant vector result"); + }; + assert_eq!(vector.values.len(), groups.len()); + for element in vector.values { + let job = element.labels.labels[0].as_str(); + let scale = groups.iter().find(|(g, _)| *g == job).unwrap().1; + // 120 values: 20 zeros, then scale * 1..=100. Rank ceil(0.9 * 120) = 108 + // is scale * 88; without the zeros it would be scale * 90. + let truth = scale * 88.0; + assert!( + (element.value - truth).abs() / truth <= alpha, + "{job}: estimate {} vs truth {truth}", + element.value + ); + } +} + /// Sliding precomputes keep their existing exact-cover composition while /// samples on every slide boundary move to the pane ending at that boundary. /// The shared 6s boundary must be counted once, not once per stored window.