From 55ec848d60211263aaebbfc67beeffd111af8d66 Mon Sep 17 00:00:00 2001 From: fdgod Date: Sun, 16 Mar 2025 22:32:30 +0530 Subject: [PATCH] change arbiter to start loop when meter receives data, stop loop when no meters or metricsEnabled is false --- metrics/init_test.go | 3 +++ metrics/meter.go | 28 +++++++++++++++++----------- metrics/meter_test.go | 27 +++++++++++++++++++++++++++ 3 files changed, 47 insertions(+), 11 deletions(-) diff --git a/metrics/init_test.go b/metrics/init_test.go index af75bee425..1487e36408 100644 --- a/metrics/init_test.go +++ b/metrics/init_test.go @@ -1,5 +1,8 @@ package metrics +import "time" + func init() { metricsEnabled = true + MeterTickerInterval = 1 * time.Second } diff --git a/metrics/meter.go b/metrics/meter.go index 194bd1f304..c0329e4212 100644 --- a/metrics/meter.go +++ b/metrics/meter.go @@ -11,6 +11,8 @@ import ( // new Meter. // Be sure to unregister the meter from the registry once it is of no use to // allow for garbage collection. +var MeterTickerInterval = time.Second * 5 + func GetOrRegisterMeter(name string, r Registry) *Meter { if r == nil { r = DefaultRegistry @@ -18,7 +20,7 @@ func GetOrRegisterMeter(name string, r Registry) *Meter { return r.GetOrRegister(name, NewMeter).(*Meter) } -// NewMeter constructs a new Meter and launches a goroutine. +// NewMeter constructs a new Meter // Be sure to call Stop() once the meter is of no use to allow for garbage collection. func NewMeter() *Meter { m := newMeter() @@ -94,7 +96,15 @@ func (m *Meter) Stop() { } // Mark records the occurrence of n events. +// It launches an arbiter goroutine to update meters if not started. func (m *Meter) Mark(n int64) { + if !metricsEnabled { + return + } + if !arbiter.started { + arbiter.started = true + go arbiter.loop() + } m.uncounted.Add(n) } @@ -129,8 +139,7 @@ var arbiter = meterTicker{meters: make(map[*Meter]struct{})} // meterTicker ticks meters every 5s from a single goroutine. // meters are references in a set for future stopping. type meterTicker struct { - mu sync.RWMutex - + mu sync.RWMutex started bool meters map[*Meter]struct{} } @@ -140,10 +149,6 @@ func (ma *meterTicker) add(m *Meter) { ma.mu.Lock() defer ma.mu.Unlock() ma.meters[m] = struct{}{} - if !ma.started { - ma.started = true - go ma.loop() - } } // remove removes a meter from the set of ticked meters. @@ -153,12 +158,13 @@ func (ma *meterTicker) remove(m *Meter) { ma.mu.Unlock() } -// loop ticks meters on a 5 second interval. +// loop ticks meters on a configured interval. func (ma *meterTicker) loop() { - ticker := time.NewTicker(5 * time.Second) + ticker := time.NewTicker(MeterTickerInterval) for range ticker.C { - if !metricsEnabled { - continue + if len(ma.meters) == 0 || !metricsEnabled { + ma.started = false + return } ma.mu.RLock() for meter := range ma.meters { diff --git a/metrics/meter_test.go b/metrics/meter_test.go index e3f39684bd..40fbfdac43 100644 --- a/metrics/meter_test.go +++ b/metrics/meter_test.go @@ -1,8 +1,11 @@ package metrics import ( + "runtime" "testing" "time" + + "go.uber.org/goleak" ) func BenchmarkMeter(b *testing.B) { @@ -81,3 +84,27 @@ func TestMeterRepeat(t *testing.T) { t.Errorf("m.Count(): 10100 != %v\n", count) } } + +func TestMeterLazyInitialization(t *testing.T) { + defer goleak.VerifyNone(t) + // Removing all meters so that goroutine stops + for meter, _ := range arbiter.meters { + arbiter.remove(meter) + } + time.Sleep(MeterTickerInterval) + initialGoroutines := runtime.NumGoroutine() + m1 := NewMeter() + afterCreateGoroutines := runtime.NumGoroutine() + if afterCreateGoroutines != initialGoroutines { + t.Errorf("Expected no new goroutines after meter creation, got: before=%d, after=%d", + initialGoroutines, afterCreateGoroutines) + } + m1.Mark(1) + afterMarkGoroutines := runtime.NumGoroutine() + if afterMarkGoroutines != initialGoroutines+1 { + t.Errorf("Expected exactly one new goroutine after Mark(), got: before=%d, after=%d", + initialGoroutines, afterMarkGoroutines) + } + m1.Stop() + time.Sleep(MeterTickerInterval) +}