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
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -88,6 +88,7 @@
* [BUGFIX] Compactor: Fix spurious `bucket operation fail after retries` error logs emitted during partial block cleanup. #7749
* [BUGFIX] Alertmanager: Fix panic in `validateAlertmanagerConfig` when receiver config traversal encounters nil interface values. #7751
* [BUGFIX] Parquet Converter: Fix `auto_forget_delay` having no effect. The ring lifecycler was created without the auto-forget delegate, so unhealthy instances were never automatically removed from the ring. #7752
* [BUGFIX] Ingester: Size the active queried series worker pool from `GOMAXPROCS` instead of `runtime.NumCPU()`, so it respects the container CPU quota (set via automaxprocs) rather than the host core count. #7671

## 1.21.1 2026-06-04

Expand Down
6 changes: 4 additions & 2 deletions pkg/ingester/active_queried_series.go
Original file line number Diff line number Diff line change
Expand Up @@ -385,8 +385,10 @@ type ActiveQueriedSeriesService struct {

// NewActiveQueriedSeriesService creates a new ActiveQueriedSeriesService service.
func NewActiveQueriedSeriesService(logger log.Logger, registerer prometheus.Registerer) *ActiveQueriedSeriesService {
// Cap at 4 workers to avoid excessive goroutines
numWorkers := max(min(runtime.NumCPU()/2, 4), 1)
// Cap at 4 workers to avoid excessive goroutines. Use GOMAXPROCS (which
// automaxprocs sets from the CPU cgroup quota) rather than NumCPU, so the
// pool is sized to the container's CPU budget instead of the host cores.
numWorkers := max(min(runtime.GOMAXPROCS(0)/2, 4), 1)

m := &ActiveQueriedSeriesService{
updateChan: make(chan activeQueriedSeriesUpdate, 10000), // Buffered channel to avoid blocking
Expand Down
28 changes: 28 additions & 0 deletions pkg/ingester/active_queried_series_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package ingester
import (
"context"
"fmt"
"runtime"
"sync"
"testing"
"time"
Expand Down Expand Up @@ -546,3 +547,30 @@ func TestActiveQueriedSeriesService_NoSendOnClosedChannelOnShutdown(t *testing.T
t.Fatalf("producer goroutine panicked during shutdown: %v", panicMsg.Load())
}
}

func TestActiveQueriedSeriesService_NumWorkersFollowsGOMAXPROCS(t *testing.T) {
// The worker pool must be sized from GOMAXPROCS (which automaxprocs derives
// from the CPU cgroup quota) and not from runtime.NumCPU (host cores, which
// ignores cgroup limits). Save and restore GOMAXPROCS around the test.
defer runtime.GOMAXPROCS(runtime.GOMAXPROCS(1))

for _, tc := range []struct {
gomaxprocs int
expectedWorkers int
}{
// GOMAXPROCS(1)/2 == 0, floored to the minimum of 1 worker. Against the
// buggy runtime.NumCPU() this yields min(NumCPU/2, 4) on a multi-core
// host and fails.
{gomaxprocs: 1, expectedWorkers: 1},
{gomaxprocs: 2, expectedWorkers: 1},
{gomaxprocs: 4, expectedWorkers: 2},
// Capped at 4 workers regardless of how many CPUs are available.
{gomaxprocs: 16, expectedWorkers: 4},
} {
t.Run(fmt.Sprintf("gomaxprocs=%d", tc.gomaxprocs), func(t *testing.T) {
runtime.GOMAXPROCS(tc.gomaxprocs)
svc := NewActiveQueriedSeriesService(log.NewNopLogger(), nil)
assert.Equal(t, tc.expectedWorkers, svc.numWorkers)
})
}
}