From 7badca597d81448e2d25714a33cf80c32c1ebc48 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Tomislav=20Brada=C4=8D?= Date: Wed, 16 Sep 2026 15:09:56 +0200 Subject: [PATCH] [Metricbeat][statsd] Reset timer and histogram metrics on every flush The statsd server registry resets counters and sets after every flush, but timers (`ms`) and histograms (`h`) were never reset. Their `count` therefore grew monotonically and `min`, `max`, `mean`, `stddev` and the percentiles aggregated every measurement ever recorded, for as long as the metric kept being updated. The values only went back down once the metric went stale and was evicted by `ttl`, which makes the reported series impossible to interpret: a consumer cannot tell whether a document updates the previous one or starts a new window. Clear the histogram sample after each flush and track the timer count separately from the meter, which is monotonic and cannot be reset. The moving average rates (`1m_rate`, `5m_rate`, `15m_rate`, `mean_rate`) are deliberately left alone, they are defined over fixed time windows and are not tied to the flush interval. A timer or histogram that received no measurements during an interval is no longer reported for that interval, so that clearing the sample does not emit zeroed out `min`/`max`/`mean` values. The metric stays in the registry and is reported again as soon as new measurements arrive, or is dropped once it exceeds its `ttl`. Closes #39987 Closes #41002 Assisted-By: Claude Code --- ...-timer-and-histogram-metrics-on-flush.yaml | 23 +++ .../metricbeat/metricbeat-module-statsd.md | 2 + x-pack/metricbeat/module/statsd/_meta/docs.md | 2 + .../module/statsd/server/data_test.go | 138 ++++++++++++++++++ .../module/statsd/server/registry.go | 53 ++++++- 5 files changed, 213 insertions(+), 5 deletions(-) create mode 100644 changelog/fragments/1789563172-statsd-reset-timer-and-histogram-metrics-on-flush.yaml diff --git a/changelog/fragments/1789563172-statsd-reset-timer-and-histogram-metrics-on-flush.yaml b/changelog/fragments/1789563172-statsd-reset-timer-and-histogram-metrics-on-flush.yaml new file mode 100644 index 000000000000..a336d93fd750 --- /dev/null +++ b/changelog/fragments/1789563172-statsd-reset-timer-and-histogram-metrics-on-flush.yaml @@ -0,0 +1,23 @@ +kind: bug-fix + +summary: Reset statsd timer and histogram metrics on every flush + +description: | + The statsd module reset counters and sets on every flush but never reset + timers (`ms`) and histograms (`h`), so their `count` grew monotonically and + their `min`, `max`, `mean`, `stddev` and percentiles aggregated every + measurement ever recorded, until the metric went stale and was dropped by + `ttl`. + + Timers and histograms now aggregate a single flush interval, as counters and + sets already did, and a timer or histogram that received no measurements + during an interval is no longer reported for that interval. The moving + average rates (`1m_rate`, `5m_rate`, `15m_rate`, `mean_rate`) are unchanged: + they are defined over fixed time windows and are not tied to the flush + interval. + + Also fixes https://github.com/elastic/beats/issues/41002. + +component: metricbeat + +issue: https://github.com/elastic/beats/issues/39987 diff --git a/docs/reference/metricbeat/metricbeat-module-statsd.md b/docs/reference/metricbeat/metricbeat-module-statsd.md index b1f51289c64c..5bf1419dffe4 100644 --- a/docs/reference/metricbeat/metricbeat-module-statsd.md +++ b/docs/reference/metricbeat/metricbeat-module-statsd.md @@ -56,6 +56,8 @@ The `statsd` module has these additional config options: **`ttl`** : It defines how long a metric will be reported after it was last recorded. Irrespective of the given ttl, metrics will be reported at least once. A ttl of zero means metrics will never expire. + {applies_to}`stack: ga 9.6+` It controls how long a metric is kept, not what it aggregates: every flush reports what was recorded since the previous one, and a timer or histogram that received no measurements is not reported for that flush. + **`statsd.mappings`** : It defines how metrics will mapped from the original metric label to the event json. Here’s an example configuration: diff --git a/x-pack/metricbeat/module/statsd/_meta/docs.md b/x-pack/metricbeat/module/statsd/_meta/docs.md index 248ac2fee25b..492894506768 100644 --- a/x-pack/metricbeat/module/statsd/_meta/docs.md +++ b/x-pack/metricbeat/module/statsd/_meta/docs.md @@ -44,6 +44,8 @@ The `statsd` module has these additional config options: **`ttl`** : It defines how long a metric will be reported after it was last recorded. Irrespective of the given ttl, metrics will be reported at least once. A ttl of zero means metrics will never expire. + {applies_to}`stack: ga 9.6+` It controls how long a metric is kept, not what it aggregates: every flush reports what was recorded since the previous one, and a timer or histogram that received no measurements is not reported for that flush. + **`statsd.mappings`** : It defines how metrics will mapped from the original metric label to the event json. Here’s an example configuration: diff --git a/x-pack/metricbeat/module/statsd/server/data_test.go b/x-pack/metricbeat/module/statsd/server/data_test.go index f2044e4d61e4..70eada7aad06 100644 --- a/x-pack/metricbeat/module/statsd/server/data_test.go +++ b/x-pack/metricbeat/module/statsd/server/data_test.go @@ -1387,6 +1387,144 @@ func TestTimerSampled(t *testing.T) { assert.True(t, actualMetric01["15m_rate"].(float64) > 10) } +// newStatsdMetricSet builds the statsd server metricset for a test. +func newStatsdMetricSet(t *testing.T, config map[string]any) *MetricSet { + t.Helper() + ms, ok := mbtest.NewMetricSet(t, config).(*MetricSet) + require.True(t, ok, "the statsd module must build a *MetricSet") + return ms +} + +// metricValues returns the values reported for a single metric. +func metricValues(t *testing.T, fields mapstr.M, name string) map[string]any { + t.Helper() + values, ok := fields[name].(map[string]any) + require.True(t, ok, "metric %q must be reported as a map of values, got %#v", name, fields[name]) + return values +} + +func TestTimerReset(t *testing.T) { + ms := newStatsdMetricSet(t, map[string]any{"module": "statsd"}) + + err := process([]string{"metric01:2|ms", "metric01:4|ms"}, ms) + require.NoError(t, err, "the timer packets must be accepted") + + events := ms.getEvents() + require.Len(t, events, 1, "the timer received measurements, so it must be reported") + + actual := metricValues(t, events[0].MetricSetFields, "metric01") + assert.Equal(t, int64(2), actual["count"], "both measurements of the first interval must be counted") + assert.Equal(t, int64(2), actual["min"], "min must be the smallest measurement of the first interval") + assert.Equal(t, int64(4), actual["max"], "max must be the largest measurement of the first interval") + assert.InDelta(t, 3.0, actual["mean"], 0.001, "mean must average the measurements of the first interval") + + err = process([]string{"metric01:10|ms"}, ms) + require.NoError(t, err, "the timer packet of the second interval must be accepted") + + events = ms.getEvents() + require.Len(t, events, 1, "the timer received a new measurement, so it must be reported again") + + actual = metricValues(t, events[0].MetricSetFields, "metric01") + assert.Equal(t, int64(1), actual["count"], "count must restart from the second interval instead of accumulating the first") + assert.Equal(t, int64(10), actual["min"], "min must forget the measurements of the first interval") + assert.Equal(t, int64(10), actual["max"], "max must forget the measurements of the first interval") + assert.InDelta(t, 10.0, actual["mean"], 0.001, "mean must average only the second interval") + + meanRate, ok := actual["mean_rate"].(float64) + require.True(t, ok, "mean_rate must be reported as a float64, got %#v", actual["mean_rate"]) + assert.Positive(t, meanRate, "the moving average rates span fixed time windows and must survive the reset") + + assert.Empty(t, ms.getEvents(), "a timer that received no new measurements must not be reported") +} + +func TestTimerSampledReset(t *testing.T) { + ms := newStatsdMetricSet(t, map[string]any{"module": "statsd"}) + + err := process([]string{"metric01:2|ms|@0.1"}, ms) + require.NoError(t, err, "the sampled timer packet must be accepted") + + events := ms.getEvents() + require.Len(t, events, 1, "the sampled timer received a measurement, so it must be reported") + assert.Equal(t, int64(10), metricValues(t, events[0].MetricSetFields, "metric01")["count"], + "count must be extrapolated from the 0.1 sample rate") + + err = process([]string{"metric01:2|ms|@0.5"}, ms) + require.NoError(t, err, "the sampled timer packet of the second interval must be accepted") + + events = ms.getEvents() + require.Len(t, events, 1, "the sampled timer received a new measurement, so it must be reported again") + assert.Equal(t, int64(2), metricValues(t, events[0].MetricSetFields, "metric01")["count"], + "the extrapolated count must restart from the second interval instead of adding to the previous 10") +} + +// TestTimerSampleRateAboveOne covers a sample rate above one, for which the +// extrapolated count rounds down to zero. The measurement was still recorded, +// so the timer must still be reported. +func TestTimerSampleRateAboveOne(t *testing.T) { + ms := newStatsdMetricSet(t, map[string]any{"module": "statsd"}) + + err := process([]string{"metric01:2|ms|@2"}, ms) + require.NoError(t, err, "a sample rate above one is accepted, only a rate of zero or less is rejected") + + events := ms.getEvents() + require.Len(t, events, 1, "the timer recorded a measurement, so it must be reported even though 1/2 truncates its count to zero") + + actual := metricValues(t, events[0].MetricSetFields, "metric01") + assert.Equal(t, int64(2), actual["min"], "min must come from the recorded measurement") + assert.Equal(t, int64(2), actual["max"], "max must come from the recorded measurement") + + // nothing new arrived, so the next flush skips the timer again + assert.Empty(t, ms.getEvents(), "an idle timer must not be reported on the following flush") +} + +func TestHistogramReset(t *testing.T) { + ms := newStatsdMetricSet(t, map[string]any{"module": "statsd"}) + + err := process([]string{"metric01:2|h", "metric01:4|h"}, ms) + require.NoError(t, err, "the histogram packets must be accepted") + + events := ms.getEvents() + require.Len(t, events, 1, "the histogram received measurements, so it must be reported") + + actual := metricValues(t, events[0].MetricSetFields, "metric01") + assert.Equal(t, int64(2), actual["count"], "both measurements of the first interval must be counted") + assert.Equal(t, int64(2), actual["min"], "min must be the smallest measurement of the first interval") + assert.Equal(t, int64(4), actual["max"], "max must be the largest measurement of the first interval") + assert.InDelta(t, 3.0, actual["mean"], 0.001, "mean must average the measurements of the first interval") + + err = process([]string{"metric01:10|h"}, ms) + require.NoError(t, err, "the histogram packet of the second interval must be accepted") + + events = ms.getEvents() + require.Len(t, events, 1, "the histogram received a new measurement, so it must be reported again") + + actual = metricValues(t, events[0].MetricSetFields, "metric01") + assert.Equal(t, int64(1), actual["count"], "count must restart from the second interval instead of accumulating the first") + assert.Equal(t, int64(10), actual["min"], "min must forget the measurements of the first interval") + assert.Equal(t, int64(10), actual["max"], "max must forget the measurements of the first interval") + assert.InDelta(t, 10.0, actual["mean"], 0.001, "mean must average only the second interval") + + assert.Empty(t, ms.getEvents(), "a histogram that received no new measurements must not be reported") +} + +// TestIdleTimerDoesNotMaskOtherMetrics makes sure that an idle timer is skipped +// without dropping the metrics it shares a tag group with. +func TestIdleTimerDoesNotMaskOtherMetrics(t *testing.T) { + ms := newStatsdMetricSet(t, map[string]any{"module": "statsd"}) + + err := process([]string{"metric01:2|ms|#k1:v1", "metric02:1.0|g|#k1:v1"}, ms) + require.NoError(t, err, "the timer and gauge packets must be accepted") + + events := ms.getEvents() + require.Len(t, events, 2, "the timer and the gauge of the tag group must each be reported") + + events = ms.getEvents() + require.Len(t, events, 1, "the gauge must still be reported once the timer goes idle, skipping the timer must not drop its tag group") + assert.NotContains(t, events[0].MetricSetFields, "metric01", "the idle timer must not be reported") + assert.InDelta(t, 1.0, metricValues(t, events[0].MetricSetFields, "metric02")["value"], 0.001, + "the gauge must keep reporting its value, it is not scoped to a flush interval") +} + func TestChangeType(t *testing.T) { ms := mbtest.NewMetricSet(t, map[string]any{"module": "statsd"}).(*MetricSet) testData := []string{ diff --git a/x-pack/metricbeat/module/statsd/server/registry.go b/x-pack/metricbeat/module/statsd/server/registry.go index b2d4c4be34cc..98e129cccfa7 100644 --- a/x-pack/metricbeat/module/statsd/server/registry.go +++ b/x-pack/metricbeat/module/statsd/server/registry.go @@ -5,6 +5,7 @@ package server import ( + "sync/atomic" "time" "github.com/rcrowley/go-metrics" @@ -71,6 +72,9 @@ type samplingTimer struct { metrics.Timer meter metrics.Meter histogram metrics.Histogram + // count is the number of measurements since the last flush, tracked here + // because the meter count is monotonic and cannot be reset. + count atomic.Int64 } // NewSamplingTimer returns a new SamplingTimer @@ -88,25 +92,37 @@ func newSamplingTimer() *samplingTimer { // SampledUpdate will update the timer a sampled measurement func (s *samplingTimer) SampledUpdate(d time.Duration, sampleRate float64) { s.histogram.Update(int64(d)) - s.meter.Mark(int64(1 / sampleRate)) + count := int64(1 / sampleRate) + s.meter.Mark(count) + s.count.Add(count) } // Snapshot gets a snapshot of the SamplingTimer func (s *samplingTimer) Snapshot() samplingTimerSnapshot { return samplingTimerSnapshot{ + count: s.count.Load(), histogram: s.histogram.Snapshot(), meter: s.meter.Snapshot(), } } +// Reset clears the values that are aggregated over a single flush interval. The +// meter's moving average rates are deliberately left alone, they span fixed +// time windows and are not tied to the interval. +func (s *samplingTimer) Reset() { + s.count.Store(0) + s.histogram.Clear() +} + type samplingTimerSnapshot struct { + count int64 histogram metrics.Histogram meter metrics.Meter } -// Count returns the number of events recorded at the time the snapshot was -// taken. -func (t *samplingTimerSnapshot) Count() int64 { return t.meter.Count() } +// Count returns the number of events recorded since the last flush at the time +// the snapshot was taken. +func (t *samplingTimerSnapshot) Count() int64 { return t.count } // Max returns the maximum value at the time the snapshot was taken. func (t *samplingTimerSnapshot) Max() int64 { return t.histogram.Max() } @@ -183,6 +199,10 @@ type metricsGroup struct { metrics mapstr.M } +// getMetric returns the values to report for a single metric and resets the +// accumulators that are scoped to one flush interval, so that every flush +// reports what happened since the previous one. It returns nil when the metric +// has nothing to report for this interval. func (r *registry) getMetric(metric any) map[string]any { values := map[string]any{} switch m := metric.(type) { @@ -193,6 +213,11 @@ func (r *registry) getMetric(metric any) map[string]any { values["value"] = m.Value() case metrics.Histogram: h := m.Snapshot() + if h.Count() == 0 { + // Reporting the aggregates of an empty sample would only add + // zeroed out min/max/mean values. + return nil + } ps := h.Percentiles([]float64{0.5, 0.75, 0.95, 0.99, 0.999}) values["count"] = h.Count() values["min"] = h.Min() @@ -204,8 +229,15 @@ func (r *registry) getMetric(metric any) map[string]any { values["p95"] = ps[2] values["p99"] = ps[3] values["p99_9"] = ps[4] + m.Clear() case *samplingTimer: t := m.Snapshot() + // Count() is extrapolated from the sample rate and rounds down to zero + // for a sample rate above one, so the raw number of measurements is + // what decides whether anything was recorded. See the histogram case. + if t.histogram.Count() == 0 { + return nil + } ps := t.Percentiles([]float64{0.5, 0.75, 0.95, 0.99, 0.999}) values["count"] = t.Count() values["min"] = t.Min() @@ -221,6 +253,7 @@ func (r *registry) getMetric(metric any) map[string]any { values["5m_rate"] = t.Rate5() values["15m_rate"] = t.Rate15() values["mean_rate"] = t.RateMean() + m.Reset() case *setMetric: values["count"] = m.Count() m.Reset() @@ -255,7 +288,12 @@ func (r *registry) GetAll() []metricsGroup { // all the .tags are the same for this metricsMap // we just need one tags = m.tags - fields[m.name] = r.getMetric(m.metric) + + values := r.getMetric(m.metric) + if len(values) == 0 { + continue + } + fields[m.name] = values } // cleanup the tag group if it's empty @@ -264,6 +302,11 @@ func (r *registry) GetAll() []metricsGroup { continue } + // every metric in this group is idle, there is nothing to report + if len(fields) == 0 { + continue + } + tagGroups = append(tagGroups, metricsGroup{ metrics: fields, tags: tags,