From 5ac415f47cc30cd4aad1b0507729bdf237d6659e Mon Sep 17 00:00:00 2001 From: Pavan Nalam Date: Wed, 7 Oct 2026 12:12:10 -0700 Subject: [PATCH] Bound memory use during scrapes to avoid OOM kills reportMonitoringMetrics starts one goroutine per metric descriptor per project with no cap, and each holds a fully decoded TimeSeries.List response while in flight. When google.projects.filter (or a long google.project-ids list) resolves to many projects, or a project has many metric descriptors, the number of in-flight responses is unbounded and can push a memory-limited container past its limit and get it OOM killed. Fix this without adding any configuration: - Cap concurrent TimeSeries.List requests process-wide at a fixed internal limit (maxConcurrentTimeSeriesRequests), shared by every project's collector, so a scrape never holds more than that many decoded responses in memory at once. - Read the container's cgroup memory limit at startup and set GOMEMLIMIT to 90% of it via automemlimit, so the garbage collector reclaims more aggressively as usage approaches the real limit. GOMEMLIMIT and AUTOMEMLIMIT still override it. Signed-off-by: Pavan Nalam --- README.md | 24 +++++ collectors/monitoring_collector.go | 7 ++ collectors/monitoring_collector_test.go | 124 ++++++++++++++++++++++++ go.mod | 2 + go.sum | 4 + stackdriver_exporter.go | 4 + 6 files changed, 165 insertions(+) diff --git a/README.md b/README.md index 17273837..ff87ce47 100644 --- a/README.md +++ b/README.md @@ -189,6 +189,30 @@ stackdriver_exporter \ --google.projects.filter='labels.monitoring="true"' ``` +### Memory limits in constrained environments + +When `google.projects.filter` (or a long, repeated `google.project-ids`) resolves to many +projects, each scrape fetches metrics for every metric descriptor of every project. In +memory-constrained environments (e.g. a GKE pod with a small memory limit), an unbounded +fan-out of concurrent Monitoring API requests and JSON decodes can grow the heap enough to +OOM the process. Two fixes address this together, and neither requires configuration: + +- Concurrent `TimeSeries.List` requests are capped process-wide at a fixed internal limit, + so a scrape can never have more than that many API responses in memory at once, no matter + how many projects or descriptors it fans out across. +- `stackdriver_exporter` also reads the container's memory limit (from the cgroup) at + startup and sets Go's [`GOMEMLIMIT`][gomemlimit] to 90% of it, via + [`automemlimit`][automemlimit]. This makes the garbage collector reclaim memory more + aggressively as usage approaches the limit. This happens automatically whenever a memory + limit is set on the container (e.g. `resources.limits.memory` in a Kubernetes pod spec). + +If you need to override the detected value, set the `GOMEMLIMIT` environment variable +directly, or `AUTOMEMLIMIT` to change the ratio (or `off` to disable auto-detection +entirely). + +[gomemlimit]: https://pkg.go.dev/runtime/debug#SetMemoryLimit +[automemlimit]: https://github.com/KimMachineGun/automemlimit + ### Filtering enabled collectors The `stackdriver_exporter` collects all metrics type prefixes by default. diff --git a/collectors/monitoring_collector.go b/collectors/monitoring_collector.go index 7d5b138d..b1c9667f 100644 --- a/collectors/monitoring_collector.go +++ b/collectors/monitoring_collector.go @@ -29,6 +29,10 @@ import ( const namespace = "stackdriver" +const maxConcurrentTimeSeriesRequests = 20 + +var timeSeriesRequestLimiter = make(chan struct{}, maxConcurrentTimeSeriesRequests) + type MetricFilter struct { TargetedMetricPrefix string FilterQuery string @@ -339,6 +343,9 @@ func (c *MonitoringCollector) reportMonitoringMetrics(ch chan<- prometheus.Metri c.logger.Debug("retrieving Google Stackdriver Monitoring metrics with filter", "filter", filter) + timeSeriesRequestLimiter <- struct{}{} + defer func() { <-timeSeriesRequestLimiter }() + timeSeriesListCall := c.monitoringService.Projects.TimeSeries.List(projectResource(c.projectID)). Filter(filter). IntervalStartTime(startTime.Format(time.RFC3339Nano)). diff --git a/collectors/monitoring_collector_test.go b/collectors/monitoring_collector_test.go index 9e391ef2..d2f798c3 100644 --- a/collectors/monitoring_collector_test.go +++ b/collectors/monitoring_collector_test.go @@ -14,8 +14,20 @@ package collectors import ( + "context" + "encoding/json" + "log/slog" + "net/http" + "net/http/httptest" "reflect" + "strings" + "sync" "testing" + "time" + + "github.com/prometheus/client_golang/prometheus" + "google.golang.org/api/monitoring/v3" + "google.golang.org/api/option" ) func TestIsGoogleMetric(t *testing.T) { @@ -113,3 +125,115 @@ func TestProjectResource(t *testing.T) { t.Fatalf("projectResource() = %q, want %q", got, "projects/fake-project-1") } } + +func TestTimeSeriesRequestLimiterBoundsConcurrency(t *testing.T) { + numDescriptors := maxConcurrentTimeSeriesRequests + 10 + + var ( + mu sync.Mutex + inFlight int + maxInFlight int + ) + + mux := http.NewServeMux() + mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) { + switch { + case strings.Contains(r.URL.Path, "metricDescriptors"): + descriptors := make([]*monitoring.MetricDescriptor, 0, numDescriptors) + for i := 0; i < numDescriptors; i++ { + descriptors = append(descriptors, &monitoring.MetricDescriptor{ + Type: "custom.googleapis.com/metric_" + strings.Repeat("a", i+1), + }) + } + writeJSONResponse(w, &monitoring.ListMetricDescriptorsResponse{MetricDescriptors: descriptors}) + + case strings.Contains(r.URL.Path, "timeSeries"): + mu.Lock() + inFlight++ + if inFlight > maxInFlight { + maxInFlight = inFlight + } + mu.Unlock() + + time.Sleep(50 * time.Millisecond) + + mu.Lock() + inFlight-- + mu.Unlock() + + writeJSONResponse(w, &monitoring.ListTimeSeriesResponse{}) + + default: + http.NotFound(w, r) + } + }) + + server := httptest.NewServer(mux) + defer server.Close() + + ctx := context.Background() + service, err := monitoring.NewService(ctx, + option.WithHTTPClient(server.Client()), + option.WithEndpoint(server.URL), + option.WithoutAuthentication(), + ) + if err != nil { + t.Fatalf("failed to create monitoring service: %v", err) + } + + logger := slog.New(slog.DiscardHandler) + + collector, err := NewMonitoringCollector( + "test-project", + service, + MonitoringCollectorOptions{ + MetricTypePrefixes: []string{"custom.googleapis.com"}, + RequestInterval: 5 * time.Minute, + DescriptorCacheTTL: 0, + }, + logger, + noopCounterStore{}, + noopHistogramStore{}, + ) + if err != nil { + t.Fatalf("failed to create collector: %v", err) + } + + ch := make(chan prometheus.Metric, numDescriptors+10) + done := make(chan struct{}) + go func() { + defer close(done) + for range ch { + } + }() + + collector.Collect(ch) + close(ch) + <-done + + mu.Lock() + observed := maxInFlight + mu.Unlock() + + if observed > maxConcurrentTimeSeriesRequests { + t.Fatalf("observed %d concurrent TimeSeries.List requests, want <= %d", observed, maxConcurrentTimeSeriesRequests) + } + if observed < maxConcurrentTimeSeriesRequests { + t.Fatalf("expected concurrency to reach the limiter's capacity %d, got max observed %d; test may not be exercising real contention", maxConcurrentTimeSeriesRequests, observed) + } +} + +type noopCounterStore struct{} + +func (noopCounterStore) Increment(*monitoring.MetricDescriptor, *ConstMetric) {} +func (noopCounterStore) ListMetrics(string) []*ConstMetric { return nil } + +type noopHistogramStore struct{} + +func (noopHistogramStore) Increment(*monitoring.MetricDescriptor, *HistogramMetric) {} +func (noopHistogramStore) ListMetrics(string) []*HistogramMetric { return nil } + +func writeJSONResponse(w http.ResponseWriter, v any) { + w.Header().Set("Content-Type", "application/json") + _ = json.NewEncoder(w).Encode(v) +} diff --git a/go.mod b/go.mod index 9dc21439..aa56b816 100644 --- a/go.mod +++ b/go.mod @@ -3,6 +3,7 @@ module github.com/prometheus-community/stackdriver_exporter go 1.26.0 require ( + github.com/KimMachineGun/automemlimit v1.0.0 github.com/PuerkitoBio/rehttp v1.4.0 github.com/alecthomas/kingpin/v2 v2.4.0 github.com/fatih/camelcase v1.0.0 @@ -39,6 +40,7 @@ require ( github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect github.com/mwitkow/go-conntrack v0.0.0-20190716064945-2f068394615f // indirect github.com/nxadm/tail v1.4.8 // indirect + github.com/pbnjay/memory v0.0.0-20210728143218-7b4eea64cf58 // indirect github.com/prometheus/client_model v0.6.2 // indirect github.com/prometheus/procfs v0.21.1 // indirect github.com/xhit/go-str2duration/v2 v2.1.0 // indirect diff --git a/go.sum b/go.sum index 2b7f2b1c..f50534ef 100644 --- a/go.sum +++ b/go.sum @@ -4,6 +4,8 @@ cloud.google.com/go/auth/oauth2adapt v0.2.8 h1:keo8NaayQZ6wimpNSmW5OPc283g65QNIi cloud.google.com/go/auth/oauth2adapt v0.2.8/go.mod h1:XQ9y31RkqZCcwJWNSx2Xvric3RrU88hAYYbjDWYDL+c= cloud.google.com/go/compute/metadata v0.9.0 h1:pDUj4QMoPejqq20dK0Pg2N4yG9zIkYGdBtwLoEkH9Zs= cloud.google.com/go/compute/metadata v0.9.0/go.mod h1:E0bWwX5wTnLPedCKqk3pJmVgCBSM6qQI1yTBdEb3C10= +github.com/KimMachineGun/automemlimit v1.0.0 h1:+MqlvDE/pkJNjk1rU+O14QsH8k10nJAD0frB0lsxyvw= +github.com/KimMachineGun/automemlimit v1.0.0/go.mod h1:n+BSXxQWDFS1DKh67Rqo0lgTsowsg6x65ak5uyngML0= github.com/PuerkitoBio/rehttp v1.4.0 h1:rIN7A2s+O9fmHUM1vUcInvlHj9Ysql4hE+Y0wcl/xk8= github.com/PuerkitoBio/rehttp v1.4.0/go.mod h1:LUwKPoDbDIA2RL5wYZCNsQ90cx4OJ4AWBmq6KzWZL1s= github.com/alecthomas/kingpin/v2 v2.4.0 h1:f48lwail6p8zpO1bC4TxtqACaGqHYA22qkHjHpqDjYY= @@ -86,6 +88,8 @@ github.com/onsi/gomega v1.7.1/go.mod h1:XdKZgCCFLUoM/7CFJVPcG8C1xQ1AJ0vpAezJrB7J github.com/onsi/gomega v1.10.1/go.mod h1:iN09h71vgCQne3DLsj+A5owkum+a2tYe+TOCB1ybHNo= github.com/onsi/gomega v1.43.0 h1:VlG/1FxqNxhSO+lq/OHBNaaqwiBK/mO8JbVkX9Y+FeU= github.com/onsi/gomega v1.43.0/go.mod h1:REff/hsDsodHoKlWsP2mAPhu1+5/6hVYNf9rIEBpeSg= +github.com/pbnjay/memory v0.0.0-20210728143218-7b4eea64cf58 h1:onHthvaw9LFnH4t2DcNVpwGmV9E1BkGknEliJkfwQj0= +github.com/pbnjay/memory v0.0.0-20210728143218-7b4eea64cf58/go.mod h1:DXv8WO4yhMYhSNPKjeNKa5WY9YCIEBRbNzFFPJbWO6Y= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= github.com/prometheus/client_golang v1.24.1 h1:JnJkREXzWxUdCuPFpIWZiPispT9xVV59uiuyR2bPlnU= github.com/prometheus/client_golang v1.24.1/go.mod h1:F+oSRECHg4sse5ucfYpYDeIv/hu68Zo0uoHKetWnzcE= diff --git a/stackdriver_exporter.go b/stackdriver_exporter.go index 85a4091d..d40a5c08 100644 --- a/stackdriver_exporter.go +++ b/stackdriver_exporter.go @@ -23,6 +23,7 @@ import ( "strconv" "strings" + "github.com/KimMachineGun/automemlimit/memlimit" "github.com/alecthomas/kingpin/v2" "github.com/prometheus/client_golang/prometheus" versioncollector "github.com/prometheus/client_golang/prometheus/collectors/version" @@ -228,6 +229,9 @@ func main() { kingpin.Parse() logger := promslog.New(promslogConfig) + + _, _ = memlimit.Set(memlimit.WithLogger(logger)) + if *projectID != "" { logger.Warn("The google.project-id flag is deprecated and will be replaced by google.project-ids.") }