use std::mem::size_of;
use hdrhistogram::Histogram;
use super::{
garnet_latency_metrics_session::GarnetLatencyMetricsSession,
latency_metrics_entry::time_stamp::TICKS_PER_MICROSECOND,
latency_metrics_type::LatencyMetricsType,
};
use crate::metrics::{metrics_item::MetricsItem, resp_write_utils::RespWriteUtils};
#[derive(Clone)]
pub struct GarnetLatencyMetrics {
pub metrics: Vec<Histogram<u64>>,
}
impl GarnetLatencyMetrics {
pub const DEFAULT_LATENCY_TYPES: &[LatencyMetricsType] = &LatencyMetricsType::ALL;
pub fn new(latency_types: &'static [LatencyMetricsType]) -> Self {
Self {
metrics: latency_types
.iter()
.map(|_| Self::new_histogram())
.collect(),
}
}
pub fn return_to_pool(&mut self) {
for hist in &mut self.metrics {
hist.reset();
}
self.metrics.clear();
}
pub fn merge(&mut self, lm: &GarnetLatencyMetricsSession) {
let Some(session_metrics) = lm.metrics_snapshot() else {
return;
};
let ver = lm.prior_version();
self.merge_session_snapshot(&session_metrics, ver);
}
pub fn merge_session_snapshot(
&mut self,
session_metrics: &[super::latency_metrics_entry_session::LatencyMetricsEntrySession],
ver: usize,
) {
for (dst, src) in self.metrics.iter_mut().zip(session_metrics.iter()) {
if !src.latency[ver].is_empty() {
let _ = dst.add(&src.latency[ver]);
}
}
}
pub fn reset(&mut self, cmd: LatencyMetricsType) {
let Some(hist) = self.metrics.get_mut(cmd.idx()) else {
return;
};
hist.reset();
}
fn get_percentiles(&self, idx: usize) -> Option<Vec<MetricsItem>> {
let hist = self.metrics.get(idx)?;
if hist.is_empty() {
return None;
}
let ticks = LatencyMetricsType::IS_TICKS[idx];
let raw = |p: f64| hist.value_at_percentile(p);
let fmt_value = |v: u64| {
if ticks {
fmt_n2(v as f64 / TICKS_PER_MICROSECOND as f64)
} else {
v.to_string()
}
};
let fmt_mean = |v: f64| {
if ticks {
fmt_n2(v / TICKS_PER_MICROSECOND as f64)
} else {
format!("{v:.2}")
}
};
Some(vec![
MetricsItem::new("calls", hist.len().to_string()),
MetricsItem::new("min", fmt_value(raw(0.0))),
MetricsItem::new("5th", fmt_value(raw(5.0))),
MetricsItem::new("50th", fmt_value(raw(50.0))),
MetricsItem::new("mean", fmt_mean(hist.mean())),
MetricsItem::new("95th", fmt_value(raw(95.0))),
MetricsItem::new("99th", fmt_value(raw(99.0))),
MetricsItem::new("99.9th", fmt_value(raw(99.9))),
])
}
pub fn get_resp_histogram(
&self,
idx: usize,
event_type: LatencyMetricsType,
response: &mut String,
) -> bool {
let Some(hist) = self.metrics.get(idx) else {
return false;
};
if hist.is_empty() {
return false;
}
let Some(p) = &self.get_percentiles(idx) else {
return false;
};
let cmd_type = event_type.cs_name();
response.push_str(&RespWriteUtils::bulk_string(cmd_type));
response.push_str("*6\r\n");
response.push_str(&RespWriteUtils::bulk_string("calls"));
response.push_str(&format!(":{}\r\n", p[0].value));
response.push_str(&RespWriteUtils::bulk_string("size"));
response.push_str(&format!(
":{}\r\n", (hist.distinct_values() as u64) * (size_of::<u64>() as u64)
));
response.push_str(&RespWriteUtils::bulk_string(
if LatencyMetricsType::IS_TICKS[idx] {
"histogram_usec"
} else {
"histogram_cnt"
},
));
response.push_str(&format!("*{}\r\n", (p.len() - 1) * 2));
for item in &p[1..] {
response.push_str(&RespWriteUtils::bulk_string(&item.name));
response.push_str(&RespWriteUtils::bulk_string(&item.value));
}
true
}
pub fn get_resp_histograms(&self, events: &[LatencyMetricsType]) -> String {
let mut cmd_count = 0;
let mut response = String::new();
for event_type in events {
let idx = event_type.idx();
let mut cmd_histogram = String::new();
if self.get_resp_histogram(idx, *event_type, &mut cmd_histogram) {
response.push_str(&cmd_histogram);
cmd_count += 1;
}
}
if cmd_count == 0 {
"*0\r\n".into()
} else {
format!("*{}\r\n", cmd_count * 2) + response.as_str()
}
}
pub fn get_latency_metrics(&self, latency_metrics_type: LatencyMetricsType) -> Vec<MetricsItem> {
self
.get_percentiles(latency_metrics_type.idx())
.unwrap_or_default()
}
pub fn get_latency_metrics_multi(
&self,
latency_metrics_types: &[LatencyMetricsType],
) -> Vec<(LatencyMetricsType, Vec<MetricsItem>)> {
latency_metrics_types
.iter()
.filter_map(|&event_type| {
self
.get_percentiles(event_type.idx())
.map(|items| (event_type, items))
})
.collect()
}
pub fn dump(&self, idx: usize) {
let Some(hist) = self.metrics.get(idx) else {
return;
};
if !hist.is_empty() {
let mean = hist.mean() / TICKS_PER_MICROSECOND as f64;
log::info!(
"min (us); 5th (us); median (us); avg (us); 95th (us); 99th (us); 99.9th (us); cnt\n{}; {}; {}; {mean:.1}; {}; {}; {}; {}",
hist.value_at_percentile(0.0) / TICKS_PER_MICROSECOND,
hist.value_at_percentile(5.0) / TICKS_PER_MICROSECOND,
hist.value_at_percentile(50.0) / TICKS_PER_MICROSECOND,
hist.value_at_percentile(95.0) / TICKS_PER_MICROSECOND,
hist.value_at_percentile(99.0) / TICKS_PER_MICROSECOND,
hist.value_at_percentile(99.9) / TICKS_PER_MICROSECOND,
hist.len()
);
}
}
fn new_histogram() -> Histogram<u64> {
Histogram::new_with_bounds(
1,
super::latency_metrics_entry::LatencyMetricsEntry::HISTOGRAM_UPPER_BOUND,
2,
)
.expect("直方图边界为编译期常量,构造必成功")
}
}
pub fn fmt_n2(v: f64) -> String {
let rounded = format!("{v:.2}");
let (int_part, frac_part) = rounded.split_once('.').unwrap_or((rounded.as_str(), "00"));
let (sign, digits) = if let Some(d) = int_part.strip_prefix('-') {
("-", d)
} else {
("", int_part)
};
let mut grouped = String::with_capacity(digits.len() + digits.len() / 3 + 4);
grouped.push_str(sign);
for (i, c) in digits.chars().enumerate() {
if i > 0 && (digits.len() - i) % 3 == 0 {
grouped.push(',');
}
grouped.push(c);
}
grouped.push('.');
grouped.push_str(frac_part);
grouped
}