From 51831d6850988a8f3d0550cbb865f2e9b951eddd Mon Sep 17 00:00:00 2001 From: Chris Lu Date: Mon, 14 Sep 2026 21:47:31 -0700 Subject: [PATCH] admin: derive interval rates and latency quantiles from scrapes Raw counters and histograms are process-lifetime cumulative, so charting them directly is meaningless. Counters now become per-second rates and histograms become p50/p95/p99, both computed from the delta against the previous scrape, with counter resets skipped. Quantiles use linear interpolation within the matching bucket, as Prometheus does. --- weed/admin/dash/admin_server.go | 6 +- weed/admin/dash/metrics_derive.go | 123 +++++++++++++++++++++++++ weed/admin/dash/metrics_derive_test.go | 105 +++++++++++++++++++++ weed/admin/dash/metrics_scraper.go | 69 ++++++++------ 4 files changed, 275 insertions(+), 28 deletions(-) create mode 100644 weed/admin/dash/metrics_derive.go create mode 100644 weed/admin/dash/metrics_derive_test.go diff --git a/weed/admin/dash/admin_server.go b/weed/admin/dash/admin_server.go index 34ca3be8c..eec4a2f55 100644 --- a/weed/admin/dash/admin_server.go +++ b/weed/admin/dash/admin_server.go @@ -128,7 +128,10 @@ type AdminServer struct { // metricsStore holds scraped per-server Prometheus series for the // monitoring pages. Filled by the scrape loop in startMetricsScraper. - metricsStore *metricsStore + // metricsDeriver turns raw counters and histograms into interval rates + // and latency quantiles before they are stored. + metricsStore *metricsStore + metricsDeriver *metricsDeriver // Filer discovery and caching cachedFilers []string @@ -215,6 +218,7 @@ func NewAdminServer(masters string, filerGroup string, templateFS http.FileSyste adminPresenceLock: presenceLock, bgCancel: bgCancel, metricsStore: newMetricsStore(), + metricsDeriver: newMetricsDeriver(), } // Initialize topic retention purger diff --git a/weed/admin/dash/metrics_derive.go b/weed/admin/dash/metrics_derive.go new file mode 100644 index 000000000..69793b9a9 --- /dev/null +++ b/weed/admin/dash/metrics_derive.go @@ -0,0 +1,123 @@ +package dash + +import ( + "math" + "sync" + "time" +) + +// Suffixes for series derived from raw scrapes. Counters become per-second +// rates and histograms become latency quantiles, both computed from the delta +// against the previous scrape so the values reflect the last interval rather +// than process lifetime totals. +const ( + suffixRate = ":rate" + suffixP50 = ":p50" + suffixP95 = ":p95" + suffixP99 = ":p99" +) + +type prevScrape struct { + t time.Time + value float64 + buckets []histogramBucket +} + +type metricsDeriver struct { + mu sync.Mutex + prev map[string]prevScrape +} + +func newMetricsDeriver() *metricsDeriver { + return &metricsDeriver{prev: make(map[string]prevScrape)} +} + +// record stores the raw value and, for counters and histograms, the derived +// rate/quantile series for this interval. +func (d *metricsDeriver) record(store *metricsStore, source string, m scrapedMetric, now time.Time) { + key := source + "/" + m.name + "/" + labelKey(m.labels) + + d.mu.Lock() + prev, hadPrev := d.prev[key] + d.prev[key] = prevScrape{t: now, value: m.value, buckets: m.buckets} + d.mu.Unlock() + + switch m.kind { + case kindGauge: + store.recordLabeled(source, m.name, m.labels, m.value, now) + case kindCounter: + if !hadPrev { + return + } + dt := now.Sub(prev.t).Seconds() + if dt <= 0 { + return + } + delta := m.value - prev.value + if delta < 0 { + // Counter reset (process restart); skip this interval. + return + } + store.recordLabeled(source, m.name+suffixRate, m.labels, delta/dt, now) + case kindHistogram: + if !hadPrev { + return + } + delta := bucketDelta(prev.buckets, m.buckets) + if len(delta) == 0 { + return + } + for suffix, q := range map[string]float64{suffixP50: 0.5, suffixP95: 0.95, suffixP99: 0.99} { + store.recordLabeled(source, m.name+suffix, m.labels, histogramQuantile(delta, q), now) + } + } +} + +// bucketDelta subtracts cumulative bucket counts, yielding the distribution +// observed during the interval. Returns nil on a reset or bucket mismatch. +func bucketDelta(prev, cur []histogramBucket) []histogramBucket { + if len(prev) != len(cur) { + return nil + } + out := make([]histogramBucket, len(cur)) + for i := range cur { + if cur[i].upperBound != prev[i].upperBound { + return nil + } + c := cur[i].count - prev[i].count + if c < 0 { + return nil + } + out[i] = histogramBucket{upperBound: cur[i].upperBound, count: c} + } + if out[len(out)-1].count == 0 { + return nil + } + return out +} + +// histogramQuantile estimates a quantile from cumulative buckets by linear +// interpolation within the matching bucket, matching Prometheus' approach. +func histogramQuantile(buckets []histogramBucket, q float64) float64 { + total := buckets[len(buckets)-1].count + if total == 0 { + return 0 + } + rank := q * total + prevCount, prevBound := 0.0, 0.0 + for _, b := range buckets { + if b.count < rank { + prevCount, prevBound = b.count, b.upperBound + continue + } + if math.IsInf(b.upperBound, 1) { + return prevBound + } + span := b.count - prevCount + if span <= 0 { + return b.upperBound + } + return prevBound + (b.upperBound-prevBound)*(rank-prevCount)/span + } + return buckets[len(buckets)-1].upperBound +} diff --git a/weed/admin/dash/metrics_derive_test.go b/weed/admin/dash/metrics_derive_test.go new file mode 100644 index 000000000..fb93c159c --- /dev/null +++ b/weed/admin/dash/metrics_derive_test.go @@ -0,0 +1,105 @@ +package dash + +import ( + "math" + "testing" + "time" +) + +func TestHistogramQuantile(t *testing.T) { + // 100 observations spread evenly across 0-1s in 10 buckets. + buckets := []histogramBucket{ + {0.1, 10}, {0.2, 20}, {0.3, 30}, {0.4, 40}, {0.5, 50}, + {0.6, 60}, {0.7, 70}, {0.8, 80}, {0.9, 90}, {1.0, 100}, + {math.Inf(1), 100}, + } + for _, tc := range []struct { + q float64 + want float64 + }{ + {0.5, 0.5}, + {0.95, 0.95}, + {0.99, 0.99}, + } { + got := histogramQuantile(buckets, tc.q) + if math.Abs(got-tc.want) > 1e-9 { + t.Errorf("quantile(%v) = %v, want %v", tc.q, got, tc.want) + } + } +} + +func TestHistogramQuantileEmpty(t *testing.T) { + if got := histogramQuantile([]histogramBucket{{math.Inf(1), 0}}, 0.99); got != 0 { + t.Errorf("empty histogram quantile = %v, want 0", got) + } +} + +func TestBucketDeltaResetAndMismatch(t *testing.T) { + prev := []histogramBucket{{0.1, 5}, {math.Inf(1), 10}} + if got := bucketDelta(prev, []histogramBucket{{0.1, 1}, {math.Inf(1), 2}}); got != nil { + t.Errorf("counter reset should yield nil, got %v", got) + } + if got := bucketDelta(prev, []histogramBucket{{0.2, 5}, {math.Inf(1), 10}}); got != nil { + t.Errorf("bound mismatch should yield nil, got %v", got) + } + if got := bucketDelta(prev, prev); got != nil { + t.Errorf("no new observations should yield nil, got %v", got) + } + got := bucketDelta(prev, []histogramBucket{{0.1, 7}, {math.Inf(1), 14}}) + if len(got) != 2 || got[0].count != 2 || got[1].count != 4 { + t.Errorf("unexpected delta %v", got) + } +} + +func TestDeriveCounterRate(t *testing.T) { + store := newMetricsStore() + d := newMetricsDeriver() + t0 := time.Now() + m := scrapedMetric{name: "reqs", kind: kindCounter, value: 100} + + d.record(store, "volume/a", m, t0) + if got := store.match("volume", "reqs"+suffixRate); len(got) != 0 { + t.Fatalf("first scrape should not emit a rate, got %d series", len(got)) + } + + m.value = 250 + d.record(store, "volume/a", m, t0.Add(15*time.Second)) + series := store.match("volume", "reqs"+suffixRate) + if len(series) != 1 { + t.Fatalf("expected 1 rate series, got %d", len(series)) + } + samples := series[0].snapshot() + if len(samples) != 1 { + t.Fatalf("expected 1 sample, got %d", len(samples)) + } + if want := 10.0; math.Abs(samples[0].values[""]-want) > 1e-9 { + t.Errorf("rate = %v, want %v", samples[0].values[""], want) + } +} + +func TestDeriveCounterResetSkipped(t *testing.T) { + store := newMetricsStore() + d := newMetricsDeriver() + t0 := time.Now() + m := scrapedMetric{name: "reqs", kind: kindCounter, value: 100} + d.record(store, "volume/a", m, t0) + m.value = 5 + d.record(store, "volume/a", m, t0.Add(15*time.Second)) + if got := store.match("volume", "reqs"+suffixRate); len(got) != 0 { + t.Errorf("counter reset should emit no rate, got %d series", len(got)) + } +} + +func TestStoreRingIsBounded(t *testing.T) { + s := newMetricsSeries() + for i := 0; i < metricsMaxSamples+50; i++ { + s.record(time.Now(), map[string]float64{"": float64(i)}) + } + got := s.snapshot() + if len(got) != metricsMaxSamples { + t.Fatalf("len = %d, want %d", len(got), metricsMaxSamples) + } + if got[len(got)-1].values[""] != float64(metricsMaxSamples+49) { + t.Errorf("newest sample not retained: %v", got[len(got)-1].values[""]) + } +} diff --git a/weed/admin/dash/metrics_scraper.go b/weed/admin/dash/metrics_scraper.go index 0ba8e587c..4d57ac7a1 100644 --- a/weed/admin/dash/metrics_scraper.go +++ b/weed/admin/dash/metrics_scraper.go @@ -16,10 +16,25 @@ import ( util_http "github.com/seaweedfs/seaweedfs/weed/util/http" ) +type metricKind int + +const ( + kindGauge metricKind = iota + kindCounter + kindHistogram +) + +type histogramBucket struct { + upperBound float64 + count float64 +} + type scrapedMetric struct { - name string - labels map[string]string - value float64 + name string + labels map[string]string + kind metricKind + value float64 + buckets []histogramBucket } func scrapeMetrics(ctx context.Context, target string) ([]scrapedMetric, error) { @@ -59,30 +74,34 @@ func parsePrometheusText(r io.Reader) ([]scrapedMetric, error) { return nil, err } for _, m := range fam.Metric { - labels := map[string]string{} - for _, l := range m.Label { - labels[l.GetName()] = l.GetValue() - } - out = append(out, scrapedMetric{name: fam.GetName(), labels: labels, value: metricValue(m)}) + out = append(out, toScrapedMetric(fam.GetName(), m)) } } return out, nil } -func metricValue(m *dto.Metric) float64 { - switch { - case m.Gauge != nil: - return m.Gauge.GetValue() - case m.Counter != nil: - return m.Counter.GetValue() - case m.Untyped != nil: - return m.Untyped.GetValue() - case m.Histogram != nil: - return m.Histogram.GetSampleSum() - case m.Summary != nil: - return m.Summary.GetSampleSum() +func toScrapedMetric(name string, m *dto.Metric) scrapedMetric { + labels := map[string]string{} + for _, l := range m.Label { + labels[l.GetName()] = l.GetValue() } - return 0 + sm := scrapedMetric{name: name, labels: labels} + switch { + case m.Counter != nil: + sm.kind, sm.value = kindCounter, m.Counter.GetValue() + case m.Histogram != nil: + sm.kind = kindHistogram + for _, b := range m.Histogram.Bucket { + sm.buckets = append(sm.buckets, histogramBucket{upperBound: b.GetUpperBound(), count: float64(b.GetCumulativeCount())}) + } + case m.Summary != nil: + sm.kind, sm.value = kindCounter, m.Summary.GetSampleSum() + case m.Gauge != nil: + sm.value = m.Gauge.GetValue() + case m.Untyped != nil: + sm.value = m.Untyped.GetValue() + } + return sm } // gatherLocalMetrics records the admin's own registry (maintenance tasks, @@ -95,11 +114,7 @@ func (s *AdminServer) gatherLocalMetrics(now time.Time) { } for _, fam := range families { for _, m := range fam.Metric { - labels := map[string]string{} - for _, l := range m.Label { - labels[l.GetName()] = l.GetValue() - } - s.metricsStore.recordLabeled("admin/local", fam.GetName(), labels, metricValue(m), now) + s.metricsDeriver.record(s.metricsStore, "admin/local", toScrapedMetric(fam.GetName(), m), now) } } } @@ -133,7 +148,7 @@ func (s *AdminServer) scrapeAllServers(ctx context.Context) { continue } for _, m := range r.metrics { - s.metricsStore.recordLabeled(r.source, m.name, m.labels, m.value, now) + s.metricsDeriver.record(s.metricsStore, r.source, m, now) } } }