diff --git a/pkg/metric/metric.go b/pkg/metric/metric.go index 4bea1b487..66f994f43 100644 --- a/pkg/metric/metric.go +++ b/pkg/metric/metric.go @@ -825,9 +825,9 @@ type DistributionMetric struct { // last (i.e. infinite) bucket. samples [][]atomicbitops.Uint64 - // sampleSum is the sum of samples. + // statistics is a set of statistics about each distribution. // It is mapped by the concatenation of the fields using `fieldsToKey`. - sampleSum []atomicbitops.Int64 + statistics []distributionStatistics } // NewDistributionMetric creates and registers a new distribution metric. @@ -867,7 +867,7 @@ func NewDistributionMetric(name string, sync bool, bucketer Bucketer, unit pb.Me exponentialBucketer: exponentialBucketer, fieldsToKey: fieldsToKey, samples: samples, - sampleSum: make([]atomicbitops.Int64, fieldsToKey.numKeys()), + statistics: make([]distributionStatistics, fieldsToKey.numKeys()), metadata: &pb.MetricMetadata{ Name: name, PrometheusName: nameToPrometheusName(name), @@ -898,6 +898,132 @@ func MustCreateNewDistributionMetric(name string, sync bool, bucketer Bucketer, return distrib } +// distributionStatistics is a set of useful statistics for a distribution. +// As metric update operations must be non-blocking, this uses a bunch of +// atomic numbers rather than a mutex. +type distributionStatistics struct { + // sampleCount is the total number of samples. + sampleCount atomicbitops.Uint64 + + // sampleSum is the sum of samples. + sampleSum atomicbitops.Int64 + + // sumOfSquaredDeviations is the running sum of squared deviations from the + // mean of each sample. + // This quantity is useful as part of Welford's online algorithm: + // https://en.wikipedia.org/wiki/Algorithms_for_calculating_variance#Welford's_online_algorithm + sumOfSquaredDeviations atomicbitops.Float64 + + // min and max are the minimum and maximum samples ever recorded. + min, max atomicbitops.Int64 +} + +// Update updates the distribution statistics with the given sample. +// This function must be non-blocking, i.e. no mutexes. +// As a result, it is not entirely accurate when it races with itself, +// though the imprecision should be fairly small and should not practically +// matter for distributions with more than a handful of records. +func (s *distributionStatistics) Update(sample int64) { + newSampleCount := s.sampleCount.Add(1) + newSampleSum := s.sampleSum.Add(sample) + + if newSampleCount > 1 { + // Not the first sample of the distribution. + floatSample := float64(sample) + oldMean := float64(newSampleSum-sample) / float64(newSampleCount-1) + newMean := float64(newSampleSum) / float64(newSampleCount) + devSquared := (floatSample - oldMean) * (floatSample - newMean) + s.sumOfSquaredDeviations.Add(devSquared) + + // Update min and max. + // We optimistically load racily here in the hope that it passes the CaS + // operation. If it doesn't, we'll load it atomically, so this is not a + // race. + sync.RaceDisable() + for oldMin := s.min.RacyLoad(); sample < oldMin && !s.min.CompareAndSwap(oldMin, sample); oldMin = s.min.Load() { + } + for oldMax := s.max.RacyLoad(); sample > oldMax && !s.max.CompareAndSwap(oldMax, sample); oldMax = s.max.Load() { + } + sync.RaceEnable() + } else { + // We are the first sample, so set the min and max to the current sample. + // See above for why disabling race detection is safe here as well. + sync.RaceDisable() + if !s.min.CompareAndSwap(0, sample) { + for oldMin := s.min.RacyLoad(); sample < oldMin && !s.min.CompareAndSwap(oldMin, sample); oldMin = s.min.Load() { + } + } + if !s.max.CompareAndSwap(0, sample) { + for oldMax := s.max.RacyLoad(); sample > oldMax && !s.max.CompareAndSwap(oldMax, sample); oldMax = s.max.Load() { + } + } + sync.RaceEnable() + } +} + +// distributionStatisticsSnapshot an atomically-loaded snapshot of +// distributionStatistics. +type distributionStatisticsSnapshot struct { + // sampleCount is the total number of samples. + sampleCount uint64 + + // sampleSum is the sum of samples. + sampleSum int64 + + // sumOfSquaredDeviations is the running sum of squared deviations from the + // mean of each sample. + // This quantity is useful as part of Welford's online algorithm: + // https://en.wikipedia.org/wiki/Algorithms_for_calculating_variance#Welford's_online_algorithm + sumOfSquaredDeviations float64 + + // min and max are the minimum and maximum samples ever recorded. + min, max int64 +} + +// Load generates a consistent snapshot of the distribution statistics. +func (s *distributionStatistics) Load() distributionStatisticsSnapshot { + // We start out reading things racily, but will verify each of them + // atomically later in this function, so this is OK. Disable the race + // checker for this part of the function. + sync.RaceDisable() + snapshot := distributionStatisticsSnapshot{ + sampleCount: s.sampleCount.RacyLoad(), + sampleSum: s.sampleSum.RacyLoad(), + sumOfSquaredDeviations: s.sumOfSquaredDeviations.RacyLoad(), + min: s.min.RacyLoad(), + max: s.max.RacyLoad(), + } + sync.RaceEnable() + + // Now verify that we loaded an atomic snapshot of the statistics. + // This relies on the fact that each update should at least change the + // count statistic, so we should be able to tell if anything changed based + // on whether we have an exact match with the currently-loaded values. + // If not, we reload that value and try again until all is consistent. +retry: + if sampleCount := s.sampleCount.Load(); sampleCount != snapshot.sampleCount { + snapshot.sampleCount = sampleCount + goto retry + } + if sampleSum := s.sampleSum.Load(); sampleSum != snapshot.sampleSum { + snapshot.sampleSum = sampleSum + goto retry + } + if ssd := s.sumOfSquaredDeviations.Load(); ssd != snapshot.sumOfSquaredDeviations { + snapshot.sumOfSquaredDeviations = ssd + goto retry + } + if min := s.min.Load(); min != snapshot.min { + snapshot.min = min + goto retry + } + if max := s.max.Load(); max != snapshot.max { + snapshot.max = max + goto retry + } + return snapshot +} + // AddSample adds a sample to the distribution. // This *must* be called with the correct number of fields, or it will panic. // +checkescape:all @@ -914,7 +1040,7 @@ func (d *DistributionMetric) AddSample(sample int64, fields ...*FieldValue) { func (d *DistributionMetric) addSampleByKey(sample int64, key int) { bucket := d.exponentialBucketer.BucketIndex(sample) d.samples[key][bucket+1].Add(1) - d.sampleSum[key].Add(sample) + d.statistics[key].Update(sample) } // Minimum number of buckets for NewDurationBucket. @@ -1067,7 +1193,7 @@ func (m *metricSet) Values() metricValues { uint64Metrics: make(map[string]any, len(m.uint64Metrics)), distributionMetrics: make(map[string][][]uint64, len(m.distributionMetrics)), distributionTotalSamples: make(map[string][]uint64, len(m.distributionMetrics)), - distributionSampleSum: make(map[string][]int64, len(m.distributionMetrics)), + distributionStatistics: make(map[string][]distributionStatisticsSnapshot, len(m.distributionMetrics)), stages: stages, } for k, v := range m.uint64Metrics { @@ -1094,7 +1220,7 @@ func (m *metricSet) Values() metricValues { for name, metric := range m.distributionMetrics { fieldKeysToValues := make([][]uint64, len(metric.samples)) fieldKeysToTotalSamples := make([]uint64, len(metric.samples)) - fieldKeysToSampleSum := make([]int64, len(metric.samples)) + fieldKeysToStatistics := make([]distributionStatisticsSnapshot, len(metric.samples)) for fieldKey, samples := range metric.samples { samplesSnapshot := snapshotDistribution(samples) totalSamples := uint64(0) @@ -1106,17 +1232,17 @@ func (m *metricSet) Values() metricValues { // the maps for this fieldKey as nil. This lessens the memory cost // of distributions with unused field combinations. fieldKeysToTotalSamples[fieldKey] = 0 - fieldKeysToSampleSum[fieldKey] = 0 + fieldKeysToStatistics[fieldKey] = distributionStatisticsSnapshot{} fieldKeysToValues[fieldKey] = nil } else { fieldKeysToTotalSamples[fieldKey] = totalSamples - fieldKeysToSampleSum[fieldKey] = metric.sampleSum[fieldKey].Load() + fieldKeysToStatistics[fieldKey] = metric.statistics[fieldKey].Load() fieldKeysToValues[fieldKey] = samplesSnapshot } } vals.distributionMetrics[name] = fieldKeysToValues vals.distributionTotalSamples[name] = fieldKeysToTotalSamples - vals.distributionSampleSum[name] = fieldKeysToSampleSum + vals.distributionStatistics[name] = fieldKeysToStatistics } return vals } @@ -1147,9 +1273,8 @@ type metricValues struct { // no new samples are not retransmitted. distributionTotalSamples map[string][]uint64 - // distributionSampleSum is the sum of samples for each distribution metric - // and field values. - distributionSampleSum map[string][]int64 + // distributionStatistics is a set of statistics about the samples. + distributionStatistics map[string][]distributionStatisticsSnapshot // Information on when initialization stages were reached. Does not include // the currently-ongoing stage, if any. @@ -1342,7 +1467,7 @@ func GetSnapshot(options SnapshotOptions) (*prometheus.Snapshot, error) { } distributionSamples := values.distributionMetrics[k] numFiniteBuckets := m.exponentialBucketer.NumFiniteBuckets() - sampleSums := values.distributionSampleSum[k] + statistics := values.distributionStatistics[k] for fieldKey := range dists { var labels map[string]string if numFields := m.fieldsToKey.numKeys(); numFields > 0 { @@ -1380,7 +1505,7 @@ func GetSnapshot(options SnapshotOptions) (*prometheus.Snapshot, error) { Metric: m.prometheusMetric, Labels: labels, HistogramValue: &prometheus.Histogram{ - Total: prometheus.Number{Int: sampleSums[fieldKey]}, + Total: prometheus.Number{Int: statistics[fieldKey].sampleSum}, Buckets: buckets, }, })