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
318 changes: 164 additions & 154 deletions docs/configuration/settings.md

Large diffs are not rendered by default.

157 changes: 150 additions & 7 deletions kyuubi-common/src/main/scala/org/apache/kyuubi/config/KyuubiConf.scala
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ import org.apache.kyuubi.engine.{EngineType, ShareLevel}
import org.apache.kyuubi.engine.deploy.DeployMode
import org.apache.kyuubi.operation.{NoneMode, PlainStyle}
import org.apache.kyuubi.service.authentication.{AuthTypes, SaslQOP}
import org.apache.kyuubi.util.ThreadUtils

case class KyuubiConf(loadSysDefault: Boolean = true) extends Logging {

Expand All @@ -59,6 +60,14 @@ case class KyuubiConf(loadSysDefault: Boolean = true) extends Logging {
this
}

private[kyuubi] def validateServerVirtualThreadConfigs(): Unit = {
if (!ThreadUtils.isVirtualThreadSupported) {
serverVirtualThreadConfigs.foreach { config =>
require(!get(config), s"${config.key}=true requires Java 21 or later")
}
}
}

def set[T](entry: ConfigEntry[T], value: T): KyuubiConf = {
require(entry != null, "entry cannot be null")
require(value != null, s"value cannot be null for key: ${entry.key}")
Expand Down Expand Up @@ -635,17 +644,29 @@ object KyuubiConf {
.immutable
.fallbackConf(FRONTEND_BIND_PORT)

val SERVER_VIRTUAL_THREADS_ENABLED: ConfigEntry[Boolean] =
buildConf("kyuubi.server.virtualThreads.enabled")
.doc("Whether to use virtual threads by default for all supported executors in the " +
Comment thread
wangzhigang1999 marked this conversation as resolved.
"Kyuubi server. Requires at least Java 21; Java 25 or later is recommended. " +
"An explicitly configured component-level virtual thread option takes precedence.")
.version("1.13.0")
.audience(SERVER)
.immutable
.booleanConf
.createWithDefault(false)

val FRONTEND_THRIFT_BINARY_VIRTUAL_THREADS_ENABLED: ConfigEntry[Boolean] =
buildConf("kyuubi.frontend.thrift.binary.virtualThreads.enabled")
.doc("Whether to use virtual threads for the Kyuubi server thrift binary frontend " +
"workers. This requires Java 21 or later. The maximum number of concurrent workers " +
"workers. Requires at least Java 21; Java 25 or later is recommended. " +
"The maximum number of concurrent workers " +
"remains limited by kyuubi.frontend.thrift.max.worker.threads. The minimum worker " +
"threads and worker keepalive configurations do not apply in this mode.")
"threads and worker keepalive configurations do not apply in this mode. " +
s"Defaults to ${SERVER_VIRTUAL_THREADS_ENABLED.key}.")
.version("1.13.0")
.audience(SERVER)
.immutable
.booleanConf
.createWithDefault(false)
.fallbackConf(SERVER_VIRTUAL_THREADS_ENABLED)

val FRONTEND_THRIFT_HTTP_BIND_HOST: ConfigEntry[Option[String]] =
buildConf("kyuubi.frontend.thrift.http.bind.host")
Expand Down Expand Up @@ -1441,6 +1462,26 @@ object KyuubiConf {
.toSequence()
.createWithDefault(Nil)

val KUBERNETES_CLIENT_DISPATCHER_VIRTUAL_THREADS_ENABLED: ConfigEntry[Boolean] =
buildConf("kyuubi.kubernetes.client.dispatcher.virtualThreads.enabled")
.doc("Whether the Kyuubi server Kubernetes HTTP client dispatcher uses virtual threads. " +
"Requires at least Java 21; Java 25 or later is recommended. " +
s"Defaults to ${SERVER_VIRTUAL_THREADS_ENABLED.key}.")
.version("1.13.0")
.audience(SERVER)
.immutable
.fallbackConf(SERVER_VIRTUAL_THREADS_ENABLED)

val KUBERNETES_APPLICATION_CLEANUP_VIRTUAL_THREADS_ENABLED: ConfigEntry[Boolean] =
buildConf("kyuubi.kubernetes.application.cleanup.virtualThreads.enabled")
.doc("Whether asynchronous Kubernetes application cleanup tasks in the Kyuubi server " +
"use virtual threads. Requires at least Java 21; Java 25 or later is recommended. " +
s"Defaults to ${SERVER_VIRTUAL_THREADS_ENABLED.key}.")
.version("1.13.0")
.audience(SERVER)
.immutable
.fallbackConf(SERVER_VIRTUAL_THREADS_ENABLED)

val KUBERNETES_MASTER: OptionalConfigEntry[String] =
buildConf("kyuubi.kubernetes.master.address")
.doc("The internal Kubernetes master (API server) address to be used for kyuubi.")
Expand Down Expand Up @@ -1601,6 +1642,16 @@ object KyuubiConf {
.checkValue(_ > 0, "must be positive number")
.createWithDefault(Duration.ofDays(1).toMillis)

val SERVER_ENGINE_LOG_CAPTURE_VIRTUAL_THREADS_ENABLED: ConfigEntry[Boolean] =
buildConf("kyuubi.session.engine.log.capture.virtualThreads.enabled")
.doc("Whether the Kyuubi server uses virtual threads to capture engine startup logs. " +
"Requires at least Java 21; Java 25 or later is recommended. " +
s"Defaults to ${SERVER_VIRTUAL_THREADS_ENABLED.key}.")
.version("1.13.0")
.audience(SERVER)
.immutable
.fallbackConf(SERVER_VIRTUAL_THREADS_ENABLED)

val ENGINE_SPARK_MAIN_RESOURCE: OptionalConfigEntry[String] =
buildConf("kyuubi.session.engine.spark.main.resource")
.doc("The package used to create Spark SQL engine remote application. If it is undefined," +
Expand Down Expand Up @@ -1805,6 +1856,17 @@ object KyuubiConf {
.timeConf
.createWithDefault(Duration.ofSeconds(15).toMillis)

val ENGINE_RPC_CLIENT_VIRTUAL_THREADS_ENABLED: ConfigEntry[Boolean] =
buildConf("kyuubi.session.engine.rpc.client.virtualThreads.enabled")
.doc("Whether the Kyuubi server uses virtual threads for per-session blocking RPC calls " +
"to SQL engines. Requires at least Java 21; Java 25 or later is recommended. " +
"RPC calls for one session remain serialized. " +
s"Defaults to ${SERVER_VIRTUAL_THREADS_ENABLED.key}.")
.version("1.13.0")
.audience(SERVER)
.immutable
.fallbackConf(SERVER_VIRTUAL_THREADS_ENABLED)

val ENGINE_ALIVE_MAX_FAILURES: ConfigEntry[Int] =
buildConf("kyuubi.session.engine.alive.max.failures")
.doc("The maximum number of failures allowed for the engine.")
Expand All @@ -1821,6 +1883,17 @@ object KyuubiConf {
.booleanConf
.createWithDefault(false)

val ENGINE_ALIVE_PROBE_VIRTUAL_THREADS_ENABLED: ConfigEntry[Boolean] =
buildConf("kyuubi.session.engine.alive.probe.virtualThreads.enabled")
.doc("Whether engine alive probes initiated by the Kyuubi server use virtual threads. " +
"Requires at least Java 21; Java 25 or later is recommended. " +
"Probes for one session remain serialized. " +
s"Defaults to ${SERVER_VIRTUAL_THREADS_ENABLED.key}.")
.version("1.13.0")
.audience(SERVER)
.immutable
.fallbackConf(SERVER_VIRTUAL_THREADS_ENABLED)

val ENGINE_ALIVE_PROBE_INTERVAL: ConfigEntry[Long] =
buildConf("kyuubi.session.engine.alive.probe.interval")
.doc("The interval for engine alive probe.")
Expand Down Expand Up @@ -2159,6 +2232,17 @@ object KyuubiConf {
.intConf
.createWithDefault(16)

val BATCH_SUBMITTER_VIRTUAL_THREADS_ENABLED: ConfigEntry[Boolean] =
buildConf("kyuubi.batch.submitter.virtualThreads.enabled")
.internal
.audience(SERVER)
.immutable
.doc("Whether batch submitter workers use virtual threads. " +
"Requires at least Java 21; Java 25 or later is recommended. " +
s"Defaults to ${SERVER_VIRTUAL_THREADS_ENABLED.key}.")
.version("1.13.0")
.fallbackConf(SERVER_VIRTUAL_THREADS_ENABLED)

val BATCH_IMPL_VERSION: ConfigEntry[String] =
buildConf("kyuubi.batch.impl.version")
.internal
Expand Down Expand Up @@ -2191,6 +2275,18 @@ object KyuubiConf {
.intConf
.createWithDefault(100)

val SERVER_EXEC_POOL_VIRTUAL_THREADS_ENABLED: ConfigEntry[Boolean] =
buildConf("kyuubi.backend.server.exec.pool.virtualThreads.enabled")
.doc("Whether the Kyuubi server uses virtual threads for asynchronous operation " +
"execution. Requires at least Java 21; Java 25 or later is recommended. " +
"The configured concurrency limit, wait queue capacity, rejection behavior, and " +
"metrics remain enforced. " +
s"Defaults to ${SERVER_VIRTUAL_THREADS_ENABLED.key}.")
.version("1.13.0")
.audience(SERVER)
.immutable
.fallbackConf(SERVER_VIRTUAL_THREADS_ENABLED)

val ENGINE_EXEC_POOL_SIZE: ConfigEntry[Int] =
buildConf("kyuubi.backend.engine.exec.pool.size")
.doc("Number of threads in the operation execution thread pool of SQL engine applications")
Expand Down Expand Up @@ -2261,19 +2357,31 @@ object KyuubiConf {
buildConf("kyuubi.metadata.recovery.threads")
.audience(SERVER)
.immutable
.doc("The number of threads for recovery from the metadata store " +
.doc("The maximum number of concurrent tasks for recovery from the metadata store " +
"when the Kyuubi server restarts.")
.version("1.6.0")
.intConf
.createWithDefault(10)

val METADATA_RECOVERY_VIRTUAL_THREADS_ENABLED: ConfigEntry[Boolean] =
buildConf("kyuubi.metadata.recovery.virtualThreads.enabled")
.audience(SERVER)
.immutable
.doc("Whether metadata recovery workers use virtual threads. " +
"Requires at least Java 21; Java 25 or later is recommended. " +
"The configured concurrency limit remains enforced. " +
s"Defaults to ${SERVER_VIRTUAL_THREADS_ENABLED.key}.")
.version("1.13.0")
.fallbackConf(SERVER_VIRTUAL_THREADS_ENABLED)

val METADATA_RECOVERY_WAIT_ENGINE_SUBMISSION: ConfigEntry[Boolean] =
buildConf("kyuubi.metadata.recovery.waitEngineSubmission")
.audience(SERVER)
.immutable
.doc("Whether a metadata recovery task should wait for its corresponding engine " +
"submission to complete before finishing. All recovery tasks are submitted to a fixed " +
s"thread pool controlled by ${METADATA_RECOVERY_THREADS.key}. If true, a task blocks " +
"submission to complete before finishing. All recovery tasks are submitted to a " +
s"concurrency-limited executor controlled by ${METADATA_RECOVERY_THREADS.key}. If true, " +
"a task blocks " +
"until the engine submission is done, helping throttle the load on the system " +
s"if ${SESSION_ENGINE_STARTUP_WAIT_COMPLETION.key} is false. " +
"If false, the task returns immediately after opening the session without waiting.")
Expand Down Expand Up @@ -2313,6 +2421,17 @@ object KyuubiConf {
.intConf
.createWithDefault(10)

val METADATA_REQUEST_ASYNC_RETRY_VIRTUAL_THREADS_ENABLED: ConfigEntry[Boolean] =
buildConf("kyuubi.metadata.request.async.retry.virtualThreads.enabled")
.audience(SERVER)
.immutable
.doc("Whether metadata asynchronous retry workers use virtual threads. " +
"Requires at least Java 21; Java 25 or later is recommended. " +
"The configured concurrency limit remains enforced. " +
s"Defaults to ${SERVER_VIRTUAL_THREADS_ENABLED.key}.")
.version("1.13.0")
.fallbackConf(SERVER_VIRTUAL_THREADS_ENABLED)

val METADATA_REQUEST_ASYNC_RETRY_QUEUE_SIZE: ConfigEntry[Int] =
buildConf("kyuubi.metadata.request.async.retry.queue.size")
.audience(SERVER)
Expand Down Expand Up @@ -4069,6 +4188,17 @@ object KyuubiConf {
.timeConf
.createWithDefaultString("PT2M")

val FRONTEND_DATA_AGENT_OPERATION_SUBMIT_VIRTUAL_THREADS_ENABLED: ConfigEntry[Boolean] =
buildConf("kyuubi.frontend.data.agent.operation.submit.virtualThreads.enabled")
.doc("Whether blocking Data Agent operation submissions in the Kyuubi server use virtual " +
"threads. Requires at least Java 21; Java 25 or later is recommended. " +
"The concurrency and queue limits remain enforced. " +
s"Defaults to ${SERVER_VIRTUAL_THREADS_ENABLED.key}.")
.version("1.13.0")
.audience(SERVER)
.immutable
.fallbackConf(SERVER_VIRTUAL_THREADS_ENABLED)

val ENGINE_JDBC_MEMORY: ConfigEntry[String] =
buildConf("kyuubi.engine.jdbc.memory")
.doc("The heap memory for the JDBC query engine")
Expand Down Expand Up @@ -4280,4 +4410,17 @@ object KyuubiConf {
.audience(SERVER)
.immutable
.fallbackConf(HIVE_SERVER2_THRIFT_RESULTSET_DEFAULT_FETCH_SIZE)

private[kyuubi] val serverVirtualThreadConfigs: Seq[ConfigEntry[Boolean]] = Seq(
FRONTEND_THRIFT_BINARY_VIRTUAL_THREADS_ENABLED,
KUBERNETES_CLIENT_DISPATCHER_VIRTUAL_THREADS_ENABLED,
KUBERNETES_APPLICATION_CLEANUP_VIRTUAL_THREADS_ENABLED,
SERVER_ENGINE_LOG_CAPTURE_VIRTUAL_THREADS_ENABLED,
ENGINE_RPC_CLIENT_VIRTUAL_THREADS_ENABLED,
ENGINE_ALIVE_PROBE_VIRTUAL_THREADS_ENABLED,
BATCH_SUBMITTER_VIRTUAL_THREADS_ENABLED,
SERVER_EXEC_POOL_VIRTUAL_THREADS_ENABLED,
METADATA_RECOVERY_VIRTUAL_THREADS_ENABLED,
METADATA_REQUEST_ASYNC_RETRY_VIRTUAL_THREADS_ENABLED,
FRONTEND_DATA_AGENT_OPERATION_SUBMIT_VIRTUAL_THREADS_ENABLED)
}
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@ package org.apache.kyuubi.session

import java.io.IOException
import java.nio.file.{Files, Paths}
import java.util.concurrent.{ConcurrentHashMap, Future, ThreadPoolExecutor, TimeUnit}
import java.util.concurrent.{ConcurrentHashMap, ExecutorService, Future, TimeUnit}

import scala.collection.JavaConverters._
import scala.concurrent.duration.Duration
Expand Down Expand Up @@ -76,7 +76,10 @@ abstract class SessionManager(name: String) extends CompositeService(name) {

protected def isServer: Boolean

private var execPool: ThreadPoolExecutor = _
private var execPool: ExecutorService = _
private var execPoolSize: () => Int = _
private var execPoolActiveCount: () => Int = _
private var execPoolQueueSize: () => Int = _

def submitBackgroundOperation(r: Runnable): Future[_] = execPool.submit(r)

Expand Down Expand Up @@ -170,17 +173,17 @@ abstract class SessionManager(name: String) extends CompositeService(name) {

def getExecPoolSize: Int = {
assert(execPool != null)
execPool.getPoolSize
execPoolSize()
}

def getActiveCount: Int = {
assert(execPool != null)
execPool.getActiveCount
execPoolActiveCount()
}

def getWorkQueueSize: Int = {
assert(execPool != null)
execPool.getQueue.size()
execPoolQueueSize()
}

private var _confRestrictList: Set[String] = _
Expand Down Expand Up @@ -282,11 +285,28 @@ abstract class SessionManager(name: String) extends CompositeService(name) {
s"${SESSION_USER_SIGN_ENABLED.key}"
_batchConfIgnoreList = conf.get(BATCH_CONF_IGNORE_LIST)

execPool = ThreadUtils.newDaemonQueuedThreadPool(
poolSize,
waitQueueSize,
keepAliveMs,
s"$name-exec-pool")
if (isServer && conf.get(SERVER_EXEC_POOL_VIRTUAL_THREADS_ENABLED)) {
val virtualThreadPool = ThreadUtils.newBoundedQueuedVirtualThreadPerTaskExecutor(
poolSize,
waitQueueSize,
s"$name-exec-pool")
execPool = virtualThreadPool
execPoolSize = () => virtualThreadPool.getPoolSize
execPoolActiveCount = () => virtualThreadPool.getActiveCount
execPoolQueueSize = () => virtualThreadPool.getQueueSize
info(s"$name-exec-pool: concurrency limit: $poolSize, wait queue size: " +
s"$waitQueueSize, virtual threads")
} else {
val platformThreadPool = ThreadUtils.newDaemonQueuedThreadPool(
poolSize,
waitQueueSize,
keepAliveMs,
s"$name-exec-pool")
execPool = platformThreadPool
execPoolSize = () => platformThreadPool.getPoolSize
execPoolActiveCount = () => platformThreadPool.getActiveCount
execPoolQueueSize = () => platformThreadPool.getQueue.size()
}
super.initialize(conf)
}

Expand Down
Loading
Loading