From 73958e7f65078c92688837a930075560f2eb4880 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Nicol=C3=A1s=20Pazos?= Date: Sun, 16 Aug 2026 16:45:44 -0300 Subject: [PATCH 1/5] compactor: report scheduler queue metrics per lane Adds a lane label to the backlog gauge and a per-lane last-empty timestamp, so a fleet serving one lane can be sized from its own backlog. --- pkg/compactor/scheduler/job_tracker_test.go | 20 ++-- pkg/compactor/scheduler/lane_policy.go | 43 ++++--- pkg/compactor/scheduler/metrics.go | 108 ++++++++++++------ .../scheduler/persistence_bbolt_test.go | 8 +- pkg/compactor/scheduler/rotator.go | 54 +++++---- pkg/compactor/scheduler/rotator_test.go | 92 +++++++++++---- pkg/compactor/scheduler/scheduler.go | 8 +- pkg/compactor/scheduler/scheduler_test.go | 7 +- 8 files changed, 227 insertions(+), 113 deletions(-) diff --git a/pkg/compactor/scheduler/job_tracker_test.go b/pkg/compactor/scheduler/job_tracker_test.go index cd3164930b8..c56206c8e34 100644 --- a/pkg/compactor/scheduler/job_tracker_test.go +++ b/pkg/compactor/scheduler/job_tracker_test.go @@ -21,9 +21,13 @@ func at(hour, minute int) time.Time { return time.Date(2026, 1, 2, hour, minute, 0, 0, time.UTC) } +func newTestSchedulerMetrics(reg prometheus.Registerer) *schedulerMetrics { + return newSchedulerMetrics(reg, newSimpleLanePolicy()) +} + func newTestJobTracker(clk clock.Clock) (*JobTracker, *prometheus.Registry) { reg := prometheus.NewPedanticRegistry() - metrics := newSchedulerMetrics(reg) + metrics := newTestSchedulerMetrics(reg) return NewJobTracker(&NopJobPersister{}, "test", clk, newSimpleLanePolicy(), infiniteLeases, infiniteLeases, metrics.newTrackerMetricsForTenant("test"), log.NewNopLogger()), reg } @@ -103,7 +107,7 @@ func TestJobTracker_Maintenance_Planning(t *testing.T) { } t.Run("returns error on persist failure", func(t *testing.T) { - metrics := newSchedulerMetrics(prometheus.NewPedanticRegistry()) + metrics := newTestSchedulerMetrics(prometheus.NewPedanticRegistry()) jt := NewJobTracker(&errJobPersister{}, "test", clock.New(), newSimpleLanePolicy(), infiniteLeases, infiniteLeases, metrics.newTrackerMetricsForTenant("test"), log.NewNopLogger()) transition, err := jt.Maintenance(leaseDuration, false, true, planningInterval, compactionWaitPeriod) @@ -113,7 +117,7 @@ func TestJobTracker_Maintenance_Planning(t *testing.T) { }) t.Run("planning skipped when plan is false", func(t *testing.T) { - metrics := newSchedulerMetrics(prometheus.NewPedanticRegistry()) + metrics := newTestSchedulerMetrics(prometheus.NewPedanticRegistry()) jt := NewJobTracker(&errJobPersister{}, "test", clock.New(), newSimpleLanePolicy(), infiniteLeases, infiniteLeases, metrics.newTrackerMetricsForTenant("test"), log.NewNopLogger()) transition, err := jt.Maintenance(leaseDuration, false, false, planningInterval, compactionWaitPeriod) require.NoError(t, err) @@ -253,8 +257,8 @@ func assertTrackerBytes(t *testing.T, reg *prometheus.Registry, msg string, spli require.NoError(t, prom_testutil.GatherAndCompare(reg, strings.NewReader(fmt.Sprintf(` # HELP cortex_compactor_scheduler_incomplete_compaction_jobs_bytes The total bytes of blocks in compaction jobs that have not yet completed (pending or active). # TYPE cortex_compactor_scheduler_incomplete_compaction_jobs_bytes gauge - cortex_compactor_scheduler_incomplete_compaction_jobs_bytes{compaction_type="merge"} %g - cortex_compactor_scheduler_incomplete_compaction_jobs_bytes{compaction_type="split"} %g + cortex_compactor_scheduler_incomplete_compaction_jobs_bytes{compaction_type="merge",lane="compaction"} %g + cortex_compactor_scheduler_incomplete_compaction_jobs_bytes{compaction_type="split",lane="compaction"} %g `, mergeBytes, splitBytes)), "cortex_compactor_scheduler_incomplete_compaction_jobs_bytes"), msg) } @@ -335,7 +339,7 @@ func TestJobTracker_PlanJobTracking(t *testing.T) { func TestJobTracker_Cleanup(t *testing.T) { clk := clock.NewMock() reg := prometheus.NewPedanticRegistry() - sm := newSchedulerMetrics(reg) + sm := newTestSchedulerMetrics(reg) // Two tenants share the same aggregate gauges (incompleteJobsBytes, pendingJobs, activeJobs). jt1 := NewJobTracker(&NopJobPersister{}, "tenant1", clk, newSimpleLanePolicy(), infiniteLeases, infiniteLeases, sm.newTrackerMetricsForTenant("tenant1"), log.NewNopLogger()) @@ -397,7 +401,7 @@ func TestJobTracker_CancelLease_PlanJobAlwaysRevives(t *testing.T) { const maxLeases = 2 clk := clock.NewMock() - metrics := newSchedulerMetrics(prometheus.NewPedanticRegistry()) + metrics := newTestSchedulerMetrics(prometheus.NewPedanticRegistry()) jt := NewJobTracker(&NopJobPersister{}, "test", clk, newSimpleLanePolicy(), maxLeases, infiniteLeases, metrics.newTrackerMetricsForTenant("test"), log.NewNopLogger()) _, err := jt.Maintenance(time.Minute, false, true, time.Hour, 15*time.Minute) @@ -447,7 +451,7 @@ func TestJobTracker_CancelLease_Interrupted(t *testing.T) { } { t.Run(tc.name, func(t *testing.T) { clk := clock.NewMock() - metrics := newSchedulerMetrics(prometheus.NewPedanticRegistry()) + metrics := newTestSchedulerMetrics(prometheus.NewPedanticRegistry()) jt := NewJobTracker(&NopJobPersister{}, "test", clk, newSimpleLanePolicy(), tc.maxLeases, tc.threshold, metrics.newTrackerMetricsForTenant("test"), log.NewNopLogger()) lane := compactionLane diff --git a/pkg/compactor/scheduler/lane_policy.go b/pkg/compactor/scheduler/lane_policy.go index 9dded750cc0..813c0396507 100644 --- a/pkg/compactor/scheduler/lane_policy.go +++ b/pkg/compactor/scheduler/lane_policy.go @@ -10,16 +10,17 @@ import ( "github.com/grafana/mimir/pkg/compactor/scheduler/compactorschedulerpb" ) -// lane is an in-memory identifier of pending work logically enqueued together -type lane uint8 +// lane is an in-memory identifier of pending work logically enqueued together. Its value is +// exported as a metric label, so renaming one changes existing metrics. +type lane string const ( - lanePolicySimple = "simple" - - planLane lane = iota - compactionLane + planLane lane = "plan" + compactionLane lane = "compaction" ) +const lanePolicySimple = "simple" + type laneTransition struct { lane lane kind rotationTransition @@ -27,9 +28,19 @@ type laneTransition struct { // Defines how to map jobs and requests into lanes type lanePolicy interface { - AllLanes() []lane // All possible lanes defined by this policy. - LaneForJob(TrackedJob) lane // The lane this job is assigned to. A job must always map to some lane. - LanesForRequest(*compactorschedulerpb.LeaseJobRequest) ([]lane, error) // The lanes this worker requested, or an error. + // AllLanes returns every lane this policy defines. The scheduler tracks exactly these, and + // serves them in this order to a worker that names no job type. + AllLanes() []lane + + // CompactionLanes returns the lanes carrying compaction jobs. + CompactionLanes() []lane + + // LaneForJob returns the job's lane. It must always return the same lane for a job for the + // lifetime of the process, as different callers may re-derive it at different times. + LaneForJob(TrackedJob) lane + + // LanesForRequest returns the lanes this worker requested. + LanesForRequest(*compactorschedulerpb.LeaseJobRequest) ([]lane, error) } type LanePolicyConfig struct { @@ -37,12 +48,12 @@ type LanePolicyConfig struct { } func (cfg *LanePolicyConfig) RegisterFlagsWithPrefix(prefix string, f *flag.FlagSet) { - f.StringVar(&cfg.Policy, prefix+".policy", "simple", "The lane policy the compactor scheduler should use. Valid values: "+lanePolicySimple) + f.StringVar(&cfg.Policy, prefix+".policy", lanePolicySimple, "The lane policy the compactor scheduler should use. Valid values: "+lanePolicySimple) } func newLanePolicy(cfg LanePolicyConfig) (lanePolicy, error) { switch cfg.Policy { - case "simple": + case lanePolicySimple: return newSimpleLanePolicy(), nil default: return nil, fmt.Errorf("unrecognized lane policy: %s", cfg.Policy) @@ -51,12 +62,14 @@ func newLanePolicy(cfg LanePolicyConfig) (lanePolicy, error) { // simpleLanePolicy assigns a lane per job type type simpleLanePolicy struct { - allLanes []lane + allLanes []lane + compactionLanes []lane } func newSimpleLanePolicy() lanePolicy { return &simpleLanePolicy{ - allLanes: []lane{planLane, compactionLane}, + allLanes: []lane{planLane, compactionLane}, + compactionLanes: []lane{compactionLane}, } } @@ -71,6 +84,10 @@ func (slp *simpleLanePolicy) AllLanes() []lane { return slp.allLanes } +func (slp *simpleLanePolicy) CompactionLanes() []lane { + return slp.compactionLanes +} + // requestedLanes maps a lease request to scheduler lanes func (slp *simpleLanePolicy) LanesForRequest(req *compactorschedulerpb.LeaseJobRequest) ([]lane, error) { numLanes := len(req.LaneRequests) diff --git a/pkg/compactor/scheduler/metrics.go b/pkg/compactor/scheduler/metrics.go index 858099fb2a0..55ffaf7d517 100644 --- a/pkg/compactor/scheduler/metrics.go +++ b/pkg/compactor/scheduler/metrics.go @@ -15,19 +15,36 @@ const ( compactionTypeMerge = "merge" ) +type incompleteBytesKey struct { + compactionType string + lane lane +} + type schedulerMetrics struct { - pendingJobs *prometheus.GaugeVec - pendingJobsByUser *prometheus.GaugeVec - pendingJobsLastEmpty prometheus.Gauge - incompleteJobsBytes *prometheus.GaugeVec - activeJobs *prometheus.GaugeVec - activeJobsByUser *prometheus.GaugeVec - jobsCompleted *prometheus.CounterVec - repeatedJobFailures prometheus.Counter + pendingJobs *prometheus.GaugeVec + pendingJobsByUser *prometheus.GaugeVec + pendingJobsLastEmpty prometheus.Gauge + lanePendingJobsLastEmpty *prometheus.GaugeVec + incompleteJobsBytes *prometheus.GaugeVec + activeJobs *prometheus.GaugeVec + activeJobsByUser *prometheus.GaugeVec + jobsCompleted *prometheus.CounterVec + repeatedJobFailures prometheus.Counter + lanePolicy lanePolicy + + // One gauge per series of incompleteJobsBytes, resolved once and shared across tenants. + incompleteBytesGauges map[incompleteBytesKey]prometheus.Gauge + + lanePendingJobsLastEmptyGauges map[lane]prometheus.Gauge } -func newSchedulerMetrics(reg prometheus.Registerer) *schedulerMetrics { +// newSchedulerMetrics labels incomplete compaction bytes by the lanes that carry compaction work, +// so that a fleet serving one lane can be sized from its own backlog. +func newSchedulerMetrics(reg prometheus.Registerer, lanePolicy lanePolicy) *schedulerMetrics { + allLanes := lanePolicy.AllLanes() + compactionLanes := lanePolicy.CompactionLanes() m := &schedulerMetrics{ + lanePolicy: lanePolicy, pendingJobs: promauto.With(reg).NewGaugeVec(prometheus.GaugeOpts{ Name: "cortex_compactor_scheduler_pending_jobs", Help: "The number of queued pending jobs.", @@ -40,10 +57,14 @@ func newSchedulerMetrics(reg prometheus.Registerer) *schedulerMetrics { Name: "cortex_compactor_scheduler_pending_jobs_last_empty_timestamp_seconds", Help: "Unix timestamp of the last time there were no pending jobs remaining.", }), + lanePendingJobsLastEmpty: promauto.With(reg).NewGaugeVec(prometheus.GaugeOpts{ + Name: "cortex_compactor_scheduler_lane_pending_jobs_last_empty_timestamp_seconds", + Help: "Unix timestamp of the last time there were no pending jobs remaining in this lane.", + }, []string{"lane"}), incompleteJobsBytes: promauto.With(reg).NewGaugeVec(prometheus.GaugeOpts{ Name: "cortex_compactor_scheduler_incomplete_compaction_jobs_bytes", Help: "The total bytes of blocks in compaction jobs that have not yet completed (pending or active).", - }, []string{"compaction_type"}), + }, []string{"compaction_type", "lane"}), activeJobs: promauto.With(reg).NewGaugeVec(prometheus.GaugeOpts{ Name: "cortex_compactor_scheduler_active_jobs", Help: "The number of jobs active in workers.", @@ -68,12 +89,25 @@ func newSchedulerMetrics(reg prometheus.Registerer) *schedulerMetrics { m.pendingJobs.WithLabelValues(jobTypeCompaction) m.activeJobs.WithLabelValues(jobTypePlan) m.activeJobs.WithLabelValues(jobTypeCompaction) - m.incompleteJobsBytes.WithLabelValues(compactionTypeSplit) - m.incompleteJobsBytes.WithLabelValues(compactionTypeMerge) + m.lanePendingJobsLastEmptyGauges = make(map[lane]prometheus.Gauge, len(allLanes)) + for _, l := range allLanes { + m.lanePendingJobsLastEmptyGauges[l] = m.lanePendingJobsLastEmpty.WithLabelValues(string(l)) + } + m.incompleteBytesGauges = make(map[incompleteBytesKey]prometheus.Gauge, 2*len(compactionLanes)) + for _, t := range []string{compactionTypeSplit, compactionTypeMerge} { + for _, l := range compactionLanes { + k := incompleteBytesKey{compactionType: t, lane: l} + m.incompleteBytesGauges[k] = m.incompleteJobsBytes.WithLabelValues(t, string(l)) + } + } return m } func (s *schedulerMetrics) newTrackerMetricsForTenant(tenant string) *trackerMetrics { + byKey := make(map[incompleteBytesKey]*incompleteBytes, len(s.incompleteBytesGauges)) + for k, g := range s.incompleteBytesGauges { + byKey[k] = &incompleteBytes{gauge: g} + } return &trackerMetrics{ queue: &queueMetrics{ pendingJobsByUser: s.pendingJobsByUser.WithLabelValues(tenant), @@ -82,8 +116,8 @@ func (s *schedulerMetrics) newTrackerMetricsForTenant(tenant string) *trackerMet pendingCompactionJobs: s.pendingJobs.WithLabelValues(jobTypeCompaction), activePlanJobs: s.activeJobs.WithLabelValues(jobTypePlan), activeCompactionJobs: s.activeJobs.WithLabelValues(jobTypeCompaction), - incompleteSplitBytes: s.incompleteJobsBytes.WithLabelValues(compactionTypeSplit), - incompleteMergeBytes: s.incompleteJobsBytes.WithLabelValues(compactionTypeMerge), + incompleteBytes: byKey, + laneForJob: s.lanePolicy.LaneForJob, clear: func() { s.pendingJobsByUser.DeleteLabelValues(tenant) s.activeJobsByUser.DeleteLabelValues(tenant) @@ -93,6 +127,12 @@ func (s *schedulerMetrics) newTrackerMetricsForTenant(tenant string) *trackerMet } } +// incompleteBytes pairs a gauge shared across tenants with this tenant's contribution to it. +type incompleteBytes struct { + gauge prometheus.Gauge + contributed uint64 +} + type trackerMetrics struct { queue *queueMetrics repeatedJobFailures prometheus.Counter @@ -102,14 +142,14 @@ type trackerMetrics struct { // shared gauges. Must be called when a tenant is removed. func (m *trackerMetrics) Clear() { q := m.queue - q.incompleteSplitBytes.Sub(float64(q.splitBytes)) - q.incompleteMergeBytes.Sub(float64(q.mergeBytes)) + for _, b := range q.incompleteBytes { + b.gauge.Sub(float64(b.contributed)) + b.contributed = 0 + } q.pendingPlanJobs.Sub(float64(q.pendingPlanCount)) q.pendingCompactionJobs.Sub(float64(q.pendingCompactionCount)) q.activePlanJobs.Sub(float64(q.activePlanCount)) q.activeCompactionJobs.Sub(float64(q.activeCompactionCount)) - q.splitBytes = 0 - q.mergeBytes = 0 q.pendingPlanCount = 0 q.pendingCompactionCount = 0 q.activePlanCount = 0 @@ -130,13 +170,11 @@ type queueMetrics struct { pendingCompactionJobs prometheus.Gauge activePlanJobs prometheus.Gauge activeCompactionJobs prometheus.Gauge - incompleteSplitBytes prometheus.Gauge - incompleteMergeBytes prometheus.Gauge + incompleteBytes map[incompleteBytesKey]*incompleteBytes + laneForJob func(TrackedJob) lane // This tenant's contribution to the shared gauges, tracked so Clear() can subtract exactly // the right amount on tenant removal. - splitBytes uint64 - mergeBytes uint64 pendingPlanCount int pendingCompactionCount int activePlanCount int @@ -237,22 +275,22 @@ func (q *queueMetrics) decActive(isPlan bool) { } } -func (q *queueMetrics) addBytes(cj *TrackedCompactionJob) { +func (q *queueMetrics) bytesFor(cj *TrackedCompactionJob) *incompleteBytes { + compactionType := compactionTypeMerge if cj.value.isSplit { - q.splitBytes += cj.totalBlockBytes - q.incompleteSplitBytes.Add(float64(cj.totalBlockBytes)) - } else { - q.mergeBytes += cj.totalBlockBytes - q.incompleteMergeBytes.Add(float64(cj.totalBlockBytes)) + compactionType = compactionTypeSplit } + return q.incompleteBytes[incompleteBytesKey{compactionType: compactionType, lane: q.laneForJob(cj)}] +} + +func (q *queueMetrics) addBytes(cj *TrackedCompactionJob) { + b := q.bytesFor(cj) + b.contributed += cj.totalBlockBytes + b.gauge.Add(float64(cj.totalBlockBytes)) } func (q *queueMetrics) subBytes(cj *TrackedCompactionJob) { - if cj.value.isSplit { - q.splitBytes -= cj.totalBlockBytes - q.incompleteSplitBytes.Sub(float64(cj.totalBlockBytes)) - } else { - q.mergeBytes -= cj.totalBlockBytes - q.incompleteMergeBytes.Sub(float64(cj.totalBlockBytes)) - } + b := q.bytesFor(cj) + b.contributed -= cj.totalBlockBytes + b.gauge.Sub(float64(cj.totalBlockBytes)) } diff --git a/pkg/compactor/scheduler/persistence_bbolt_test.go b/pkg/compactor/scheduler/persistence_bbolt_test.go index 8a66561fd63..d18c0877ba9 100644 --- a/pkg/compactor/scheduler/persistence_bbolt_test.go +++ b/pkg/compactor/scheduler/persistence_bbolt_test.go @@ -74,7 +74,7 @@ func TestBboltJobPersistenceManager_RecoverAll(t *testing.T) { }) allowedTenants := util.NewAllowList(nil, nil) - metrics := newSchedulerMetrics(prometheus.NewPedanticRegistry()) + metrics := newTestSchedulerMetrics(prometheus.NewPedanticRegistry()) jobTrackerFactory := func(tenant string, persister JobPersister) *JobTracker { return NewJobTracker(persister, tenant, clock.New(), newSimpleLanePolicy(), infiniteLeases, infiniteLeases, metrics.newTrackerMetricsForTenant(tenant), log.NewNopLogger()) } @@ -108,7 +108,7 @@ func TestBboltJobPersistenceManager_RecoverAll_Cleanup(t *testing.T) { require.NoError(t, mgr.Close()) }) - metrics := newSchedulerMetrics(prometheus.NewPedanticRegistry()) + metrics := newTestSchedulerMetrics(prometheus.NewPedanticRegistry()) jobTrackerFactory := func(tenant string, persister JobPersister) *JobTracker { return NewJobTracker(persister, tenant, clock.New(), newSimpleLanePolicy(), infiniteLeases, infiniteLeases, metrics.newTrackerMetricsForTenant(tenant), log.NewNopLogger()) } @@ -256,7 +256,7 @@ func TestRunMigration_ScaleUp(t *testing.T) { require.Len(t, mgr.dbs, 2) allowedTenants := util.NewAllowList(nil, nil) - metrics := newSchedulerMetrics(prometheus.NewPedanticRegistry()) + metrics := newTestSchedulerMetrics(prometheus.NewPedanticRegistry()) trackers, err := mgr.RecoverAll(allowedTenants, func(tenant string, persister JobPersister) *JobTracker { return NewJobTracker(persister, tenant, clock.New(), newSimpleLanePolicy(), infiniteLeases, infiniteLeases, metrics.newTrackerMetricsForTenant(tenant), log.NewNopLogger()) }) @@ -318,7 +318,7 @@ func TestRunMigration_ScaleDown(t *testing.T) { require.True(t, os.IsNotExist(err), "extra shard file should be deleted") allowedTenants := util.NewAllowList(nil, nil) - metrics := newSchedulerMetrics(prometheus.NewPedanticRegistry()) + metrics := newTestSchedulerMetrics(prometheus.NewPedanticRegistry()) trackers, err := mgr.RecoverAll(allowedTenants, func(tenant string, persister JobPersister) *JobTracker { return NewJobTracker(persister, tenant, clock.New(), newSimpleLanePolicy(), infiniteLeases, infiniteLeases, metrics.newTrackerMetricsForTenant(tenant), log.NewNopLogger()) }) diff --git a/pkg/compactor/scheduler/rotator.go b/pkg/compactor/scheduler/rotator.go index fd461f6a001..521491112a5 100644 --- a/pkg/compactor/scheduler/rotator.go +++ b/pkg/compactor/scheduler/rotator.go @@ -45,6 +45,7 @@ type Rotator struct { intervalsBeforeColdStartPlanning int clock clock.Clock pendingJobsLastEmpty prometheus.Gauge + lanePendingJobsLastEmpty map[lane]prometheus.Gauge logger log.Logger mtx sync.RWMutex @@ -63,7 +64,7 @@ type tenantRotationState struct { elements map[lane]*list.Element // tenant's slot in each lane's rotation (if they are present in that lane) } -func NewRotator(leaseDuration, planningInterval, compactionWaitPeriod, maintenanceInterval time.Duration, intervalsBeforeLeaseExpiration, intervalsBeforeColdStartPlanning int, lanePolicy lanePolicy, pendingJobsLastEmpty prometheus.Gauge, logger log.Logger) *Rotator { +func NewRotator(leaseDuration, planningInterval, compactionWaitPeriod, maintenanceInterval time.Duration, intervalsBeforeLeaseExpiration, intervalsBeforeColdStartPlanning int, lanePolicy lanePolicy, pendingJobsLastEmpty prometheus.Gauge, lanePendingJobsLastEmpty map[lane]prometheus.Gauge, logger log.Logger) *Rotator { laneRotations := make(map[lane]*laneRotation) for _, lane := range lanePolicy.AllLanes() { laneRotations[lane] = &laneRotation{rotation: list.New()} @@ -78,6 +79,7 @@ func NewRotator(leaseDuration, planningInterval, compactionWaitPeriod, maintenan intervalsBeforeColdStartPlanning: intervalsBeforeColdStartPlanning, clock: clock.New(), pendingJobsLastEmpty: pendingJobsLastEmpty, + lanePendingJobsLastEmpty: lanePendingJobsLastEmpty, mtx: sync.RWMutex{}, tenantStateMap: make(map[string]*tenantRotationState), laneRotations: laneRotations, @@ -159,9 +161,7 @@ func (r *Rotator) RecoverFrom(jobTrackers map[string]*JobTracker, creationTime t } } - if r.allRotationsEmpty() { - r.pendingJobsLastEmpty.Set(float64(r.clock.Now().Unix())) - } + r.recordEmptyQueues() } func (r *Rotator) LeaseJob(ctx context.Context, lanes []lane) (*compactorschedulerpb.LeaseJobResponse, bool, error) { @@ -399,18 +399,19 @@ func (r *Rotator) Maintenance(ctx context.Context, enforceLeaseExpiration, plan addRotationFor = append(addRotationFor, tenantLane{tenant: tenant, lane: l}) } } - stayingEmpty := len(addRotationFor) == 0 && r.allRotationsEmpty() - r.mtx.RUnlock() - - if len(addRotationFor) == 0 || ctx.Err() != nil { - // No tenant needs to be moved into the rotation or we're shutting down and don't care - if stayingEmpty && ctx.Err() == nil { - // Ensures periodic updates for the time we last saw an empty queue rather than only on transition. - // The context check is to ensure we don't update this when mid-shutdown. - r.pendingJobsLastEmpty.Set(float64(r.clock.Now().Unix())) - } + if ctx.Err() != nil { + // Shutting down, don't update any metrics + r.mtx.RUnlock() + return + } + if len(addRotationFor) == 0 { + // No tenant needs to be moved into the rotation. Ensures periodic updates for the time we last + // saw an empty queue rather than only on transition. + r.recordEmptyQueues() + r.mtx.RUnlock() return } + r.mtx.RUnlock() r.mtx.Lock() defer r.mtx.Unlock() @@ -420,6 +421,10 @@ func (r *Rotator) Maintenance(ctx context.Context, enforceLeaseExpiration, plan r.addToRotation(tl.lane, tl.tenant, tenantState) } } + // Ensures the same periodic updates, whichever way the maintenance tick ended. + if ctx.Err() == nil { + r.recordEmptyQueues() + } } func (r *Rotator) possiblyRemoveFromRotation(l lane, tenant string) { @@ -457,17 +462,22 @@ func (r *Rotator) removeFromRotation(l lane, tenantState *tenantRotationState) { lr.rotation.Remove(elem) delete(tenantState.elements, l) - if r.allRotationsEmpty() { - r.pendingJobsLastEmpty.Set(float64(r.clock.Now().Unix())) - } + r.recordEmptyQueues() } -// allRotationsEmpty reports whether no tenant has pending work in any lane. A write lock must be held in order to call this function. -func (r *Rotator) allRotationsEmpty() bool { - for _, lr := range r.laneRotations { +// recordEmptyQueues records now as the last time each empty lane, and the queue as a whole, was seen +// empty. The rotations must be up to date, and a lock must be held, in order to call this function. +func (r *Rotator) recordEmptyQueues() { + now := float64(r.clock.Now().Unix()) + allEmpty := true + for l, lr := range r.laneRotations { if lr.rotation.Len() > 0 { - return false + allEmpty = false + continue } + r.lanePendingJobsLastEmpty[l].Set(now) + } + if allEmpty { + r.pendingJobsLastEmpty.Set(now) } - return true } diff --git a/pkg/compactor/scheduler/rotator_test.go b/pkg/compactor/scheduler/rotator_test.go index d9c975f7780..4a5cd78ac40 100644 --- a/pkg/compactor/scheduler/rotator_test.go +++ b/pkg/compactor/scheduler/rotator_test.go @@ -6,6 +6,8 @@ import ( "container/list" "context" "fmt" + "maps" + "slices" "strings" "sync" "testing" @@ -72,8 +74,8 @@ func TestRotator_RecoverFrom_ColdStartDelay(t *testing.T) { for _, tc := range tests { t.Run(tc.name, func(t *testing.T) { reg := prometheus.NewPedanticRegistry() - metrics := newSchedulerMetrics(reg) - r := NewRotator(0, 0, 0, maintenanceInterval, 0, intervalsBeforeColdStartPlanning, newSimpleLanePolicy(), metrics.pendingJobsLastEmpty, log.NewNopLogger()) + metrics := newTestSchedulerMetrics(reg) + r := NewRotator(0, 0, 0, maintenanceInterval, 0, intervalsBeforeColdStartPlanning, newSimpleLanePolicy(), metrics.pendingJobsLastEmpty, metrics.lanePendingJobsLastEmptyGauges, log.NewNopLogger()) r.clock = clock r.RecoverFrom(tc.jobTrackers, tc.creationTime) @@ -84,14 +86,14 @@ func TestRotator_RecoverFrom_ColdStartDelay(t *testing.T) { func newRotatorForTest() *Rotator { reg := prometheus.NewPedanticRegistry() - metrics := newSchedulerMetrics(reg) - return NewRotator(0, 0, 0, time.Minute, 0, 0, newSimpleLanePolicy(), metrics.pendingJobsLastEmpty, log.NewNopLogger()) + metrics := newTestSchedulerMetrics(reg) + return NewRotator(0, 0, 0, time.Minute, 0, 0, newSimpleLanePolicy(), metrics.pendingJobsLastEmpty, metrics.lanePendingJobsLastEmptyGauges, log.NewNopLogger()) } // newTrackerWithPendingJobs builds a JobTracker for the named tenant holding numJobs pending // compaction jobs (numJobs == 0 yields an empty tracker). func newTrackerWithPendingJobs(clk clock.Clock, name string, numJobs int) *JobTracker { - metrics := newSchedulerMetrics(prometheus.NewPedanticRegistry()) + metrics := newTestSchedulerMetrics(prometheus.NewPedanticRegistry()) jt := NewJobTracker(&NopJobPersister{}, name, clk, newSimpleLanePolicy(), infiniteLeases, infiniteLeases, metrics.newTrackerMetricsForTenant(name), log.NewNopLogger()) for j := range numJobs { id := fmt.Sprintf("%s-%d", name, j) @@ -228,8 +230,8 @@ func TestRotator_LeaseJob_ConcurrentLeasersDoNotSkipPendingWork(t *testing.T) { func TestRotator_LeaseJob_LanePriority(t *testing.T) { clk := clock.New() lanePolicy := newSimpleLanePolicy() - metrics := newSchedulerMetrics(prometheus.NewPedanticRegistry()) - r := NewRotator(0, 0, 0, time.Minute, 0, 0, lanePolicy, metrics.pendingJobsLastEmpty, log.NewNopLogger()) + metrics := newTestSchedulerMetrics(prometheus.NewPedanticRegistry()) + r := NewRotator(0, 0, 0, time.Minute, 0, 0, lanePolicy, metrics.pendingJobsLastEmpty, metrics.lanePendingJobsLastEmptyGauges, log.NewNopLogger()) // Add a tenant with a plan job and a compaction job jt := NewJobTracker(&NopJobPersister{}, "t1", clk, lanePolicy, infiniteLeases, infiniteLeases, metrics.newTrackerMetricsForTenant("t1"), log.NewNopLogger()) @@ -286,29 +288,65 @@ func TestRotator_PendingJobsLastEmpty(t *testing.T) { } tests := map[string]struct { - action func(r *Rotator, clk *clock.Mock) - expected int64 + action func(r *Rotator, clk *clock.Mock) + // expected is for the whole queue, expectedByLane for each lane on its own. + expected int64 + expectedByLane map[lane]int64 }{ "recovery with no pending jobs": { - action: func(r *Rotator, _ *clock.Mock) { r.RecoverFrom(nil, now) }, - expected: now.Unix(), + action: func(r *Rotator, _ *clock.Mock) { r.RecoverFrom(nil, now) }, + expected: now.Unix(), + expectedByLane: map[lane]int64{planLane: now.Unix(), compactionLane: now.Unix()}, }, - "recovery with pending jobs leaves metric at 0": { + "recovery with pending jobs leaves only their lane unset": { action: func(r *Rotator, clk *clock.Mock) { r.RecoverFrom(map[string]*JobTracker{"t1": pendingTracker(clk)}, now) }, - expected: 0, + expected: 0, + expectedByLane: map[lane]int64{planLane: 0, compactionLane: now.Unix()}, + }, + // The point of tracking this per lane: a compaction backlog must not hide that planning is idle. + "recovery with a pending compaction job leaves only the compaction lane unset": { + action: func(r *Rotator, clk *clock.Mock) { + jt, _ := newTestJobTracker(clk) + jt.toPendingBack(NewTrackedCompactionJob("c1", &CompactionJob{}, 1, 0, clk.Now())) + r.RecoverFrom(map[string]*JobTracker{"t1": jt}, now) + }, + expected: 0, + expectedByLane: map[lane]int64{planLane: now.Unix(), compactionLane: 0}, }, "maintenance over empty rotation": { - action: func(r *Rotator, _ *clock.Mock) { r.Maintenance(context.Background(), false, false) }, - expected: now.Unix(), + action: func(r *Rotator, _ *clock.Mock) { r.Maintenance(context.Background(), false, false) }, + expected: now.Unix(), + expectedByLane: map[lane]int64{planLane: now.Unix(), compactionLane: now.Unix()}, }, + // The clock advances so the assertion pins when the lanes were stamped, not merely that they were. "removing the last tenant in rotation": { action: func(r *Rotator, clk *clock.Mock) { r.AddTenant("t1", pendingTracker(clk)) + clk.Add(time.Minute) _, _ = r.RemoveTenant("t1") }, - expected: now.Unix(), + expected: now.Add(time.Minute).Unix(), + expectedByLane: map[lane]int64{planLane: now.Add(time.Minute).Unix(), compactionLane: now.Add(time.Minute).Unix()}, + }, + "maintenance mid-shutdown records nothing": { + action: func(r *Rotator, _ *clock.Mock) { + ctx, cancel := context.WithCancel(context.Background()) + cancel() + r.Maintenance(ctx, false, false) + }, + expected: 0, + expectedByLane: map[lane]int64{planLane: 0, compactionLane: 0}, + }, + "maintenance moving a tenant into a rotation still records the idle lanes": { + action: func(r *Rotator, clk *clock.Mock) { + jt, _ := newTestJobTracker(clk) + r.AddTenant("t1", jt) + r.Maintenance(context.Background(), false, true) + }, + expected: 0, + expectedByLane: map[lane]int64{planLane: 0, compactionLane: now.Unix()}, }, } @@ -317,18 +355,24 @@ func TestRotator_PendingJobsLastEmpty(t *testing.T) { clk := clock.NewMock() clk.Set(now) reg := prometheus.NewPedanticRegistry() - metrics := newSchedulerMetrics(reg) - r := NewRotator(0, 0, 0, 0, 0, 0, newSimpleLanePolicy(), metrics.pendingJobsLastEmpty, log.NewNopLogger()) + metrics := newTestSchedulerMetrics(reg) + r := NewRotator(0, 0, 0, 0, 0, 0, newSimpleLanePolicy(), metrics.pendingJobsLastEmpty, metrics.lanePendingJobsLastEmptyGauges, log.NewNopLogger()) r.clock = clk tc.action(r, clk) - metricName := "cortex_compactor_scheduler_pending_jobs_last_empty_timestamp_seconds" - require.NoError(t, prom_testutil.GatherAndCompare(reg, strings.NewReader(fmt.Sprintf(` - # HELP %s Unix timestamp of the last time there were no pending jobs remaining. - # TYPE %s gauge - %s %d - `, metricName, metricName, metricName, tc.expected)), metricName)) + const ( + laneMetric = "cortex_compactor_scheduler_lane_pending_jobs_last_empty_timestamp_seconds" + anyMetric = "cortex_compactor_scheduler_pending_jobs_last_empty_timestamp_seconds" + ) + var expected strings.Builder + fmt.Fprintf(&expected, "# HELP %s Unix timestamp of the last time there were no pending jobs remaining in this lane.\n# TYPE %s gauge\n", laneMetric, laneMetric) + for _, l := range slices.Sorted(maps.Keys(tc.expectedByLane)) { + fmt.Fprintf(&expected, "%s{lane=%q} %d\n", laneMetric, l, tc.expectedByLane[l]) + } + fmt.Fprintf(&expected, "# HELP %s Unix timestamp of the last time there were no pending jobs remaining.\n# TYPE %s gauge\n%s %d\n", anyMetric, anyMetric, anyMetric, tc.expected) + + require.NoError(t, prom_testutil.GatherAndCompare(reg, strings.NewReader(expected.String()), laneMetric, anyMetric)) }) } } diff --git a/pkg/compactor/scheduler/scheduler.go b/pkg/compactor/scheduler/scheduler.go index d4c8e3f2eaf..aaf8b7423a1 100644 --- a/pkg/compactor/scheduler/scheduler.go +++ b/pkg/compactor/scheduler/scheduler.go @@ -119,7 +119,6 @@ func NewCompactorScheduler( logger log.Logger, registerer prometheus.Registerer) (*Scheduler, error) { - metrics := newSchedulerMetrics(registerer) allowList := util.NewAllowList(compactorCfg.EnabledTenants, compactorCfg.DisabledTenants) jpm, err := jobPersistenceManagerFactory(cfg, logger) @@ -132,7 +131,7 @@ func NewCompactorScheduler( return nil, err } - return newCompactorScheduler(compactorCfg, cfg, allowList, bkt, jpm, metrics, logger) + return newCompactorScheduler(compactorCfg, cfg, allowList, bkt, jpm, registerer, logger) } func newCompactorScheduler( @@ -141,7 +140,7 @@ func newCompactorScheduler( allowList *util.AllowList, bkt objstore.Bucket, jpm JobPersistenceManager, - metrics *schedulerMetrics, + registerer prometheus.Registerer, logger log.Logger) (*Scheduler, error) { lanePolicy, err := newLanePolicy(cfg.LanePolicy) @@ -149,6 +148,8 @@ func newCompactorScheduler( return nil, err } + metrics := newSchedulerMetrics(registerer, lanePolicy) + rotator := NewRotator( cfg.LeaseDuration, cfg.PlanningInterval, @@ -158,6 +159,7 @@ func newCompactorScheduler( cfg.MaintenanceIntervalsBeforeColdStartPlanning, lanePolicy, metrics.pendingJobsLastEmpty, + metrics.lanePendingJobsLastEmptyGauges, logger, ) diff --git a/pkg/compactor/scheduler/scheduler_test.go b/pkg/compactor/scheduler/scheduler_test.go index a3d25af02c1..1643b3cbf0b 100644 --- a/pkg/compactor/scheduler/scheduler_test.go +++ b/pkg/compactor/scheduler/scheduler_test.go @@ -44,8 +44,8 @@ func TestScheduler_JobLifecycleMetrics(t *testing.T) { require.NoError(t, prom_testutil.GatherAndCompare(reg, strings.NewReader(fmt.Sprintf(` # HELP cortex_compactor_scheduler_incomplete_compaction_jobs_bytes The total bytes of blocks in compaction jobs that have not yet completed (pending or active). # TYPE cortex_compactor_scheduler_incomplete_compaction_jobs_bytes gauge - cortex_compactor_scheduler_incomplete_compaction_jobs_bytes{compaction_type="merge"} %g - cortex_compactor_scheduler_incomplete_compaction_jobs_bytes{compaction_type="split"} %g + cortex_compactor_scheduler_incomplete_compaction_jobs_bytes{compaction_type="merge",lane="compaction"} %g + cortex_compactor_scheduler_incomplete_compaction_jobs_bytes{compaction_type="split",lane="compaction"} %g `, mergeBytes, splitBytes)), "cortex_compactor_scheduler_incomplete_compaction_jobs_bytes"), msg) } @@ -187,10 +187,9 @@ func newTestScheduler(t *testing.T, bkt objstore.Bucket, cfg Config) (*Scheduler t.Helper() reg := prometheus.NewPedanticRegistry() - metrics := newSchedulerMetrics(reg) compactorCfg := compactor.Config{CompactionWaitPeriod: 15 * time.Minute} - scheduler, err := newCompactorScheduler(compactorCfg, cfg, util.NewAllowList(nil, nil), bkt, &NopJobPersistenceManager{}, metrics, log.NewNopLogger()) + scheduler, err := newCompactorScheduler(compactorCfg, cfg, util.NewAllowList(nil, nil), bkt, &NopJobPersistenceManager{}, reg, log.NewNopLogger()) require.NoError(t, err) ctx, cancel := context.WithCancel(context.Background()) From f1c35e206c0197b9a760fcba63dbd28fb97c874f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Nicol=C3=A1s=20Pazos?= Date: Thu, 27 Aug 2026 11:06:49 -0300 Subject: [PATCH 2/5] compactor: keep the lane metrics change closer to the original code --- pkg/compactor/scheduler/lane_policy.go | 25 +++++---------- pkg/compactor/scheduler/metrics.go | 44 +++++++++++--------------- pkg/compactor/scheduler/rotator.go | 18 ++++------- 3 files changed, 34 insertions(+), 53 deletions(-) diff --git a/pkg/compactor/scheduler/lane_policy.go b/pkg/compactor/scheduler/lane_policy.go index 813c0396507..506a625c67f 100644 --- a/pkg/compactor/scheduler/lane_policy.go +++ b/pkg/compactor/scheduler/lane_policy.go @@ -15,12 +15,12 @@ import ( type lane string const ( + lanePolicySimple = "simple" + planLane lane = "plan" compactionLane lane = "compaction" ) -const lanePolicySimple = "simple" - type laneTransition struct { lane lane kind rotationTransition @@ -28,19 +28,10 @@ type laneTransition struct { // Defines how to map jobs and requests into lanes type lanePolicy interface { - // AllLanes returns every lane this policy defines. The scheduler tracks exactly these, and - // serves them in this order to a worker that names no job type. - AllLanes() []lane - - // CompactionLanes returns the lanes carrying compaction jobs. - CompactionLanes() []lane - - // LaneForJob returns the job's lane. It must always return the same lane for a job for the - // lifetime of the process, as different callers may re-derive it at different times. - LaneForJob(TrackedJob) lane - - // LanesForRequest returns the lanes this worker requested. - LanesForRequest(*compactorschedulerpb.LeaseJobRequest) ([]lane, error) + AllLanes() []lane // All possible lanes defined by this policy. + CompactionLanes() []lane // The lanes that carry compaction jobs. + LaneForJob(TrackedJob) lane // The lane this job is assigned to. A job must always map to some lane. + LanesForRequest(*compactorschedulerpb.LeaseJobRequest) ([]lane, error) // The lanes this worker requested, or an error. } type LanePolicyConfig struct { @@ -48,12 +39,12 @@ type LanePolicyConfig struct { } func (cfg *LanePolicyConfig) RegisterFlagsWithPrefix(prefix string, f *flag.FlagSet) { - f.StringVar(&cfg.Policy, prefix+".policy", lanePolicySimple, "The lane policy the compactor scheduler should use. Valid values: "+lanePolicySimple) + f.StringVar(&cfg.Policy, prefix+".policy", "simple", "The lane policy the compactor scheduler should use. Valid values: "+lanePolicySimple) } func newLanePolicy(cfg LanePolicyConfig) (lanePolicy, error) { switch cfg.Policy { - case lanePolicySimple: + case "simple": return newSimpleLanePolicy(), nil default: return nil, fmt.Errorf("unrecognized lane policy: %s", cfg.Policy) diff --git a/pkg/compactor/scheduler/metrics.go b/pkg/compactor/scheduler/metrics.go index 55ffaf7d517..616ccf3758c 100644 --- a/pkg/compactor/scheduler/metrics.go +++ b/pkg/compactor/scheduler/metrics.go @@ -21,25 +21,20 @@ type incompleteBytesKey struct { } type schedulerMetrics struct { - pendingJobs *prometheus.GaugeVec - pendingJobsByUser *prometheus.GaugeVec - pendingJobsLastEmpty prometheus.Gauge - lanePendingJobsLastEmpty *prometheus.GaugeVec - incompleteJobsBytes *prometheus.GaugeVec - activeJobs *prometheus.GaugeVec - activeJobsByUser *prometheus.GaugeVec - jobsCompleted *prometheus.CounterVec - repeatedJobFailures prometheus.Counter - lanePolicy lanePolicy - - // One gauge per series of incompleteJobsBytes, resolved once and shared across tenants. - incompleteBytesGauges map[incompleteBytesKey]prometheus.Gauge - + pendingJobs *prometheus.GaugeVec + pendingJobsByUser *prometheus.GaugeVec + pendingJobsLastEmpty prometheus.Gauge + lanePendingJobsLastEmpty *prometheus.GaugeVec lanePendingJobsLastEmptyGauges map[lane]prometheus.Gauge + incompleteJobsBytes *prometheus.GaugeVec + incompleteBytesGauges map[incompleteBytesKey]prometheus.Gauge + activeJobs *prometheus.GaugeVec + activeJobsByUser *prometheus.GaugeVec + jobsCompleted *prometheus.CounterVec + repeatedJobFailures prometheus.Counter + lanePolicy lanePolicy } -// newSchedulerMetrics labels incomplete compaction bytes by the lanes that carry compaction work, -// so that a fleet serving one lane can be sized from its own backlog. func newSchedulerMetrics(reg prometheus.Registerer, lanePolicy lanePolicy) *schedulerMetrics { allLanes := lanePolicy.AllLanes() compactionLanes := lanePolicy.CompactionLanes() @@ -127,7 +122,6 @@ func (s *schedulerMetrics) newTrackerMetricsForTenant(tenant string) *trackerMet } } -// incompleteBytes pairs a gauge shared across tenants with this tenant's contribution to it. type incompleteBytes struct { gauge prometheus.Gauge contributed uint64 @@ -275,14 +269,6 @@ func (q *queueMetrics) decActive(isPlan bool) { } } -func (q *queueMetrics) bytesFor(cj *TrackedCompactionJob) *incompleteBytes { - compactionType := compactionTypeMerge - if cj.value.isSplit { - compactionType = compactionTypeSplit - } - return q.incompleteBytes[incompleteBytesKey{compactionType: compactionType, lane: q.laneForJob(cj)}] -} - func (q *queueMetrics) addBytes(cj *TrackedCompactionJob) { b := q.bytesFor(cj) b.contributed += cj.totalBlockBytes @@ -294,3 +280,11 @@ func (q *queueMetrics) subBytes(cj *TrackedCompactionJob) { b.contributed -= cj.totalBlockBytes b.gauge.Sub(float64(cj.totalBlockBytes)) } + +func (q *queueMetrics) bytesFor(cj *TrackedCompactionJob) *incompleteBytes { + compactionType := compactionTypeMerge + if cj.value.isSplit { + compactionType = compactionTypeSplit + } + return q.incompleteBytes[incompleteBytesKey{compactionType: compactionType, lane: q.laneForJob(cj)}] +} diff --git a/pkg/compactor/scheduler/rotator.go b/pkg/compactor/scheduler/rotator.go index 521491112a5..1783878eb56 100644 --- a/pkg/compactor/scheduler/rotator.go +++ b/pkg/compactor/scheduler/rotator.go @@ -399,20 +399,17 @@ func (r *Rotator) Maintenance(ctx context.Context, enforceLeaseExpiration, plan addRotationFor = append(addRotationFor, tenantLane{tenant: tenant, lane: l}) } } - if ctx.Err() != nil { - // Shutting down, don't update any metrics - r.mtx.RUnlock() - return - } - if len(addRotationFor) == 0 { - // No tenant needs to be moved into the rotation. Ensures periodic updates for the time we last - // saw an empty queue rather than only on transition. + if len(addRotationFor) == 0 && ctx.Err() == nil { + // Ensures periodic updates for the time we last saw an empty queue rather than only on transition. r.recordEmptyQueues() - r.mtx.RUnlock() - return } r.mtx.RUnlock() + if len(addRotationFor) == 0 || ctx.Err() != nil { + // No tenant needs to be moved into the rotation or we're shutting down and don't care + return + } + r.mtx.Lock() defer r.mtx.Unlock() for _, tl := range addRotationFor { @@ -421,7 +418,6 @@ func (r *Rotator) Maintenance(ctx context.Context, enforceLeaseExpiration, plan r.addToRotation(tl.lane, tl.tenant, tenantState) } } - // Ensures the same periodic updates, whichever way the maintenance tick ended. if ctx.Err() == nil { r.recordEmptyQueues() } From 27ee3208cf1246fe8e08cba49f8232f2d6faf02d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Nicol=C3=A1s=20Pazos?= Date: Thu, 27 Aug 2026 11:24:31 -0300 Subject: [PATCH 3/5] compactor: add changelog entry for per-lane queue metrics --- CHANGELOG.md | 1 + 1 file changed, 1 insertion(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index 5710c32a0c9..e0b6ddfd963 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -19,6 +19,7 @@ * [ENHANCEMENT] Store-gateway: Add `force_attempt_http2` option to block storage client config to enable HTTP/2. #16385, #16422 * [ENHANCEMENT] Validation: Add an optional `reason` field to `limited_queries` rules, aligning them with `blocked_queries`. When set, the reason is included in the client-facing error and the query-frontend's `"query limited"` log line. #16407 * [ENHANCEMENT] Compactor: Add the experimental `-compactor.scheduler-client.enable-ring-based-cleanup` option, which when disabled stops a scheduler-mode compactor from running the ring-based background blocks cleaner. #16457 +* [ENHANCEMENT] Compactor scheduler: Add per-lane queue metrics. #16489 * [FEATURE] Querier: Add experimental per-tenant limit `-querier.max-blocks-per-store-request` to cap the number of blocks a single store-gateway request may reference. Disabled by default. #16292 * [FEATURE] Validation: Add optional `id`, `note`, `created_by`, `created_at`, and `expires_at` fields to `blocked_queries` and `limited_queries` rules, for tooling to attach ownership/context metadata to a rule. For rules with `expires_at` set, the earliest `expires_at` per tenant and `id` (rules without an `id` are grouped together) is exported as the `cortex_blocked_query_rule_expires_at`/`cortex_limited_query_rule_expires_at` metrics, so an alert can fire on stale rules; this is informational only and never affects enforcement. The query-frontend's `"query blocked"` log line now also includes the matched rule's `id` and whether it is expired, and rate-limited queries are now logged with a new `"query limited"` line carrying the same fields. #16395 * [BUGFIX] Query-frontend: Wait for the querier ring to be populated during startup, up to 30 seconds, before reporting the query-frontend as ready. Previously a query-frontend could become ready before it had seen any querier in the ring and fail every query it received until the ring was populated. Only applies when remote execution is enabled, and can be disabled with the experimental `-query-frontend.wait-for-querier-ring-on-startup=false`. #16333 From 68237be2236659b75ffc07b11ef4c5e87f0a9c74 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Nicol=C3=A1s=20Pazos?= Date: Thu, 27 Aug 2026 14:39:39 -0300 Subject: [PATCH 4/5] compactor: track incomplete bytes per lane without a composite key --- pkg/compactor/scheduler/metrics.go | 74 +++++++++++++++--------------- 1 file changed, 36 insertions(+), 38 deletions(-) diff --git a/pkg/compactor/scheduler/metrics.go b/pkg/compactor/scheduler/metrics.go index 616ccf3758c..98ad77ab14b 100644 --- a/pkg/compactor/scheduler/metrics.go +++ b/pkg/compactor/scheduler/metrics.go @@ -15,11 +15,6 @@ const ( compactionTypeMerge = "merge" ) -type incompleteBytesKey struct { - compactionType string - lane lane -} - type schedulerMetrics struct { pendingJobs *prometheus.GaugeVec pendingJobsByUser *prometheus.GaugeVec @@ -27,7 +22,8 @@ type schedulerMetrics struct { lanePendingJobsLastEmpty *prometheus.GaugeVec lanePendingJobsLastEmptyGauges map[lane]prometheus.Gauge incompleteJobsBytes *prometheus.GaugeVec - incompleteBytesGauges map[incompleteBytesKey]prometheus.Gauge + incompleteSplitBytes map[lane]prometheus.Gauge + incompleteMergeBytes map[lane]prometheus.Gauge activeJobs *prometheus.GaugeVec activeJobsByUser *prometheus.GaugeVec jobsCompleted *prometheus.CounterVec @@ -88,21 +84,16 @@ func newSchedulerMetrics(reg prometheus.Registerer, lanePolicy lanePolicy) *sche for _, l := range allLanes { m.lanePendingJobsLastEmptyGauges[l] = m.lanePendingJobsLastEmpty.WithLabelValues(string(l)) } - m.incompleteBytesGauges = make(map[incompleteBytesKey]prometheus.Gauge, 2*len(compactionLanes)) - for _, t := range []string{compactionTypeSplit, compactionTypeMerge} { - for _, l := range compactionLanes { - k := incompleteBytesKey{compactionType: t, lane: l} - m.incompleteBytesGauges[k] = m.incompleteJobsBytes.WithLabelValues(t, string(l)) - } + m.incompleteSplitBytes = make(map[lane]prometheus.Gauge, len(compactionLanes)) + m.incompleteMergeBytes = make(map[lane]prometheus.Gauge, len(compactionLanes)) + for _, l := range compactionLanes { + m.incompleteSplitBytes[l] = m.incompleteJobsBytes.WithLabelValues(compactionTypeSplit, string(l)) + m.incompleteMergeBytes[l] = m.incompleteJobsBytes.WithLabelValues(compactionTypeMerge, string(l)) } return m } func (s *schedulerMetrics) newTrackerMetricsForTenant(tenant string) *trackerMetrics { - byKey := make(map[incompleteBytesKey]*incompleteBytes, len(s.incompleteBytesGauges)) - for k, g := range s.incompleteBytesGauges { - byKey[k] = &incompleteBytes{gauge: g} - } return &trackerMetrics{ queue: &queueMetrics{ pendingJobsByUser: s.pendingJobsByUser.WithLabelValues(tenant), @@ -111,7 +102,10 @@ func (s *schedulerMetrics) newTrackerMetricsForTenant(tenant string) *trackerMet pendingCompactionJobs: s.pendingJobs.WithLabelValues(jobTypeCompaction), activePlanJobs: s.activeJobs.WithLabelValues(jobTypePlan), activeCompactionJobs: s.activeJobs.WithLabelValues(jobTypeCompaction), - incompleteBytes: byKey, + incompleteSplitBytes: s.incompleteSplitBytes, + incompleteMergeBytes: s.incompleteMergeBytes, + splitBytes: make(map[lane]uint64, len(s.incompleteSplitBytes)), + mergeBytes: make(map[lane]uint64, len(s.incompleteMergeBytes)), laneForJob: s.lanePolicy.LaneForJob, clear: func() { s.pendingJobsByUser.DeleteLabelValues(tenant) @@ -122,11 +116,6 @@ func (s *schedulerMetrics) newTrackerMetricsForTenant(tenant string) *trackerMet } } -type incompleteBytes struct { - gauge prometheus.Gauge - contributed uint64 -} - type trackerMetrics struct { queue *queueMetrics repeatedJobFailures prometheus.Counter @@ -136,9 +125,13 @@ type trackerMetrics struct { // shared gauges. Must be called when a tenant is removed. func (m *trackerMetrics) Clear() { q := m.queue - for _, b := range q.incompleteBytes { - b.gauge.Sub(float64(b.contributed)) - b.contributed = 0 + for l, bytes := range q.splitBytes { + q.incompleteSplitBytes[l].Sub(float64(bytes)) + delete(q.splitBytes, l) + } + for l, bytes := range q.mergeBytes { + q.incompleteMergeBytes[l].Sub(float64(bytes)) + delete(q.mergeBytes, l) } q.pendingPlanJobs.Sub(float64(q.pendingPlanCount)) q.pendingCompactionJobs.Sub(float64(q.pendingCompactionCount)) @@ -164,11 +157,14 @@ type queueMetrics struct { pendingCompactionJobs prometheus.Gauge activePlanJobs prometheus.Gauge activeCompactionJobs prometheus.Gauge - incompleteBytes map[incompleteBytesKey]*incompleteBytes + incompleteSplitBytes map[lane]prometheus.Gauge + incompleteMergeBytes map[lane]prometheus.Gauge laneForJob func(TrackedJob) lane // This tenant's contribution to the shared gauges, tracked so Clear() can subtract exactly // the right amount on tenant removal. + splitBytes map[lane]uint64 + mergeBytes map[lane]uint64 pendingPlanCount int pendingCompactionCount int activePlanCount int @@ -270,21 +266,23 @@ func (q *queueMetrics) decActive(isPlan bool) { } func (q *queueMetrics) addBytes(cj *TrackedCompactionJob) { - b := q.bytesFor(cj) - b.contributed += cj.totalBlockBytes - b.gauge.Add(float64(cj.totalBlockBytes)) + l := q.laneForJob(cj) + if cj.value.isSplit { + q.splitBytes[l] += cj.totalBlockBytes + q.incompleteSplitBytes[l].Add(float64(cj.totalBlockBytes)) + } else { + q.mergeBytes[l] += cj.totalBlockBytes + q.incompleteMergeBytes[l].Add(float64(cj.totalBlockBytes)) + } } func (q *queueMetrics) subBytes(cj *TrackedCompactionJob) { - b := q.bytesFor(cj) - b.contributed -= cj.totalBlockBytes - b.gauge.Sub(float64(cj.totalBlockBytes)) -} - -func (q *queueMetrics) bytesFor(cj *TrackedCompactionJob) *incompleteBytes { - compactionType := compactionTypeMerge + l := q.laneForJob(cj) if cj.value.isSplit { - compactionType = compactionTypeSplit + q.splitBytes[l] -= cj.totalBlockBytes + q.incompleteSplitBytes[l].Sub(float64(cj.totalBlockBytes)) + } else { + q.mergeBytes[l] -= cj.totalBlockBytes + q.incompleteMergeBytes[l].Sub(float64(cj.totalBlockBytes)) } - return q.incompleteBytes[incompleteBytesKey{compactionType: compactionType, lane: q.laneForJob(cj)}] } From f7bef828616c5da1994e6146aa030509fb81a347 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Nicol=C3=A1s=20Pazos?= Date: Thu, 27 Aug 2026 14:53:06 -0300 Subject: [PATCH 5/5] compactor: drop the lane gauge vec field and simplify its test assertion --- pkg/compactor/scheduler/metrics.go | 38 ++++++++++++------------- pkg/compactor/scheduler/rotator_test.go | 35 +++++++++++------------ pkg/compactor/scheduler/scheduler.go | 2 +- 3 files changed, 37 insertions(+), 38 deletions(-) diff --git a/pkg/compactor/scheduler/metrics.go b/pkg/compactor/scheduler/metrics.go index 98ad77ab14b..8773c059995 100644 --- a/pkg/compactor/scheduler/metrics.go +++ b/pkg/compactor/scheduler/metrics.go @@ -16,19 +16,18 @@ const ( ) type schedulerMetrics struct { - pendingJobs *prometheus.GaugeVec - pendingJobsByUser *prometheus.GaugeVec - pendingJobsLastEmpty prometheus.Gauge - lanePendingJobsLastEmpty *prometheus.GaugeVec - lanePendingJobsLastEmptyGauges map[lane]prometheus.Gauge - incompleteJobsBytes *prometheus.GaugeVec - incompleteSplitBytes map[lane]prometheus.Gauge - incompleteMergeBytes map[lane]prometheus.Gauge - activeJobs *prometheus.GaugeVec - activeJobsByUser *prometheus.GaugeVec - jobsCompleted *prometheus.CounterVec - repeatedJobFailures prometheus.Counter - lanePolicy lanePolicy + pendingJobs *prometheus.GaugeVec + pendingJobsByUser *prometheus.GaugeVec + pendingJobsLastEmpty prometheus.Gauge + lanePendingJobsLastEmpty map[lane]prometheus.Gauge + incompleteJobsBytes *prometheus.GaugeVec + incompleteSplitBytes map[lane]prometheus.Gauge + incompleteMergeBytes map[lane]prometheus.Gauge + activeJobs *prometheus.GaugeVec + activeJobsByUser *prometheus.GaugeVec + jobsCompleted *prometheus.CounterVec + repeatedJobFailures prometheus.Counter + lanePolicy lanePolicy } func newSchedulerMetrics(reg prometheus.Registerer, lanePolicy lanePolicy) *schedulerMetrics { @@ -48,10 +47,6 @@ func newSchedulerMetrics(reg prometheus.Registerer, lanePolicy lanePolicy) *sche Name: "cortex_compactor_scheduler_pending_jobs_last_empty_timestamp_seconds", Help: "Unix timestamp of the last time there were no pending jobs remaining.", }), - lanePendingJobsLastEmpty: promauto.With(reg).NewGaugeVec(prometheus.GaugeOpts{ - Name: "cortex_compactor_scheduler_lane_pending_jobs_last_empty_timestamp_seconds", - Help: "Unix timestamp of the last time there were no pending jobs remaining in this lane.", - }, []string{"lane"}), incompleteJobsBytes: promauto.With(reg).NewGaugeVec(prometheus.GaugeOpts{ Name: "cortex_compactor_scheduler_incomplete_compaction_jobs_bytes", Help: "The total bytes of blocks in compaction jobs that have not yet completed (pending or active).", @@ -73,6 +68,11 @@ func newSchedulerMetrics(reg prometheus.Registerer, lanePolicy lanePolicy) *sche Help: "Total number of failures for jobs that exceeded the repeated failure threshold.", }), } + laneLastEmpty := promauto.With(reg).NewGaugeVec(prometheus.GaugeOpts{ + Name: "cortex_compactor_scheduler_lane_pending_jobs_last_empty_timestamp_seconds", + Help: "Unix timestamp of the last time there were no pending jobs remaining in this lane.", + }, []string{"lane"}) + // Pre-initialize job type labels so we get zeros instead of no data. m.jobsCompleted.WithLabelValues(jobTypePlan) m.jobsCompleted.WithLabelValues(jobTypeCompaction) @@ -80,9 +80,9 @@ func newSchedulerMetrics(reg prometheus.Registerer, lanePolicy lanePolicy) *sche m.pendingJobs.WithLabelValues(jobTypeCompaction) m.activeJobs.WithLabelValues(jobTypePlan) m.activeJobs.WithLabelValues(jobTypeCompaction) - m.lanePendingJobsLastEmptyGauges = make(map[lane]prometheus.Gauge, len(allLanes)) + m.lanePendingJobsLastEmpty = make(map[lane]prometheus.Gauge, len(allLanes)) for _, l := range allLanes { - m.lanePendingJobsLastEmptyGauges[l] = m.lanePendingJobsLastEmpty.WithLabelValues(string(l)) + m.lanePendingJobsLastEmpty[l] = laneLastEmpty.WithLabelValues(string(l)) } m.incompleteSplitBytes = make(map[lane]prometheus.Gauge, len(compactionLanes)) m.incompleteMergeBytes = make(map[lane]prometheus.Gauge, len(compactionLanes)) diff --git a/pkg/compactor/scheduler/rotator_test.go b/pkg/compactor/scheduler/rotator_test.go index 4a5cd78ac40..80092327b3a 100644 --- a/pkg/compactor/scheduler/rotator_test.go +++ b/pkg/compactor/scheduler/rotator_test.go @@ -6,8 +6,6 @@ import ( "container/list" "context" "fmt" - "maps" - "slices" "strings" "sync" "testing" @@ -75,7 +73,7 @@ func TestRotator_RecoverFrom_ColdStartDelay(t *testing.T) { t.Run(tc.name, func(t *testing.T) { reg := prometheus.NewPedanticRegistry() metrics := newTestSchedulerMetrics(reg) - r := NewRotator(0, 0, 0, maintenanceInterval, 0, intervalsBeforeColdStartPlanning, newSimpleLanePolicy(), metrics.pendingJobsLastEmpty, metrics.lanePendingJobsLastEmptyGauges, log.NewNopLogger()) + r := NewRotator(0, 0, 0, maintenanceInterval, 0, intervalsBeforeColdStartPlanning, newSimpleLanePolicy(), metrics.pendingJobsLastEmpty, metrics.lanePendingJobsLastEmpty, log.NewNopLogger()) r.clock = clock r.RecoverFrom(tc.jobTrackers, tc.creationTime) @@ -87,7 +85,7 @@ func TestRotator_RecoverFrom_ColdStartDelay(t *testing.T) { func newRotatorForTest() *Rotator { reg := prometheus.NewPedanticRegistry() metrics := newTestSchedulerMetrics(reg) - return NewRotator(0, 0, 0, time.Minute, 0, 0, newSimpleLanePolicy(), metrics.pendingJobsLastEmpty, metrics.lanePendingJobsLastEmptyGauges, log.NewNopLogger()) + return NewRotator(0, 0, 0, time.Minute, 0, 0, newSimpleLanePolicy(), metrics.pendingJobsLastEmpty, metrics.lanePendingJobsLastEmpty, log.NewNopLogger()) } // newTrackerWithPendingJobs builds a JobTracker for the named tenant holding numJobs pending @@ -231,7 +229,7 @@ func TestRotator_LeaseJob_LanePriority(t *testing.T) { clk := clock.New() lanePolicy := newSimpleLanePolicy() metrics := newTestSchedulerMetrics(prometheus.NewPedanticRegistry()) - r := NewRotator(0, 0, 0, time.Minute, 0, 0, lanePolicy, metrics.pendingJobsLastEmpty, metrics.lanePendingJobsLastEmptyGauges, log.NewNopLogger()) + r := NewRotator(0, 0, 0, time.Minute, 0, 0, lanePolicy, metrics.pendingJobsLastEmpty, metrics.lanePendingJobsLastEmpty, log.NewNopLogger()) // Add a tenant with a plan job and a compaction job jt := NewJobTracker(&NopJobPersister{}, "t1", clk, lanePolicy, infiniteLeases, infiniteLeases, metrics.newTrackerMetricsForTenant("t1"), log.NewNopLogger()) @@ -356,23 +354,24 @@ func TestRotator_PendingJobsLastEmpty(t *testing.T) { clk.Set(now) reg := prometheus.NewPedanticRegistry() metrics := newTestSchedulerMetrics(reg) - r := NewRotator(0, 0, 0, 0, 0, 0, newSimpleLanePolicy(), metrics.pendingJobsLastEmpty, metrics.lanePendingJobsLastEmptyGauges, log.NewNopLogger()) + r := NewRotator(0, 0, 0, 0, 0, 0, newSimpleLanePolicy(), metrics.pendingJobsLastEmpty, metrics.lanePendingJobsLastEmpty, log.NewNopLogger()) r.clock = clk tc.action(r, clk) - const ( - laneMetric = "cortex_compactor_scheduler_lane_pending_jobs_last_empty_timestamp_seconds" - anyMetric = "cortex_compactor_scheduler_pending_jobs_last_empty_timestamp_seconds" - ) - var expected strings.Builder - fmt.Fprintf(&expected, "# HELP %s Unix timestamp of the last time there were no pending jobs remaining in this lane.\n# TYPE %s gauge\n", laneMetric, laneMetric) - for _, l := range slices.Sorted(maps.Keys(tc.expectedByLane)) { - fmt.Fprintf(&expected, "%s{lane=%q} %d\n", laneMetric, l, tc.expectedByLane[l]) - } - fmt.Fprintf(&expected, "# HELP %s Unix timestamp of the last time there were no pending jobs remaining.\n# TYPE %s gauge\n%s %d\n", anyMetric, anyMetric, anyMetric, tc.expected) - - require.NoError(t, prom_testutil.GatherAndCompare(reg, strings.NewReader(expected.String()), laneMetric, anyMetric)) + expected := fmt.Sprintf(` + # HELP cortex_compactor_scheduler_lane_pending_jobs_last_empty_timestamp_seconds Unix timestamp of the last time there were no pending jobs remaining in this lane. + # TYPE cortex_compactor_scheduler_lane_pending_jobs_last_empty_timestamp_seconds gauge + cortex_compactor_scheduler_lane_pending_jobs_last_empty_timestamp_seconds{lane="compaction"} %d + cortex_compactor_scheduler_lane_pending_jobs_last_empty_timestamp_seconds{lane="plan"} %d + # HELP cortex_compactor_scheduler_pending_jobs_last_empty_timestamp_seconds Unix timestamp of the last time there were no pending jobs remaining. + # TYPE cortex_compactor_scheduler_pending_jobs_last_empty_timestamp_seconds gauge + cortex_compactor_scheduler_pending_jobs_last_empty_timestamp_seconds %d + `, tc.expectedByLane[compactionLane], tc.expectedByLane[planLane], tc.expected) + + require.NoError(t, prom_testutil.GatherAndCompare(reg, strings.NewReader(expected), + "cortex_compactor_scheduler_lane_pending_jobs_last_empty_timestamp_seconds", + "cortex_compactor_scheduler_pending_jobs_last_empty_timestamp_seconds")) }) } } diff --git a/pkg/compactor/scheduler/scheduler.go b/pkg/compactor/scheduler/scheduler.go index aaf8b7423a1..55f690a9e4a 100644 --- a/pkg/compactor/scheduler/scheduler.go +++ b/pkg/compactor/scheduler/scheduler.go @@ -159,7 +159,7 @@ func newCompactorScheduler( cfg.MaintenanceIntervalsBeforeColdStartPlanning, lanePolicy, metrics.pendingJobsLastEmpty, - metrics.lanePendingJobsLastEmptyGauges, + metrics.lanePendingJobsLastEmpty, logger, )