From 9da332c47d2792f9e64ff91e8238d38b32283f99 Mon Sep 17 00:00:00 2001 From: evan slack <51209817+evanofslack@users.noreply.github.com> Date: Tue, 17 Feb 2026 12:03:35 -0500 Subject: [PATCH] perf(metrics): Cache metric handles instead of creating each call (#1852) Co-authored-by: heihutu <30542132+heihutu@users.noreply.github.com> --- crates/obs/src/recorder.rs | 153 ++++++++++++++++++++++++++++++++++--- 1 file changed, 144 insertions(+), 9 deletions(-) diff --git a/crates/obs/src/recorder.rs b/crates/obs/src/recorder.rs index 92d69ead5..934959665 100644 --- a/crates/obs/src/recorder.rs +++ b/crates/obs/src/recorder.rs @@ -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>>, + // cache metric handlers as to not reregister on each call + cached_counters: Arc>>, + cached_gauges: Arc>>, + cached_histograms: Arc>>, } 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"); + } }