From 5571942ce76c8b1e70e5e44eec820eaa07ed5c0e Mon Sep 17 00:00:00 2001 From: "nilesh.rathi" Date: Fri, 7 Aug 2026 09:08:07 +0530 Subject: [PATCH] [LIVY-1063] Fix resource leak in Kubernetes service account token refresh --- .../livy/utils/SparkKubernetesApp.scala | 37 ++++++++++--------- 1 file changed, 19 insertions(+), 18 deletions(-) diff --git a/server/src/main/scala/org/apache/livy/utils/SparkKubernetesApp.scala b/server/src/main/scala/org/apache/livy/utils/SparkKubernetesApp.scala index 22ab6b8f8..7ba366160 100644 --- a/server/src/main/scala/org/apache/livy/utils/SparkKubernetesApp.scala +++ b/server/src/main/scala/org/apache/livy/utils/SparkKubernetesApp.scala @@ -77,26 +77,27 @@ object SparkKubernetesApp extends Logging { } val RefreshServiceAccountTokenThread = new Thread() { + // Token expires after 1 hour by default; refresh every 5 minutes. + private val tokenRefreshIntervalMs = 300000L + override def run(): Unit = { - while (true) { - var currentContext = new Context() - var currentContextName = new String - val config = kubernetesClient.getConfiguration - if (config.getCurrentContext != null) { - currentContext = config.getCurrentContext.getContext - currentContextName = config.getCurrentContext.getName + while (!Thread.currentThread().isInterrupted) { + try { + val config = kubernetesClient.getConfiguration + val contextName = Option(config.getCurrentContext).map(_.getName).getOrElse("") + val newestConfig = Config.autoConfigure(contextName) + config.setOauthToken(newestConfig.getOauthToken) + info("Refreshed Kubernetes service account token.") + Thread.sleep(tokenRefreshIntervalMs) + } catch { + case e: InterruptedException => + Thread.currentThread().interrupt() + warn("RefreshServiceAccountTokenThread was interrupted, exiting.", e) + return + case NonFatal(e) => + warn(s"Failed to refresh Kubernetes service account token, will retry.", e) + Thread.sleep(tokenRefreshIntervalMs) } - - var newAccessToken = new String - val newestConfig = Config.autoConfigure(currentContextName) - newAccessToken = newestConfig.getOauthToken - info("Refreshed Kubernetes service account token.") - - config.setOauthToken(newAccessToken) - kubernetesClient = new DefaultKubernetesClient(config) - - // Token will expire 1 hour default, community recommend to update every 5 minutes - Thread.sleep(300000) } } }