Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -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
2 changes: 2 additions & 0 deletions docs/reference/metricbeat/metricbeat-module-statsd.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:

Expand Down
2 changes: 2 additions & 0 deletions x-pack/metricbeat/module/statsd/_meta/docs.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:

Expand Down
138 changes: 138 additions & 0 deletions x-pack/metricbeat/module/statsd/server/data_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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{
Expand Down
53 changes: 48 additions & 5 deletions x-pack/metricbeat/module/statsd/server/registry.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
package server

import (
"sync/atomic"
"time"

"github.com/rcrowley/go-metrics"
Expand Down Expand Up @@ -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
Expand All @@ -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() }
Expand Down Expand Up @@ -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) {
Expand All @@ -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()
Expand All @@ -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()
Expand All @@ -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()
Expand Down Expand Up @@ -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
Expand All @@ -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,
Expand Down
Loading