diff --git a/docs/monitor/metrics.md b/docs/monitor/metrics.md
index c19a25a22b2..404105f8243 100644
--- a/docs/monitor/metrics.md
+++ b/docs/monitor/metrics.md
@@ -63,6 +63,7 @@ These metrics include:
| `kyuubi.operation.exec_time` | `${operationType}` | histogram | 1.7.0 |
execution time histogram for the operation `${operationType}`, now only `ExecuteStatement` is enabled.
|
| `kyuubi.engine.total` | | counter | 1.2.0 | cumulative created engines
|
| `kyuubi.engine.startup.time` | | histogram | 1.12.0 | startup time histogram for successfully created engines, from submitting a new engine to the engine being registered and discovered by Kyuubi server
|
+| `kyuubi.engine.kubernetes.pod.discovery.latency` | | histogram | 1.13.0 | Time in milliseconds from Spark engine Pod creation to the informer ADD callback. Recorded for each ADD event, including existing Pods discovered during initial listing. UPDATE and DELETE events do not contribute samples. |
| `kyuubi.engine.timeout` | | counter | 1.2.0 | cumulative timeout engines
|
| `kyuubi.engine.failed` | `${user}` | counter | 1.2.0 | cumulative explicitly failed engine count for a `${user}`
|
| `kyuubi.engine.failed` | `${errorType}` | counter | 1.2.0 | cumulative explicitly failed engine count for a particular `${errorType}`, e.g. `ClassNotFoundException`
|
diff --git a/kyuubi-metrics/src/main/scala/org/apache/kyuubi/metrics/MetricsConstants.scala b/kyuubi-metrics/src/main/scala/org/apache/kyuubi/metrics/MetricsConstants.scala
index 285bb2992b4..c2f9589ecae 100644
--- a/kyuubi-metrics/src/main/scala/org/apache/kyuubi/metrics/MetricsConstants.scala
+++ b/kyuubi-metrics/src/main/scala/org/apache/kyuubi/metrics/MetricsConstants.scala
@@ -58,6 +58,8 @@ object MetricsConstants {
final private val ENGINE = KYUUBI + "engine."
final val ENGINE_FAIL: String = ENGINE + "failed"
final val ENGINE_STARTUP_TIME: String = ENGINE + "startup.time"
+ final val ENGINE_KUBERNETES_POD_DISCOVERY_LATENCY: String =
+ ENGINE + "kubernetes.pod.discovery.latency"
final val ENGINE_TIMEOUT: String = ENGINE + "timeout"
final val ENGINE_TOTAL: String = ENGINE + "total"
diff --git a/kyuubi-server/src/main/scala/org/apache/kyuubi/engine/KubernetesApplicationOperation.scala b/kyuubi-server/src/main/scala/org/apache/kyuubi/engine/KubernetesApplicationOperation.scala
index 85723bff0fd..186da6bb132 100644
--- a/kyuubi-server/src/main/scala/org/apache/kyuubi/engine/KubernetesApplicationOperation.scala
+++ b/kyuubi-server/src/main/scala/org/apache/kyuubi/engine/KubernetesApplicationOperation.scala
@@ -17,6 +17,8 @@
package org.apache.kyuubi.engine
+import java.time.Instant
+import java.time.format.DateTimeParseException
import java.util.Locale
import java.util.concurrent.{ConcurrentHashMap, ExecutorService, ScheduledExecutorService, TimeUnit}
@@ -37,6 +39,8 @@ import org.apache.kyuubi.config.KyuubiConf.KubernetesCleanupDriverPodStrategy.{A
import org.apache.kyuubi.engine.ApplicationState.{isTerminated, ApplicationState, FAILED, FINISHED, KILLED, NOT_FOUND, PENDING, RUNNING, UNKNOWN}
import org.apache.kyuubi.engine.KubernetesApplicationUrlSource._
import org.apache.kyuubi.engine.KubernetesResourceEventTypes.KubernetesResourceEventType
+import org.apache.kyuubi.metrics.MetricsConstants.ENGINE_KUBERNETES_POD_DISCOVERY_LATENCY
+import org.apache.kyuubi.metrics.MetricsSystem
import org.apache.kyuubi.operation.OperationState
import org.apache.kyuubi.server.metadata.MetadataManager
import org.apache.kyuubi.server.metadata.api.KubernetesEngineInfo
@@ -359,7 +363,26 @@ class KubernetesApplicationOperation extends ApplicationOperation with Logging {
extends ResourceEventHandler[Pod] {
override def onAdd(pod: Pod): Unit = {
+ val observedAt = System.currentTimeMillis()
if (isSparkEnginePod(pod)) {
+ Option(pod.getMetadata.getCreationTimestamp) match {
+ case Some(creationTimestamp) =>
+ try {
+ val latency = observedAt - Instant.parse(creationTimestamp).toEpochMilli
+ if (latency >= 0) {
+ MetricsSystem.tracing(_.updateHistogram(
+ ENGINE_KUBERNETES_POD_DISCOVERY_LATENCY,
+ latency))
+ }
+ } catch {
+ case e: DateTimeParseException =>
+ warn(s"Invalid creation timestamp for engine pod ${pod.getMetadata.getName}", e)
+ }
+ case None =>
+ warn(s"[$kubernetesInfo] Missing creation timestamp for engine pod " +
+ s"${pod.getMetadata.getNamespace}/${pod.getMetadata.getName}, " +
+ "skipping pod discovery latency metric")
+ }
val eventType = KubernetesResourceEventTypes.ADD
updateApplicationState(kubernetesInfo, pod, eventType)
val appState = toApplicationState(pod, appStateSource, appStateContainer, eventType)
diff --git a/kyuubi-server/src/test/scala/org/apache/kyuubi/engine/KubernetesApplicationOperationSuite.scala b/kyuubi-server/src/test/scala/org/apache/kyuubi/engine/KubernetesApplicationOperationSuite.scala
index 4cc79b583db..e167af97ebb 100644
--- a/kyuubi-server/src/test/scala/org/apache/kyuubi/engine/KubernetesApplicationOperationSuite.scala
+++ b/kyuubi-server/src/test/scala/org/apache/kyuubi/engine/KubernetesApplicationOperationSuite.scala
@@ -17,15 +17,113 @@
package org.apache.kyuubi.engine
+import java.time.Instant
+
import io.fabric8.kubernetes.api.model.{ContainerState, ContainerStateWaiting, PodBuilder}
import org.apache.kyuubi.{KyuubiException, KyuubiFunSuite}
import org.apache.kyuubi.config.KyuubiConf
import org.apache.kyuubi.engine.ApplicationState.{FAILED, FINISHED, PENDING}
import org.apache.kyuubi.engine.KubernetesApplicationOperation.LABEL_KYUUBI_UNIQUE_KEY
+import org.apache.kyuubi.metrics.MetricsConstants.ENGINE_KUBERNETES_POD_DISCOVERY_LATENCY
+import org.apache.kyuubi.metrics.MetricsSystem
class KubernetesApplicationOperationSuite extends KyuubiFunSuite {
+ test("record pod discovery latency only for ADD events") {
+ val operation = new KubernetesApplicationOperation()
+ val metricsSystem = new MetricsSystem
+ val conf = KyuubiConf()
+ operation.initialize(conf, None)
+ metricsSystem.initialize(conf)
+ metricsSystem.start()
+ val createdAt = Instant.now().minusSeconds(180)
+ val pod = new PodBuilder()
+ .withNewMetadata()
+ .withName("delayed-driver")
+ .withCreationTimestamp(createdAt.toString)
+ .addToLabels(LABEL_KYUUBI_UNIQUE_KEY, "delayed-app")
+ .addToLabels("spark-app-selector", "spark-application")
+ .endMetadata()
+ .withNewStatus()
+ .withPhase("Running")
+ .endStatus()
+ .build()
+ val handler = new operation.SparkEnginePodEventHandler(KubernetesInfo())
+ val histogram = MetricsSystem.getMetricsRegistry.get
+ .histogram(ENGINE_KUBERNETES_POD_DISCOVERY_LATENCY)
+ try {
+ handler.onUpdate(pod, pod)
+ assert(histogram.getCount == 0)
+ val before = System.currentTimeMillis()
+ handler.onAdd(pod)
+ val after = System.currentTimeMillis()
+ assert(histogram.getCount == 1)
+ val latency = histogram.getSnapshot.getMax
+ assert(latency >= before - createdAt.toEpochMilli)
+ assert(latency <= after - createdAt.toEpochMilli)
+ handler.onAdd(pod)
+ assert(histogram.getCount == 2)
+ handler.onUpdate(pod, pod)
+ handler.onDelete(pod, false)
+ assert(histogram.getCount == 2)
+ } finally {
+ operation.stop()
+ metricsSystem.stop()
+ }
+ }
+
+ test("skip invalid pod discovery samples without preventing pod event handling") {
+ val operation = new KubernetesApplicationOperation()
+ val metricsSystem = new MetricsSystem
+ val conf = KyuubiConf()
+ operation.initialize(conf, None)
+ metricsSystem.initialize(conf)
+ metricsSystem.start()
+ val histogram = MetricsSystem.getMetricsRegistry.get
+ .histogram(ENGINE_KUBERNETES_POD_DISCOVERY_LATENCY)
+ val handler = new operation.SparkEnginePodEventHandler(KubernetesInfo())
+ try {
+ Seq(null, "invalid", Instant.now().plusSeconds(3600).toString).zipWithIndex.foreach {
+ case (creationTimestamp, index) =>
+ val pod = new PodBuilder()
+ .withNewMetadata()
+ .withName(s"driver-$index")
+ .withCreationTimestamp(creationTimestamp)
+ .addToLabels(LABEL_KYUUBI_UNIQUE_KEY, s"app-$index")
+ .addToLabels("spark-app-selector", "spark-application")
+ .endMetadata()
+ .withNewStatus()
+ .withPhase("Succeeded")
+ .endStatus()
+ .build()
+ handler.onAdd(pod)
+ assert(histogram.getCount == 0)
+ assert(operation.cleanupTerminatedAppInfoTrigger.getIfPresent(s"app-$index") == FINISHED)
+ }
+ val deletedPod = new PodBuilder()
+ .withNewMetadata()
+ .withName("deleted-driver")
+ .withCreationTimestamp(Instant.now().minusSeconds(180).toString)
+ .addToLabels(LABEL_KYUUBI_UNIQUE_KEY, "deleted-app")
+ .addToLabels("spark-app-selector", "spark-application")
+ .endMetadata()
+ .withNewStatus()
+ .withPhase("Running")
+ .endStatus()
+ .build()
+ handler.onDelete(deletedPod, false)
+ assert(histogram.getCount == 0)
+ deletedPod.getMetadata.getLabels.put(LABEL_KYUUBI_UNIQUE_KEY, "non-engine-app")
+ deletedPod.getMetadata.getLabels.remove("spark-app-selector")
+ handler.onAdd(deletedPod)
+ assert(histogram.getCount == 0)
+ } finally {
+ operation.stop()
+ metricsSystem.stop()
+ }
+ }
+
test("mark terminated application received from pod add event") {
val operation = new KubernetesApplicationOperation()
operation.initialize(KyuubiConf(), None)