From 05bc0496d596e9dd10dc28cd774caa7a1735e0d6 Mon Sep 17 00:00:00 2001 From: ruanwenjun Date: Thu, 10 Sep 2026 18:14:24 +0800 Subject: [PATCH 1/3] [SERVER] Add Kubernetes engine Pod discovery latency metric --- docs/monitor/metrics.md | 1 + .../kyuubi/metrics/MetricsConstants.scala | 2 + .../KubernetesApplicationOperation.scala | 20 ++++ .../KubernetesApplicationOperationSuite.scala | 99 +++++++++++++++++++ 4 files changed, 122 insertions(+) diff --git a/docs/monitor/metrics.md b/docs/monitor/metrics.md index c19a25a22b2..8509b090a58 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 its first ADD or UPDATE populating the local application cache. Includes existing Pods discovered during initial listing; may be recorded again after cache eviction or server restart. | | `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..6354d3caeae 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 @@ -447,6 +451,7 @@ class KubernetesApplicationOperation extends ApplicationOperation with Logging { kubernetesInfo: KubernetesInfo, pod: Pod, eventType: KubernetesResourceEventType): Unit = { + val observedAt = System.currentTimeMillis() val (appState, appError) = toApplicationStateAndError(pod, appStateSource, appStateContainer, eventType) debug(s"Driver Informer changes pod: ${pod.getMetadata.getName} to state: $appState") @@ -474,6 +479,21 @@ class KubernetesApplicationOperation extends ApplicationOperation with Logging { url = getPodAppUrl(sparkAppUrlSource, sparkAppUrlPattern, kubernetesInfo, pod), error = appError, podName = Some(pod.getMetadata.getName))) + if (eventType != KubernetesResourceEventTypes.DELETE) { + Option(pod.getMetadata.getCreationTimestamp).foreach { 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) + } + } + } } } } 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..00b2d1d7bfc 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,114 @@ 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 { + for (firstEvent <- Seq("ADD", "UPDATE")) { + test(s"record pod discovery latency once for first $firstEvent") { + 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 { + val before = System.currentTimeMillis() + if (firstEvent == "ADD") handler.onAdd(pod) else handler.onUpdate(pod, 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) + handler.onUpdate(pod, pod) + handler.onDelete(pod, false) + assert(histogram.getCount == 1) + } finally { + operation.stop() + metricsSystem.stop() + } + } + } + + test("skip invalid pod discovery samples without preventing cache updates") { + 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("Running") + .endStatus() + .build() + handler.onAdd(pod) + // A valid repeated ADD must not sample an already cached Pod. + pod.getMetadata.setCreationTimestamp(Instant.now().minusSeconds(180).toString) + handler.onAdd(pod) + assert(histogram.getCount == 0) + } + 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) From 38cac831bba583f27e3aa8540b341bed01a4cbcd Mon Sep 17 00:00:00 2001 From: ruanwenjun Date: Thu, 10 Sep 2026 20:15:10 +0800 Subject: [PATCH 2/3] [SERVER] Record Pod discovery latency in informer ADD callback --- docs/monitor/metrics.md | 2 +- .../KubernetesApplicationOperation.scala | 30 +++---- .../KubernetesApplicationOperationSuite.scala | 87 +++++++++---------- 3 files changed, 58 insertions(+), 61 deletions(-) diff --git a/docs/monitor/metrics.md b/docs/monitor/metrics.md index 8509b090a58..404105f8243 100644 --- a/docs/monitor/metrics.md +++ b/docs/monitor/metrics.md @@ -63,7 +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 its first ADD or UPDATE populating the local application cache. Includes existing Pods discovered during initial listing; may be recorded again after cache eviction or server restart. | +| `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-server/src/main/scala/org/apache/kyuubi/engine/KubernetesApplicationOperation.scala b/kyuubi-server/src/main/scala/org/apache/kyuubi/engine/KubernetesApplicationOperation.scala index 6354d3caeae..1e8c49eae4d 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 @@ -363,7 +363,21 @@ 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).foreach { 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) + } + } val eventType = KubernetesResourceEventTypes.ADD updateApplicationState(kubernetesInfo, pod, eventType) val appState = toApplicationState(pod, appStateSource, appStateContainer, eventType) @@ -451,7 +465,6 @@ class KubernetesApplicationOperation extends ApplicationOperation with Logging { kubernetesInfo: KubernetesInfo, pod: Pod, eventType: KubernetesResourceEventType): Unit = { - val observedAt = System.currentTimeMillis() val (appState, appError) = toApplicationStateAndError(pod, appStateSource, appStateContainer, eventType) debug(s"Driver Informer changes pod: ${pod.getMetadata.getName} to state: $appState") @@ -479,21 +492,6 @@ class KubernetesApplicationOperation extends ApplicationOperation with Logging { url = getPodAppUrl(sparkAppUrlSource, sparkAppUrlPattern, kubernetesInfo, pod), error = appError, podName = Some(pod.getMetadata.getName))) - if (eventType != KubernetesResourceEventTypes.DELETE) { - Option(pod.getMetadata.getCreationTimestamp).foreach { 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) - } - } - } } } } 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 00b2d1d7bfc..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 @@ -30,49 +30,50 @@ import org.apache.kyuubi.metrics.MetricsSystem class KubernetesApplicationOperationSuite extends KyuubiFunSuite { - for (firstEvent <- Seq("ADD", "UPDATE")) { - test(s"record pod discovery latency once for first $firstEvent") { - 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 { - val before = System.currentTimeMillis() - if (firstEvent == "ADD") handler.onAdd(pod) else handler.onUpdate(pod, 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) - handler.onUpdate(pod, pod) - handler.onDelete(pod, false) - assert(histogram.getCount == 1) - } finally { - operation.stop() - metricsSystem.stop() - } + 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 cache updates") { + test("skip invalid pod discovery samples without preventing pod event handling") { val operation = new KubernetesApplicationOperation() val metricsSystem = new MetricsSystem val conf = KyuubiConf() @@ -93,14 +94,12 @@ class KubernetesApplicationOperationSuite extends KyuubiFunSuite { .addToLabels("spark-app-selector", "spark-application") .endMetadata() .withNewStatus() - .withPhase("Running") + .withPhase("Succeeded") .endStatus() .build() handler.onAdd(pod) - // A valid repeated ADD must not sample an already cached Pod. - pod.getMetadata.setCreationTimestamp(Instant.now().minusSeconds(180).toString) - handler.onAdd(pod) assert(histogram.getCount == 0) + assert(operation.cleanupTerminatedAppInfoTrigger.getIfPresent(s"app-$index") == FINISHED) } val deletedPod = new PodBuilder() .withNewMetadata() From b4007d2657062bc33a0210db392930dcf85a3c65 Mon Sep 17 00:00:00 2001 From: ruanwenjun Date: Thu, 10 Sep 2026 21:03:18 +0800 Subject: [PATCH 3/3] [SERVER] Warn when engine Pod creation timestamp is missing --- .../KubernetesApplicationOperation.scala | 27 +++++++++++-------- 1 file changed, 16 insertions(+), 11 deletions(-) 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 1e8c49eae4d..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 @@ -365,18 +365,23 @@ class KubernetesApplicationOperation extends ApplicationOperation with Logging { override def onAdd(pod: Pod): Unit = { val observedAt = System.currentTimeMillis() if (isSparkEnginePod(pod)) { - Option(pod.getMetadata.getCreationTimestamp).foreach { creationTimestamp => - try { - val latency = observedAt - Instant.parse(creationTimestamp).toEpochMilli - if (latency >= 0) { - MetricsSystem.tracing(_.updateHistogram( - ENGINE_KUBERNETES_POD_DISCOVERY_LATENCY, - latency)) + 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) } - } 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)