This is an automated email from the ASF dual-hosted git repository.
chengpan pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/kyuubi.git
The following commit(s) were added to refs/heads/master by this push:
new 7562a9711 [KYUUBI #6156] Remove `flink.` prefix for create session
configurations
7562a9711 is described below
commit 7562a9711506ee922729a73336ff6c7b17c21904
Author: wforget <[email protected]>
AuthorDate: Tue Mar 12 19:39:38 2024 +0800
[KYUUBI #6156] Remove `flink.` prefix for create session configurations
# :mag: Description
## Issue References ๐
This pull request fixes #6156
## Describe Your Solution ๐ง
Remove `flink.` prefix for open flink session configurations.
## Types of changes :bookmark:
- [x] Bugfix (non-breaking change which fixes an issue)
- [ ] New feature (non-breaking change which adds functionality)
- [ ] Breaking change (fix or feature that would cause existing
functionality to change)
## Test Plan ๐งช
#### Behavior Without This Pull Request :coffin:
#### Behavior With This Pull Request :tada:
#### Related Unit Tests
---
# Checklist ๐
- [X] This patch was not authored or co-authored using [Generative
Tooling](https://www.apache.org/legal/generative-tooling.html)
**Be nice. Be informative.**
Closes #6157 from wForget/KYUUBI-6156.
Closes #6156
fc750dc5a [wforget] comment
f6134919c [wforget] Remove `flink.` prefix for create session configurations
Authored-by: wforget <[email protected]>
Signed-off-by: Cheng Pan <[email protected]>
---
.../flink/session/FlinkSQLSessionManager.scala | 35 +++++++++++-----------
1 file changed, 18 insertions(+), 17 deletions(-)
diff --git
a/externals/kyuubi-flink-sql-engine/src/main/scala/org/apache/kyuubi/engine/flink/session/FlinkSQLSessionManager.scala
b/externals/kyuubi-flink-sql-engine/src/main/scala/org/apache/kyuubi/engine/flink/session/FlinkSQLSessionManager.scala
index bbdcdc437..73908bef8 100644
---
a/externals/kyuubi-flink-sql-engine/src/main/scala/org/apache/kyuubi/engine/flink/session/FlinkSQLSessionManager.scala
+++
b/externals/kyuubi-flink-sql-engine/src/main/scala/org/apache/kyuubi/engine/flink/session/FlinkSQLSessionManager.scala
@@ -54,23 +54,24 @@ class FlinkSQLSessionManager(engineContext: DefaultContext)
password: String,
ipAddress: String,
conf: Map[String, String]): Session = {
- conf.get(KYUUBI_SESSION_HANDLE_KEY).map(SessionHandle.fromUUID).flatMap(
- getSessionOption).getOrElse {
- val flinkInternalSession = sessionManager.openSession(
- SessionEnvironment.newBuilder
- .setSessionEndpointVersion(SqlGatewayRestAPIVersion.V1)
- .addSessionConfig(mapAsJavaMap(conf))
- .build)
- val session = new FlinkSessionImpl(
- protocol,
- user,
- password,
- ipAddress,
- conf,
- this,
- flinkInternalSession)
- session
- }
+ val normalizedConf = conf.map { case (k, v) => k.stripPrefix("flink.") ->
v }
+ normalizedConf.get(KYUUBI_SESSION_HANDLE_KEY).map(SessionHandle.fromUUID)
+ .flatMap(getSessionOption).getOrElse {
+ val flinkInternalSession = sessionManager.openSession(
+ SessionEnvironment.newBuilder
+ .setSessionEndpointVersion(SqlGatewayRestAPIVersion.V1)
+ .addSessionConfig(mapAsJavaMap(normalizedConf))
+ .build)
+ val session = new FlinkSessionImpl(
+ protocol,
+ user,
+ password,
+ ipAddress,
+ normalizedConf,
+ this,
+ flinkInternalSession)
+ session
+ }
}
override def getSessionOption(sessionHandle: SessionHandle): Option[Session]
= {