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]