From 938cfcd9248b34bc459edcde38cc0bc09d669056 Mon Sep 17 00:00:00 2001 From: Zeying Zhu Date: Mon, 5 Oct 2026 19:03:54 +0000 Subject: [PATCH 1/6] feat(query-engine): add DDSketch aggregation type Add a DDSketch quantile aggregation backed by asap_sketchlib::DDSketch, the implementation sketch-bench measures for the planner's cost table. - AggregationType::DDSketch (aliases: DDSketchAccumulator, ddsketch, dd), single-population, served for Statistic::Quantile by capability matching. - DDSketchAccumulator: update, quantile query, merge, msgpack/JSON output. Non-positive values are dropped, as DDSketch indexes positive values only. - Accumulator factory updater; config validation requires an `alpha` parameter in (0, 1). - Planner: DDSketch parameter building (default alpha 0.01, overridable via sketch_parameters.DDSketch.alpha) and sketch properties. The MILP optimizer skips DDSketch until it has a parameter grid and cost rows. Co-Authored-By: Claude Opus 5.5 --- .../rs/asap_types/src/aggregation_config.rs | 24 ++ .../rs/asap_types/src/capability_matching.rs | 30 +- .../rs/asap_types/src/streaming_config.rs | 57 ++++ .../src/query_logics/enums.rs | 4 + .../src/query_logics/logics.rs | 3 +- asap-planner-rs/src/config/input.rs | 7 + .../src/optimizer/candidate_gen.rs | 2 + .../src/optimizer/sketch_properties.rs | 2 + asap-planner-rs/src/planner/sketch.rs | 37 +++ .../precompute_engine/accumulator_factory.rs | 140 ++++++++- .../ddsketch_accumulator.rs | 289 ++++++++++++++++++ .../src/precompute_operators/mod.rs | 2 + .../tests/e2e_precompute_equivalence.rs | 72 +++++ 13 files changed, 663 insertions(+), 6 deletions(-) create mode 100644 asap-query-engine/src/precompute_operators/ddsketch_accumulator.rs 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..2d21567a 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; the optimizer side lands with #762. + 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/precompute_engine/accumulator_factory.rs b/asap-query-engine/src/precompute_engine/accumulator_factory.rs index 9a1f1008..58bad75a 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,50 @@ 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 dense store. + std::mem::size_of::() + + std::mem::size_of_val(self.acc.inner.store_counts()) + } +} + // --------------------------------------------------------------------------- // HllAccumulatorUpdater // --------------------------------------------------------------------------- @@ -820,6 +864,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 +969,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 +1183,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 +1500,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..9100b0db --- /dev/null +++ b/asap-query-engine/src/precompute_operators/ddsketch_accumulator.rs @@ -0,0 +1,289 @@ +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 is within a +/// factor of `(1 + alpha) / (1 - alpha)` of the true value. DDSketch indexes +/// positive values only, so zero, negative and non-finite inputs 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) + } +} + +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 non_positive_values_are_dropped() { + let mut dd = DDSketchAccumulator::new(ALPHA); + dd.update(0.0); + dd.update(-5.0); + dd.update(f64::NAN); + dd.update(42.0); + assert_eq!(dd.count(), 1); + assert!(rel_err(dd.get_quantile(0.5), 42.0) <= ALPHA); + } + + #[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_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..2eebc775 100644 --- a/asap-query-engine/tests/e2e_precompute_equivalence.rs +++ b/asap-query-engine/tests/e2e_precompute_equivalence.rs @@ -1701,6 +1701,78 @@ async fn e2e_grouped_quantile_preserves_output_label_shape() { ); } +/// A grouped quantile served by a DDSketch aggregation: samples flow through +/// remote write and precompute into one DDSketch per group, and each returned +/// quantile stays within the sketch's relative-accuracy bound. +#[tokio::test] +async fn e2e_grouped_quantile_over_ddsketch_is_within_alpha() { + let port = 19421u16; + 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)); + + // Each group gets 1..=100 (scaled per group) inside the (1s, 2s] window, + // plus one later sample to close the window. + let groups = [("frontend", 1.0), ("backend", 10.0)]; + let samples = groups + .iter() + .flat_map(|&(job, scale)| { + (1..=100) + .map(move |i| { + make_timeseries(metric, vec![("job", job)], 1_000 + 5 * i, scale * i as f64) + }) + .chain(std::iter::once(make_timeseries( + metric, + vec![("job", job)], + 3_500, + scale, + ))) + }) + .collect(); + let (engine, query) = NativeDagScenario { + port, + metric, + query, + aggregation_configs: vec![config], + schema_labels: vec!["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; + // The 0.9 quantile of scale * {1..=100} is scale * 90. + let truth = scale * 90.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. From 3302e14f1d6ab067a9d94b0d8b0e8d171498808f Mon Sep 17 00:00:00 2001 From: zz_y Date: Tue, 6 Oct 2026 19:23:35 +0000 Subject: [PATCH 2/6] build(deps): bump asap_sketchlib to 5158bf2 for DDSketch zeros and negatives asap_sketchlib#141 gives DDSketch a zero bucket and a negative store, as ClickHouse's quantilesDD has. At the old pin, zero and negative samples were dropped, so quantiles over data containing them were computed over the positive part only. No release includes the fix yet, so the git pin moves. The accumulator's doc comment and drop test now describe the kept values. Co-Authored-By: Claude Opus 5.5 --- Cargo.lock | 4 ++-- Cargo.toml | 7 ++++--- .../precompute_operators/ddsketch_accumulator.rs | 14 ++++++++------ 3 files changed, 14 insertions(+), 11 deletions(-) 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..034cde3c 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 +# (ProjectASAP/asap_sketchlib#76) and DDSketch's zero bucket and negative store +# (ProjectASAP/asap_sketchlib#141). 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-query-engine/src/precompute_operators/ddsketch_accumulator.rs b/asap-query-engine/src/precompute_operators/ddsketch_accumulator.rs index 9100b0db..471d4957 100644 --- a/asap-query-engine/src/precompute_operators/ddsketch_accumulator.rs +++ b/asap-query-engine/src/precompute_operators/ddsketch_accumulator.rs @@ -14,8 +14,8 @@ use promql_utilities::query_logics::enums::Statistic; /// implementation sketch-bench measures for the planner's cost table. /// /// `alpha` is the relative-accuracy bound: a returned quantile is within a -/// factor of `(1 + alpha) / (1 - alpha)` of the true value. DDSketch indexes -/// positive values only, so zero, negative and non-finite inputs are dropped. +/// factor of `(1 + alpha) / (1 - alpha)` of the true value. Negative values go +/// to a mirrored store and zeros to a zero bucket; non-finite inputs are dropped. #[derive(Clone, Debug)] pub struct DDSketchAccumulator { pub inner: DDSketch, @@ -192,14 +192,16 @@ mod tests { } #[test] - fn non_positive_values_are_dropped() { + fn zero_and_negative_values_are_counted() { let mut dd = DDSketchAccumulator::new(ALPHA); - dd.update(0.0); dd.update(-5.0); + dd.update(0.0); dd.update(f64::NAN); dd.update(42.0); - assert_eq!(dd.count(), 1); - assert!(rel_err(dd.get_quantile(0.5), 42.0) <= ALPHA); + 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] From b43075543fdb77a4a86690fb6fc276376dd0c703 Mon Sep 17 00:00:00 2001 From: zz_y Date: Tue, 6 Oct 2026 19:23:35 +0000 Subject: [PATCH 3/6] test(query-engine): DDSketch e2e with two instances per job and zeros MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Each job now has two instances feeding one DDSketch, sent in time order since the watermark is per group, plus ten zeros per instance. The zeros move the true p90 from 90·scale to 88·scale, outside the alpha bound, so the test fails if zeros are dropped. Co-Authored-By: Claude Opus 5.5 --- .../tests/e2e_precompute_equivalence.rs | 57 +++++++++++-------- 1 file changed, 34 insertions(+), 23 deletions(-) diff --git a/asap-query-engine/tests/e2e_precompute_equivalence.rs b/asap-query-engine/tests/e2e_precompute_equivalence.rs index 2eebc775..ace8e8b8 100644 --- a/asap-query-engine/tests/e2e_precompute_equivalence.rs +++ b/asap-query-engine/tests/e2e_precompute_equivalence.rs @@ -1701,9 +1701,11 @@ async fn e2e_grouped_quantile_preserves_output_label_shape() { ); } -/// A grouped quantile served by a DDSketch aggregation: samples flow through -/// remote write and precompute into one DDSketch per group, and each returned -/// quantile stays within the sketch's relative-accuracy bound. +/// 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 = 19421u16; @@ -1721,30 +1723,38 @@ async fn e2e_grouped_quantile_over_ddsketch_is_within_alpha() { ); config.parameters.insert("alpha".to_string(), json!(alpha)); - // Each group gets 1..=100 (scaled per group) inside the (1s, 2s] window, - // plus one later sample to close the window. + // 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 samples = groups - .iter() - .flat_map(|&(job, scale)| { - (1..=100) - .map(move |i| { - make_timeseries(metric, vec![("job", job)], 1_000 + 5 * i, scale * i as f64) - }) - .chain(std::iter::once(make_timeseries( - metric, - vec![("job", job)], - 3_500, - scale, - ))) - }) - .collect(); + 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!["job".to_string()], + schema_labels: vec!["instance".to_string(), "job".to_string()], samples, evaluation_time_seconds: 2.0, base_interval_ms: 1_000, @@ -1763,8 +1773,9 @@ async fn e2e_grouped_quantile_over_ddsketch_is_within_alpha() { for element in vector.values { let job = element.labels.labels[0].as_str(); let scale = groups.iter().find(|(g, _)| *g == job).unwrap().1; - // The 0.9 quantile of scale * {1..=100} is scale * 90. - let truth = scale * 90.0; + // 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}", From 848546eeae89b95e7696a7e6002199e23684e0e0 Mon Sep 17 00:00:00 2001 From: zz_y Date: Tue, 6 Oct 2026 19:40:57 +0000 Subject: [PATCH 4/6] fix(query-engine): address DDSketch review findings MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Move the DDSketch e2e test off port 19421, which the legacy-feature topk range test also binds. - State the accuracy bound as |x̂ - x| <= alpha·|x|; (1+alpha)/(1-alpha) is the bucket width, not the guarantee. - Count the negative store in memory_usage_bytes. - Drop issue numbers from code comments, per AGENTS.md. Co-Authored-By: Claude Opus 5.5 --- Cargo.toml | 4 ++-- asap-planner-rs/src/optimizer/candidate_gen.rs | 2 +- .../src/precompute_engine/accumulator_factory.rs | 3 ++- .../src/precompute_operators/ddsketch_accumulator.rs | 4 ++-- asap-query-engine/tests/e2e_precompute_equivalence.rs | 2 +- 5 files changed, 8 insertions(+), 7 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index 034cde3c..06e9fca7 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -36,8 +36,8 @@ 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 the published releases to pick up the CMS estimate i32::MAX clamp fix -# (ProjectASAP/asap_sketchlib#76) and DDSketch's zero bucket and negative store -# (ProjectASAP/asap_sketchlib#141). Revert to a published version once released. +# 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"] } diff --git a/asap-planner-rs/src/optimizer/candidate_gen.rs b/asap-planner-rs/src/optimizer/candidate_gen.rs index 2d21567a..42edfdec 100644 --- a/asap-planner-rs/src/optimizer/candidate_gen.rs +++ b/asap-planner-rs/src/optimizer/candidate_gen.rs @@ -25,7 +25,7 @@ const OPTIMIZER_SKIPPED_AGG_TYPES: &[AggregationType] = &[ AggregationType::Sum, AggregationType::MinMax, AggregationType::Increase, - // No parameter grid or cost rows yet; the optimizer side lands with #762. + // No parameter grid or cost rows yet. AggregationType::DDSketch, ]; diff --git a/asap-query-engine/src/precompute_engine/accumulator_factory.rs b/asap-query-engine/src/precompute_engine/accumulator_factory.rs index 58bad75a..e55149b5 100644 --- a/asap-query-engine/src/precompute_engine/accumulator_factory.rs +++ b/asap-query-engine/src/precompute_engine/accumulator_factory.rs @@ -337,9 +337,10 @@ impl AccumulatorUpdater for DDSketchAccumulatorUpdater { } fn memory_usage_bytes(&self) -> usize { - // One u64 counter per bucket in the dense store. + // 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()) } } diff --git a/asap-query-engine/src/precompute_operators/ddsketch_accumulator.rs b/asap-query-engine/src/precompute_operators/ddsketch_accumulator.rs index 471d4957..d6e20b97 100644 --- a/asap-query-engine/src/precompute_operators/ddsketch_accumulator.rs +++ b/asap-query-engine/src/precompute_operators/ddsketch_accumulator.rs @@ -13,8 +13,8 @@ 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 is within a -/// factor of `(1 + alpha) / (1 - alpha)` of the true value. Negative values go +/// `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 to a zero bucket; non-finite inputs are dropped. #[derive(Clone, Debug)] pub struct DDSketchAccumulator { diff --git a/asap-query-engine/tests/e2e_precompute_equivalence.rs b/asap-query-engine/tests/e2e_precompute_equivalence.rs index ace8e8b8..7ff11222 100644 --- a/asap-query-engine/tests/e2e_precompute_equivalence.rs +++ b/asap-query-engine/tests/e2e_precompute_equivalence.rs @@ -1708,7 +1708,7 @@ async fn e2e_grouped_quantile_preserves_output_label_shape() { /// that dropped them would answer outside the bound. #[tokio::test] async fn e2e_grouped_quantile_over_ddsketch_is_within_alpha() { - let port = 19421u16; + let port = 19428u16; let metric = "dd_latency"; let query = "quantile by (job) (0.9, dd_latency)"; let alpha = 0.01; From 6c128bcbdc6907be42e29d65e0c27639907c9b2a Mon Sep 17 00:00:00 2001 From: zz_y Date: Tue, 6 Oct 2026 19:44:20 +0000 Subject: [PATCH 5/6] perf(query-engine): batch-merge DDSketch window buckets merge_accumulators_batch had no DDSketch arm, so merging N buckets went through merge_with, which clones the running sketch at every step. DDSketchAccumulator::merge_multiple clones the first sketch once and merges the rest into it in place, as the KLL and CMS arms do. Co-Authored-By: Claude Opus 5.5 --- asap-query-engine/src/engines/merge_utils.rs | 11 +++- .../src/engines/simple_engine/mod.rs | 41 +++++++++++++- .../ddsketch_accumulator.rs | 55 +++++++++++++++++++ 3 files changed, 105 insertions(+), 2 deletions(-) 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_operators/ddsketch_accumulator.rs b/asap-query-engine/src/precompute_operators/ddsketch_accumulator.rs index d6e20b97..4b21f6de 100644 --- a/asap-query-engine/src/precompute_operators/ddsketch_accumulator.rs +++ b/asap-query-engine/src/precompute_operators/ddsketch_accumulator.rs @@ -52,6 +52,28 @@ impl DDSketchAccumulator { 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 { @@ -255,6 +277,39 @@ mod tests { 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); From 505f9ba443050f69ffe42aab1c4a1bb293a02890 Mon Sep 17 00:00:00 2001 From: zz_y Date: Tue, 6 Oct 2026 19:50:53 +0000 Subject: [PATCH 6/6] docs(query-engine): list every input DDSketch drops or files as zero Co-Authored-By: Claude Opus 5.5 --- .../src/precompute_operators/ddsketch_accumulator.rs | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/asap-query-engine/src/precompute_operators/ddsketch_accumulator.rs b/asap-query-engine/src/precompute_operators/ddsketch_accumulator.rs index 4b21f6de..b8041eff 100644 --- a/asap-query-engine/src/precompute_operators/ddsketch_accumulator.rs +++ b/asap-query-engine/src/precompute_operators/ddsketch_accumulator.rs @@ -14,8 +14,9 @@ use promql_utilities::query_logics::enums::Statistic; /// 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 to a zero bucket; non-finite inputs are dropped. +/// 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,