mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-07 23:07:48 +02:00
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.
This commit is contained in:
1 parent
d235dd280b
commit
51831d6850
4 files changed
+275
-28
No files matched your search
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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[""])
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in new issue
Block a user