diff --git a/pkg/metric/profiling_metric.go b/pkg/metric/profiling_metric.go index 7fff92719..1ff910410 100644 --- a/pkg/metric/profiling_metric.go +++ b/pkg/metric/profiling_metric.go @@ -30,6 +30,7 @@ import ( "gvisor.dev/gvisor/pkg/atomicbitops" "gvisor.dev/gvisor/pkg/log" "gvisor.dev/gvisor/pkg/prometheus" + "gvisor.dev/gvisor/pkg/sync" ) const ( @@ -59,6 +60,53 @@ const ( MetricsStartTimeIndicator = "START_TIME\t" ) +// CollectionStats contains statistics about the profiling metrics collection +// process itself. +type CollectionStats struct { + // mu protects the fields below. + mu sync.Mutex `json:"-"` + + // CollectionRate is the rate at which the metrics are meant to be + // collected. + CollectionRateNanos uint64 `json:"collection_rate"` + + // CheapStartNanos is the time at which the collector started in nanoseconds, + // as returned by CheapNowNano. + CheapStartNanos uint64 `json:"cheap_start_nanos"` + + // CheapLastCollectionNanos is the time at which the last collection was + // meant to be performed, in nanoseconds as returned by CheapNowNano. + CheapLastCollectionNanos uint64 `json:"cheap_last_collection_nanos"` + + // TotalSnapshots is the total number of snapshots successfully taken. + TotalSnapshots uint64 `json:"total_snapshots"` + + // TotalSleepTimingError is the running sum of absolute difference in timing + // between when the collector was meant to start collecting a metric snapshot + // vs when it actually collected it. It should be divided by TotalSnapshots + // to get the average sleep timing error. + TotalSleepTimingErrorNanos uint64 `json:"total_sleep_timing_error"` + + // TotalCollectionTimingError is the running sum of time spent doing actual + // metric collection, i.e. retrieving numerical values from metrics. + // The larger this duration is, the more time gap there is between the first + // metric being collected and the last within a single metric collection + // cycle. This can cause the metric data to be less accurate because all of + // these points will be recorded as having the same timestamp despite this + // not actually being the case. + // It should be divided by TotalSnapshots to get the average + // per-collection-cycle collection timing error. + TotalCollectionTimingErrorNanos uint64 `json:"total_collection_timing_error"` + + // NumBackoffSleeps is the number of times the collector had to back off + // because the writer was too slow. + NumBackoffSleeps uint64 `json:"num_backoff_sleeps"` + + // TotalBackoffSleep is the running sum of time the collector had to back + // off because the writer was too slow. + TotalBackoffSleepNanos uint64 `json:"total_backoff_sleep"` +} + var ( // profilingMetricsStarted indicates whether StartProfilingMetrics has // been called. @@ -67,8 +115,9 @@ var ( // goroutine to stop recording and writing metrics. stopProfilingMetrics atomicbitops.Bool // doneProfilingMetrics is used to signal that the profiling metrics - // goroutines are finished. - doneProfilingMetrics chan bool + // goroutines are finished. It carries information about the stats of + // the profiling metrics collection process. + doneProfilingMetrics chan *CollectionStats // definedProfilingMetrics is the set of metrics known to be created for // profiling (see condmetric_profiling.go). definedProfilingMetrics []string @@ -209,17 +258,21 @@ func StartProfilingMetrics[T ProfilingMetricsWriter](opts ProfilingMetricsOption _ = opts.Sink.Truncate(0) stopProfilingMetrics = atomicbitops.FromBool(false) - doneProfilingMetrics = make(chan bool, 1) + doneProfilingMetrics = make(chan *CollectionStats, 1) writeCh := make(chan writeReq, snapshotRingbufferSize) s.startTime = time.Now().UnixNano() cheapStartTime := CheapNowNano() - go collectProfilingMetrics(&s, values, cheapStartTime, opts.Rate, writeCh) + stats := CollectionStats{ + CollectionRateNanos: uint64(opts.Rate.Nanoseconds()), + CheapStartNanos: uint64(cheapStartTime), + } + go collectProfilingMetrics(&s, values, cheapStartTime, opts.Rate, writeCh, &stats) if opts.Lossy { lossySink := newLossyBufferedWriter(opts.Sink) - go writeProfilingMetrics[*lossyBufferedWriter[T]](lossySink, &s, headers, writeCh) + go writeProfilingMetrics[*lossyBufferedWriter[T]](lossySink, &s, headers, writeCh, &stats) } else { bufferedSink := newBufferedWriter(opts.Sink) - go writeProfilingMetrics[*bufferedWriter[T]](bufferedSink, &s, headers, writeCh) + go writeProfilingMetrics[*bufferedWriter[T]](bufferedSink, &s, headers, writeCh, &stats) } log.Infof("Profiling metrics started.") @@ -228,12 +281,16 @@ func StartProfilingMetrics[T ProfilingMetricsWriter](opts ProfilingMetricsOption // collectProfilingMetrics will send metrics to the writeCh until it receives a // signal via the stopProfilingMetrics channel. -func collectProfilingMetrics(s *snapshots, values []func(fieldValues ...*FieldValue) uint64, cheapStartTime int64, profilingRate time.Duration, writeCh chan<- writeReq) { +func collectProfilingMetrics(s *snapshots, values []func(fieldValues ...*FieldValue) uint64, cheapStartTime int64, profilingRate time.Duration, writeCh chan<- writeReq, stats *CollectionStats) { defer close(writeCh) + stats.mu.Lock() + defer stats.mu.Unlock() + numEntries := s.numMetrics + 1 // to account for the timestamp ringbufferIdx := 0 curSnapshot := 0 + var beforeCollectionTimestamp int64 // If we write faster than the writer can keep up, we back off. // The backoff factor starts small but increases exponentially @@ -254,16 +311,21 @@ func collectProfilingMetrics(s *snapshots, values []func(fieldValues ...*FieldVa stopCollecting := false for nextCollection := cheapStartTime; !stopCollecting; nextCollection += profilingRate.Nanoseconds() { + if stopProfilingMetrics.Load() { + stopCollecting = true + stats.CheapLastCollectionNanos = uint64(nextCollection) + // Collect one last time before stopping. + } - // For small durations, just spin. Otherwise sleep. + // For small durations, just spin (and maybe yield). Otherwise sleep. for { const ( wakeUpNanos = 10 spinMaxNanos = 250 yieldMaxNanos = 1_000 ) - now := CheapNowNano() - nanosToNextCollection := nextCollection - now + beforeCollectionTimestamp = CheapNowNano() + nanosToNextCollection := nextCollection - beforeCollectionTimestamp if nanosToNextCollection <= 0 { // Collect now. break @@ -280,19 +342,19 @@ func collectProfilingMetrics(s *snapshots, values []func(fieldValues ...*FieldVa time.Sleep(time.Duration(nanosToNextCollection-wakeUpNanos) * time.Nanosecond) } - if stopProfilingMetrics.Load() { - stopCollecting = true - // Collect one last time before stopping. - } - - timestamp := time.Duration(CheapNowNano() - cheapStartTime) - base := curSnapshot * numEntries ringBuf := s.ringbuffer[ringbufferIdx] - ringBuf[base] = uint64(timestamp) + base := curSnapshot * numEntries for i := 1; i < numEntries; i++ { ringBuf[base+i] = values[i-1]() } + afterCollectionTimestamp := CheapNowNano() + middleCollectionTimestamp := (beforeCollectionTimestamp + afterCollectionTimestamp) / 2 + relativeCollectionTimestamp := time.Duration(middleCollectionTimestamp - cheapStartTime) + ringBuf[base] = uint64(relativeCollectionTimestamp.Nanoseconds()) curSnapshot++ + stats.TotalSnapshots++ + stats.TotalSleepTimingErrorNanos += uint64(max(nextCollection-beforeCollectionTimestamp, beforeCollectionTimestamp-nextCollection)) + stats.TotalCollectionTimingErrorNanos += uint64(afterCollectionTimestamp - beforeCollectionTimestamp) if curSnapshot == snapshotBufferSize { writeCh <- writeReq{ringbufferIdx: ringbufferIdx, numLines: curSnapshot} @@ -303,6 +365,8 @@ func collectProfilingMetrics(s *snapshots, values []func(fieldValues ...*FieldVa // Going too fast, stop collecting for a bit. backoffSleep := profilingRate * time.Duration(backoffFactor) log.Warningf("Profiling metrics collector exhausted the entire ringbuffer... backing off for %v to let writer catch up.", backoffSleep) + stats.NumBackoffSleeps++ + stats.TotalBackoffSleepNanos += uint64(backoffSleep.Nanoseconds()) time.Sleep(backoffSleep) backoffFactor = min(backoffFactor*backoffFactorGrowth, backoffFactorMax) } @@ -504,7 +568,7 @@ func (w *lossyBufferedWriter[T]) Close() error { // writeProfilingMetrics will write to the ProfilingMetricsWriter on every // request via writeReqs, until writeReqs is closed. -func writeProfilingMetrics[T bufferedMetricsWriter](sink T, s *snapshots, headers []string, writeReqs <-chan writeReq) { +func writeProfilingMetrics[T bufferedMetricsWriter](sink T, s *snapshots, headers []string, writeReqs <-chan writeReq, stats *CollectionStats) { numEntries := s.numMetrics + 1 for _, header := range headers { sink.WriteString(header) @@ -527,11 +591,73 @@ func writeProfilingMetrics[T bufferedMetricsWriter](sink T, s *snapshots, header } sink.Close() - doneProfilingMetrics <- true + doneProfilingMetrics <- stats close(doneProfilingMetrics) profilingMetricsStarted.Store(false) } +// Log logs some statistics about the profiling metrics collection process. +func (s *CollectionStats) Log() { + s.mu.Lock() + defer s.mu.Unlock() + + captureDuration := time.Duration(s.CheapLastCollectionNanos-s.CheapStartNanos) * time.Nanosecond + collectionRate := time.Duration(s.CollectionRateNanos) * time.Nanosecond + totalSleepTimingError := time.Duration(s.TotalSleepTimingErrorNanos) * time.Nanosecond + totalCollectionTimingError := time.Duration(s.TotalCollectionTimingErrorNanos) * time.Nanosecond + totalBackoffSleep := time.Duration(s.TotalBackoffSleepNanos) * time.Nanosecond + + expectedSnapshots := uint64(captureDuration / collectionRate) + if s.TotalSnapshots == expectedSnapshots+1 { + // Depending on the timing of when the stop signal was sent, the + // collection goroutine is expected to do an extra collection cycle, + // so we add one to the expected number of snapshots here to make it not + // look like the capture rate is >100%. + expectedSnapshots = s.TotalSnapshots + } + captureRate := 0.0 + if expectedSnapshots > 0 { + captureRate = float64(s.TotalSnapshots) / float64(expectedSnapshots) + } + if captureRate < .99 { + log.Warningf("Profiling metrics: Captured %d snapshots out of %d expected (%.2f%% capture rate) over %v.", s.TotalSnapshots, expectedSnapshots, captureRate*100.0, captureDuration) + log.Warningf("This indicates that the profiling metrics writer is not keeping up with the metrics collection rate.") + log.Warningf("Ensure that the profiling metrics log is stored on a fast storage device, or consider reducing the metric profiling rate.") + } else { + log.Infof("Profiling metrics: Captured %d snapshots out of %d expected (%.2f%% capture rate) over %v. This is OK.", s.TotalSnapshots, expectedSnapshots, captureRate*100.0, captureDuration) + } + averageSleepTimingError := totalSleepTimingError / time.Duration(s.TotalSnapshots) + sleepTimingErrorVsRate := float64(averageSleepTimingError) / float64(collectionRate) + if sleepTimingErrorVsRate > .1 { + log.Warningf("Profiling metrics: Average sleep timing error is high: %v (%.2f%% of the collection interval).", averageSleepTimingError, sleepTimingErrorVsRate*100.0) + log.Warningf("This means the profiling metrics collector is not waking up at the correct time to collect the next snapshot.") + log.Warningf("This may mean that the CPU is overloaded (e.g. from other processes running on the same machine or from the workload itself taking up all the cores).") + log.Warningf("Consider using a slower profiling rate, removing other background processes, or tweaking the sleep consts in profiling_metric.go.") + } else { + log.Infof("Profiling metrics: Average sleep timing error: %v (%.2f%% of the collection interval). This is OK.", averageSleepTimingError, sleepTimingErrorVsRate*100.0) + } + averageCollectionTimingError := totalCollectionTimingError / time.Duration(s.TotalSnapshots) + collectionTimingErrorVsRate := float64(averageCollectionTimingError) / float64(collectionRate) + if collectionTimingErrorVsRate > .1 { + log.Warningf("Profiling metrics: Average collection timing error is high: %v (%.2f%% of the collection interval).", averageCollectionTimingError, collectionTimingErrorVsRate*100.0) + log.Warningf("This means the time between getting the value of the first metric vs the last metric within a single collection cycle is a too large fraction of the profiling rate.") + log.Warningf("This means there is significant drift between the time a value is reported as having vs the time it was actually scraped, relative to the profiling interval.") + log.Warningf("Consider using a slower profiling rate or profiling fewer metrics at a time.") + } else { + log.Infof("Profiling metrics: Average collection timing error: %v (%.2f%% of the collection interval). This is OK.", averageCollectionTimingError, collectionTimingErrorVsRate*100.0) + } + if s.NumBackoffSleeps > 0 { + ratioLostToBackoff := float64(totalBackoffSleep) / float64(captureDuration) + if ratioLostToBackoff > .05 { + log.Warningf("Profiling metrics: Backed off %d times due to slow writer; total %v spent in backoff sleep (%.2f%% of the capture duration).") + log.Warningf("This indicates that the profiling metrics writer is not keeping up with the metrics collection rate.") + log.Warningf("Ensure that the profiling metrics log is stored on a fast storage device, or consider reducing the metric profiling rate.") + } else { + log.Infof("Profiling metrics: Backed off %d times due to slow writer; total %v spent in backoff sleep (%.2f%% of the capture duration). This is OK.", s.NumBackoffSleeps, totalBackoffSleep, ratioLostToBackoff*100.0) + } + } +} + // StopProfilingMetrics stops the profiling metrics goroutines. Call to make sure // all metric data has been flushed. // Note that calling this function prior to StartProfilingMetrics has no effect. @@ -539,9 +665,12 @@ func StopProfilingMetrics() { if !profilingMetricsStarted.Load() { return } - if stopProfilingMetrics.CompareAndSwap(false, true) { - <-doneProfilingMetrics + if !stopProfilingMetrics.CompareAndSwap(false, true) { + // If the CAS fails, this means the signal was already sent, + // so don't wait on doneProfilingMetrics. + return } - // If the CAS fails, this means the signal was already sent, - // so don't wait on doneProfilingMetrics. + stats := <-doneProfilingMetrics + log.Infof("Profiling metrics stopped.") + stats.Log() }