mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-26 08:18:18 +00:00
perf(metrics): Cache metric handles instead of creating each call (#1852)
Co-authored-by: heihutu <30542132+heihutu@users.noreply.github.com>
This commit is contained in:
+144
-9
@@ -24,7 +24,7 @@ use std::{
|
||||
collections::HashMap,
|
||||
ops::Deref,
|
||||
sync::{
|
||||
Arc, Mutex,
|
||||
Arc, Mutex, RwLock,
|
||||
atomic::{AtomicU64, Ordering},
|
||||
},
|
||||
};
|
||||
@@ -72,6 +72,9 @@ impl Builder {
|
||||
Recorder {
|
||||
meter,
|
||||
metrics_metadata: Arc::new(Mutex::new(HashMap::new())),
|
||||
cached_counters: Arc::new(RwLock::new(HashMap::new())),
|
||||
cached_gauges: Arc::new(RwLock::new(HashMap::new())),
|
||||
cached_histograms: Arc::new(RwLock::new(HashMap::new())),
|
||||
},
|
||||
)
|
||||
}
|
||||
@@ -113,6 +116,10 @@ struct MetricMetadata {
|
||||
pub struct Recorder {
|
||||
meter: Meter,
|
||||
metrics_metadata: Arc<Mutex<HashMap<KeyName, MetricMetadata>>>,
|
||||
// cache metric handlers as to not reregister on each call
|
||||
cached_counters: Arc<RwLock<HashMap<Key, Counter>>>,
|
||||
cached_gauges: Arc<RwLock<HashMap<Key, Gauge>>>,
|
||||
cached_histograms: Arc<RwLock<HashMap<Key, Histogram>>>,
|
||||
}
|
||||
|
||||
impl Recorder {
|
||||
@@ -129,6 +136,9 @@ impl Recorder {
|
||||
Recorder {
|
||||
meter,
|
||||
metrics_metadata: Arc::new(Mutex::new(HashMap::new())),
|
||||
cached_counters: Arc::new(RwLock::new(HashMap::new())),
|
||||
cached_gauges: Arc::new(RwLock::new(HashMap::new())),
|
||||
cached_histograms: Arc::new(RwLock::new(HashMap::new())),
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -158,8 +168,12 @@ impl metrics::Recorder for Recorder {
|
||||
}
|
||||
|
||||
fn register_counter(&self, key: &Key, _metadata: &Metadata<'_>) -> Counter {
|
||||
if let Some(cached) = self.cached_counters.read().unwrap().get(key) {
|
||||
return cached.clone();
|
||||
}
|
||||
|
||||
let mut builder = self.meter.u64_counter(key.name().to_owned());
|
||||
if let Some(metadata) = self.metrics_metadata.lock().unwrap().remove(key.name()) {
|
||||
if let Some(metadata) = self.metrics_metadata.lock().unwrap().get(key.name()) {
|
||||
if let Some(unit) = metadata.unit {
|
||||
builder = builder.with_unit(unit.as_canonical_label());
|
||||
}
|
||||
@@ -172,16 +186,30 @@ impl metrics::Recorder for Recorder {
|
||||
.map(|label| KeyValue::new(label.key().to_owned(), label.value().to_owned()))
|
||||
.collect();
|
||||
|
||||
Counter::from_arc(Arc::new(WrappedCounter {
|
||||
let handle = Counter::from_arc(Arc::new(WrappedCounter {
|
||||
counter,
|
||||
labels,
|
||||
value: AtomicU64::new(0),
|
||||
}))
|
||||
}));
|
||||
|
||||
// Take write lock only after building
|
||||
let mut cache = self.cached_counters.write().unwrap();
|
||||
// Quick check it wasn't created between last read and finish building
|
||||
if let Some(cached) = cache.get(key) {
|
||||
return cached.clone();
|
||||
}
|
||||
|
||||
cache.insert(key.clone(), handle.clone());
|
||||
handle
|
||||
}
|
||||
|
||||
fn register_gauge(&self, key: &Key, _metadata: &Metadata<'_>) -> Gauge {
|
||||
if let Some(cached) = self.cached_gauges.read().unwrap().get(key) {
|
||||
return cached.clone();
|
||||
}
|
||||
|
||||
let mut builder = self.meter.f64_gauge(key.name().to_owned());
|
||||
if let Some(metadata) = self.metrics_metadata.lock().unwrap().remove(key.name()) {
|
||||
if let Some(metadata) = self.metrics_metadata.lock().unwrap().get(key.name()) {
|
||||
if let Some(unit) = metadata.unit {
|
||||
builder = builder.with_unit(unit.as_canonical_label());
|
||||
}
|
||||
@@ -194,16 +222,30 @@ impl metrics::Recorder for Recorder {
|
||||
.map(|label| KeyValue::new(label.key().to_owned(), label.value().to_owned()))
|
||||
.collect();
|
||||
|
||||
Gauge::from_arc(Arc::new(WrappedGauge {
|
||||
let handle = Gauge::from_arc(Arc::new(WrappedGauge {
|
||||
gauge,
|
||||
labels,
|
||||
value: AtomicU64::new(0),
|
||||
}))
|
||||
}));
|
||||
|
||||
// Take write lock only after building
|
||||
let mut cache = self.cached_gauges.write().unwrap();
|
||||
// Quick check it wasn't created between last read and finish building
|
||||
if let Some(cached) = cache.get(key) {
|
||||
return cached.clone();
|
||||
}
|
||||
|
||||
cache.insert(key.clone(), handle.clone());
|
||||
handle
|
||||
}
|
||||
|
||||
fn register_histogram(&self, key: &Key, _metadata: &Metadata<'_>) -> Histogram {
|
||||
if let Some(cached) = self.cached_histograms.read().unwrap().get(key) {
|
||||
return cached.clone();
|
||||
}
|
||||
|
||||
let mut builder = self.meter.f64_histogram(key.name().to_owned());
|
||||
if let Some(metadata) = self.metrics_metadata.lock().unwrap().remove(key.name()) {
|
||||
if let Some(metadata) = self.metrics_metadata.lock().unwrap().get(key.name()) {
|
||||
if let Some(unit) = metadata.unit {
|
||||
builder = builder.with_unit(unit.as_canonical_label());
|
||||
}
|
||||
@@ -216,7 +258,17 @@ impl metrics::Recorder for Recorder {
|
||||
.map(|label| KeyValue::new(label.key().to_owned(), label.value().to_owned()))
|
||||
.collect();
|
||||
|
||||
Histogram::from_arc(Arc::new(WrappedHistogram { histogram, labels }))
|
||||
let handle = Histogram::from_arc(Arc::new(WrappedHistogram { histogram, labels }));
|
||||
|
||||
// Take write lock only after building
|
||||
let mut cache = self.cached_histograms.write().unwrap();
|
||||
// Quick check it wasn't created between last read and finish building
|
||||
if let Some(cached) = cache.get(key) {
|
||||
return cached.clone();
|
||||
}
|
||||
|
||||
cache.insert(key.clone(), handle.clone());
|
||||
handle
|
||||
}
|
||||
}
|
||||
|
||||
@@ -300,8 +352,25 @@ impl HistogramFn for WrappedHistogram {
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use metrics::Recorder as _;
|
||||
use opentelemetry_sdk::metrics::Temporality;
|
||||
|
||||
fn test_recorder() -> Recorder {
|
||||
let exporter = opentelemetry_stdout::MetricExporterBuilder::default()
|
||||
.with_temporality(Temporality::Cumulative)
|
||||
.build();
|
||||
|
||||
let (_provider, recorder) = Recorder::builder("test")
|
||||
.with_meter_provider(|b| b.with_periodic_exporter(exporter))
|
||||
.build();
|
||||
|
||||
recorder
|
||||
}
|
||||
|
||||
fn test_metadata() -> Metadata<'static> {
|
||||
Metadata::new(module_path!(), metrics::Level::INFO, None)
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn standard_usage() {
|
||||
let exporter = opentelemetry_stdout::MetricExporterBuilder::default()
|
||||
@@ -320,4 +389,70 @@ mod tests {
|
||||
|
||||
provider.force_flush().unwrap();
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn counter_cached_on_repeated_registration() {
|
||||
let recorder = test_recorder();
|
||||
let key = Key::from_name("requests_total");
|
||||
let meta = test_metadata();
|
||||
|
||||
let _first = recorder.register_counter(&key, &meta);
|
||||
let _second = recorder.register_counter(&key, &meta);
|
||||
|
||||
let cache = recorder.cached_counters.read().unwrap();
|
||||
assert_eq!(cache.len(), 1, "counter should be cached and inserted only once");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn gauge_cached_on_repeated_registration() {
|
||||
let recorder = test_recorder();
|
||||
let key = Key::from_name("active_connections");
|
||||
let meta = test_metadata();
|
||||
|
||||
let _first = recorder.register_gauge(&key, &meta);
|
||||
let _second = recorder.register_gauge(&key, &meta);
|
||||
|
||||
let cache = recorder.cached_gauges.read().unwrap();
|
||||
assert_eq!(cache.len(), 1, "gauge should be cached and inserted only once");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn histogram_cached_on_repeated_registration() {
|
||||
let recorder = test_recorder();
|
||||
let key = Key::from_name("request_duration");
|
||||
let meta = test_metadata();
|
||||
|
||||
let _first = recorder.register_histogram(&key, &meta);
|
||||
let _second = recorder.register_histogram(&key, &meta);
|
||||
|
||||
let cache = recorder.cached_histograms.read().unwrap();
|
||||
assert_eq!(cache.len(), 1, "histogram should be cached and inserted only once");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn concurrent_register_counter_inserts_once() {
|
||||
let recorder = test_recorder();
|
||||
let key = Key::from_name("concurrent_counter");
|
||||
let shared = Arc::new(recorder);
|
||||
let barrier = Arc::new(std::sync::Barrier::new(10));
|
||||
|
||||
let handles: Vec<_> = (0..10)
|
||||
.map(|_| {
|
||||
let r = Arc::clone(&shared);
|
||||
let k = key.clone();
|
||||
let b = Arc::clone(&barrier);
|
||||
std::thread::spawn(move || {
|
||||
b.wait();
|
||||
let _ = r.register_counter(&k, &test_metadata());
|
||||
})
|
||||
})
|
||||
.collect();
|
||||
|
||||
for h in handles {
|
||||
h.join().unwrap();
|
||||
}
|
||||
|
||||
let cache = shared.cached_counters.read().unwrap();
|
||||
assert_eq!(cache.len(), 1, "concurrent registrations should produce exactly one cache entry");
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user