This is an automated email from the ASF dual-hosted git repository.

wilfred-s pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/yunikorn-core.git


The following commit(s) were added to refs/heads/master by this push:
     new be3bb487 [YUNIKORN-3310] Add Counter metric for terminal app states 
(#1097)
be3bb487 is described below

commit be3bb487298956dd70906fca15e577bbc12d23cb
Author: Guruprasad Veerannavaru <[email protected]>
AuthorDate: Thu Jul 2 14:14:21 2026 +1000

    [YUNIKORN-3310] Add Counter metric for terminal app states (#1097)
    
    Queue removal (YUNIKORN-2908) resets terminal state gauges
    (completed/failed/rejected) to zero on queue recreation. Add a
    CounterVec (yunikorn_queue_app_total) that persists across queue
    lifecycle by reusing the existing Prometheus registration.
    
    Committed-By-Agent: claude
    
    Closes: #1097
    
    Signed-off-by: Wilfred Spiegelenburg <[email protected]>
---
 pkg/metrics/queue.go                       | 42 +++++++++++++++++++++++++++
 pkg/metrics/queue_test.go                  | 46 ++++++++++++++++++++++++++++++
 pkg/scheduler/objects/application_state.go | 12 ++++++--
 3 files changed, 97 insertions(+), 3 deletions(-)

diff --git a/pkg/metrics/queue.go b/pkg/metrics/queue.go
index 34237a62..7506ad80 100644
--- a/pkg/metrics/queue.go
+++ b/pkg/metrics/queue.go
@@ -19,6 +19,8 @@
 package metrics
 
 import (
+       "errors"
+
        "github.com/prometheus/client_golang/prometheus"
        dto "github.com/prometheus/client_model/go"
        "go.uber.org/zap"
@@ -54,6 +56,7 @@ const (
 // QueueMetrics to declare queue metrics
 type QueueMetrics struct {
        appMetrics           *prometheus.GaugeVec
+       appTerminalMetrics   *prometheus.CounterVec
        containerMetrics     *prometheus.CounterVec
        resourceMetricsLabel *prometheus.GaugeVec
        // Track known resource types
@@ -91,6 +94,14 @@ func InitQueueMetrics(name string) *QueueMetrics {
                        Help:        "Queue resource metrics. State of the 
resource includes `guaranteed`, `max`, `allocated`, `pending`, `preempting`, 
`maxRunningApps`.",
                }, []string{"state", "resource"})
 
+       q.appTerminalMetrics = prometheus.NewCounterVec(
+               prometheus.CounterOpts{
+                       Namespace:   Namespace,
+                       Name:        "queue_app_total",
+                       ConstLabels: prometheus.Labels{"queue": name},
+                       Help:        "Total number of applications that reached 
a terminal state. State includes `completed`, `failed`, `rejected`.",
+               }, []string{"state"})
+
        var queueMetricsList = []prometheus.Collector{
                q.appMetrics,
                q.containerMetrics,
@@ -107,6 +118,23 @@ func InitQueueMetrics(name string) *QueueMetrics {
                }
        }
 
+       // Register the terminal counter separately — if it already exists in 
the
+       // registry (from a previous queue lifecycle), reuse the existing one so
+       // accumulated counts are preserved.
+       if err := prometheus.Register(q.appTerminalMetrics); err != nil {
+               are := &prometheus.AlreadyRegisteredError{}
+               if errors.As(err, are) {
+                       existing, ok := 
are.ExistingCollector.(*prometheus.CounterVec)
+                       if !ok {
+                               log.Log(log.Metrics).Warn("existing terminal 
metrics collector has unexpected type", zap.Error(err))
+                       } else {
+                               q.appTerminalMetrics = existing
+                       }
+               } else {
+                       log.Log(log.Metrics).Warn("failed to register terminal 
metrics collector", zap.Error(err))
+               }
+       }
+
        q.knownResourceTypes = make(map[string]struct{})
        return q
 }
@@ -119,6 +147,8 @@ func (m *QueueMetrics) UnregisterMetrics() {
        }
 
        // Unregister the metrics
+       // appTerminalMetrics is intentionally excluded — it persists in the
+       // Prometheus registry so accumulated counts survive queue recreation.
        for _, metric := range queueMetricsList {
                prometheus.Unregister(metric)
        }
@@ -330,3 +360,15 @@ func (m *QueueMetrics) 
SetQueuePreemptingResourceMetrics(resourceName string, va
 func (m *QueueMetrics) SetQueueMaxRunningAppsMetrics(value uint64) {
        m.setQueueResource(QueueMaxRunningApps, "apps", float64(value))
 }
+
+func (m *QueueMetrics) IncQueueApplicationsCompletedTotal() {
+       m.appTerminalMetrics.WithLabelValues(AppCompleted).Inc()
+}
+
+func (m *QueueMetrics) IncQueueApplicationsFailedTotal() {
+       m.appTerminalMetrics.WithLabelValues(AppFailed).Inc()
+}
+
+func (m *QueueMetrics) IncQueueApplicationsRejectedTotal() {
+       m.appTerminalMetrics.WithLabelValues(AppRejected).Inc()
+}
diff --git a/pkg/metrics/queue_test.go b/pkg/metrics/queue_test.go
index b77772b3..30fb7834 100644
--- a/pkg/metrics/queue_test.go
+++ b/pkg/metrics/queue_test.go
@@ -352,5 +352,51 @@ func unregisterQueueMetrics() {
        prometheus.Unregister(qm.appMetrics)
        prometheus.Unregister(qm.containerMetrics)
        prometheus.Unregister(qm.resourceMetricsLabel)
+       prometheus.Unregister(qm.appTerminalMetrics)
        qm.knownResourceTypes = make(map[string]struct{})
 }
+
+func TestTerminalMetricsReuse(t *testing.T) {
+       // First registration
+       qm1 := InitQueueMetrics("root.reuse")
+       qm1.IncQueueApplicationsCompletedTotal()
+       qm1.IncQueueApplicationsCompletedTotal()
+
+       // Second registration — simulates queue recreation
+       // Should reuse the existing counter with accumulated values
+       qm2 := InitQueueMetrics("root.reuse")
+
+       metricDto := &dto.Metric{}
+       err := 
qm2.appTerminalMetrics.WithLabelValues(AppCompleted).Write(metricDto)
+       assert.NilError(t, err)
+       assert.Equal(t, 2, int(*metricDto.Counter.Value))
+
+       // Cleanup
+       prometheus.Unregister(qm2.appMetrics)
+       prometheus.Unregister(qm2.containerMetrics)
+       prometheus.Unregister(qm2.resourceMetricsLabel)
+       prometheus.Unregister(qm2.appTerminalMetrics)
+}
+
+func TestApplicationsTerminalTotal(t *testing.T) {
+       qm = getQueueMetrics()
+       defer unregisterQueueMetrics()
+
+       qm.IncQueueApplicationsCompletedTotal()
+       qm.IncQueueApplicationsCompletedTotal()
+       qm.IncQueueApplicationsFailedTotal()
+       qm.IncQueueApplicationsRejectedTotal()
+
+       metricDto := &dto.Metric{}
+       err := 
qm.appTerminalMetrics.WithLabelValues(AppCompleted).Write(metricDto)
+       assert.NilError(t, err)
+       assert.Equal(t, 2, int(*metricDto.Counter.Value))
+
+       err = qm.appTerminalMetrics.WithLabelValues(AppFailed).Write(metricDto)
+       assert.NilError(t, err)
+       assert.Equal(t, 1, int(*metricDto.Counter.Value))
+
+       err = 
qm.appTerminalMetrics.WithLabelValues(AppRejected).Write(metricDto)
+       assert.NilError(t, err)
+       assert.Equal(t, 1, int(*metricDto.Counter.Value))
+}
diff --git a/pkg/scheduler/objects/application_state.go 
b/pkg/scheduler/objects/application_state.go
index 3360f1f5..db85f908 100644
--- a/pkg/scheduler/objects/application_state.go
+++ b/pkg/scheduler/objects/application_state.go
@@ -187,7 +187,9 @@ func callbacks() fsm.Callbacks {
                },
                fmt.Sprintf("enter_%s", Rejected.String()): func(_ 
context.Context, event *fsm.Event) {
                        app := event.Args[0].(*Application) //nolint:errcheck
-                       
metrics.GetQueueMetrics(app.queuePath).IncQueueApplicationsRejected()
+                       qm := metrics.GetQueueMetrics(app.queuePath)
+                       qm.IncQueueApplicationsRejected()
+                       qm.IncQueueApplicationsRejectedTotal()
                        
metrics.GetSchedulerMetrics().IncTotalApplicationsRejected()
                        app.setStateTimer(terminatedTimeout, 
app.stateMachine.Current(), ExpireApplication)
                        app.finishedTime = time.Now()
@@ -248,7 +250,9 @@ func callbacks() fsm.Callbacks {
                fmt.Sprintf("enter_%s", Completed.String()): func(_ 
context.Context, event *fsm.Event) {
                        app := event.Args[0].(*Application) //nolint:errcheck
                        
metrics.GetSchedulerMetrics().IncTotalApplicationsCompleted()
-                       
metrics.GetQueueMetrics(app.queuePath).IncQueueApplicationsCompleted()
+                       qm := metrics.GetQueueMetrics(app.queuePath)
+                       qm.IncQueueApplicationsCompleted()
+                       qm.IncQueueApplicationsCompletedTotal()
                        app.setStateTimer(terminatedTimeout, 
app.stateMachine.Current(), ExpireApplication)
                        app.executeTerminatedCallback()
                        app.clearPlaceholderTimer()
@@ -257,7 +261,9 @@ func callbacks() fsm.Callbacks {
                fmt.Sprintf("enter_%s", Failed.String()): func(_ 
context.Context, event *fsm.Event) {
                        app := event.Args[0].(*Application) //nolint:errcheck
                        
metrics.GetSchedulerMetrics().IncTotalApplicationsFailed()
-                       
metrics.GetQueueMetrics(app.queuePath).IncQueueApplicationsFailed()
+                       qm := metrics.GetQueueMetrics(app.queuePath)
+                       qm.IncQueueApplicationsFailed()
+                       qm.IncQueueApplicationsFailedTotal()
                        app.setStateTimer(terminatedTimeout, 
app.stateMachine.Current(), ExpireApplication)
                        app.executeTerminatedCallback()
                        app.cleanupAsks()


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to