Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion charts/kyuubi/values.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -77,7 +77,10 @@ rbac:
rules:
- apiGroups: [""]
resources: ["pods"]
verbs: ["create", "list", "delete"]
verbs: ["create", "list", "watch", "delete"]
- apiGroups: [""]
resources: ["services"]
verbs: ["list", "watch"]

service:
# configuration of the headless service
Expand Down
46 changes: 23 additions & 23 deletions docs/configuration/settings.md

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
Expand Up @@ -289,7 +289,7 @@ object SparkSQLEngine extends Logging {
kyuubiConf.setIfMissing(FRONTEND_THRIFT_BINARY_BIND_PORT, 0)
kyuubiConf.setIfMissing(HA_ZK_CONN_RETRY_POLICY, RetryPolicies.N_TIME.toString)

if (Utils.isOnK8s) {
if (Utils.isOnK8s()) {
kyuubiConf.setIfMissing(FRONTEND_CONNECTION_URL_USE_HOSTNAME, false)

// https://github.com/apache/kyuubi/issues/3385
Expand Down Expand Up @@ -463,7 +463,7 @@ object SparkSQLEngine extends Logging {

private def isOnK8sClusterMode: Boolean = {
// only spark driver pod will build with `SPARK_APPLICATION_ID` env.
Utils.isOnK8s && sys.env.contains("SPARK_APPLICATION_ID")
Utils.isOnK8s() && sys.env.contains("SPARK_APPLICATION_ID")
}

@VisibleForTesting
Expand Down
4 changes: 3 additions & 1 deletion kyuubi-common/src/main/scala/org/apache/kyuubi/Utils.scala
Original file line number Diff line number Diff line change
Expand Up @@ -381,7 +381,9 @@ object Utils extends Logging {
def getContextOrKyuubiClassLoader: ClassLoader =
Option(Thread.currentThread().getContextClassLoader).getOrElse(getKyuubiClassLoader)

def isOnK8s: Boolean = Files.exists(Paths.get("/var/run/secrets/kubernetes.io"))
def isOnK8s(env: Map[String, String] = sys.env): Boolean =
env.get("KUBERNETES_SERVICE_HOST").exists(_.nonEmpty) &&
env.get("KUBERNETES_SERVICE_PORT").exists(_.nonEmpty)

/**
* Return a nice string representation of the exception. It will call "printStackTrace" to
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1406,7 +1406,10 @@ object KyuubiConf {

val KUBERNETES_NAMESPACE: ConfigEntry[String] =
buildConf("kyuubi.kubernetes.namespace")
.doc("The namespace that will be used for running the kyuubi pods and find engines.")
.doc("The default namespace used by the Kyuubi server's Kubernetes client to discover" +
" and manage engine pods, when the engine submission does not specify one (e.g." +
" `spark.kubernetes.namespace`). It does not control the namespace where the Kyuubi" +
" server itself runs.")
.version("1.7.0")
.stringConf
.createWithDefault("default")
Expand All @@ -1425,9 +1428,12 @@ object KyuubiConf {
val KUBERNETES_CLIENT_INITIALIZE_LIST: ConfigEntry[Seq[String]] =
buildConf("kyuubi.kubernetes.client.initialize.list")
.doc("The kubernetes client initialize list to register kubernetes resource informers" +
" during Kyuubi server startup. This ensure the Kyuubi server is promptly informed for" +
" any Kubernetes resource changes after startup. It is highly recommend to set it for" +
" multiple Kyuubi instances mode. The format is `context1:namespace1,context2:namespace2`.")
" during Kyuubi server startup. This ensures the Kyuubi server is promptly informed for" +
" any Kubernetes resource changes after startup. It is highly recommended to set it for" +
" multiple Kyuubi instances mode. The format is" +
" `context1:namespace1,context2:namespace2`." +
" When the list is empty and Kyuubi runs in Kubernetes, the client for the configured" +
" namespace is initialized automatically with the in-cluster configuration.")
Comment thread
miaht94 marked this conversation as resolved.
.version("1.11.0")

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Note for other reviewers: this entry shipped in v1.10.2 via the branch-1.10 backport of KYUUBI #7027, though version() says 1.11.0 - pre-existing, not related to this PR.

.audience(SERVER)
.immutable
Expand Down
11 changes: 11 additions & 0 deletions kyuubi-common/src/test/scala/org/apache/kyuubi/UtilsSuite.scala
Original file line number Diff line number Diff line change
Expand Up @@ -191,6 +191,17 @@ class UtilsSuite extends KyuubiFunSuite {
assertResult(false)(Utils.isCommandAvailable("un_exist_cmd"))
}

test("is on Kubernetes") {
val kubernetesEnv = Map(
"KUBERNETES_SERVICE_HOST" -> "kubernetes.default.svc",
"KUBERNETES_SERVICE_PORT" -> "443")

assert(Utils.isOnK8s(kubernetesEnv))
assert(!Utils.isOnK8s(Map.empty))
assert(!Utils.isOnK8s(kubernetesEnv.updated("KUBERNETES_SERVICE_HOST", "")))
assert(!Utils.isOnK8s(kubernetesEnv.updated("KUBERNETES_SERVICE_PORT", "")))
}

test("writeToTempFile rejects illegal filenames") {
val dir = Utils.createTempDir()
def stream: ByteArrayInputStream = new ByteArrayInputStream("data".getBytes)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -217,14 +217,26 @@ class KubernetesApplicationOperation extends ApplicationOperation with Logging {
}

private[kyuubi] def getKubernetesClientInitializeInfo(
kyuubiConf: KyuubiConf): Seq[KubernetesInfo] = {
kyuubiConf.get(KyuubiConf.KUBERNETES_CLIENT_INITIALIZE_LIST).map { init =>
val (context, namespace) = init.split(":") match {
case Array(ctx, ns) => (Some(ctx).filterNot(_.isEmpty), Some(ns).filterNot(_.isEmpty))
case Array(ctx) => (Some(ctx).filterNot(_.isEmpty), None)
case _ => (None, None)
kyuubiConf: KyuubiConf,
environment: Map[String, String] = sys.env): Seq[KubernetesInfo] = {
val configuredInitializeInfo =
kyuubiConf.get(KyuubiConf.KUBERNETES_CLIENT_INITIALIZE_LIST).map { init =>
val (context, namespace) = init.split(":") match {
case Array(ctx, ns) => (Some(ctx).filterNot(_.isEmpty), Some(ns).filterNot(_.isEmpty))
case Array(ctx) => (Some(ctx).filterNot(_.isEmpty), None)
case _ => (None, None)
}
KubernetesInfo(context, namespace)
}
KubernetesInfo(context, namespace)
if (configuredInitializeInfo.nonEmpty) {
configuredInitializeInfo
} else if (Utils.isOnK8s(environment)) {
Seq(KubernetesInfo(
// Kube context is not applicable to in-cluster configuration.
None,
Comment thread
miaht94 marked this conversation as resolved.
Some(kyuubiConf.get(KyuubiConf.KUBERNETES_NAMESPACE))))
} else {
Nil
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -116,12 +116,25 @@ class KubernetesApplicationOperationSuite extends KyuubiFunSuite {

test("get kubernetes client initialization info") {
val kyuubiConf = KyuubiConf()
val kubernetesEnv = Map(
KubernetesApplicationOperation.KUBERNETES_SERVICE_HOST -> "kubernetes.default.svc",
KubernetesApplicationOperation.KUBERNETES_SERVICE_PORT -> "443")
val operation = new KubernetesApplicationOperation()

assert(operation.getKubernetesClientInitializeInfo(kyuubiConf, Map.empty) === Nil)
assert(operation.getKubernetesClientInitializeInfo(
kyuubiConf,
kubernetesEnv.updated(KubernetesApplicationOperation.KUBERNETES_SERVICE_HOST, "")) === Nil)

kyuubiConf.set(KyuubiConf.KUBERNETES_NAMESPACE, "kyuubi")
assert(operation.getKubernetesClientInitializeInfo(kyuubiConf, kubernetesEnv) ===
Seq(KubernetesInfo(None, Some("kyuubi"))))

kyuubiConf.set(
KyuubiConf.KUBERNETES_CLIENT_INITIALIZE_LIST.key,
"c1:ns1,c1:ns2,c2:ns1,c2:ns2,c1:,:ns1")

val operation = new KubernetesApplicationOperation()
assert(operation.getKubernetesClientInitializeInfo(kyuubiConf) ===
assert(operation.getKubernetesClientInitializeInfo(kyuubiConf, kubernetesEnv) ===
Array(
KubernetesInfo(Some("c1"), Some("ns1")),
KubernetesInfo(Some("c1"), Some("ns2")),
Expand Down
Loading