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 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..506a625c67f 100644 --- a/pkg/compactor/scheduler/lane_policy.go +++ b/pkg/compactor/scheduler/lane_policy.go @@ -10,14 +10,15 @@ 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" ) type laneTransition struct { @@ -28,6 +29,7 @@ type laneTransition struct { // Defines how to map jobs and requests into lanes type lanePolicy interface { 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. } @@ -51,12 +53,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 +75,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..8773c059995 100644 --- a/pkg/compactor/scheduler/metrics.go +++ b/pkg/compactor/scheduler/metrics.go @@ -16,18 +16,25 @@ const ( ) 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 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) *schedulerMetrics { +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.", @@ -43,7 +50,7 @@ func newSchedulerMetrics(reg prometheus.Registerer) *schedulerMetrics { 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.", @@ -61,6 +68,11 @@ func newSchedulerMetrics(reg prometheus.Registerer) *schedulerMetrics { 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) @@ -68,8 +80,16 @@ 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.lanePendingJobsLastEmpty = make(map[lane]prometheus.Gauge, len(allLanes)) + for _, l := range allLanes { + 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)) + for _, l := range compactionLanes { + m.incompleteSplitBytes[l] = m.incompleteJobsBytes.WithLabelValues(compactionTypeSplit, string(l)) + m.incompleteMergeBytes[l] = m.incompleteJobsBytes.WithLabelValues(compactionTypeMerge, string(l)) + } return m } @@ -82,8 +102,11 @@ 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), + 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) s.activeJobsByUser.DeleteLabelValues(tenant) @@ -102,14 +125,18 @@ 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 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)) 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 +157,14 @@ type queueMetrics struct { pendingCompactionJobs prometheus.Gauge activePlanJobs prometheus.Gauge activeCompactionJobs prometheus.Gauge - incompleteSplitBytes prometheus.Gauge - incompleteMergeBytes prometheus.Gauge + 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 uint64 - mergeBytes uint64 + splitBytes map[lane]uint64 + mergeBytes map[lane]uint64 pendingPlanCount int pendingCompactionCount int activePlanCount int @@ -238,21 +266,23 @@ func (q *queueMetrics) decActive(isPlan bool) { } func (q *queueMetrics) addBytes(cj *TrackedCompactionJob) { + l := q.laneForJob(cj) if cj.value.isSplit { - q.splitBytes += cj.totalBlockBytes - q.incompleteSplitBytes.Add(float64(cj.totalBlockBytes)) + q.splitBytes[l] += cj.totalBlockBytes + q.incompleteSplitBytes[l].Add(float64(cj.totalBlockBytes)) } else { - q.mergeBytes += cj.totalBlockBytes - q.incompleteMergeBytes.Add(float64(cj.totalBlockBytes)) + q.mergeBytes[l] += cj.totalBlockBytes + q.incompleteMergeBytes[l].Add(float64(cj.totalBlockBytes)) } } func (q *queueMetrics) subBytes(cj *TrackedCompactionJob) { + l := q.laneForJob(cj) if cj.value.isSplit { - q.splitBytes -= cj.totalBlockBytes - q.incompleteSplitBytes.Sub(float64(cj.totalBlockBytes)) + q.splitBytes[l] -= cj.totalBlockBytes + q.incompleteSplitBytes[l].Sub(float64(cj.totalBlockBytes)) } else { - q.mergeBytes -= cj.totalBlockBytes - q.incompleteMergeBytes.Sub(float64(cj.totalBlockBytes)) + q.mergeBytes[l] -= cj.totalBlockBytes + q.incompleteMergeBytes[l].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..1783878eb56 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,16 +399,14 @@ func (r *Rotator) Maintenance(ctx context.Context, enforceLeaseExpiration, plan addRotationFor = append(addRotationFor, tenantLane{tenant: tenant, lane: l}) } } - stayingEmpty := len(addRotationFor) == 0 && r.allRotationsEmpty() + 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() 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())) - } return } @@ -420,6 +418,9 @@ func (r *Rotator) Maintenance(ctx context.Context, enforceLeaseExpiration, plan r.addToRotation(tl.lane, tl.tenant, tenantState) } } + if ctx.Err() == nil { + r.recordEmptyQueues() + } } func (r *Rotator) possiblyRemoveFromRotation(l lane, tenant string) { @@ -457,17 +458,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..80092327b3a 100644 --- a/pkg/compactor/scheduler/rotator_test.go +++ b/pkg/compactor/scheduler/rotator_test.go @@ -72,8 +72,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.lanePendingJobsLastEmpty, log.NewNopLogger()) r.clock = clock r.RecoverFrom(tc.jobTrackers, tc.creationTime) @@ -84,14 +84,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.lanePendingJobsLastEmpty, 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 +228,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.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()) @@ -286,29 +286,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 +353,25 @@ 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.lanePendingJobsLastEmpty, 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)) + 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 d4c8e3f2eaf..55f690a9e4a 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.lanePendingJobsLastEmpty, 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())