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,