diff --git a/docs/configuration/settings.md b/docs/configuration/settings.md
index e23143d15b6..f18e6ba7da0 100644
--- a/docs/configuration/settings.md
+++ b/docs/configuration/settings.md
@@ -60,21 +60,22 @@ You can configure the Kyuubi properties in `$KYUUBI_HOME/conf/kyuubi-defaults.co
### Backend
-| Key | Default | Meaning | Type | Since |
-|--------------------------------------------------|---------------------------|------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|----------|--------|
-| kyuubi.backend.engine.exec.pool.keepalive.time | PT1M | Time(ms) that an idle async thread of the operation execution thread pool will wait for a new task to arrive before terminating in SQL engine applications | duration | 1.0.0 |
-| kyuubi.backend.engine.exec.pool.shutdown.timeout | PT10S | Timeout(ms) for the operation execution thread pool to terminate in SQL engine applications | duration | 1.0.0 |
-| kyuubi.backend.engine.exec.pool.size | 100 | Number of threads in the operation execution thread pool of SQL engine applications | int | 1.0.0 |
-| kyuubi.backend.engine.exec.pool.wait.queue.size | 100 | Size of the wait queue for the operation execution thread pool in SQL engine applications | int | 1.0.0 |
-| kyuubi.backend.server.event.async.enabled | false | Whether backend server event logging is asynchronous. | boolean | 1.11.0 |
-| kyuubi.backend.server.event.json.log.path | file:///tmp/kyuubi/events | The location of server events go for the built-in JSON logger | string | 1.4.0 |
-| kyuubi.backend.server.event.kafka.close.timeout | PT5S | Period to wait for Kafka producer of server event handlers to close. | duration | 1.8.0 |
-| kyuubi.backend.server.event.kafka.topic | <undefined> | The topic of server events go for the built-in Kafka logger | string | 1.8.0 |
-| kyuubi.backend.server.event.loggers || A comma-separated list of server history loggers, where session/operation etc events go.
- JSON: the events will be written to the location of kyuubi.backend.server.event.json.log.path
- KAFKA: the events will be serialized in JSON format and sent to topic of `kyuubi.backend.server.event.kafka.topic`. Note: For the configs of Kafka producer, please specify them with the prefix: `kyuubi.backend.server.event.kafka.`. For example, `kyuubi.backend.server.event.kafka.bootstrap.servers=127.0.0.1:9092`
- JDBC: to be done
- CUSTOM: User-defined event handlers.
Note that: Kyuubi supports custom event handlers with the Java SPI. To register a custom event handler, the user needs to implement a class which is a child of org.apache.kyuubi.events.handler.CustomEventHandlerProvider which has a zero-arg constructor. | seq | 1.4.0 |
-| kyuubi.backend.server.exec.pool.keepalive.time | PT1M | Time(ms) that an idle async thread of the operation execution thread pool will wait for a new task to arrive before terminating in Kyuubi server | duration | 1.0.0 |
-| kyuubi.backend.server.exec.pool.shutdown.timeout | PT10S | Timeout(ms) for the operation execution thread pool to terminate in Kyuubi server | duration | 1.0.0 |
-| kyuubi.backend.server.exec.pool.size | 100 | Number of threads in the operation execution thread pool of Kyuubi server | int | 1.0.0 |
-| kyuubi.backend.server.exec.pool.wait.queue.size | 100 | Size of the wait queue for the operation execution thread pool of Kyuubi server | int | 1.0.0 |
+| Key | Default | Meaning | Type | Since |
+|--------------------------------------------------------|---------------------------|------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|----------|--------|
+| kyuubi.backend.engine.exec.pool.keepalive.time | PT1M | Time(ms) that an idle async thread of the operation execution thread pool will wait for a new task to arrive before terminating in SQL engine applications | duration | 1.0.0 |
+| kyuubi.backend.engine.exec.pool.shutdown.timeout | PT10S | Timeout(ms) for the operation execution thread pool to terminate in SQL engine applications | duration | 1.0.0 |
+| kyuubi.backend.engine.exec.pool.size | 100 | Number of threads in the operation execution thread pool of SQL engine applications | int | 1.0.0 |
+| kyuubi.backend.engine.exec.pool.wait.queue.size | 100 | Size of the wait queue for the operation execution thread pool in SQL engine applications | int | 1.0.0 |
+| kyuubi.backend.server.event.async.enabled | false | Whether backend server event logging is asynchronous. | boolean | 1.11.0 |
+| kyuubi.backend.server.event.json.log.path | file:///tmp/kyuubi/events | The location of server events go for the built-in JSON logger | string | 1.4.0 |
+| kyuubi.backend.server.event.kafka.close.timeout | PT5S | Period to wait for Kafka producer of server event handlers to close. | duration | 1.8.0 |
+| kyuubi.backend.server.event.kafka.topic | <undefined> | The topic of server events go for the built-in Kafka logger | string | 1.8.0 |
+| kyuubi.backend.server.event.loggers || A comma-separated list of server history loggers, where session/operation etc events go. - JSON: the events will be written to the location of kyuubi.backend.server.event.json.log.path
- KAFKA: the events will be serialized in JSON format and sent to topic of `kyuubi.backend.server.event.kafka.topic`. Note: For the configs of Kafka producer, please specify them with the prefix: `kyuubi.backend.server.event.kafka.`. For example, `kyuubi.backend.server.event.kafka.bootstrap.servers=127.0.0.1:9092`
- JDBC: to be done
- CUSTOM: User-defined event handlers.
Note that: Kyuubi supports custom event handlers with the Java SPI. To register a custom event handler, the user needs to implement a class which is a child of org.apache.kyuubi.events.handler.CustomEventHandlerProvider which has a zero-arg constructor. | seq | 1.4.0 |
+| kyuubi.backend.server.exec.pool.keepalive.time | PT1M | Time(ms) that an idle async thread of the operation execution thread pool will wait for a new task to arrive before terminating in Kyuubi server | duration | 1.0.0 |
+| kyuubi.backend.server.exec.pool.shutdown.timeout | PT10S | Timeout(ms) for the operation execution thread pool to terminate in Kyuubi server | duration | 1.0.0 |
+| kyuubi.backend.server.exec.pool.size | 100 | Number of threads in the operation execution thread pool of Kyuubi server | int | 1.0.0 |
+| kyuubi.backend.server.exec.pool.virtualThreads.enabled | false | 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. Defaults to kyuubi.server.virtualThreads.enabled. | boolean | 1.13.0 |
+| kyuubi.backend.server.exec.pool.wait.queue.size | 100 | Size of the wait queue for the operation execution thread pool of Kyuubi server | int | 1.0.0 |
### Batch
@@ -241,71 +242,72 @@ You can configure the Kyuubi properties in `$KYUUBI_HOME/conf/kyuubi-defaults.co
### Frontend
-| Key | Default | Meaning | Type | Since |
-|------------------------------------------------------------|--------------------|-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|----------|--------|
-| kyuubi.frontend.advertised.host | <undefined> | Hostname or IP of the Kyuubi server's frontend services to publish to external systems such as the service discovery ensemble and metadata store. Use it when you want to advertise a different hostname or IP than the bind host. | string | 1.8.0 |
-| kyuubi.frontend.bind.host | <undefined> | Hostname or IP of the machine on which to run the frontend services. | string | 1.0.0 |
-| kyuubi.frontend.bind.port | 10009 | (deprecated) Port of the machine on which to run the thrift frontend service via the binary protocol. | int | 1.0.0 |
-| kyuubi.frontend.connection.url.use.hostname | true | When true, frontend services prefer hostname, otherwise, ip address. Note that, the default value is set to `false` when engine running on Kubernetes to prevent potential network issues. | boolean | 1.5.0 |
-| kyuubi.frontend.data.agent.operation.timeout | PT2M | Timeout for waiting on data agent engine launch and operation start in the REST frontend. | duration | 1.12.0 |
-| kyuubi.frontend.jetty.sendVersion.enabled | true | Whether to send Jetty version in HTTP response. | boolean | 1.9.3 |
-| kyuubi.frontend.max.message.size | 104857600 | (deprecated) Maximum message size in bytes a Kyuubi server will accept. | int | 1.0.0 |
-| kyuubi.frontend.max.worker.threads | 999 | (deprecated) Maximum number of threads in the frontend worker thread pool for the thrift frontend service | int | 1.0.0 |
-| kyuubi.frontend.min.worker.threads | 9 | (deprecated) Minimum number of threads in the frontend worker thread pool for the thrift frontend service | int | 1.0.0 |
-| kyuubi.frontend.protocols | THRIFT_BINARY,REST | A comma-separated list for all frontend protocols - THRIFT_BINARY - HiveServer2 compatible thrift binary protocol.
- THRIFT_HTTP - HiveServer2 compatible thrift http protocol.
- REST - Kyuubi defined REST API(experimental).
- TRINO - Trino compatible http protocol(experimental).
| seq | 1.4.0 |
-| kyuubi.frontend.proxy.http.client.ip.header | X-Real-IP | The HTTP header to record the real client IP address. If your server is behind a load balancer or other proxy, the server will see this load balancer or proxy IP address as the client IP address, to get around this common issue, most load balancers or proxies offer the ability to record the real remote IP address in an HTTP header that will be added to the request for other devices to use. Note that, because the header value can be specified to any IP address, so it will not be used for authentication. | string | 1.6.0 |
-| kyuubi.frontend.rest.bind.host | <undefined> | Hostname or IP of the machine on which to run the REST frontend service. | string | 1.4.0 |
-| kyuubi.frontend.rest.bind.port | 10099 | Port of the machine on which to run the REST frontend service. | int | 1.4.0 |
-| kyuubi.frontend.rest.engine.ui.proxy.enabled | false | Whether to route Engine UI traffic via Kyuubi REST frontend proxy. When disabled, the Web UI links directly to the Kyuubi engine URL. | boolean | 1.12.0 |
-| kyuubi.frontend.rest.engine.ui.proxy.hosts || A comma-separated list of hosts that Engine UI proxy requests can route to when `kyuubi.frontend.rest.engine.ui.proxy.enabled` is enabled. Host matching is case-insensitive and supports `*` as a wildcard, for example, `*.example.com`. If empty, all Engine UI proxy requests are denied. | seq | 1.12.0 |
-| kyuubi.frontend.rest.jetty.stopTimeout | PT5S | Stop timeout for Jetty server used by the RESTful frontend service. | duration | 1.8.1 |
-| kyuubi.frontend.rest.legacy.v1.sessionsReturnAllUsers | false | When true, GET /api/v1/sessions returns all sessions on the server regardless of the calling user (legacy behavior). When false (default), only sessions owned by the authenticated user are returned. This flag is provided for backward compatibility and will be removed in a future release. | boolean | 1.12.0 |
-| kyuubi.frontend.rest.max.worker.threads | 999 | Maximum number of threads in the frontend worker thread pool for the rest frontend service | int | 1.6.2 |
-| kyuubi.frontend.rest.proxy.jetty.client.idleTimeout | PT30S | The idle timeout in milliseconds for Jetty server used by the RESTful frontend service. | duration | 1.10.0 |
-| kyuubi.frontend.rest.proxy.jetty.client.maxConnections | 32768 | The max number of connections per destination for Jetty server used by the RESTful frontend service. | int | 1.10.0 |
-| kyuubi.frontend.rest.proxy.jetty.client.maxThreads | 256 | The max number of threads of HttpClient's Executor for Jetty server used by the RESTful frontend service. | int | 1.10.0 |
-| kyuubi.frontend.rest.proxy.jetty.client.requestBufferSize | 4096 | Size of the buffer in bytes used to write requests for Jetty server used by the RESTful frontend service. | int | 1.10.0 |
-| kyuubi.frontend.rest.proxy.jetty.client.responseBufferSize | 4096 | Size of the buffer in bytes used to read response for Jetty server used by the RESTful frontend service. | int | 1.10.0 |
-| kyuubi.frontend.rest.proxy.jetty.client.timeout | PT60S | The total timeout in milliseconds for Jetty server used by the RESTful frontend service. | duration | 1.10.0 |
-| kyuubi.frontend.rest.ui.enabled | true | Whether to enable Web UI when RESTful protocol is enabled | boolean | 1.10.0 |
-| kyuubi.frontend.ssl.keystore.algorithm | <undefined> | SSL certificate keystore algorithm. | string | 1.7.0 |
-| kyuubi.frontend.ssl.keystore.password | <undefined> | SSL certificate keystore password. | string | 1.7.0 |
-| kyuubi.frontend.ssl.keystore.path | <undefined> | SSL certificate keystore location. | string | 1.7.0 |
-| kyuubi.frontend.ssl.keystore.type | <undefined> | SSL certificate keystore type. | string | 1.7.0 |
-| kyuubi.frontend.thrift.binary.bind.host | <undefined> | Hostname or IP of the machine on which to run the thrift frontend service via the binary protocol. | string | 1.4.0 |
-| kyuubi.frontend.thrift.binary.bind.port | 10009 | Port of the machine on which to run the thrift frontend service via the binary protocol. | int | 1.4.0 |
-| kyuubi.frontend.thrift.binary.ssl.disallowed.protocols | SSLv2,SSLv3 | SSL versions to disallow for Kyuubi thrift binary frontend. | set | 1.7.0 |
-| kyuubi.frontend.thrift.binary.ssl.enabled | false | Set this to true for using SSL encryption in thrift binary frontend server. | boolean | 1.7.0 |
-| kyuubi.frontend.thrift.binary.ssl.include.ciphersuites || A comma-separated list of include SSL cipher suite names for thrift binary frontend. | seq | 1.7.0 |
-| kyuubi.frontend.thrift.binary.virtualThreads.enabled | false | 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 remains limited by kyuubi.frontend.thrift.max.worker.threads. The minimum worker threads and worker keepalive configurations do not apply in this mode. | boolean | 1.13.0 |
-| kyuubi.frontend.thrift.client.max.message.size | 1073741824 | Maximum message size in bytes a thrift client will receive. | int | 1.9.3 |
-| kyuubi.frontend.thrift.http.bind.host | <undefined> | Hostname or IP of the machine on which to run the thrift frontend service via http protocol. | string | 1.6.0 |
-| kyuubi.frontend.thrift.http.bind.port | 10010 | Port of the machine on which to run the thrift frontend service via http protocol. | int | 1.6.0 |
-| kyuubi.frontend.thrift.http.compression.enabled | true | Enable thrift http compression via Jetty compression support | boolean | 1.6.0 |
-| kyuubi.frontend.thrift.http.cookie.auth.enabled | true | When true, Kyuubi in HTTP transport mode, will use cookie-based authentication mechanism | boolean | 1.6.0 |
-| kyuubi.frontend.thrift.http.cookie.domain | <undefined> | Domain for the Kyuubi generated cookies | string | 1.6.0 |
-| kyuubi.frontend.thrift.http.cookie.is.httponly | true | HttpOnly attribute of the Kyuubi generated cookie. | boolean | 1.6.0 |
-| kyuubi.frontend.thrift.http.cookie.max.age | 86400 | Maximum age in seconds for server side cookie used by Kyuubi in HTTP mode. | int | 1.6.0 |
-| kyuubi.frontend.thrift.http.cookie.path | <undefined> | Path for the Kyuubi generated cookies | string | 1.6.0 |
-| kyuubi.frontend.thrift.http.max.idle.time | PT30M | Maximum idle time for a connection on the server when in HTTP mode. | duration | 1.6.0 |
-| kyuubi.frontend.thrift.http.path | cliservice | Path component of URL endpoint when in HTTP mode. | string | 1.6.0 |
-| kyuubi.frontend.thrift.http.request.header.size | 6144 | Request header size in bytes, when using HTTP transport mode. Jetty defaults used. | int | 1.6.0 |
-| kyuubi.frontend.thrift.http.response.header.size | 6144 | Response header size in bytes, when using HTTP transport mode. Jetty defaults used. | int | 1.6.0 |
-| kyuubi.frontend.thrift.http.ssl.exclude.ciphersuites || A comma-separated list of exclude SSL cipher suite names for thrift http frontend. | seq | 1.7.0 |
-| kyuubi.frontend.thrift.http.ssl.keystore.password | <undefined> | SSL certificate keystore password. | string | 1.6.0 |
-| kyuubi.frontend.thrift.http.ssl.keystore.path | <undefined> | SSL certificate keystore location. | string | 1.6.0 |
-| kyuubi.frontend.thrift.http.ssl.protocol.blacklist | SSLv2,SSLv3 | SSL Versions to disable when using HTTP transport mode. | seq | 1.6.0 |
-| kyuubi.frontend.thrift.http.use.SSL | false | Set this to true for using SSL encryption in http mode. | boolean | 1.6.0 |
-| kyuubi.frontend.thrift.http.xsrf.filter.enabled | false | If enabled, Kyuubi will block any requests made to it over HTTP if an X-XSRF-HEADER header is not present | boolean | 1.6.0 |
-| kyuubi.frontend.thrift.max.message.size | 104857600 | Maximum message size in bytes a Kyuubi server will accept. | int | 1.4.0 |
-| kyuubi.frontend.thrift.max.worker.threads | 999 | Maximum number of threads in the frontend worker thread pool for the thrift frontend service | int | 1.4.0 |
-| kyuubi.frontend.thrift.min.worker.threads | 9 | Minimum number of threads in the frontend worker thread pool for the thrift frontend service | int | 1.4.0 |
-| kyuubi.frontend.thrift.worker.keepalive.time | PT1M | Keep-alive time (in milliseconds) for an idle worker thread | duration | 1.4.0 |
-| kyuubi.frontend.trino.bind.host | <undefined> | Hostname or IP of the machine on which to run the TRINO frontend service. | string | 1.7.0 |
-| kyuubi.frontend.trino.bind.port | 10999 | Port of the machine on which to run the TRINO frontend service. | int | 1.7.0 |
-| kyuubi.frontend.trino.jetty.stopTimeout | PT5S | Stop timeout for Jetty server used by the Trino frontend service. | duration | 1.8.1 |
-| kyuubi.frontend.trino.max.worker.threads | 999 | Maximum number of threads in the frontend worker thread pool for the Trino frontend service | int | 1.7.0 |
-| kyuubi.frontend.worker.keepalive.time | PT1M | (deprecated) Keep-alive time (in milliseconds) for an idle worker thread | duration | 1.0.0 |
+| Key | Default | Meaning | Type | Since |
+|--------------------------------------------------------------------|--------------------|-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|----------|--------|
+| kyuubi.frontend.advertised.host | <undefined> | Hostname or IP of the Kyuubi server's frontend services to publish to external systems such as the service discovery ensemble and metadata store. Use it when you want to advertise a different hostname or IP than the bind host. | string | 1.8.0 |
+| kyuubi.frontend.bind.host | <undefined> | Hostname or IP of the machine on which to run the frontend services. | string | 1.0.0 |
+| kyuubi.frontend.bind.port | 10009 | (deprecated) Port of the machine on which to run the thrift frontend service via the binary protocol. | int | 1.0.0 |
+| kyuubi.frontend.connection.url.use.hostname | true | When true, frontend services prefer hostname, otherwise, ip address. Note that, the default value is set to `false` when engine running on Kubernetes to prevent potential network issues. | boolean | 1.5.0 |
+| kyuubi.frontend.data.agent.operation.submit.virtualThreads.enabled | false | 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. Defaults to kyuubi.server.virtualThreads.enabled. | boolean | 1.13.0 |
+| kyuubi.frontend.data.agent.operation.timeout | PT2M | Timeout for waiting on data agent engine launch and operation start in the REST frontend. | duration | 1.12.0 |
+| kyuubi.frontend.jetty.sendVersion.enabled | true | Whether to send Jetty version in HTTP response. | boolean | 1.9.3 |
+| kyuubi.frontend.max.message.size | 104857600 | (deprecated) Maximum message size in bytes a Kyuubi server will accept. | int | 1.0.0 |
+| kyuubi.frontend.max.worker.threads | 999 | (deprecated) Maximum number of threads in the frontend worker thread pool for the thrift frontend service | int | 1.0.0 |
+| kyuubi.frontend.min.worker.threads | 9 | (deprecated) Minimum number of threads in the frontend worker thread pool for the thrift frontend service | int | 1.0.0 |
+| kyuubi.frontend.protocols | THRIFT_BINARY,REST | A comma-separated list for all frontend protocols - THRIFT_BINARY - HiveServer2 compatible thrift binary protocol.
- THRIFT_HTTP - HiveServer2 compatible thrift http protocol.
- REST - Kyuubi defined REST API(experimental).
- TRINO - Trino compatible http protocol(experimental).
| seq | 1.4.0 |
+| kyuubi.frontend.proxy.http.client.ip.header | X-Real-IP | The HTTP header to record the real client IP address. If your server is behind a load balancer or other proxy, the server will see this load balancer or proxy IP address as the client IP address, to get around this common issue, most load balancers or proxies offer the ability to record the real remote IP address in an HTTP header that will be added to the request for other devices to use. Note that, because the header value can be specified to any IP address, so it will not be used for authentication. | string | 1.6.0 |
+| kyuubi.frontend.rest.bind.host | <undefined> | Hostname or IP of the machine on which to run the REST frontend service. | string | 1.4.0 |
+| kyuubi.frontend.rest.bind.port | 10099 | Port of the machine on which to run the REST frontend service. | int | 1.4.0 |
+| kyuubi.frontend.rest.engine.ui.proxy.enabled | false | Whether to route Engine UI traffic via Kyuubi REST frontend proxy. When disabled, the Web UI links directly to the Kyuubi engine URL. | boolean | 1.12.0 |
+| kyuubi.frontend.rest.engine.ui.proxy.hosts || A comma-separated list of hosts that Engine UI proxy requests can route to when `kyuubi.frontend.rest.engine.ui.proxy.enabled` is enabled. Host matching is case-insensitive and supports `*` as a wildcard, for example, `*.example.com`. If empty, all Engine UI proxy requests are denied. | seq | 1.12.0 |
+| kyuubi.frontend.rest.jetty.stopTimeout | PT5S | Stop timeout for Jetty server used by the RESTful frontend service. | duration | 1.8.1 |
+| kyuubi.frontend.rest.legacy.v1.sessionsReturnAllUsers | false | When true, GET /api/v1/sessions returns all sessions on the server regardless of the calling user (legacy behavior). When false (default), only sessions owned by the authenticated user are returned. This flag is provided for backward compatibility and will be removed in a future release. | boolean | 1.12.0 |
+| kyuubi.frontend.rest.max.worker.threads | 999 | Maximum number of threads in the frontend worker thread pool for the rest frontend service | int | 1.6.2 |
+| kyuubi.frontend.rest.proxy.jetty.client.idleTimeout | PT30S | The idle timeout in milliseconds for Jetty server used by the RESTful frontend service. | duration | 1.10.0 |
+| kyuubi.frontend.rest.proxy.jetty.client.maxConnections | 32768 | The max number of connections per destination for Jetty server used by the RESTful frontend service. | int | 1.10.0 |
+| kyuubi.frontend.rest.proxy.jetty.client.maxThreads | 256 | The max number of threads of HttpClient's Executor for Jetty server used by the RESTful frontend service. | int | 1.10.0 |
+| kyuubi.frontend.rest.proxy.jetty.client.requestBufferSize | 4096 | Size of the buffer in bytes used to write requests for Jetty server used by the RESTful frontend service. | int | 1.10.0 |
+| kyuubi.frontend.rest.proxy.jetty.client.responseBufferSize | 4096 | Size of the buffer in bytes used to read response for Jetty server used by the RESTful frontend service. | int | 1.10.0 |
+| kyuubi.frontend.rest.proxy.jetty.client.timeout | PT60S | The total timeout in milliseconds for Jetty server used by the RESTful frontend service. | duration | 1.10.0 |
+| kyuubi.frontend.rest.ui.enabled | true | Whether to enable Web UI when RESTful protocol is enabled | boolean | 1.10.0 |
+| kyuubi.frontend.ssl.keystore.algorithm | <undefined> | SSL certificate keystore algorithm. | string | 1.7.0 |
+| kyuubi.frontend.ssl.keystore.password | <undefined> | SSL certificate keystore password. | string | 1.7.0 |
+| kyuubi.frontend.ssl.keystore.path | <undefined> | SSL certificate keystore location. | string | 1.7.0 |
+| kyuubi.frontend.ssl.keystore.type | <undefined> | SSL certificate keystore type. | string | 1.7.0 |
+| kyuubi.frontend.thrift.binary.bind.host | <undefined> | Hostname or IP of the machine on which to run the thrift frontend service via the binary protocol. | string | 1.4.0 |
+| kyuubi.frontend.thrift.binary.bind.port | 10009 | Port of the machine on which to run the thrift frontend service via the binary protocol. | int | 1.4.0 |
+| kyuubi.frontend.thrift.binary.ssl.disallowed.protocols | SSLv2,SSLv3 | SSL versions to disallow for Kyuubi thrift binary frontend. | set | 1.7.0 |
+| kyuubi.frontend.thrift.binary.ssl.enabled | false | Set this to true for using SSL encryption in thrift binary frontend server. | boolean | 1.7.0 |
+| kyuubi.frontend.thrift.binary.ssl.include.ciphersuites || A comma-separated list of include SSL cipher suite names for thrift binary frontend. | seq | 1.7.0 |
+| kyuubi.frontend.thrift.binary.virtualThreads.enabled | false | Whether to use virtual threads for the Kyuubi server thrift binary frontend 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. Defaults to kyuubi.server.virtualThreads.enabled. | boolean | 1.13.0 |
+| kyuubi.frontend.thrift.client.max.message.size | 1073741824 | Maximum message size in bytes a thrift client will receive. | int | 1.9.3 |
+| kyuubi.frontend.thrift.http.bind.host | <undefined> | Hostname or IP of the machine on which to run the thrift frontend service via http protocol. | string | 1.6.0 |
+| kyuubi.frontend.thrift.http.bind.port | 10010 | Port of the machine on which to run the thrift frontend service via http protocol. | int | 1.6.0 |
+| kyuubi.frontend.thrift.http.compression.enabled | true | Enable thrift http compression via Jetty compression support | boolean | 1.6.0 |
+| kyuubi.frontend.thrift.http.cookie.auth.enabled | true | When true, Kyuubi in HTTP transport mode, will use cookie-based authentication mechanism | boolean | 1.6.0 |
+| kyuubi.frontend.thrift.http.cookie.domain | <undefined> | Domain for the Kyuubi generated cookies | string | 1.6.0 |
+| kyuubi.frontend.thrift.http.cookie.is.httponly | true | HttpOnly attribute of the Kyuubi generated cookie. | boolean | 1.6.0 |
+| kyuubi.frontend.thrift.http.cookie.max.age | 86400 | Maximum age in seconds for server side cookie used by Kyuubi in HTTP mode. | int | 1.6.0 |
+| kyuubi.frontend.thrift.http.cookie.path | <undefined> | Path for the Kyuubi generated cookies | string | 1.6.0 |
+| kyuubi.frontend.thrift.http.max.idle.time | PT30M | Maximum idle time for a connection on the server when in HTTP mode. | duration | 1.6.0 |
+| kyuubi.frontend.thrift.http.path | cliservice | Path component of URL endpoint when in HTTP mode. | string | 1.6.0 |
+| kyuubi.frontend.thrift.http.request.header.size | 6144 | Request header size in bytes, when using HTTP transport mode. Jetty defaults used. | int | 1.6.0 |
+| kyuubi.frontend.thrift.http.response.header.size | 6144 | Response header size in bytes, when using HTTP transport mode. Jetty defaults used. | int | 1.6.0 |
+| kyuubi.frontend.thrift.http.ssl.exclude.ciphersuites || A comma-separated list of exclude SSL cipher suite names for thrift http frontend. | seq | 1.7.0 |
+| kyuubi.frontend.thrift.http.ssl.keystore.password | <undefined> | SSL certificate keystore password. | string | 1.6.0 |
+| kyuubi.frontend.thrift.http.ssl.keystore.path | <undefined> | SSL certificate keystore location. | string | 1.6.0 |
+| kyuubi.frontend.thrift.http.ssl.protocol.blacklist | SSLv2,SSLv3 | SSL Versions to disable when using HTTP transport mode. | seq | 1.6.0 |
+| kyuubi.frontend.thrift.http.use.SSL | false | Set this to true for using SSL encryption in http mode. | boolean | 1.6.0 |
+| kyuubi.frontend.thrift.http.xsrf.filter.enabled | false | If enabled, Kyuubi will block any requests made to it over HTTP if an X-XSRF-HEADER header is not present | boolean | 1.6.0 |
+| kyuubi.frontend.thrift.max.message.size | 104857600 | Maximum message size in bytes a Kyuubi server will accept. | int | 1.4.0 |
+| kyuubi.frontend.thrift.max.worker.threads | 999 | Maximum number of threads in the frontend worker thread pool for the thrift frontend service | int | 1.4.0 |
+| kyuubi.frontend.thrift.min.worker.threads | 9 | Minimum number of threads in the frontend worker thread pool for the thrift frontend service | int | 1.4.0 |
+| kyuubi.frontend.thrift.worker.keepalive.time | PT1M | Keep-alive time (in milliseconds) for an idle worker thread | duration | 1.4.0 |
+| kyuubi.frontend.trino.bind.host | <undefined> | Hostname or IP of the machine on which to run the TRINO frontend service. | string | 1.7.0 |
+| kyuubi.frontend.trino.bind.port | 10999 | Port of the machine on which to run the TRINO frontend service. | int | 1.7.0 |
+| kyuubi.frontend.trino.jetty.stopTimeout | PT5S | Stop timeout for Jetty server used by the Trino frontend service. | duration | 1.8.1 |
+| kyuubi.frontend.trino.max.worker.threads | 999 | Maximum number of threads in the frontend worker thread pool for the Trino frontend service | int | 1.7.0 |
+| kyuubi.frontend.worker.keepalive.time | PT1M | (deprecated) Keep-alive time (in milliseconds) for an idle worker thread | duration | 1.0.0 |
### Ha
@@ -365,6 +367,7 @@ You can configure the Kyuubi properties in `$KYUUBI_HOME/conf/kyuubi-defaults.co
| Key | Default | Meaning | Type | Since |
|----------------------------------------------------------------------|----------------------------------------------------------------------------|--------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|----------|--------|
+| kyuubi.kubernetes.application.cleanup.virtualThreads.enabled | false | Whether asynchronous Kubernetes application cleanup tasks in the Kyuubi server use virtual threads. Requires at least Java 21; Java 25 or later is recommended. Defaults to kyuubi.server.virtualThreads.enabled. | boolean | 1.13.0 |
| kyuubi.kubernetes.application.state.container | spark-kubernetes-driver | The container name to retrieve the application state from. | string | 1.8.1 |
| kyuubi.kubernetes.application.state.source | POD | The source to retrieve the application state from. The valid values are pod and container. When the pod is in a terminated state, the container state will be ignored, and the application state will be determined based on the pod state. If the source is container and there is container inside the pod with the name of kyuubi.kubernetes.application.state.container, the application state will be from the matched container state. Otherwise, the application state will be from the pod state. | string | 1.8.1 |
| kyuubi.kubernetes.authenticate.caCertFile | <undefined> | Path to the CA cert file for connecting to the Kubernetes API server over TLS from the kyuubi. Specify this as a path as opposed to a URI (i.e. do not provide a scheme) | string | 1.7.0 |
@@ -372,6 +375,7 @@ You can configure the Kyuubi properties in `$KYUUBI_HOME/conf/kyuubi-defaults.co
| kyuubi.kubernetes.authenticate.clientKeyFile | <undefined> | Path to the client key file for connecting to the Kubernetes API server over TLS from the kyuubi. Specify this as a path as opposed to a URI (i.e. do not provide a scheme) | string | 1.7.0 |
| kyuubi.kubernetes.authenticate.oauthToken | <undefined> | The OAuth token to use when authenticating against the Kubernetes API server. Note that unlike, the other authentication options, this must be the exact string value of the token to use for the authentication. | string | 1.7.0 |
| kyuubi.kubernetes.authenticate.oauthTokenFile | <undefined> | Path to the file containing the OAuth token to use when authenticating against the Kubernetes API server. Specify this as a path as opposed to a URI (i.e. do not provide a scheme) | string | 1.7.0 |
+| kyuubi.kubernetes.client.dispatcher.virtualThreads.enabled | false | Whether the Kyuubi server Kubernetes HTTP client dispatcher uses virtual threads. Requires at least Java 21; Java 25 or later is recommended. Defaults to kyuubi.server.virtualThreads.enabled. | boolean | 1.13.0 |
| kyuubi.kubernetes.client.initialize.list || The kubernetes client initialize list to register kubernetes resource informers 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. | seq | 1.11.0 |
| kyuubi.kubernetes.context | <undefined> | The desired context from your kubernetes config file used to configure the K8s client for interacting with the cluster. | string | 1.6.0 |
| kyuubi.kubernetes.context.allow.list || The allowed kubernetes context list, if it is empty, there is no kubernetes context limitation. | set | 1.8.0 |
@@ -395,27 +399,29 @@ You can configure the Kyuubi properties in `$KYUUBI_HOME/conf/kyuubi-defaults.co
### Metadata
-| Key | Default | Meaning | Type | Since |
-|-------------------------------------------------|----------------------------------------------------------|--------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|----------|--------|
-| kyuubi.metadata.cleaner.batch.size | 1000 | The batch size for cleaning expired metadata. This is used to avoid holding the delete lock for a long time when there are too many expired metadata to be cleaned. | int | 1.11.0 |
-| kyuubi.metadata.cleaner.enabled | true | Whether to clean the metadata periodically. If it is enabled, Kyuubi will clean the metadata that is in the terminate state with max age limitation. | boolean | 1.6.0 |
-| kyuubi.metadata.cleaner.interval | PT30M | The interval to check and clean expired metadata. | duration | 1.6.0 |
-| kyuubi.metadata.max.age | PT72H | The maximum age of metadata, the metadata exceeding the age will be cleaned. | duration | 1.6.0 |
-| kyuubi.metadata.recovery.threads | 10 | The number of threads for recovery from the metadata store when the Kyuubi server restarts. | int | 1.6.0 |
-| kyuubi.metadata.recovery.waitEngineSubmission | false | Whether a metadata recovery task should wait for its corresponding engine submission to complete before finishing. All recovery tasks are submitted to a fixed thread pool controlled by kyuubi.metadata.recovery.threads. If true, a task blocks until the engine submission is done, helping throttle the load on the system if kyuubi.session.engine.startup.waitCompletion is false. If false, the task returns immediately after opening the session without waiting. | boolean | 1.10.3 |
-| kyuubi.metadata.request.async.retry.enabled | true | Whether to retry in async when metadata request failed. When true, return success response immediately even the metadata request failed, and schedule it in background until success, to tolerate long-time metadata store outages w/o blocking the submission request. | boolean | 1.7.0 |
-| kyuubi.metadata.request.async.retry.queue.size | 65536 | The maximum queue size for buffering metadata requests in memory when the external metadata storage is down. Requests will be dropped if the queue exceeds. Only take affect when kyuubi.metadata.request.async.retry.enabled is `true`. | int | 1.6.0 |
-| kyuubi.metadata.request.async.retry.threads | 10 | Number of threads in the metadata request async retry manager thread pool. Only take affect when kyuubi.metadata.request.async.retry.enabled is `true`. | int | 1.6.0 |
-| kyuubi.metadata.request.retry.interval | PT5S | The interval to check and trigger the metadata request retry tasks. | duration | 1.6.0 |
-| kyuubi.metadata.search.window | <undefined> | The time window to restrict user queries to metadata within a specific period, starting from the current time to the past. It only affects `GET /api/v1/batches` API. You may want to set this to short period to improve query performance and reduce load on the metadata store when administer want to reserve the metadata for long time. The side-effects is that, the metadata created outside the window will not be invisible to users. If it is undefined, all metadata will be visible for users. | duration | 1.10.1 |
-| kyuubi.metadata.store.class | org.apache.kyuubi.server.metadata.jdbc.JDBCMetadataStore | Fully qualified class name for server metadata store. | string | 1.6.0 |
-| kyuubi.metadata.store.jdbc.database.schema.init | true | Whether to init the JDBC metadata store database schema. | boolean | 1.6.0 |
-| kyuubi.metadata.store.jdbc.database.type | SQLITE | The database type for server jdbc metadata store. - SQLITE: SQLite3, JDBC driver `org.sqlite.JDBC`.
- MYSQL: MySQL, JDBC driver `com.mysql.cj.jdbc.Driver` (fallback `com.mysql.jdbc.Driver`).
- POSTGRESQL: PostgreSQL, JDBC driver `org.postgresql.Driver`.
- CUSTOM: User-defined database type, need to specify corresponding JDBC driver.
Note that: The JDBC datasource is powered by HiKariCP, for datasource properties, please specify them with the prefix: kyuubi.metadata.store.jdbc.datasource. For example, kyuubi.metadata.store.jdbc.datasource.connectionTimeout=10000. | string | 1.6.0 |
-| kyuubi.metadata.store.jdbc.driver | <undefined> | JDBC driver class name for server jdbc metadata store. | string | 1.6.0 |
-| kyuubi.metadata.store.jdbc.password || The password for server JDBC metadata store. | string | 1.6.0 |
-| kyuubi.metadata.store.jdbc.priority.enabled | false | Whether to enable the priority scheduling for batch impl v2. When false, ignore kyuubi.batch.priority and use the FIFO ordering strategy for batch job scheduling. Note: this feature may cause significant performance issues when using MySQL 5.7 as the metastore backend due to the lack of support for mixed order index. See more details at KYUUBI #5329. | boolean | 1.8.0 |
-| kyuubi.metadata.store.jdbc.url | jdbc:sqlite:{{KYUUBI_HOME}}/kyuubi_state_store.db | The JDBC url for server JDBC metadata store. By default, it is a SQLite database url, and the state information is not shared across Kyuubi instances. To enable high availability for multiple kyuubi instances, please specify a production JDBC url. Note: this value support the variables substitution: `{{KYUUBI_HOME}}`, `{{KYUUBI_WORK_DIR_ROOT}}`. | string | 1.6.0 |
-| kyuubi.metadata.store.jdbc.user || The username for server JDBC metadata store. | string | 1.6.0 |
+| Key | Default | Meaning | Type | Since |
+|------------------------------------------------------------|----------------------------------------------------------|--------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|----------|--------|
+| kyuubi.metadata.cleaner.batch.size | 1000 | The batch size for cleaning expired metadata. This is used to avoid holding the delete lock for a long time when there are too many expired metadata to be cleaned. | int | 1.11.0 |
+| kyuubi.metadata.cleaner.enabled | true | Whether to clean the metadata periodically. If it is enabled, Kyuubi will clean the metadata that is in the terminate state with max age limitation. | boolean | 1.6.0 |
+| kyuubi.metadata.cleaner.interval | PT30M | The interval to check and clean expired metadata. | duration | 1.6.0 |
+| kyuubi.metadata.max.age | PT72H | The maximum age of metadata, the metadata exceeding the age will be cleaned. | duration | 1.6.0 |
+| kyuubi.metadata.recovery.threads | 10 | The maximum number of concurrent tasks for recovery from the metadata store when the Kyuubi server restarts. | int | 1.6.0 |
+| kyuubi.metadata.recovery.virtualThreads.enabled | false | Whether metadata recovery workers use virtual threads. Requires at least Java 21; Java 25 or later is recommended. The configured concurrency limit remains enforced. Defaults to kyuubi.server.virtualThreads.enabled. | boolean | 1.13.0 |
+| kyuubi.metadata.recovery.waitEngineSubmission | false | Whether a metadata recovery task should wait for its corresponding engine submission to complete before finishing. All recovery tasks are submitted to a concurrency-limited executor controlled by kyuubi.metadata.recovery.threads. If true, a task blocks until the engine submission is done, helping throttle the load on the system if kyuubi.session.engine.startup.waitCompletion is false. If false, the task returns immediately after opening the session without waiting. | boolean | 1.10.3 |
+| kyuubi.metadata.request.async.retry.enabled | true | Whether to retry in async when metadata request failed. When true, return success response immediately even the metadata request failed, and schedule it in background until success, to tolerate long-time metadata store outages w/o blocking the submission request. | boolean | 1.7.0 |
+| kyuubi.metadata.request.async.retry.queue.size | 65536 | The maximum queue size for buffering metadata requests in memory when the external metadata storage is down. Requests will be dropped if the queue exceeds. Only take affect when kyuubi.metadata.request.async.retry.enabled is `true`. | int | 1.6.0 |
+| kyuubi.metadata.request.async.retry.threads | 10 | Number of threads in the metadata request async retry manager thread pool. Only take affect when kyuubi.metadata.request.async.retry.enabled is `true`. | int | 1.6.0 |
+| kyuubi.metadata.request.async.retry.virtualThreads.enabled | false | 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. Defaults to kyuubi.server.virtualThreads.enabled. | boolean | 1.13.0 |
+| kyuubi.metadata.request.retry.interval | PT5S | The interval to check and trigger the metadata request retry tasks. | duration | 1.6.0 |
+| kyuubi.metadata.search.window | <undefined> | The time window to restrict user queries to metadata within a specific period, starting from the current time to the past. It only affects `GET /api/v1/batches` API. You may want to set this to short period to improve query performance and reduce load on the metadata store when administer want to reserve the metadata for long time. The side-effects is that, the metadata created outside the window will not be invisible to users. If it is undefined, all metadata will be visible for users. | duration | 1.10.1 |
+| kyuubi.metadata.store.class | org.apache.kyuubi.server.metadata.jdbc.JDBCMetadataStore | Fully qualified class name for server metadata store. | string | 1.6.0 |
+| kyuubi.metadata.store.jdbc.database.schema.init | true | Whether to init the JDBC metadata store database schema. | boolean | 1.6.0 |
+| kyuubi.metadata.store.jdbc.database.type | SQLITE | The database type for server jdbc metadata store. - SQLITE: SQLite3, JDBC driver `org.sqlite.JDBC`.
- MYSQL: MySQL, JDBC driver `com.mysql.cj.jdbc.Driver` (fallback `com.mysql.jdbc.Driver`).
- POSTGRESQL: PostgreSQL, JDBC driver `org.postgresql.Driver`.
- CUSTOM: User-defined database type, need to specify corresponding JDBC driver.
Note that: The JDBC datasource is powered by HiKariCP, for datasource properties, please specify them with the prefix: kyuubi.metadata.store.jdbc.datasource. For example, kyuubi.metadata.store.jdbc.datasource.connectionTimeout=10000. | string | 1.6.0 |
+| kyuubi.metadata.store.jdbc.driver | <undefined> | JDBC driver class name for server jdbc metadata store. | string | 1.6.0 |
+| kyuubi.metadata.store.jdbc.password || The password for server JDBC metadata store. | string | 1.6.0 |
+| kyuubi.metadata.store.jdbc.priority.enabled | false | Whether to enable the priority scheduling for batch impl v2. When false, ignore kyuubi.batch.priority and use the FIFO ordering strategy for batch job scheduling. Note: this feature may cause significant performance issues when using MySQL 5.7 as the metastore backend due to the lack of support for mixed order index. See more details at KYUUBI #5329. | boolean | 1.8.0 |
+| kyuubi.metadata.store.jdbc.url | jdbc:sqlite:{{KYUUBI_HOME}}/kyuubi_state_store.db | The JDBC url for server JDBC metadata store. By default, it is a SQLite database url, and the state information is not shared across Kyuubi instances. To enable high availability for multiple kyuubi instances, please specify a production JDBC url. Note: this value support the variables substitution: `{{KYUUBI_HOME}}`, `{{KYUUBI_WORK_DIR_ROOT}}`. | string | 1.6.0 |
+| kyuubi.metadata.store.jdbc.user || The username for server JDBC metadata store. | string | 1.6.0 |
### Metrics
@@ -482,62 +488,66 @@ You can configure the Kyuubi properties in `$KYUUBI_HOME/conf/kyuubi-defaults.co
| kyuubi.server.tempFile.expireTime | P30D | Expiration timout for cleanup server-side temporary files, e.g. operation logs. | duration | 1.10.0 |
| kyuubi.server.tempFile.maxCount | <undefined> | The upper threshold size of server-side temporary file paths to cleanup | int | 1.10.0 |
| kyuubi.server.thrift.resultset.default.fetch.size | 1000 | The number of rows sent in one Fetch RPC call by the server to the client, if not specified by the client. Respect `hive.server2.thrift.resultset.default.fetch.size` hive conf. | int | 1.9.1 |
+| kyuubi.server.virtualThreads.enabled | false | Whether to use virtual threads by default for all supported executors in the Kyuubi server. Requires at least Java 21; Java 25 or later is recommended. An explicitly configured component-level virtual thread option takes precedence. | boolean | 1.13.0 |
### Session
-| Key | Default | Meaning | Type | Since |
-|---------------------------------------------------------|-------------------------|------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|----------|--------|
-| kyuubi.session.check.interval | PT5M | The check interval for session timeout. | duration | 1.0.0 |
-| kyuubi.session.close.on.disconnect | true | Session will be closed when client disconnects from kyuubi gateway. Set this to false to have session outlive its parent connection. | boolean | 1.8.0 |
-| kyuubi.session.conf.advisor | <undefined> | A config advisor plugin for Kyuubi Server. This plugin can provide a list of custom configs for different users or session configs and overwrite the session configs before opening a new session. This config value should be a subclass of `org.apache.kyuubi.plugin.SessionConfAdvisor` which has a zero-arg constructor. | seq | 1.5.0 |
-| kyuubi.session.conf.file.reload.interval | PT10M | When `FileSessionConfAdvisor` is used, this configuration defines the expired time of `$KYUUBI_CONF_DIR/kyuubi-session-.conf` in the cache. After exceeding this value, the file will be reloaded. | duration | 1.7.0 |
-| kyuubi.session.conf.ignore.list || A comma-separated list of ignored keys. If the client connection contains any of them, the key and the corresponding value will be removed silently during engine bootstrap and connection setup. Note that this rule is for server-side protection defined via administrators to prevent some essential configs from tampering but will not forbid users to set dynamic configurations via SET syntax. | set | 1.2.0 |
-| kyuubi.session.conf.profile | <undefined> | Specify a profile to load session-level configurations from multiple `$KYUUBI_CONF_DIR/kyuubi-session-.conf` files. This configuration will be ignored if the file does not exist. This configuration only takes effect when `kyuubi.session.conf.advisor` is set as `org.apache.kyuubi.session.FileSessionConfAdvisor`. | seq | 1.7.0 |
-| kyuubi.session.conf.restrict.list || A comma-separated list of restricted keys. If the client connection contains any of them, the connection will be rejected explicitly during engine bootstrap and connection setup. Note that this rule is for server-side protection defined via administrators to prevent some essential configs from tampering but will not forbid users to set dynamic configurations via SET syntax. | set | 1.2.0 |
-| kyuubi.session.engine.alive.max.failures | 3 | The maximum number of failures allowed for the engine. | int | 1.8.1 |
-| kyuubi.session.engine.alive.probe.enabled | false | Whether to enable the engine alive probe, it true, we will create a companion thrift client that keeps sending simple requests to check whether the engine is alive. | boolean | 1.6.0 |
-| kyuubi.session.engine.alive.probe.interval | PT10S | The interval for engine alive probe. | duration | 1.6.0 |
-| kyuubi.session.engine.alive.timeout | PT2M | The timeout for engine alive. If there is no alive probe success in the last timeout window, the engine will be marked as no-alive. | duration | 1.6.0 |
-| kyuubi.session.engine.check.interval | PT1M | The check interval for engine timeout | duration | 1.0.0 |
-| kyuubi.session.engine.flink.fetch.timeout | <undefined> | Result fetch timeout for Flink engine. If the timeout is reached, the result fetch would be stopped and the current fetched would be returned. If no data are fetched, a TimeoutException would be thrown. | duration | 1.8.0 |
-| kyuubi.session.engine.flink.initialize.sql || The initialize sql for Flink session. It fallback to `kyuubi.engine.session.initialize.sql` | seq | 1.8.1 |
-| kyuubi.session.engine.flink.main.resource | <undefined> | The package used to create Flink SQL engine remote job. If it is undefined, Kyuubi will use the default | string | 1.4.0 |
-| kyuubi.session.engine.flink.max.rows | 1000000 | Max rows of Flink query results. For batch queries, rows exceeding the limit would be ignored. For streaming queries, the query would be canceled if the limit is reached. | int | 1.5.0 |
-| kyuubi.session.engine.hive.main.resource | <undefined> | The package used to create Hive engine remote job. If it is undefined, Kyuubi will use the default | string | 1.6.0 |
-| kyuubi.session.engine.idle.timeout | PT30M | engine timeout, the engine will self-terminate when it's not accessed for this duration. 0 or negative means not to self-terminate. | duration | 1.0.0 |
-| kyuubi.session.engine.initialize.timeout | PT3M | Timeout for starting the background engine, e.g. SparkSQLEngine. | duration | 1.0.0 |
-| kyuubi.session.engine.launch.async | true | When opening kyuubi session, whether to launch the backend engine asynchronously. When true, the Kyuubi server will set up the connection with the client without delay as the backend engine will be created asynchronously. | boolean | 1.4.0 |
-| kyuubi.session.engine.log.timeout | PT24H | If we use Spark as the engine then the session submit log is the console output of spark-submit. We will retain the session submit log until over the config value. | duration | 1.1.0 |
-| kyuubi.session.engine.login.timeout | PT15S | The timeout of creating the connection to remote sql query engine | duration | 1.0.0 |
-| kyuubi.session.engine.open.max.attempts | 9 | The number of times an open engine will retry when encountering a special error. | int | 1.7.0 |
-| kyuubi.session.engine.open.onFailure | RETRY | The behavior when opening engine failed: - RETRY: retry to open engine for kyuubi.session.engine.open.max.attempts times.
- DEREGISTER_IMMEDIATELY: deregister the engine immediately.
- DEREGISTER_AFTER_RETRY: deregister the engine after retry to open engine for kyuubi.session.engine.open.max.attempts times.
| string | 1.8.1 |
-| kyuubi.session.engine.open.retry.wait | PT10S | How long to wait before retrying to open the engine after failure. | duration | 1.7.0 |
-| kyuubi.session.engine.share.level | USER | (deprecated) - Using kyuubi.engine.share.level instead | string | 1.0.0 |
-| kyuubi.session.engine.spark.initialize.sql || The initialize sql for Spark session. It fallback to `kyuubi.engine.session.initialize.sql` | seq | 1.8.1 |
-| kyuubi.session.engine.spark.main.resource | <undefined> | The package used to create Spark SQL engine remote application. If it is undefined, Kyuubi will use the default | string | 1.0.0 |
-| kyuubi.session.engine.spark.max.initial.wait | PT1M | Max wait time for the initial connection to Spark engine. The engine will self-terminate no new incoming connection is established within this time. This setting only applies at the CONNECTION share level. 0 or negative means not to self-terminate. | duration | 1.8.0 |
-| kyuubi.session.engine.spark.max.lifetime | PT0S | Max lifetime for Spark engine, the engine will self-terminate when it reaches the end of life. 0 or negative means not to self-terminate. | duration | 1.6.0 |
-| kyuubi.session.engine.spark.max.lifetime.gracefulPeriod | PT0S | Graceful period for Spark engine to wait the connections disconnected after reaching the end of life. After the graceful period, all the connections without running operations will be forcibly disconnected. 0 or negative means always waiting the connections disconnected. | duration | 1.8.1 |
-| kyuubi.session.engine.spark.progress.timeFormat | yyyy-MM-dd HH:mm:ss.SSS | The time format of the progress bar | string | 1.6.0 |
-| kyuubi.session.engine.spark.progress.update.interval | PT1S | Update period of progress bar. | duration | 1.6.0 |
-| kyuubi.session.engine.spark.showProgress | false | When true, show the progress bar in the Spark's engine log. | boolean | 1.6.0 |
-| kyuubi.session.engine.startup.destroy.timeout | PT5S | Engine startup process destroy wait time, if the process does not stop after this time, force destroy instead. This configuration only takes effect when `kyuubi.session.engine.startup.waitCompletion=false`. | duration | 1.8.0 |
-| kyuubi.session.engine.startup.error.max.size | 8192 | During engine bootstrapping, if an error occurs, using this config to limit the length of error message(characters). | int | 1.1.0 |
-| kyuubi.session.engine.startup.maxLogLines | 10 | The maximum number of engine log lines when errors occur during the engine startup phase. Note that this config effects on client-side to help track engine startup issues. | int | 1.4.0 |
-| kyuubi.session.engine.startup.waitCompletion | true | Whether to wait for completion after the engine starts. If false, the startup process will be destroyed after the engine is started. Note that only use it when the driver is not running locally, such as in yarn-cluster mode; Otherwise, the engine will be killed. | boolean | 1.5.0 |
-| kyuubi.session.engine.trino.connection.catalog | <undefined> | The default catalog that Trino engine will connect to | string | 1.5.0 |
-| kyuubi.session.engine.trino.connection.url | <undefined> | The server url that Trino engine will connect to | string | 1.5.0 |
-| kyuubi.session.engine.trino.main.resource | <undefined> | The package used to create Trino engine remote job. If it is undefined, Kyuubi will use the default | string | 1.5.0 |
-| kyuubi.session.engine.trino.progress.update.interval | PT1S | Update period of progress bar. | duration | 1.10.0 |
-| kyuubi.session.engine.trino.showProgress | true | When true, show the progress bar and final info in the Trino engine log. | boolean | 1.6.0 |
-| kyuubi.session.engine.trino.showProgress.debug | false | When true, show the progress debug info in the Trino engine log. | boolean | 1.6.0 |
-| kyuubi.session.group.provider | hadoop | A group provider plugin for Kyuubi Server. This plugin can provide primary group and groups information for different users or session configs. This config value should be a subclass of `org.apache.kyuubi.plugin.GroupProvider` which has a zero-arg constructor. Kyuubi provides the following built-in implementations: - hadoop: delegate the user group mapping to hadoop UserGroupInformation.
| string | 1.7.0 |
-| kyuubi.session.idle.timeout | PT6H | session idle timeout, it will be closed when it's not accessed for this duration | duration | 1.2.0 |
-| kyuubi.session.local.dir.allow.list || The local dir list that are allowed to access by the kyuubi session application. End-users might set some parameters such as `spark.files` and it will upload some local files when launching the kyuubi engine, if the local dir allow list is defined, kyuubi will check whether the path to upload is in the allow list. Note that, if it is empty, there is no limitation for that. And please use absolute paths. Also, currently this config takes effect only for Spark engine. | set | 1.6.0 |
-| kyuubi.session.name | <undefined> | A human readable name of the session and we use empty string by default. This name will be recorded in the event. Note that, we only apply this value from session conf. | string | 1.4.0 |
-| kyuubi.session.proxy.user | <undefined> | An alternative to hive.server2.proxy.user. The current behavior is consistent with hive.server2.proxy.user and now only takes effect in RESTFul API. When both parameters are set, kyuubi.session.proxy.user takes precedence. | string | 1.9.0 |
-| kyuubi.session.timeout | PT6H | (deprecated)session timeout, it will be closed when it's not accessed for this duration | duration | 1.0.0 |
-| kyuubi.session.user.sign.enabled | false | Whether to verify the integrity of session user name on the engine side, e.g. Authz plugin in Spark. | boolean | 1.7.0 |
+| Key | Default | Meaning | Type | Since |
+|----------------------------------------------------------|-------------------------|------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|----------|--------|
+| kyuubi.session.check.interval | PT5M | The check interval for session timeout. | duration | 1.0.0 |
+| kyuubi.session.close.on.disconnect | true | Session will be closed when client disconnects from kyuubi gateway. Set this to false to have session outlive its parent connection. | boolean | 1.8.0 |
+| kyuubi.session.conf.advisor | <undefined> | A config advisor plugin for Kyuubi Server. This plugin can provide a list of custom configs for different users or session configs and overwrite the session configs before opening a new session. This config value should be a subclass of `org.apache.kyuubi.plugin.SessionConfAdvisor` which has a zero-arg constructor. | seq | 1.5.0 |
+| kyuubi.session.conf.file.reload.interval | PT10M | When `FileSessionConfAdvisor` is used, this configuration defines the expired time of `$KYUUBI_CONF_DIR/kyuubi-session-.conf` in the cache. After exceeding this value, the file will be reloaded. | duration | 1.7.0 |
+| kyuubi.session.conf.ignore.list || A comma-separated list of ignored keys. If the client connection contains any of them, the key and the corresponding value will be removed silently during engine bootstrap and connection setup. Note that this rule is for server-side protection defined via administrators to prevent some essential configs from tampering but will not forbid users to set dynamic configurations via SET syntax. | set | 1.2.0 |
+| kyuubi.session.conf.profile | <undefined> | Specify a profile to load session-level configurations from multiple `$KYUUBI_CONF_DIR/kyuubi-session-.conf` files. This configuration will be ignored if the file does not exist. This configuration only takes effect when `kyuubi.session.conf.advisor` is set as `org.apache.kyuubi.session.FileSessionConfAdvisor`. | seq | 1.7.0 |
+| kyuubi.session.conf.restrict.list || A comma-separated list of restricted keys. If the client connection contains any of them, the connection will be rejected explicitly during engine bootstrap and connection setup. Note that this rule is for server-side protection defined via administrators to prevent some essential configs from tampering but will not forbid users to set dynamic configurations via SET syntax. | set | 1.2.0 |
+| kyuubi.session.engine.alive.max.failures | 3 | The maximum number of failures allowed for the engine. | int | 1.8.1 |
+| kyuubi.session.engine.alive.probe.enabled | false | Whether to enable the engine alive probe, it true, we will create a companion thrift client that keeps sending simple requests to check whether the engine is alive. | boolean | 1.6.0 |
+| kyuubi.session.engine.alive.probe.interval | PT10S | The interval for engine alive probe. | duration | 1.6.0 |
+| kyuubi.session.engine.alive.probe.virtualThreads.enabled | false | 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. Defaults to kyuubi.server.virtualThreads.enabled. | boolean | 1.13.0 |
+| kyuubi.session.engine.alive.timeout | PT2M | The timeout for engine alive. If there is no alive probe success in the last timeout window, the engine will be marked as no-alive. | duration | 1.6.0 |
+| kyuubi.session.engine.check.interval | PT1M | The check interval for engine timeout | duration | 1.0.0 |
+| kyuubi.session.engine.flink.fetch.timeout | <undefined> | Result fetch timeout for Flink engine. If the timeout is reached, the result fetch would be stopped and the current fetched would be returned. If no data are fetched, a TimeoutException would be thrown. | duration | 1.8.0 |
+| kyuubi.session.engine.flink.initialize.sql || The initialize sql for Flink session. It fallback to `kyuubi.engine.session.initialize.sql` | seq | 1.8.1 |
+| kyuubi.session.engine.flink.main.resource | <undefined> | The package used to create Flink SQL engine remote job. If it is undefined, Kyuubi will use the default | string | 1.4.0 |
+| kyuubi.session.engine.flink.max.rows | 1000000 | Max rows of Flink query results. For batch queries, rows exceeding the limit would be ignored. For streaming queries, the query would be canceled if the limit is reached. | int | 1.5.0 |
+| kyuubi.session.engine.hive.main.resource | <undefined> | The package used to create Hive engine remote job. If it is undefined, Kyuubi will use the default | string | 1.6.0 |
+| kyuubi.session.engine.idle.timeout | PT30M | engine timeout, the engine will self-terminate when it's not accessed for this duration. 0 or negative means not to self-terminate. | duration | 1.0.0 |
+| kyuubi.session.engine.initialize.timeout | PT3M | Timeout for starting the background engine, e.g. SparkSQLEngine. | duration | 1.0.0 |
+| kyuubi.session.engine.launch.async | true | When opening kyuubi session, whether to launch the backend engine asynchronously. When true, the Kyuubi server will set up the connection with the client without delay as the backend engine will be created asynchronously. | boolean | 1.4.0 |
+| kyuubi.session.engine.log.capture.virtualThreads.enabled | false | Whether the Kyuubi server uses virtual threads to capture engine startup logs. Requires at least Java 21; Java 25 or later is recommended. Defaults to kyuubi.server.virtualThreads.enabled. | boolean | 1.13.0 |
+| kyuubi.session.engine.log.timeout | PT24H | If we use Spark as the engine then the session submit log is the console output of spark-submit. We will retain the session submit log until over the config value. | duration | 1.1.0 |
+| kyuubi.session.engine.login.timeout | PT15S | The timeout of creating the connection to remote sql query engine | duration | 1.0.0 |
+| kyuubi.session.engine.open.max.attempts | 9 | The number of times an open engine will retry when encountering a special error. | int | 1.7.0 |
+| kyuubi.session.engine.open.onFailure | RETRY | The behavior when opening engine failed: - RETRY: retry to open engine for kyuubi.session.engine.open.max.attempts times.
- DEREGISTER_IMMEDIATELY: deregister the engine immediately.
- DEREGISTER_AFTER_RETRY: deregister the engine after retry to open engine for kyuubi.session.engine.open.max.attempts times.
| string | 1.8.1 |
+| kyuubi.session.engine.open.retry.wait | PT10S | How long to wait before retrying to open the engine after failure. | duration | 1.7.0 |
+| kyuubi.session.engine.rpc.client.virtualThreads.enabled | false | 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. Defaults to kyuubi.server.virtualThreads.enabled. | boolean | 1.13.0 |
+| kyuubi.session.engine.share.level | USER | (deprecated) - Using kyuubi.engine.share.level instead | string | 1.0.0 |
+| kyuubi.session.engine.spark.initialize.sql || The initialize sql for Spark session. It fallback to `kyuubi.engine.session.initialize.sql` | seq | 1.8.1 |
+| kyuubi.session.engine.spark.main.resource | <undefined> | The package used to create Spark SQL engine remote application. If it is undefined, Kyuubi will use the default | string | 1.0.0 |
+| kyuubi.session.engine.spark.max.initial.wait | PT1M | Max wait time for the initial connection to Spark engine. The engine will self-terminate no new incoming connection is established within this time. This setting only applies at the CONNECTION share level. 0 or negative means not to self-terminate. | duration | 1.8.0 |
+| kyuubi.session.engine.spark.max.lifetime | PT0S | Max lifetime for Spark engine, the engine will self-terminate when it reaches the end of life. 0 or negative means not to self-terminate. | duration | 1.6.0 |
+| kyuubi.session.engine.spark.max.lifetime.gracefulPeriod | PT0S | Graceful period for Spark engine to wait the connections disconnected after reaching the end of life. After the graceful period, all the connections without running operations will be forcibly disconnected. 0 or negative means always waiting the connections disconnected. | duration | 1.8.1 |
+| kyuubi.session.engine.spark.progress.timeFormat | yyyy-MM-dd HH:mm:ss.SSS | The time format of the progress bar | string | 1.6.0 |
+| kyuubi.session.engine.spark.progress.update.interval | PT1S | Update period of progress bar. | duration | 1.6.0 |
+| kyuubi.session.engine.spark.showProgress | false | When true, show the progress bar in the Spark's engine log. | boolean | 1.6.0 |
+| kyuubi.session.engine.startup.destroy.timeout | PT5S | Engine startup process destroy wait time, if the process does not stop after this time, force destroy instead. This configuration only takes effect when `kyuubi.session.engine.startup.waitCompletion=false`. | duration | 1.8.0 |
+| kyuubi.session.engine.startup.error.max.size | 8192 | During engine bootstrapping, if an error occurs, using this config to limit the length of error message(characters). | int | 1.1.0 |
+| kyuubi.session.engine.startup.maxLogLines | 10 | The maximum number of engine log lines when errors occur during the engine startup phase. Note that this config effects on client-side to help track engine startup issues. | int | 1.4.0 |
+| kyuubi.session.engine.startup.waitCompletion | true | Whether to wait for completion after the engine starts. If false, the startup process will be destroyed after the engine is started. Note that only use it when the driver is not running locally, such as in yarn-cluster mode; Otherwise, the engine will be killed. | boolean | 1.5.0 |
+| kyuubi.session.engine.trino.connection.catalog | <undefined> | The default catalog that Trino engine will connect to | string | 1.5.0 |
+| kyuubi.session.engine.trino.connection.url | <undefined> | The server url that Trino engine will connect to | string | 1.5.0 |
+| kyuubi.session.engine.trino.main.resource | <undefined> | The package used to create Trino engine remote job. If it is undefined, Kyuubi will use the default | string | 1.5.0 |
+| kyuubi.session.engine.trino.progress.update.interval | PT1S | Update period of progress bar. | duration | 1.10.0 |
+| kyuubi.session.engine.trino.showProgress | true | When true, show the progress bar and final info in the Trino engine log. | boolean | 1.6.0 |
+| kyuubi.session.engine.trino.showProgress.debug | false | When true, show the progress debug info in the Trino engine log. | boolean | 1.6.0 |
+| kyuubi.session.group.provider | hadoop | A group provider plugin for Kyuubi Server. This plugin can provide primary group and groups information for different users or session configs. This config value should be a subclass of `org.apache.kyuubi.plugin.GroupProvider` which has a zero-arg constructor. Kyuubi provides the following built-in implementations: - hadoop: delegate the user group mapping to hadoop UserGroupInformation.
| string | 1.7.0 |
+| kyuubi.session.idle.timeout | PT6H | session idle timeout, it will be closed when it's not accessed for this duration | duration | 1.2.0 |
+| kyuubi.session.local.dir.allow.list || The local dir list that are allowed to access by the kyuubi session application. End-users might set some parameters such as `spark.files` and it will upload some local files when launching the kyuubi engine, if the local dir allow list is defined, kyuubi will check whether the path to upload is in the allow list. Note that, if it is empty, there is no limitation for that. And please use absolute paths. Also, currently this config takes effect only for Spark engine. | set | 1.6.0 |
+| kyuubi.session.name | <undefined> | A human readable name of the session and we use empty string by default. This name will be recorded in the event. Note that, we only apply this value from session conf. | string | 1.4.0 |
+| kyuubi.session.proxy.user | <undefined> | An alternative to hive.server2.proxy.user. The current behavior is consistent with hive.server2.proxy.user and now only takes effect in RESTFul API. When both parameters are set, kyuubi.session.proxy.user takes precedence. | string | 1.9.0 |
+| kyuubi.session.timeout | PT6H | (deprecated)session timeout, it will be closed when it's not accessed for this duration | duration | 1.0.0 |
+| kyuubi.session.user.sign.enabled | false | Whether to verify the integrity of session user name on the engine side, e.g. Authz plugin in Spark. | boolean | 1.7.0 |
### Spnego
diff --git a/kyuubi-common/src/main/scala/org/apache/kyuubi/config/KyuubiConf.scala b/kyuubi-common/src/main/scala/org/apache/kyuubi/config/KyuubiConf.scala
index 30eec5fefb6..ba68133b8f0 100644
--- a/kyuubi-common/src/main/scala/org/apache/kyuubi/config/KyuubiConf.scala
+++ b/kyuubi-common/src/main/scala/org/apache/kyuubi/config/KyuubiConf.scala
@@ -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 {
@@ -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}")
@@ -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 " +
+ "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")
@@ -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.")
@@ -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," +
@@ -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.")
@@ -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.")
@@ -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
@@ -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")
@@ -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.")
@@ -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)
@@ -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")
@@ -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)
}
diff --git a/kyuubi-common/src/main/scala/org/apache/kyuubi/session/SessionManager.scala b/kyuubi-common/src/main/scala/org/apache/kyuubi/session/SessionManager.scala
index 35848dfa9e7..31e54189e41 100644
--- a/kyuubi-common/src/main/scala/org/apache/kyuubi/session/SessionManager.scala
+++ b/kyuubi-common/src/main/scala/org/apache/kyuubi/session/SessionManager.scala
@@ -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
@@ -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)
@@ -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] = _
@@ -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)
}
diff --git a/kyuubi-common/src/main/scala/org/apache/kyuubi/util/ThreadUtils.scala b/kyuubi-common/src/main/scala/org/apache/kyuubi/util/ThreadUtils.scala
index 831c9b63e0a..21f27663558 100644
--- a/kyuubi-common/src/main/scala/org/apache/kyuubi/util/ThreadUtils.scala
+++ b/kyuubi-common/src/main/scala/org/apache/kyuubi/util/ThreadUtils.scala
@@ -18,6 +18,7 @@
package org.apache.kyuubi.util
import java.util.concurrent._
+import java.util.concurrent.atomic.AtomicInteger
import scala.concurrent.Awaitable
import scala.concurrent.duration.{Duration, FiniteDuration}
@@ -26,29 +27,70 @@ import org.apache.kyuubi.{KyuubiException, Logging}
object ThreadUtils extends Logging {
- def newBoundedVirtualThreadPerTaskExecutor(
- maxConcurrentTasks: Int,
- threadNamePrefix: String): ExecutorService = {
- require(maxConcurrentTasks > 0, "maxConcurrentTasks must be positive")
- new BoundedExecutorService(
- newVirtualThreadPerTaskExecutor(threadNamePrefix),
- maxConcurrentTasks)
+ trait ExecutorServiceWithMetrics extends ExecutorService {
+ def getPoolSize: Int
+ def getActiveCount: Int
+ def getQueueSize: Int
+ }
+
+ lazy val isVirtualThreadSupported: Boolean = {
+ try {
+ classOf[Thread].getMethod("ofVirtual")
+ classOf[Executors].getMethod("newThreadPerTaskExecutor", classOf[ThreadFactory])
+ true
+ } catch {
+ case _: ReflectiveOperationException => false
+ }
}
- private def newVirtualThreadPerTaskExecutor(threadNamePrefix: String): ExecutorService =
+ def newVirtualThreadFactory(threadNamePrefix: String): ThreadFactory =
try {
val builder = classOf[Thread].getMethod("ofVirtual").invoke(null)
val builderClass = Class.forName("java.lang.Thread$Builder")
val namedBuilder = builderClass
.getMethod("name", classOf[String], java.lang.Long.TYPE)
.invoke(builder, s"$threadNamePrefix-", Long.box(0L))
- val threadFactory = builderClass
+ val configuredBuilder = builderClass
+ .getMethod("uncaughtExceptionHandler", classOf[Thread.UncaughtExceptionHandler])
+ .invoke(namedBuilder, NamedThreadFactory.kyuubiUncaughtExceptionHandler)
+ builderClass
.getMethod("factory")
- .invoke(namedBuilder)
+ .invoke(configuredBuilder)
.asInstanceOf[ThreadFactory]
+ } catch {
+ case e: ReflectiveOperationException =>
+ throw new IllegalStateException("Virtual threads require Java 21 or later", e)
+ }
+
+ def newBoundedVirtualThreadPerTaskExecutor(
+ maxConcurrentTasks: Int,
+ threadNamePrefix: String): ExecutorService = {
+ require(maxConcurrentTasks > 0, "maxConcurrentTasks must be positive")
+ new BoundedExecutorService(
+ newVirtualThreadPerTaskExecutor(threadNamePrefix),
+ maxConcurrentTasks)
+ }
+
+ def newBoundedQueuedVirtualThreadPerTaskExecutor(
+ maxConcurrentTasks: Int,
+ maxQueuedTasks: Int,
+ threadNamePrefix: String): ExecutorServiceWithMetrics = {
+ require(maxConcurrentTasks > 0, "maxConcurrentTasks must be positive")
+ require(maxQueuedTasks >= 0, "maxQueuedTasks must not be negative")
+ require(
+ maxQueuedTasks <= Int.MaxValue - maxConcurrentTasks,
+ "maxConcurrentTasks plus maxQueuedTasks must not exceed Int.MaxValue")
+ new BoundedQueuedExecutorService(
+ newVirtualThreadPerTaskExecutor(threadNamePrefix),
+ maxConcurrentTasks,
+ maxQueuedTasks)
+ }
+
+ def newVirtualThreadPerTaskExecutor(threadNamePrefix: String): ExecutorService =
+ try {
classOf[Executors]
.getMethod("newThreadPerTaskExecutor", classOf[ThreadFactory])
- .invoke(null, threadFactory)
+ .invoke(null, newVirtualThreadFactory(threadNamePrefix))
.asInstanceOf[ExecutorService]
} catch {
case e: ReflectiveOperationException =>
@@ -66,6 +108,16 @@ object ThreadUtils extends Logging {
executor
}
+ def newVirtualThreadSingleThreadScheduledExecutor(
+ threadName: String,
+ executeExistingDelayedTasksAfterShutdown: Boolean = true): ScheduledExecutorService = {
+ val executor = new ScheduledThreadPoolExecutor(1, newVirtualThreadFactory(threadName))
+ executor.setRemoveOnCancelPolicy(true)
+ executor
+ .setExecuteExistingDelayedTasksAfterShutdownPolicy(executeExistingDelayedTasksAfterShutdown)
+ executor
+ }
+
def newDaemonQueuedThreadPool(
poolSize: Int,
poolQueueSize: Int,
@@ -220,4 +272,79 @@ object ThreadUtils extends Logging {
}
}
}
+
+ // The fair semaphore does not guarantee FIFO execution: virtual threads may reach it
+ // out of submission order.
+ private class BoundedQueuedExecutorService(
+ delegate: ExecutorService,
+ maxConcurrentTasks: Int,
+ maxQueuedTasks: Int)
+ extends AbstractExecutorService with ExecutorServiceWithMetrics {
+
+ private val admittedTasks = new Semaphore(maxConcurrentTasks + maxQueuedTasks)
+ private val runningTasks = new Semaphore(maxConcurrentTasks, true)
+ private val activeCount = new AtomicInteger()
+ private val queueSize = new AtomicInteger()
+
+ override def getPoolSize: Int = getActiveCount
+
+ override def getActiveCount: Int = activeCount.get()
+
+ override def getQueueSize: Int = queueSize.get()
+
+ override def shutdown(): Unit = delegate.shutdown()
+
+ override def shutdownNow(): java.util.List[Runnable] = delegate.shutdownNow()
+
+ override def isShutdown: Boolean = delegate.isShutdown
+
+ override def isTerminated: Boolean = delegate.isTerminated
+
+ override def awaitTermination(timeout: Long, unit: TimeUnit): Boolean =
+ delegate.awaitTermination(timeout, unit)
+
+ override def execute(command: Runnable): Unit = {
+ if (!admittedTasks.tryAcquire()) {
+ throw new RejectedExecutionException(
+ s"Maximum running task limit $maxConcurrentTasks and queue limit " +
+ s"$maxQueuedTasks reached")
+ }
+ queueSize.incrementAndGet()
+
+ try {
+ delegate.execute(new Runnable {
+ override def run(): Unit = {
+ var active = false
+ try {
+ runningTasks.acquire()
+ queueSize.decrementAndGet()
+ activeCount.incrementAndGet()
+ active = true
+ command.run()
+ } catch {
+ case _: InterruptedException =>
+ command match {
+ case future: java.util.concurrent.Future[_] => future.cancel(true)
+ case _ =>
+ }
+ Thread.currentThread().interrupt()
+ } finally {
+ if (active) {
+ activeCount.decrementAndGet()
+ runningTasks.release()
+ } else {
+ queueSize.decrementAndGet()
+ }
+ admittedTasks.release()
+ }
+ }
+ })
+ } catch {
+ case t: Throwable =>
+ queueSize.decrementAndGet()
+ admittedTasks.release()
+ throw t
+ }
+ }
+ }
}
diff --git a/kyuubi-common/src/test/scala/org/apache/kyuubi/config/KyuubiConfSuite.scala b/kyuubi-common/src/test/scala/org/apache/kyuubi/config/KyuubiConfSuite.scala
index 9ee691773c3..af506f4e94f 100644
--- a/kyuubi-common/src/test/scala/org/apache/kyuubi/config/KyuubiConfSuite.scala
+++ b/kyuubi-common/src/test/scala/org/apache/kyuubi/config/KyuubiConfSuite.scala
@@ -21,6 +21,7 @@ import java.time.Duration
import org.apache.kyuubi.KyuubiFunSuite
import org.apache.kyuubi.engine.EngineType
+import org.apache.kyuubi.util.ThreadUtils
class KyuubiConfSuite extends KyuubiFunSuite {
@@ -352,14 +353,59 @@ class KyuubiConfSuite extends KyuubiFunSuite {
}
}
- test("getEngineConf excludes the server thrift binary virtual thread config") {
+ test("getEngineConf excludes server virtual thread configs") {
val kyuubiConf = KyuubiConf(false)
- kyuubiConf.set(FRONTEND_THRIFT_BINARY_VIRTUAL_THREADS_ENABLED, true)
+ val serverConfigs = SERVER_VIRTUAL_THREADS_ENABLED +: serverVirtualThreadConfigs
+ serverConfigs.foreach(kyuubiConf.set(_, true))
EngineType.values.foreach { engineType =>
- assert(!kyuubiConf.getEngineConf(engineType)
- .contains(FRONTEND_THRIFT_BINARY_VIRTUAL_THREADS_ENABLED.key))
+ val engineConf = kyuubiConf.getEngineConf(engineType)
+ serverConfigs.foreach { config =>
+ assert(!engineConf.contains(config.key), s"$engineType should not receive ${config.key}")
+ }
+ }
+ }
+
+ test("server virtual thread switch and component overrides") {
+ val kyuubiConf = KyuubiConf(false)
+
+ assert(!kyuubiConf.get(SERVER_VIRTUAL_THREADS_ENABLED))
+ serverVirtualThreadConfigs.foreach(config => assert(!kyuubiConf.get(config)))
+
+ kyuubiConf.set(SERVER_VIRTUAL_THREADS_ENABLED, true)
+ serverVirtualThreadConfigs.foreach(config => assert(kyuubiConf.get(config)))
+
+ kyuubiConf.set(ENGINE_RPC_CLIENT_VIRTUAL_THREADS_ENABLED, false)
+ assert(!kyuubiConf.get(ENGINE_RPC_CLIENT_VIRTUAL_THREADS_ENABLED))
+ serverVirtualThreadConfigs.filterNot(
+ _ == ENGINE_RPC_CLIENT_VIRTUAL_THREADS_ENABLED).foreach {
+ config => assert(kyuubiConf.get(config))
+ }
+ }
+
+ test("validate server virtual thread configs") {
+ KyuubiConf(false).validateServerVirtualThreadConfigs()
+
+ (SERVER_VIRTUAL_THREADS_ENABLED +: serverVirtualThreadConfigs).foreach { config =>
+ val conf = KyuubiConf(false).set(config, true)
+ if (ThreadUtils.isVirtualThreadSupported) {
+ conf.validateServerVirtualThreadConfigs()
+ } else {
+ val error = intercept[IllegalArgumentException] {
+ conf.validateServerVirtualThreadConfigs()
+ }
+ val resolvedConfig = if (config == SERVER_VIRTUAL_THREADS_ENABLED) {
+ serverVirtualThreadConfigs.head
+ } else {
+ config
+ }
+ assert(error.getMessage.contains(s"${resolvedConfig.key}=true requires Java 21"))
+ }
}
+
+ val conf = KyuubiConf(false).set(SERVER_VIRTUAL_THREADS_ENABLED, true)
+ serverVirtualThreadConfigs.foreach(config => conf.set(config, false))
+ conf.validateServerVirtualThreadConfigs()
}
test("getEngineConf passes through reserved keys") {
diff --git a/kyuubi-common/src/test/scala/org/apache/kyuubi/session/SessionManagerValidationSuite.scala b/kyuubi-common/src/test/scala/org/apache/kyuubi/session/SessionManagerValidationSuite.scala
index 2da0ef455c6..fd6bc738f58 100644
--- a/kyuubi-common/src/test/scala/org/apache/kyuubi/session/SessionManagerValidationSuite.scala
+++ b/kyuubi-common/src/test/scala/org/apache/kyuubi/session/SessionManagerValidationSuite.scala
@@ -17,9 +17,13 @@
package org.apache.kyuubi.session
+import java.util.concurrent.TimeUnit
+import java.util.concurrent.atomic.AtomicBoolean
+
import org.apache.kyuubi.KyuubiFunSuite
import org.apache.kyuubi.config.KyuubiConf
import org.apache.kyuubi.config.KyuubiConf._
+import org.apache.kyuubi.util.ThreadUtils
class SessionManagerValidationSuite extends KyuubiFunSuite {
@@ -106,4 +110,34 @@ class SessionManagerValidationSuite extends KyuubiFunSuite {
assert(result("spark.executor.memory") === "4g")
}
+ test("server operation pool uses virtual threads when enabled") {
+ val conf = KyuubiConf(false)
+ .set(SERVER_EXEC_POOL_VIRTUAL_THREADS_ENABLED, true)
+ val manager = new NoopSessionManager() {
+ override protected def isServer: Boolean = true
+ }
+
+ if (!ThreadUtils.isVirtualThreadSupported) {
+ val error = intercept[IllegalStateException](manager.initialize(conf))
+ assert(error.getMessage.contains("Java 21"))
+ } else {
+ val isVirtual = classOf[Thread].getMethod("isVirtual")
+ val taskWasVirtual = new AtomicBoolean()
+
+ manager.initialize(conf)
+ manager.start()
+ try {
+ val task = manager.submitBackgroundOperation(new Runnable {
+ override def run(): Unit = {
+ taskWasVirtual.set(isVirtual.invoke(Thread.currentThread()).asInstanceOf[Boolean])
+ }
+ })
+ task.get(10, TimeUnit.SECONDS)
+ assert(taskWasVirtual.get())
+ } finally {
+ manager.stop()
+ }
+ }
+ }
+
}
diff --git a/kyuubi-common/src/test/scala/org/apache/kyuubi/util/ThreadUtilsSuite.scala b/kyuubi-common/src/test/scala/org/apache/kyuubi/util/ThreadUtilsSuite.scala
index d3a8140e613..afc9c384401 100644
--- a/kyuubi-common/src/test/scala/org/apache/kyuubi/util/ThreadUtilsSuite.scala
+++ b/kyuubi-common/src/test/scala/org/apache/kyuubi/util/ThreadUtilsSuite.scala
@@ -19,6 +19,8 @@ package org.apache.kyuubi.util
import java.util.concurrent.{ConcurrentLinkedQueue, CountDownLatch, RejectedExecutionException, TimeUnit}
+import scala.concurrent.duration._
+
import org.apache.kyuubi.KyuubiFunSuite
class ThreadUtilsSuite extends KyuubiFunSuite {
@@ -62,6 +64,23 @@ class ThreadUtilsSuite extends KyuubiFunSuite {
assert(threadName startsWith "")
}
+ test("New virtual thread single thread scheduled executor") {
+ if (ThreadUtils.isVirtualThreadSupported) {
+ val service = ThreadUtils.newVirtualThreadSingleThreadScheduledExecutor(
+ "ThreadUtilsVirtualScheduledTest")
+ val isVirtual = classOf[Thread].getMethod("isVirtual")
+ try {
+ val task = service.schedule(
+ () => isVirtual.invoke(Thread.currentThread()).asInstanceOf[Boolean],
+ 10,
+ TimeUnit.MILLISECONDS)
+ assert(task.get(10, TimeUnit.SECONDS))
+ } finally {
+ ThreadUtils.shutdown(service)
+ }
+ }
+ }
+
test("New daemon scheduled thread pool") {
val pool = ThreadUtils.newDaemonScheduledThreadPool(2, 10, "ThreadUtilsSchedTest")
// submit a task to ensure pool operational
@@ -78,15 +97,7 @@ class ThreadUtilsSuite extends KyuubiFunSuite {
}
test("New bounded virtual thread per task executor") {
- val virtualThreadsSupported =
- try {
- classOf[Thread].getMethod("isVirtual")
- true
- } catch {
- case _: NoSuchMethodException => false
- }
-
- if (!virtualThreadsSupported) {
+ if (!ThreadUtils.isVirtualThreadSupported) {
val error = intercept[IllegalStateException] {
ThreadUtils.newBoundedVirtualThreadPerTaskExecutor(2, "ThreadUtilsVirtualTest")
}
@@ -120,6 +131,8 @@ class ThreadUtilsSuite extends KyuubiFunSuite {
val last = executor.submit(new Runnable {
override def run(): Unit = {
+ assert(Thread.currentThread().getUncaughtExceptionHandler eq
+ NamedThreadFactory.kyuubiUncaughtExceptionHandler)
threadNames.add(Thread.currentThread().getName)
tasksAreVirtual.add(isVirtual.invoke(Thread.currentThread()).asInstanceOf[Boolean])
}
@@ -134,4 +147,52 @@ class ThreadUtilsSuite extends KyuubiFunSuite {
}
}
}
+
+ test("New bounded queued virtual thread per task executor") {
+ if (ThreadUtils.isVirtualThreadSupported) {
+ val executor =
+ ThreadUtils.newBoundedQueuedVirtualThreadPerTaskExecutor(
+ 2,
+ 1,
+ "ThreadUtilsQueuedTest")
+ val ready = new CountDownLatch(2)
+ val release = new CountDownLatch(1)
+ val isVirtual = classOf[Thread].getMethod("isVirtual")
+
+ def blockingTask: Runnable = new Runnable {
+ override def run(): Unit = {
+ assert(isVirtual.invoke(Thread.currentThread()).asInstanceOf[Boolean])
+ ready.countDown()
+ release.await()
+ }
+ }
+
+ try {
+ val first = executor.submit(blockingTask)
+ val second = executor.submit(blockingTask)
+ assert(ready.await(10, TimeUnit.SECONDS))
+ val third = executor.submit(new Runnable {
+ override def run(): Unit = ()
+ })
+
+ assert(executor.getPoolSize === 2)
+ assert(executor.getActiveCount === 2)
+ assert(executor.getQueueSize === 1)
+ intercept[RejectedExecutionException](executor.submit(blockingTask))
+
+ release.countDown()
+ first.get(10, TimeUnit.SECONDS)
+ second.get(10, TimeUnit.SECONDS)
+ third.get(10, TimeUnit.SECONDS)
+ eventually(timeout(10.seconds), interval(10.millis)) {
+ assert(executor.getPoolSize === 0)
+ assert(executor.getActiveCount === 0)
+ assert(executor.getQueueSize === 0)
+ }
+ } finally {
+ release.countDown()
+ ThreadUtils.shutdown(executor)
+ }
+ }
+ }
}
diff --git a/kyuubi-server/src/main/scala/org/apache/kyuubi/client/KyuubiSyncThriftClient.scala b/kyuubi-server/src/main/scala/org/apache/kyuubi/client/KyuubiSyncThriftClient.scala
index c36e9ec06c7..8d8d03b4d72 100644
--- a/kyuubi-server/src/main/scala/org/apache/kyuubi/client/KyuubiSyncThriftClient.scala
+++ b/kyuubi-server/src/main/scala/org/apache/kyuubi/client/KyuubiSyncThriftClient.scala
@@ -29,7 +29,7 @@ import com.google.common.annotations.VisibleForTesting
import org.apache.kyuubi.{KyuubiSQLException, Logging, Utils}
import org.apache.kyuubi.config.KyuubiConf
-import org.apache.kyuubi.config.KyuubiConf.ENGINE_LOGIN_TIMEOUT
+import org.apache.kyuubi.config.KyuubiConf.{ENGINE_ALIVE_PROBE_VIRTUAL_THREADS_ENABLED, ENGINE_LOGIN_TIMEOUT, ENGINE_RPC_CLIENT_VIRTUAL_THREADS_ENABLED}
import org.apache.kyuubi.config.KyuubiReservedKeys._
import org.apache.kyuubi.operation.FetchOrientation
import org.apache.kyuubi.operation.FetchOrientation.FetchOrientation
@@ -47,7 +47,9 @@ class KyuubiSyncThriftClient private (
protocol: TProtocol,
engineAliveProbeProtocol: Option[TProtocol],
engineAliveProbeInterval: Long,
- engineAliveTimeout: Long)
+ engineAliveTimeout: Long,
+ useVirtualThreadsForAliveProbe: Boolean,
+ useVirtualThreadsForAsyncRequests: Boolean)
extends TCLIService.Client(protocol) with Logging {
@volatile private var _remoteSessionHandle: TSessionHandle = _
@@ -73,8 +75,15 @@ class KyuubiSyncThriftClient private (
@volatile private var asyncRequestExecutorInitialized: Boolean = false
private lazy val asyncRequestExecutor: ExecutorService = {
asyncRequestExecutorInitialized = true
- ThreadUtils.newDaemonSingleThreadScheduledExecutor(
- "async-request-executor-" + SessionHandle(_remoteSessionHandle))
+ val threadName = "async-request-executor-" + SessionHandle(_remoteSessionHandle)
+ if (useVirtualThreadsForAsyncRequests) {
+ // The session lock serializes submissions, so this is not a normal request backlog.
+ // Retain the platform executor's effectively unbounded queue: a completed or cancelled
+ // Future does not imply that the previous task has released its execution permit.
+ ThreadUtils.newBoundedQueuedVirtualThreadPerTaskExecutor(1, Int.MaxValue - 1, threadName)
+ } else {
+ ThreadUtils.newDaemonSingleThreadScheduledExecutor(threadName)
+ }
}
@VisibleForTesting
@@ -91,8 +100,12 @@ class KyuubiSyncThriftClient private (
}
private def startEngineAliveProbe(): Unit = {
- engineAliveThreadPool = ThreadUtils.newDaemonSingleThreadScheduledExecutor(
- "engine-alive-probe-" + _aliveProbeSessionHandle)
+ val threadName = "engine-alive-probe-" + _aliveProbeSessionHandle
+ engineAliveThreadPool = if (useVirtualThreadsForAliveProbe) {
+ ThreadUtils.newVirtualThreadSingleThreadScheduledExecutor(threadName)
+ } else {
+ ThreadUtils.newDaemonSingleThreadScheduledExecutor(threadName)
+ }
def closeClient(): Unit = {
warn(s"Removing Clients for ${_remoteSessionHandle}")
@@ -491,6 +504,8 @@ private[kyuubi] object KyuubiSyncThriftClient extends Logging {
val aliveProbeEnabled = conf.get(KyuubiConf.ENGINE_ALIVE_PROBE_ENABLED)
val aliveProbeInterval = conf.get(KyuubiConf.ENGINE_ALIVE_PROBE_INTERVAL).toInt
val aliveTimeout = conf.get(KyuubiConf.ENGINE_ALIVE_TIMEOUT)
+ val useVirtualThreadsForAliveProbe = conf.get(ENGINE_ALIVE_PROBE_VIRTUAL_THREADS_ENABLED)
+ val useVirtualThreadsForAsyncRequests = conf.get(ENGINE_RPC_CLIENT_VIRTUAL_THREADS_ENABLED)
val tProtocol = createTProtocol(user, passwd, host, port, 0, loginTimeout, maxMessageSize)
@@ -512,6 +527,8 @@ private[kyuubi] object KyuubiSyncThriftClient extends Logging {
tProtocol,
aliveProbeProtocol,
aliveProbeInterval,
- aliveTimeout)
+ aliveTimeout,
+ useVirtualThreadsForAliveProbe,
+ useVirtualThreadsForAsyncRequests)
}
}
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 d5deeeca967..85723bff0fd 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
@@ -18,7 +18,7 @@
package org.apache.kyuubi.engine
import java.util.Locale
-import java.util.concurrent.{ConcurrentHashMap, ScheduledExecutorService, ThreadPoolExecutor, TimeUnit}
+import java.util.concurrent.{ConcurrentHashMap, ExecutorService, ScheduledExecutorService, TimeUnit}
import scala.collection.JavaConverters._
import scala.util.control.NonFatal
@@ -82,7 +82,7 @@ class KubernetesApplicationOperation extends ApplicationOperation with Logging {
private var expireCleanUpTriggerCacheExecutor: ScheduledExecutorService = _
- private var cleanupCanceledAppPodExecutor: ThreadPoolExecutor = _
+ private var cleanupCanceledAppPodExecutor: ExecutorService = _
private def getOrCreateKubernetesClient(kubernetesInfo: KubernetesInfo): KubernetesClient = {
checkKubernetesInfo(kubernetesInfo)
@@ -173,8 +173,13 @@ class KubernetesApplicationOperation extends ApplicationOperation with Logging {
cleanupDriverPodCheckInterval,
cleanupDriverPodCheckInterval,
TimeUnit.MILLISECONDS)
- cleanupCanceledAppPodExecutor = ThreadUtils.newDaemonCachedThreadPool(
- "cleanup-canceled-app-pod-thread")
+ val useVirtualThreads =
+ conf.get(KyuubiConf.KUBERNETES_APPLICATION_CLEANUP_VIRTUAL_THREADS_ENABLED)
+ cleanupCanceledAppPodExecutor = if (useVirtualThreads) {
+ ThreadUtils.newVirtualThreadPerTaskExecutor("cleanup-canceled-app-pod-thread")
+ } else {
+ ThreadUtils.newDaemonCachedThreadPool("cleanup-canceled-app-pod-thread")
+ }
initializeKubernetesClient(kyuubiConf)
}
diff --git a/kyuubi-server/src/main/scala/org/apache/kyuubi/engine/ProcBuilder.scala b/kyuubi-server/src/main/scala/org/apache/kyuubi/engine/ProcBuilder.scala
index db083579253..fab5ed82ad5 100644
--- a/kyuubi-server/src/main/scala/org/apache/kyuubi/engine/ProcBuilder.scala
+++ b/kyuubi-server/src/main/scala/org/apache/kyuubi/engine/ProcBuilder.scala
@@ -32,7 +32,7 @@ import org.apache.kyuubi._
import org.apache.kyuubi.config.{KyuubiConf, KyuubiReservedKeys}
import org.apache.kyuubi.config.KyuubiConf.KYUUBI_HOME_ENV_VAR_NAME
import org.apache.kyuubi.operation.log.OperationLog
-import org.apache.kyuubi.util.{JavaUtils, NamedThreadFactory}
+import org.apache.kyuubi.util.{JavaUtils, NamedThreadFactory, ThreadUtils}
trait ProcBuilder {
@@ -254,7 +254,12 @@ trait ProcBuilder {
}
logCaptureThreadReleased = false
- logCaptureThread = PROC_BUILD_LOGGER.newThread(redirect)
+ logCaptureThread =
+ if (conf.get(KyuubiConf.SERVER_ENGINE_LOG_CAPTURE_VIRTUAL_THREADS_ENABLED)) {
+ PROC_BUILD_VIRTUAL_LOGGER.newThread(redirect)
+ } else {
+ PROC_BUILD_LOGGER.newThread(redirect)
+ }
logCaptureThread.start()
process
}
@@ -363,6 +368,8 @@ trait ProcBuilder {
object ProcBuilder extends Logging {
private val PROC_BUILD_LOGGER = new NamedThreadFactory("process-logger-capture", daemon = true)
+ private lazy val PROC_BUILD_VIRTUAL_LOGGER =
+ ThreadUtils.newVirtualThreadFactory("process-logger-capture")
private val UNCAUGHT_ERROR = new RuntimeException("Uncaught error")
diff --git a/kyuubi-server/src/main/scala/org/apache/kyuubi/server/KyuubiBatchService.scala b/kyuubi-server/src/main/scala/org/apache/kyuubi/server/KyuubiBatchService.scala
index 38bb999c342..ee6b5a904b3 100644
--- a/kyuubi-server/src/main/scala/org/apache/kyuubi/server/KyuubiBatchService.scala
+++ b/kyuubi-server/src/main/scala/org/apache/kyuubi/server/KyuubiBatchService.scala
@@ -19,6 +19,7 @@ package org.apache.kyuubi.server
import java.util.concurrent.atomic.AtomicBoolean
+import org.apache.kyuubi.config.KyuubiConf
import org.apache.kyuubi.config.KyuubiConf.BATCH_SUBMITTER_THREADS
import org.apache.kyuubi.engine.ApplicationState
import org.apache.kyuubi.operation.OperationState
@@ -40,8 +41,17 @@ class KyuubiBatchService(
private lazy val metadataManager: MetadataManager = sessionManager.metadataManager.get
private val running: AtomicBoolean = new AtomicBoolean(false)
- private lazy val batchExecutor = ThreadUtils
- .newDaemonFixedThreadPool(conf.get(BATCH_SUBMITTER_THREADS), "kyuubi-batch-submitter")
+ private lazy val batchExecutor = {
+ val poolSize = conf.get(BATCH_SUBMITTER_THREADS)
+ if (conf.get(KyuubiConf.BATCH_SUBMITTER_VIRTUAL_THREADS_ENABLED)) {
+ ThreadUtils.newBoundedQueuedVirtualThreadPerTaskExecutor(
+ poolSize,
+ Int.MaxValue - poolSize,
+ "kyuubi-batch-submitter")
+ } else {
+ ThreadUtils.newDaemonFixedThreadPool(poolSize, "kyuubi-batch-submitter")
+ }
+ }
def cancelUnscheduledBatch(batchId: String): Boolean = {
metadataManager.cancelUnscheduledBatch(batchId)
@@ -114,7 +124,7 @@ class KyuubiBatchService(
}
}
}
- (0 until batchExecutor.getCorePoolSize).foreach(_ => batchExecutor.submit(submitTask))
+ (0 until conf.get(BATCH_SUBMITTER_THREADS)).foreach(_ => batchExecutor.submit(submitTask))
super.start()
}
diff --git a/kyuubi-server/src/main/scala/org/apache/kyuubi/server/KyuubiRestFrontendService.scala b/kyuubi-server/src/main/scala/org/apache/kyuubi/server/KyuubiRestFrontendService.scala
index e71ac256c6b..72fa9edd00b 100644
--- a/kyuubi-server/src/main/scala/org/apache/kyuubi/server/KyuubiRestFrontendService.scala
+++ b/kyuubi-server/src/main/scala/org/apache/kyuubi/server/KyuubiRestFrontendService.scala
@@ -179,12 +179,23 @@ class KyuubiRestFrontendService(override val serverable: Serverable)
}
}
+ private def newBatchRecoveryExecutor(poolSize: Int, name: String) = {
+ if (conf.get(METADATA_RECOVERY_VIRTUAL_THREADS_ENABLED)) {
+ ThreadUtils.newBoundedQueuedVirtualThreadPerTaskExecutor(
+ poolSize,
+ Int.MaxValue - poolSize,
+ name)
+ } else {
+ ThreadUtils.newDaemonFixedThreadPool(poolSize, name)
+ }
+ }
+
@VisibleForTesting
private[kyuubi] def recoverBatchSessions(): Unit = withBatchRecoveryLockRequired {
val recoveryNumThreads = conf.get(METADATA_RECOVERY_THREADS)
val recoveryWaitEngineSubmission = conf.get(METADATA_RECOVERY_WAIT_ENGINE_SUBMISSION)
val batchRecoveryExecutor =
- ThreadUtils.newDaemonFixedThreadPool(recoveryNumThreads, "batch-recovery-executor")
+ newBatchRecoveryExecutor(recoveryNumThreads, "batch-recovery-executor")
try {
val batchSessionsToRecover = sessionManager.getBatchSessionsToRecover(connectionUrl)
val pendingRecoveryTasksCount = new AtomicInteger(0)
@@ -234,7 +245,7 @@ class KyuubiRestFrontendService(override val serverable: Serverable)
withBatchRecoveryLockRequired {
val recoveryNumThreads = conf.get(METADATA_RECOVERY_THREADS)
val batchRecoveryExecutor =
- ThreadUtils.newDaemonFixedThreadPool(recoveryNumThreads, "batch-reassign-recovery-executor")
+ newBatchRecoveryExecutor(recoveryNumThreads, "batch-reassign-recovery-executor")
try {
val batchSessionsToRecover =
sessionManager.getSpecificBatchSessionsToRecover(batchIds, connectionUrl)
diff --git a/kyuubi-server/src/main/scala/org/apache/kyuubi/server/KyuubiServer.scala b/kyuubi-server/src/main/scala/org/apache/kyuubi/server/KyuubiServer.scala
index 5e7405be238..20bc01279eb 100644
--- a/kyuubi-server/src/main/scala/org/apache/kyuubi/server/KyuubiServer.scala
+++ b/kyuubi-server/src/main/scala/org/apache/kyuubi/server/KyuubiServer.scala
@@ -43,6 +43,7 @@ object KyuubiServer extends Logging {
private var commandArgs: Array[String] = Array.empty[String]
def startServer(conf: KyuubiConf): KyuubiServer = {
+ conf.validateServerVirtualThreadConfigs()
hadoopConf = KyuubiHadoopUtils.newHadoopConf(conf)
var embeddedZkServer: Option[EmbeddedZookeeper] = None
if (!ServiceDiscovery.supportServiceDiscovery(conf)) {
diff --git a/kyuubi-server/src/main/scala/org/apache/kyuubi/server/api/v1/DataAgentResource.scala b/kyuubi-server/src/main/scala/org/apache/kyuubi/server/api/v1/DataAgentResource.scala
index 60aea983c00..1788a4dc59e 100644
--- a/kyuubi-server/src/main/scala/org/apache/kyuubi/server/api/v1/DataAgentResource.scala
+++ b/kyuubi-server/src/main/scala/org/apache/kyuubi/server/api/v1/DataAgentResource.scala
@@ -162,7 +162,8 @@ private[v1] class DataAgentResource extends ApiRequestContext with Logging {
try {
val f = CompletableFuture.supplyAsync(
() => client.executeStatement(text, confOverlay, true, 0L),
- opSubmitter)
+ opSubmitter(fe.getConf.get(
+ KyuubiConf.FRONTEND_DATA_AGENT_OPERATION_SUBMIT_VIRTUAL_THREADS_ENABLED)))
f.whenComplete((handle, _) => {
if (handle != null && timedOut.get() && closed.compareAndSet(false, true)) {
info(s"Closing orphaned op for session $sessionHandleStr (servlet already timed out)")
@@ -458,21 +459,29 @@ private[v1] class DataAgentResource extends ApiRequestContext with Logging {
private[server] object DataAgentResource {
// Bounded pool for blocking executeStatement submissions; rebuildable after service restart.
- @volatile private var opSubmitExecutor: ExecutorService = newOpSubmitExecutor()
-
- private def newOpSubmitExecutor(): ExecutorService =
- ThreadUtils.newDaemonQueuedThreadPool(
- poolSize = 8,
- poolQueueSize = 64,
- keepAliveMs = 60000L,
- threadPoolName = "data-agent-op-submit")
+ @volatile private var opSubmitExecutor: ExecutorService = _
+
+ private def newOpSubmitExecutor(useVirtualThreads: Boolean): ExecutorService = {
+ if (useVirtualThreads) {
+ ThreadUtils.newBoundedQueuedVirtualThreadPerTaskExecutor(
+ maxConcurrentTasks = 8,
+ maxQueuedTasks = 64,
+ threadNamePrefix = "data-agent-op-submit")
+ } else {
+ ThreadUtils.newDaemonQueuedThreadPool(
+ poolSize = 8,
+ poolQueueSize = 64,
+ keepAliveMs = 60000L,
+ threadPoolName = "data-agent-op-submit")
+ }
+ }
- private def opSubmitter: ExecutorService = {
+ private def opSubmitter(useVirtualThreads: Boolean): ExecutorService = {
val current = opSubmitExecutor
if (current != null && !current.isShutdown) current
else synchronized {
if (opSubmitExecutor == null || opSubmitExecutor.isShutdown) {
- opSubmitExecutor = newOpSubmitExecutor()
+ opSubmitExecutor = newOpSubmitExecutor(useVirtualThreads)
}
opSubmitExecutor
}
diff --git a/kyuubi-server/src/main/scala/org/apache/kyuubi/server/metadata/MetadataManager.scala b/kyuubi-server/src/main/scala/org/apache/kyuubi/server/metadata/MetadataManager.scala
index c5182979ee7..567dccfbc94 100644
--- a/kyuubi-server/src/main/scala/org/apache/kyuubi/server/metadata/MetadataManager.scala
+++ b/kyuubi-server/src/main/scala/org/apache/kyuubi/server/metadata/MetadataManager.scala
@@ -17,7 +17,7 @@
package org.apache.kyuubi.server.metadata
-import java.util.concurrent.{ConcurrentHashMap, ThreadPoolExecutor, TimeUnit}
+import java.util.concurrent.{ConcurrentHashMap, ExecutorService, TimeUnit}
import java.util.concurrent.atomic.AtomicInteger
import scala.collection.JavaConverters._
@@ -59,10 +59,17 @@ class MetadataManager extends AbstractService("MetadataManager") {
private lazy val requestsAsyncRetryTrigger =
ThreadUtils.newDaemonSingleThreadScheduledExecutor("metadata-requests-async-retry-trigger")
- private lazy val requestsAsyncRetryExecutor: ThreadPoolExecutor =
- ThreadUtils.newDaemonFixedThreadPool(
- conf.get(KyuubiConf.METADATA_REQUEST_ASYNC_RETRY_THREADS),
- "metadata-requests-async-retry")
+ private lazy val requestsAsyncRetryExecutor: ExecutorService = {
+ val poolSize = conf.get(KyuubiConf.METADATA_REQUEST_ASYNC_RETRY_THREADS)
+ if (conf.get(KyuubiConf.METADATA_REQUEST_ASYNC_RETRY_VIRTUAL_THREADS_ENABLED)) {
+ ThreadUtils.newBoundedQueuedVirtualThreadPerTaskExecutor(
+ poolSize,
+ Int.MaxValue - poolSize,
+ "metadata-requests-async-retry")
+ } else {
+ ThreadUtils.newDaemonFixedThreadPool(poolSize, "metadata-requests-async-retry")
+ }
+ }
private lazy val cleanerEnabled = conf.get(KyuubiConf.METADATA_CLEANER_ENABLED)
diff --git a/kyuubi-server/src/main/scala/org/apache/kyuubi/util/KubernetesUtils.scala b/kyuubi-server/src/main/scala/org/apache/kyuubi/util/KubernetesUtils.scala
index e3bfd13aea6..de9a38d14b9 100644
--- a/kyuubi-server/src/main/scala/org/apache/kyuubi/util/KubernetesUtils.scala
+++ b/kyuubi-server/src/main/scala/org/apache/kyuubi/util/KubernetesUtils.scala
@@ -96,8 +96,13 @@ object KubernetesUtils extends Logging {
}.build()
// https://github.com/fabric8io/kubernetes-client/issues/3547
- val dispatcher = new Dispatcher(
- ThreadUtils.newDaemonCachedThreadPool("kubernetes-dispatcher"))
+ val dispatcherExecutor =
+ if (conf.get(KUBERNETES_CLIENT_DISPATCHER_VIRTUAL_THREADS_ENABLED)) {
+ ThreadUtils.newVirtualThreadPerTaskExecutor("kubernetes-dispatcher")
+ } else {
+ ThreadUtils.newDaemonCachedThreadPool("kubernetes-dispatcher")
+ }
+ val dispatcher = new Dispatcher(dispatcherExecutor)
val factoryWithCustomDispatcher = new OkHttpClientFactory() {
override protected def additionalConfig(builder: OkHttpClient.Builder): Unit = {
builder.dispatcher(dispatcher)