This is an automated email from the ASF dual-hosted git repository.
rexxiong pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/incubator-celeborn.git
The following commit(s) were added to refs/heads/main by this push:
new d89dcf0e0 [CELEBORN-1054] Support db based dynamic config service
d89dcf0e0 is described below
commit d89dcf0e064e2844607f8691f49620e334421ef8
Author: Shuang <[email protected]>
AuthorDate: Mon Feb 5 13:23:25 2024 +0800
[CELEBORN-1054] Support db based dynamic config service
### What changes were proposed in this pull request?
Support database based store backend implementation for dynamic
configuration management
### Why are the changes needed?
Currently celeborn provides `FsConfigServiceImpl` implementation for
dynamic config service which is based on file system, We cloud Support database
based store backend implementation.
### Does this PR introduce _any_ user-facing change?
No
### How was this patch tested?
- `ConfigServiceSuiteJ#testDbConfig`
Closes #2273 from RexXiong/CELEBORN-1054.
Authored-by: Shuang <[email protected]>
Signed-off-by: Shuang <[email protected]>
---
LICENSE-binary | 2 +
NOTICE-binary | 72 +++++++++++
build/make-distribution.sh | 4 +
.../org/apache/celeborn/common/CelebornConf.scala | 111 ++++++++++++++++-
dev/deps/dependencies-server | 2 +
docs/configuration/master.md | 13 +-
docs/configuration/worker.md | 13 +-
pom.xml | 23 ++++
project/CelebornBuild.scala | 11 +-
service/pom.xml | 15 +++
.../service/config/BaseConfigServiceImpl.java | 88 +++++++++++++
.../common/service/config/ConfigService.java | 18 ++-
.../common/service/config/DbConfigServiceImpl.java | 54 ++++++++
.../config/DynamicConfigServiceFactory.java | 9 +-
.../common/service/config/FsConfigServiceImpl.java | 48 ++-----
.../server/common/service/config/SystemConfig.java | 7 ++
.../server/common/service/config/TenantConfig.java | 25 +++-
.../server/common/service/model/ClusterInfo.java | 78 ++++++++++++
.../common/service/model/ClusterSystemConfig.java | 86 +++++++++++++
.../common/service/model/ClusterTenantConfig.java | 138 +++++++++++++++++++++
.../IServiceManager.java} | 25 ++--
.../common/service/store/db/DBSessionFactory.java | 95 ++++++++++++++
.../service/store/db/DbServiceManagerImpl.java | 138 +++++++++++++++++++++
.../db/mapper/ClusterInfoMapper.java} | 29 ++---
.../db/mapper/ClusterSystemConfigMapper.java} | 24 ++--
.../store/db/mapper/ClusterTenantConfigMapper.java | 42 +++++++
.../server/common/service/utils/JsonUtils.java | 50 ++++++++
.../resources/sql/mysql/celeborn-0.5.0-mysql.sql | 56 +++++++++
.../common/service/config/ConfigServiceSuiteJ.java | 102 ++++++++++-----
.../test/resources/celeborn-0.5.0-h2-ut-data.sql | 37 ++++++
service/src/test/resources/celeborn-0.5.0-h2.sql | 64 ++++++++++
31 files changed, 1346 insertions(+), 133 deletions(-)
diff --git a/LICENSE-binary b/LICENSE-binary
index 945016f5b..a33986495 100644
--- a/LICENSE-binary
+++ b/LICENSE-binary
@@ -214,6 +214,7 @@ com.thoughtworks.paranamer:paranamer
commons-cli:commons-cli
commons-io:commons-io
commons-logging:commons-logging
+com.zaxxer:HikariCP
io.dropwizard.metrics:metrics-core
io.dropwizard.metrics:metrics-graphite
io.dropwizard.metrics:metrics-jvm
@@ -251,6 +252,7 @@ org.apache.commons:commons-crypto
org.apache.commons:commons-lang3
org.apache.hadoop:hadoop-client-api
org.apache.hadoop:hadoop-client-runtime
+org.apache.ibatis:mybatis
org.apache.logging.log4j:log4j-1.2-api
org.apache.logging.log4j:log4j-api
org.apache.logging.log4j:log4j-core
diff --git a/NOTICE-binary b/NOTICE-binary
index 8e59fff41..4a3bb196f 100644
--- a/NOTICE-binary
+++ b/NOTICE-binary
@@ -123,3 +123,75 @@ This software includes projects with other licenses -- see
`doc/LICENSE.md`.
This product includes software developed by Google
Snappy: http://code.google.com/p/snappy/ (New BSD License)
+
+MyBatis
+Copyright 2010-2023
+
+This product includes software developed by
+The MyBatis Team (https://www.mybatis.org/).
+
+iBATIS
+Copyright 2010 The Apache Software Foundation
+
+Licensed under the Apache License, Version 2.0 (the "License");
+you may not use this file except in compliance with the License.
+You may obtain a copy of the License at
+
+ https://www.apache.org/licenses/LICENSE-2.0
+
+Unless required by applicable law or agreed to in writing, software
+distributed under the License is distributed on an "AS IS" BASIS,
+WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+See the License for the specific language governing permissions and
+limitations under the License.
+
+OGNL
+//--------------------------------------------------------------------------
+// Copyright (c) 2004, Drew Davidson and Luke Blanshard
+// All rights reserved.
+//
+// Redistribution and use in source and binary forms, with or without
+// modification, are permitted provided that the following conditions are
+// met:
+//
+// Redistributions of source code must retain the above copyright notice,
+// this list of conditions and the following disclaimer.
+// Redistributions in binary form must reproduce the above copyright
+// notice, this list of conditions and the following disclaimer in the
+// documentation and/or other materials provided with the distribution.
+// Neither the name of the Drew Davidson nor the names of its contributors
+// may be used to endorse or promote products derived from this software
+// without specific prior written permission.
+//
+// THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS
+// "AS IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT
+// LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS
+// FOR A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE
+// COPYRIGHT OWNER OR CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT,
+// INCIDENTAL, SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING,
+// BUT NOT LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS
+// OF USE, DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED
+// AND ON ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY,
+// OR TORT (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF
+// THE USE OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH
+// DAMAGE.
+//--------------------------------------------------------------------------
+
+Refactored SqlBuilder class (SQL, AbstractSQL)
+
+This product includes software developed by
+Adam Gent (https://gist.github.com/3650165)
+
+Copyright 2010 Adam Gent
+
+Licensed under the Apache License, Version 2.0 (the "License");
+you may not use this file except in compliance with the License.
+You may obtain a copy of the License at
+
+ https://www.apache.org/licenses/LICENSE-2.0
+
+Unless required by applicable law or agreed to in writing, software
+distributed under the License is distributed on an "AS IS" BASIS,
+WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+See the License for the specific language governing permissions and
+limitations under the License.
diff --git a/build/make-distribution.sh b/build/make-distribution.sh
index b063942ca..962245b0a 100755
--- a/build/make-distribution.sh
+++ b/build/make-distribution.sh
@@ -382,6 +382,10 @@ cp "$PROJECT_DIR"/conf/*.template "$DIST_DIR/conf"
cp -r "$PROJECT_DIR/bin" "$DIST_DIR"
cp -r "$PROJECT_DIR/sbin" "$DIST_DIR"
+# Copy db scripts
+mkdir "$DIST_DIR/db-scripts"
+cp -r "$PROJECT_DIR/service/src/main/resources/sql/" "$DIST_DIR/db-scripts"
+
# Copy container related resources
mkdir "$DIST_DIR/docker"
cp "$PROJECT_DIR/docker/Dockerfile" "$DIST_DIR/docker"
diff --git
a/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala
b/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala
index 240fb24f9..822fe5459 100644
--- a/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala
+++ b/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala
@@ -367,7 +367,23 @@ class CelebornConf(loadDefaults: Boolean) extends
Cloneable with Logging with Se
}
def dynamicConfigStoreBackend: String = get(DYNAMIC_CONFIG_STORE_BACKEND)
+ def dynamicConfigEnabled: Boolean = get(DYNAMIC_CONFIG_ENABLED)
def dynamicConfigRefreshInterval: Long = get(DYNAMIC_CONFIG_REFRESH_INTERVAL)
+ def dynamicConfigStoreDbFetchPageSize: Int =
get(DYNAMIC_CONFIG_STORE_DB_FETCH_PAGE_SIZE)
+ def dynamicConfigStoreDbHikariDriverClassName: String =
+ get(DYNAMIC_CONFIG_STORE_DB_HIKARI_DRIVER_CLASS_NAME)
+ def dynamicConfigStoreDbHikariJdbcUrl: String =
get(DYNAMIC_CONFIG_STORE_DB_HIKARI_JDBC_URL)
+ def dynamicConfigStoreDbHikariUsername: String =
get(DYNAMIC_CONFIG_STORE_DB_HIKARI_USERNAME)
+ def dynamicConfigStoreDbHikariPassword: String =
get(DYNAMIC_CONFIG_STORE_DB_HIKARI_PASSWORD)
+ def dynamicConfigStoreDbHikariConnectionTimeout: Long =
+ get(DYNAMIC_CONFIG_STORE_DB_HIKARI_CONNECTION_TIMEOUT)
+ def dynamicConfigStoreDbHikariIdleTimeout: Long =
get(DYNAMIC_CONFIG_STORE_DB_HIKARI_IDLE_TIMEOUT)
+ def dynamicConfigStoreDbHikariMaxLifetime: Long =
get(DYNAMIC_CONFIG_STORE_DB_HIKARI_MAX_LIFETIME)
+ def dynamicConfigStoreDbHikariMaximumPoolSize: Int =
+ get(DYNAMIC_CONFIG_STORE_DB_HIKARI_MAXIMUM_POOL_SIZE)
+ def dynamicConfigStoreDbHikariCustomConfigs: JMap[String, String] = {
+
settings.asScala.filter(_._1.startsWith("celeborn.dynamicConfig.store.db.hikari")).toMap.asJava
+ }
// //////////////////////////////////////////////////////
// Network //
@@ -543,6 +559,7 @@ class CelebornConf(loadDefaults: Boolean) extends Cloneable
with Logging with Se
def estimatedPartitionSizeForEstimationUpdateInterval: Long =
get(ESTIMATED_PARTITION_SIZE_UPDATE_INTERVAL)
def masterResourceConsumptionInterval: Long =
get(MASTER_RESOURCE_CONSUMPTION_INTERVAL)
+ def clusterName: String = get(CLUSTER_NAME)
// //////////////////////////////////////////////////////
// Address && HA && RATIS //
@@ -2206,6 +2223,14 @@ object CelebornConf extends Logging {
.timeConf(TimeUnit.MILLISECONDS)
.createWithDefaultString("30s")
+ val CLUSTER_NAME: ConfigEntry[String] =
+ buildConf("celeborn.cluster.name")
+ .categories("master", "worker")
+ .version("0.5.0")
+ .doc("Celeborn cluster name.")
+ .stringConf
+ .createWithDefaultString("default")
+
val SHUFFLE_CHUNK_SIZE: ConfigEntry[Long] =
buildConf("celeborn.shuffle.chunk.size")
.categories("worker")
@@ -4361,12 +4386,20 @@ object CelebornConf extends Logging {
val DYNAMIC_CONFIG_STORE_BACKEND: ConfigEntry[String] =
buildConf("celeborn.dynamicConfig.store.backend")
.categories("master", "worker")
- .doc("Store backend for dynamic config. Available options: NONE, FS.
Note: NONE means disabling dynamic config store.")
+ .doc("Store backend for dynamic config service. Available options: FS,
DB.")
.version("0.4.0")
.stringConf
.transform(_.toUpperCase(Locale.ROOT))
- .checkValues(Set("NONE", "FS"))
- .createWithDefault("NONE")
+ .checkValues(Set("FS", "DB"))
+ .createWithDefault("FS")
+
+ val DYNAMIC_CONFIG_ENABLED: ConfigEntry[Boolean] =
+ buildConf("celeborn.dynamicConfig.enabled")
+ .categories("master", "worker")
+ .version("0.5.0")
+ .doc("Whether to enable dynamic configuration.")
+ .booleanConf
+ .createWithDefault(false)
val DYNAMIC_CONFIG_REFRESH_INTERVAL: ConfigEntry[Long] =
buildConf("celeborn.dynamicConfig.refresh.interval")
@@ -4376,6 +4409,78 @@ object CelebornConf extends Logging {
.timeConf(TimeUnit.MILLISECONDS)
.createWithDefaultString("120s")
+ val DYNAMIC_CONFIG_STORE_DB_FETCH_PAGE_SIZE: ConfigEntry[Int] =
+ buildConf("celeborn.dynamicConfig.store.db.fetch.pageSize")
+ .categories("master", "worker")
+ .version("0.5.0")
+ .doc("The page size for db store to query configurations.")
+ .intConf
+ .createWithDefaultString("1000")
+
+ val DYNAMIC_CONFIG_STORE_DB_HIKARI_DRIVER_CLASS_NAME: ConfigEntry[String] =
+ buildConf("celeborn.dynamicConfig.store.db.hikari.driverClassName")
+ .categories("master", "worker")
+ .version("0.5.0")
+ .doc("The jdbc driver class name of db store backend.")
+ .stringConf
+ .createWithDefaultString("")
+
+ val DYNAMIC_CONFIG_STORE_DB_HIKARI_JDBC_URL: ConfigEntry[String] =
+ buildConf("celeborn.dynamicConfig.store.db.hikari.jdbcUrl")
+ .categories("master", "worker")
+ .version("0.5.0")
+ .doc("The jdbc url of db store backend.")
+ .stringConf
+ .createWithDefaultString("")
+
+ val DYNAMIC_CONFIG_STORE_DB_HIKARI_USERNAME: ConfigEntry[String] =
+ buildConf("celeborn.dynamicConfig.store.db.hikari.username")
+ .categories("master", "worker")
+ .version("0.5.0")
+ .doc("The username of db store backend.")
+ .stringConf
+ .createWithDefaultString("")
+
+ val DYNAMIC_CONFIG_STORE_DB_HIKARI_PASSWORD: ConfigEntry[String] =
+ buildConf("celeborn.dynamicConfig.store.db.hikari.password")
+ .categories("master", "worker")
+ .version("0.5.0")
+ .doc("The password of db store backend.")
+ .stringConf
+ .createWithDefaultString("")
+
+ val DYNAMIC_CONFIG_STORE_DB_HIKARI_CONNECTION_TIMEOUT: ConfigEntry[Long] =
+ buildConf("celeborn.dynamicConfig.store.db.hikari.connectionTimeout")
+ .categories("master", "worker")
+ .version("0.5.0")
+ .doc("The connection timeout that a client will wait for a connection
from the pool for db store backend.")
+ .timeConf(TimeUnit.MILLISECONDS)
+ .createWithDefaultString("30s")
+
+ val DYNAMIC_CONFIG_STORE_DB_HIKARI_IDLE_TIMEOUT: ConfigEntry[Long] =
+ buildConf("celeborn.dynamicConfig.store.db.hikari.idleTimeout")
+ .categories("master", "worker")
+ .version("0.5.0")
+ .doc("The idle timeout that a connection is allowed to sit idle in the
pool for db store backend.")
+ .timeConf(TimeUnit.MILLISECONDS)
+ .createWithDefaultString("600s")
+
+ val DYNAMIC_CONFIG_STORE_DB_HIKARI_MAX_LIFETIME: ConfigEntry[Long] =
+ buildConf("celeborn.dynamicConfig.store.db.hikari.maxLifetime")
+ .categories("master", "worker")
+ .version("0.5.0")
+ .doc("The maximum lifetime of a connection in the pool for db store
backend.")
+ .timeConf(TimeUnit.MILLISECONDS)
+ .createWithDefaultString("1800s")
+
+ val DYNAMIC_CONFIG_STORE_DB_HIKARI_MAXIMUM_POOL_SIZE: ConfigEntry[Int] =
+ buildConf("celeborn.dynamicConfig.store.db.hikari.maximumPoolSize")
+ .categories("master", "worker")
+ .version("0.5.0")
+ .doc("The maximum pool size of db store backend.")
+ .intConf
+ .createWithDefaultString("2")
+
val REGISTER_SHUFFLE_FILTER_EXCLUDED_WORKER_ENABLED: ConfigEntry[Boolean] =
buildConf("celeborn.client.shuffle.register.filterExcludedWorker.enabled")
.categories("client")
diff --git a/dev/deps/dependencies-server b/dev/deps/dependencies-server
index a8324cabe..f507f96dd 100644
--- a/dev/deps/dependencies-server
+++ b/dev/deps/dependencies-server
@@ -15,6 +15,7 @@
# limitations under the License.
#
+HikariCP/4.0.3//HikariCP-4.0.3.jar
RoaringBitmap/0.9.32//RoaringBitmap-0.9.32.jar
commons-cli/1.5.0//commons-cli-1.5.0.jar
commons-crypto/1.0.0//commons-crypto-1.0.0.jar
@@ -44,6 +45,7 @@ maven-jdk-tools-wrapper/0.1//maven-jdk-tools-wrapper-0.1.jar
metrics-core/3.2.6//metrics-core-3.2.6.jar
metrics-graphite/3.2.6//metrics-graphite-3.2.6.jar
metrics-jvm/3.2.6//metrics-jvm-3.2.6.jar
+mybatis/3.5.15//mybatis-3.5.15.jar
netty-all/4.1.101.Final//netty-all-4.1.101.Final.jar
netty-buffer/4.1.101.Final//netty-buffer-4.1.101.Final.jar
netty-codec-dns/4.1.101.Final//netty-codec-dns-4.1.101.Final.jar
diff --git a/docs/configuration/master.md b/docs/configuration/master.md
index 65d9079e3..091c362f5 100644
--- a/docs/configuration/master.md
+++ b/docs/configuration/master.md
@@ -19,8 +19,19 @@ license: |
<!--begin-include-->
| Key | Default | Description | Since | Deprecated |
| --- | ------- | ----------- | ----- | ---------- |
+| celeborn.cluster.name | default | Celeborn cluster name. | 0.5.0 | |
+| celeborn.dynamicConfig.enabled | false | Whether to enable dynamic
configuration. | 0.5.0 | |
| celeborn.dynamicConfig.refresh.interval | 120s | Interval for refreshing the
corresponding dynamic config periodically. | 0.4.0 | |
-| celeborn.dynamicConfig.store.backend | NONE | Store backend for dynamic
config. Available options: NONE, FS. Note: NONE means disabling dynamic config
store. | 0.4.0 | |
+| celeborn.dynamicConfig.store.backend | FS | Store backend for dynamic config
service. Available options: FS, DB. | 0.4.0 | |
+| celeborn.dynamicConfig.store.db.fetch.pageSize | 1000 | The page size for db
store to query configurations. | 0.5.0 | |
+| celeborn.dynamicConfig.store.db.hikari.connectionTimeout | 30s | The
connection timeout that a client will wait for a connection from the pool for
db store backend. | 0.5.0 | |
+| celeborn.dynamicConfig.store.db.hikari.driverClassName | | The jdbc driver
class name of db store backend. | 0.5.0 | |
+| celeborn.dynamicConfig.store.db.hikari.idleTimeout | 600s | The idle timeout
that a connection is allowed to sit idle in the pool for db store backend. |
0.5.0 | |
+| celeborn.dynamicConfig.store.db.hikari.jdbcUrl | | The jdbc url of db store
backend. | 0.5.0 | |
+| celeborn.dynamicConfig.store.db.hikari.maxLifetime | 1800s | The maximum
lifetime of a connection in the pool for db store backend. | 0.5.0 | |
+| celeborn.dynamicConfig.store.db.hikari.maximumPoolSize | 2 | The maximum
pool size of db store backend. | 0.5.0 | |
+| celeborn.dynamicConfig.store.db.hikari.password | | The password of db
store backend. | 0.5.0 | |
+| celeborn.dynamicConfig.store.db.hikari.username | | The username of db
store backend. | 0.5.0 | |
| celeborn.internal.port.enabled | false | Whether to create a internal port
on Masters/Workers for inter-Masters/Workers communication. This is beneficial
when SASL authentication is enforced for all interactions between clients and
Celeborn Services, but the services can exchange messages without being subject
to SASL authentication. | 0.5.0 | |
| celeborn.master.estimatedPartitionSize.initialSize | 64mb | Initial
partition size for estimation, it will change according to runtime stats. |
0.3.0 | celeborn.shuffle.initialEstimatedPartitionSize |
| celeborn.master.estimatedPartitionSize.update.initialDelay | 5min | Initial
delay time before start updating partition size for estimation. | 0.3.0 |
celeborn.shuffle.estimatedPartitionSize.update.initialDelay |
diff --git a/docs/configuration/worker.md b/docs/configuration/worker.md
index 7cef6bdbe..c792840f0 100644
--- a/docs/configuration/worker.md
+++ b/docs/configuration/worker.md
@@ -19,8 +19,19 @@ license: |
<!--begin-include-->
| Key | Default | Description | Since | Deprecated |
| --- | ------- | ----------- | ----- | ---------- |
+| celeborn.cluster.name | default | Celeborn cluster name. | 0.5.0 | |
+| celeborn.dynamicConfig.enabled | false | Whether to enable dynamic
configuration. | 0.5.0 | |
| celeborn.dynamicConfig.refresh.interval | 120s | Interval for refreshing the
corresponding dynamic config periodically. | 0.4.0 | |
-| celeborn.dynamicConfig.store.backend | NONE | Store backend for dynamic
config. Available options: NONE, FS. Note: NONE means disabling dynamic config
store. | 0.4.0 | |
+| celeborn.dynamicConfig.store.backend | FS | Store backend for dynamic config
service. Available options: FS, DB. | 0.4.0 | |
+| celeborn.dynamicConfig.store.db.fetch.pageSize | 1000 | The page size for db
store to query configurations. | 0.5.0 | |
+| celeborn.dynamicConfig.store.db.hikari.connectionTimeout | 30s | The
connection timeout that a client will wait for a connection from the pool for
db store backend. | 0.5.0 | |
+| celeborn.dynamicConfig.store.db.hikari.driverClassName | | The jdbc driver
class name of db store backend. | 0.5.0 | |
+| celeborn.dynamicConfig.store.db.hikari.idleTimeout | 600s | The idle timeout
that a connection is allowed to sit idle in the pool for db store backend. |
0.5.0 | |
+| celeborn.dynamicConfig.store.db.hikari.jdbcUrl | | The jdbc url of db store
backend. | 0.5.0 | |
+| celeborn.dynamicConfig.store.db.hikari.maxLifetime | 1800s | The maximum
lifetime of a connection in the pool for db store backend. | 0.5.0 | |
+| celeborn.dynamicConfig.store.db.hikari.maximumPoolSize | 2 | The maximum
pool size of db store backend. | 0.5.0 | |
+| celeborn.dynamicConfig.store.db.hikari.password | | The password of db
store backend. | 0.5.0 | |
+| celeborn.dynamicConfig.store.db.hikari.username | | The username of db
store backend. | 0.5.0 | |
| celeborn.internal.port.enabled | false | Whether to create a internal port
on Masters/Workers for inter-Masters/Workers communication. This is beneficial
when SASL authentication is enforced for all interactions between clients and
Celeborn Services, but the services can exchange messages without being subject
to SASL authentication. | 0.5.0 | |
| celeborn.master.endpoints | <localhost>:9097 | Endpoints of master
nodes for celeborn client to connect, allowed pattern is:
`<host1>:<port1>[,<host2>:<port2>]*`, e.g. `clb1:9097,clb2:9098,clb3:9099`. If
the port is omitted, 9097 will be used. | 0.2.0 | |
| celeborn.master.estimatedPartitionSize.minSize | 8mb | Ignore partition size
smaller than this configuration of partition size for estimation. | 0.3.0 |
celeborn.shuffle.minPartitionSizeToEstimate |
diff --git a/pom.xml b/pom.xml
index a0af02c6b..a84445057 100644
--- a/pom.xml
+++ b/pom.xml
@@ -100,6 +100,11 @@
<jackson.version>2.15.3</jackson.version>
<snappy.version>1.1.10.5</snappy.version>
+ <!-- Db dependencies -->
+ <mybatis.version>3.5.15</mybatis.version>
+ <hikaricp.version>4.0.3</hikaricp.version>
+ <h2.version>2.2.224</h2.version>
+
<shading.prefix>org.apache.celeborn.shaded</shading.prefix>
<maven.plugin.antrun.version>3.0.0</maven.plugin.antrun.version>
@@ -426,6 +431,24 @@
</exclusions>
</dependency>
+ <!-- Db dependencies -->
+ <dependency>
+ <groupId>org.mybatis</groupId>
+ <artifactId>mybatis</artifactId>
+ <version>${mybatis.version}</version>
+ </dependency>
+ <dependency>
+ <groupId>com.zaxxer</groupId>
+ <artifactId>HikariCP</artifactId>
+ <version>${hikaricp.version}</version>
+ </dependency>
+ <dependency>
+ <groupId>com.h2database</groupId>
+ <artifactId>h2</artifactId>
+ <version>${h2.version}</version>
+ <scope>test</scope>
+ </dependency>
+
<dependency>
<groupId>junit</groupId>
<artifactId>junit</artifactId>
diff --git a/project/CelebornBuild.scala b/project/CelebornBuild.scala
index 4c3f34626..a280b922c 100644
--- a/project/CelebornBuild.scala
+++ b/project/CelebornBuild.scala
@@ -63,6 +63,9 @@ object Dependencies {
val slf4jVersion = "1.7.36"
val snakeyamlVersion = "2.2"
val snappyVersion = "1.1.10.5"
+ val mybatisVersion = "3.5.15"
+ val hikaricpVersion = "4.0.3"
+ val h2Version = "2.2.224"
// Versions for proto
val protocVersion = "3.21.7"
@@ -123,6 +126,8 @@ object Dependencies {
val snakeyaml = "org.yaml" % "snakeyaml" % snakeyamlVersion
val snappyJava = "org.xerial.snappy" % "snappy-java" % snappyVersion
val zstdJni = "com.github.luben" % "zstd-jni" % zstdJniVersion
+ val mybatis = "org.mybatis" % "mybatis" % mybatisVersion
+ val hikaricp = "com.zaxxer" % "HikariCP" % hikaricpVersion
// Test dependencies
// https://www.scala-sbt.org/1.x/docs/Testing.html
@@ -132,6 +137,7 @@ object Dependencies {
val mockitoInline = "org.mockito" % "mockito-inline" % mockitoVersion
val scalatestMockito = "org.mockito" %% "mockito-scala-scalatest" %
scalatestMockitoVersion
val scalatest = "org.scalatest" %% "scalatest" % scalatestVersion
+ val h2 = "com.h2database" % "h2" % h2Version
}
object CelebornCommonSettings {
@@ -443,8 +449,11 @@ object CelebornService {
Dependencies.javaxServletApi,
Dependencies.commonsCrypto,
Dependencies.slf4jApi,
+ Dependencies.mybatis,
+ Dependencies.hikaricp,
Dependencies.log4jSlf4jImpl % "test",
- Dependencies.log4j12Api % "test"
+ Dependencies.log4j12Api % "test",
+ Dependencies.h2 % "test"
) ++ commonUnitTestDependencies
)
}
diff --git a/service/pom.xml b/service/pom.xml
index 5a8639658..faf80a142 100644
--- a/service/pom.xml
+++ b/service/pom.xml
@@ -60,6 +60,21 @@
<artifactId>jsr305</artifactId>
</dependency>
+ <!-- Db dependencies -->
+ <dependency>
+ <groupId>org.mybatis</groupId>
+ <artifactId>mybatis</artifactId>
+ </dependency>
+ <dependency>
+ <groupId>com.zaxxer</groupId>
+ <artifactId>HikariCP</artifactId>
+ </dependency>
+ <dependency>
+ <groupId>com.h2database</groupId>
+ <artifactId>h2</artifactId>
+ <scope>test</scope>
+ </dependency>
+
<!-- Test dependencies -->
<dependency>
<groupId>org.mockito</groupId>
diff --git
a/service/src/main/java/org/apache/celeborn/server/common/service/config/BaseConfigServiceImpl.java
b/service/src/main/java/org/apache/celeborn/server/common/service/config/BaseConfigServiceImpl.java
new file mode 100644
index 000000000..5e7c27009
--- /dev/null
+++
b/service/src/main/java/org/apache/celeborn/server/common/service/config/BaseConfigServiceImpl.java
@@ -0,0 +1,88 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.celeborn.server.common.service.config;
+
+import java.io.IOException;
+import java.util.HashMap;
+import java.util.Map;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicReference;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import org.apache.celeborn.common.CelebornConf;
+import org.apache.celeborn.common.util.ThreadUtils;
+
+public abstract class BaseConfigServiceImpl implements ConfigService {
+ private static final Logger LOG =
LoggerFactory.getLogger(BaseConfigServiceImpl.class);
+
+ protected final CelebornConf celebornConf;
+ protected final AtomicReference<SystemConfig> systemConfigAtomicReference =
+ new AtomicReference<>();
+ protected final AtomicReference<Map<String, TenantConfig>>
tenantConfigAtomicReference =
+ new AtomicReference<>(new HashMap<>());
+
+ private final ScheduledExecutorService configRefreshService =
+
ThreadUtils.newDaemonSingleThreadScheduledExecutor("celeborn-config-refresher");
+
+ public BaseConfigServiceImpl(CelebornConf celebornConf) throws IOException {
+ this.celebornConf = celebornConf;
+ this.systemConfigAtomicReference.set(new SystemConfig(celebornConf));
+ this.refreshAllCache();
+ boolean dynamicConfigEnabled = celebornConf.dynamicConfigEnabled();
+ if (dynamicConfigEnabled) {
+ LOG.info("Celeborn dynamic config is enabled.");
+ long dynamicConfigRefreshInterval =
celebornConf.dynamicConfigRefreshInterval();
+ this.configRefreshService.scheduleWithFixedDelay(
+ () -> {
+ try {
+ refreshAllCache();
+ } catch (Throwable e) {
+ LOG.error("Refresh config encounter exception: {}",
e.getMessage(), e);
+ }
+ },
+ dynamicConfigRefreshInterval,
+ dynamicConfigRefreshInterval,
+ TimeUnit.MILLISECONDS);
+ } else {
+ LOG.info("Celeborn dynamic config is disabled, config can not be
refreshed after updated.");
+ }
+ }
+
+ @Override
+ public CelebornConf getCelebornConf() {
+ return celebornConf;
+ }
+
+ @Override
+ public SystemConfig getSystemConfigFromCache() {
+ return systemConfigAtomicReference.get();
+ }
+
+ @Override
+ public TenantConfig getRawTenantConfigFromCache(String tenantId) {
+ return tenantConfigAtomicReference.get().get(tenantId);
+ }
+
+ @Override
+ public void shutdown() {
+ ThreadUtils.shutdown(configRefreshService);
+ }
+}
diff --git
a/service/src/main/java/org/apache/celeborn/server/common/service/config/ConfigService.java
b/service/src/main/java/org/apache/celeborn/server/common/service/config/ConfigService.java
index 80e02dc23..fc918b0a5 100644
---
a/service/src/main/java/org/apache/celeborn/server/common/service/config/ConfigService.java
+++
b/service/src/main/java/org/apache/celeborn/server/common/service/config/ConfigService.java
@@ -17,22 +17,28 @@
package org.apache.celeborn.server.common.service.config;
+import java.io.IOException;
+
+import org.apache.celeborn.common.CelebornConf;
+
public interface ConfigService {
- SystemConfig getSystemConfig();
+ CelebornConf getCelebornConf();
+
+ SystemConfig getSystemConfigFromCache();
- TenantConfig getRawTenantConfig(String tenantId);
+ TenantConfig getRawTenantConfigFromCache(String tenantId);
- default DynamicConfig getTenantConfig(String tenantId) {
- TenantConfig tenantConfig = getRawTenantConfig(tenantId);
+ default DynamicConfig getTenantConfigFromCache(String tenantId) {
+ TenantConfig tenantConfig = getRawTenantConfigFromCache(tenantId);
if (tenantConfig == null || tenantConfig.getConfigs().isEmpty()) {
- return getSystemConfig();
+ return getSystemConfigFromCache();
} else {
return tenantConfig;
}
}
- void refreshAllCache();
+ void refreshAllCache() throws IOException;
void shutdown();
}
diff --git
a/service/src/main/java/org/apache/celeborn/server/common/service/config/DbConfigServiceImpl.java
b/service/src/main/java/org/apache/celeborn/server/common/service/config/DbConfigServiceImpl.java
new file mode 100644
index 000000000..6217545b5
--- /dev/null
+++
b/service/src/main/java/org/apache/celeborn/server/common/service/config/DbConfigServiceImpl.java
@@ -0,0 +1,54 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.celeborn.server.common.service.config;
+
+import java.io.IOException;
+import java.util.List;
+import java.util.Map;
+import java.util.function.Function;
+import java.util.stream.Collectors;
+
+import org.apache.celeborn.common.CelebornConf;
+import org.apache.celeborn.server.common.service.store.IServiceManager;
+import org.apache.celeborn.server.common.service.store.db.DbServiceManagerImpl;
+
+public class DbConfigServiceImpl extends BaseConfigServiceImpl implements
ConfigService {
+ private volatile IServiceManager iServiceManager;
+
+ public DbConfigServiceImpl(CelebornConf celebornConf) throws IOException {
+ super(celebornConf);
+ }
+
+ @Override
+ public void refreshAllCache() throws IOException {
+ if (iServiceManager == null) {
+ synchronized (this) {
+ if (iServiceManager == null) {
+ iServiceManager = new DbServiceManagerImpl(celebornConf, this);
+ }
+ }
+ }
+
+ systemConfigAtomicReference.set(iServiceManager.getSystemConfig());
+ List<TenantConfig> allTenantConfigs =
iServiceManager.getAllTenantConfigs();
+ Map<String, TenantConfig> tenantConfigMap =
+ allTenantConfigs.stream()
+ .collect(Collectors.toMap(TenantConfig::getTenantId,
Function.identity()));
+ tenantConfigAtomicReference.set(tenantConfigMap);
+ }
+}
diff --git
a/service/src/main/java/org/apache/celeborn/server/common/service/config/DynamicConfigServiceFactory.java
b/service/src/main/java/org/apache/celeborn/server/common/service/config/DynamicConfigServiceFactory.java
index 2d0154966..5466d718a 100644
---
a/service/src/main/java/org/apache/celeborn/server/common/service/config/DynamicConfigServiceFactory.java
+++
b/service/src/main/java/org/apache/celeborn/server/common/service/config/DynamicConfigServiceFactory.java
@@ -17,16 +17,21 @@
package org.apache.celeborn.server.common.service.config;
+import java.io.IOException;
+
import org.apache.celeborn.common.CelebornConf;
public class DynamicConfigServiceFactory {
- public static ConfigService getConfigService(CelebornConf celebornConf) {
+ public static ConfigService getConfigService(CelebornConf celebornConf)
throws IOException {
String configStoreBackend = celebornConf.dynamicConfigStoreBackend();
if ("FS".equals(configStoreBackend)) {
return new FsConfigServiceImpl(celebornConf);
+ } else if ("DB".equals(configStoreBackend)) {
+ return new DbConfigServiceImpl(celebornConf);
}
- return null;
+ throw new UnsupportedOperationException(
+ "Unsupported dynamic config store backend:" + configStoreBackend);
}
}
diff --git
a/service/src/main/java/org/apache/celeborn/server/common/service/config/FsConfigServiceImpl.java
b/service/src/main/java/org/apache/celeborn/server/common/service/config/FsConfigServiceImpl.java
index cd519c14d..decf061af 100644
---
a/service/src/main/java/org/apache/celeborn/server/common/service/config/FsConfigServiceImpl.java
+++
b/service/src/main/java/org/apache/celeborn/server/common/service/config/FsConfigServiceImpl.java
@@ -19,13 +19,11 @@ package org.apache.celeborn.server.common.service.config;
import java.io.File;
import java.io.FileInputStream;
+import java.io.IOException;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Optional;
-import java.util.concurrent.ScheduledExecutorService;
-import java.util.concurrent.TimeUnit;
-import java.util.concurrent.atomic.AtomicReference;
import java.util.stream.Collectors;
import org.slf4j.Logger;
@@ -33,31 +31,19 @@ import org.slf4j.LoggerFactory;
import org.yaml.snakeyaml.Yaml;
import org.apache.celeborn.common.CelebornConf;
-import org.apache.celeborn.common.util.ThreadUtils;
-public class FsConfigServiceImpl implements ConfigService {
+public class FsConfigServiceImpl extends BaseConfigServiceImpl implements
ConfigService {
private static final Logger LOG =
LoggerFactory.getLogger(FsConfigServiceImpl.class);
- private final CelebornConf celebornConf;
- private final AtomicReference<SystemConfig> systemConfigAtomicReference =
new AtomicReference<>();
- private final AtomicReference<Map<String, TenantConfig>>
tenantConfigAtomicReference =
- new AtomicReference<>(new HashMap<>());
private static final String CONF_TENANT_ID = "tenantId";
private static final String CONF_LEVEL = "level";
private static final String CONF_CONFIG = "config";
- private final ScheduledExecutorService configRefreshService =
-
ThreadUtils.newDaemonSingleThreadScheduledExecutor("celeborn-config-refresher");
-
- public FsConfigServiceImpl(CelebornConf celebornConf) {
- this.celebornConf = celebornConf;
- this.systemConfigAtomicReference.set(new SystemConfig(celebornConf));
- this.refresh();
- long dynamicConfigRefreshTime =
celebornConf.dynamicConfigRefreshInterval();
- this.configRefreshService.scheduleWithFixedDelay(
- this::refresh, dynamicConfigRefreshTime, dynamicConfigRefreshTime,
TimeUnit.MILLISECONDS);
+ public FsConfigServiceImpl(CelebornConf celebornConf) throws IOException {
+ super(celebornConf);
}
- private synchronized void refresh() {
+ @Override
+ public synchronized void refreshAllCache() {
File configurationFile = getConfigurationFile(System.getenv());
if (!configurationFile.exists()) {
return;
@@ -76,7 +62,7 @@ public class FsConfigServiceImpl implements ConfigService {
.entrySet().stream()
.collect(Collectors.toMap(Map.Entry::getKey, a ->
a.getValue().toString()));
if (ConfigLevel.TENANT.name().equals(level)) {
- TenantConfig tenantConfig = new TenantConfig(this, tenantId, config);
+ TenantConfig tenantConfig = new TenantConfig(this, tenantId, null,
config);
tenantConfs.put(tenantId, tenantConfig);
} else {
systemConfig = new SystemConfig(celebornConf, config);
@@ -93,26 +79,6 @@ public class FsConfigServiceImpl implements ConfigService {
}
}
- @Override
- public SystemConfig getSystemConfig() {
- return systemConfigAtomicReference.get();
- }
-
- @Override
- public TenantConfig getRawTenantConfig(String tenantId) {
- return tenantConfigAtomicReference.get().get(tenantId);
- }
-
- @Override
- public void refreshAllCache() {
- this.refresh();
- }
-
- @Override
- public void shutdown() {
- ThreadUtils.shutdown(configRefreshService);
- }
-
private File getConfigurationFile(Map<String, String> env) {
if (!this.celebornConf.quotaConfigurationPath().isEmpty()) {
return new File(this.celebornConf.quotaConfigurationPath().get());
diff --git
a/service/src/main/java/org/apache/celeborn/server/common/service/config/SystemConfig.java
b/service/src/main/java/org/apache/celeborn/server/common/service/config/SystemConfig.java
index 82b5714b8..c9f627e3e 100644
---
a/service/src/main/java/org/apache/celeborn/server/common/service/config/SystemConfig.java
+++
b/service/src/main/java/org/apache/celeborn/server/common/service/config/SystemConfig.java
@@ -18,10 +18,12 @@
package org.apache.celeborn.server.common.service.config;
import java.util.HashMap;
+import java.util.List;
import java.util.Map;
import org.apache.celeborn.common.CelebornConf;
import org.apache.celeborn.common.internal.config.ConfigEntry;
+import org.apache.celeborn.server.common.service.model.ClusterSystemConfig;
public class SystemConfig extends DynamicConfig {
private final CelebornConf celebornConf;
@@ -36,6 +38,11 @@ public class SystemConfig extends DynamicConfig {
this.configs = new HashMap<>();
}
+ public SystemConfig(CelebornConf celebornConf, List<ClusterSystemConfig>
systemConfigs) {
+ this.celebornConf = celebornConf;
+ systemConfigs.forEach(t -> configs.put(t.getConfigKey(),
t.getConfigValue()));
+ }
+
@Override
public DynamicConfig getParentLevelConfig() {
return null;
diff --git
a/service/src/main/java/org/apache/celeborn/server/common/service/config/TenantConfig.java
b/service/src/main/java/org/apache/celeborn/server/common/service/config/TenantConfig.java
index c5e4f983c..bea2f5960 100644
---
a/service/src/main/java/org/apache/celeborn/server/common/service/config/TenantConfig.java
+++
b/service/src/main/java/org/apache/celeborn/server/common/service/config/TenantConfig.java
@@ -17,24 +17,45 @@
package org.apache.celeborn.server.common.service.config;
+import java.util.List;
import java.util.Map;
+import org.apache.celeborn.server.common.service.model.ClusterTenantConfig;
+
public class TenantConfig extends DynamicConfig {
private final String tenantId;
private final ConfigService configService;
+ private final String name;
- public TenantConfig(ConfigService configService, String tenantId,
Map<String, String> configs) {
+ public TenantConfig(
+ ConfigService configService, String tenantId, String name, Map<String,
String> configs) {
this.configService = configService;
+ this.tenantId = tenantId;
+ this.name = name;
this.configs.putAll(configs);
+ }
+
+ public TenantConfig(
+ ConfigService configService,
+ String tenantId,
+ String name,
+ List<ClusterTenantConfig> tenantConfigs) {
+ this.configService = configService;
this.tenantId = tenantId;
+ this.name = name;
+ tenantConfigs.forEach(t -> configs.put(t.getConfigKey(),
t.getConfigValue()));
}
public String getTenantId() {
return tenantId;
}
+ public String getName() {
+ return name;
+ }
+
@Override
public DynamicConfig getParentLevelConfig() {
- return configService.getSystemConfig();
+ return configService.getSystemConfigFromCache();
}
}
diff --git
a/service/src/main/java/org/apache/celeborn/server/common/service/model/ClusterInfo.java
b/service/src/main/java/org/apache/celeborn/server/common/service/model/ClusterInfo.java
new file mode 100644
index 000000000..0adb34f59
--- /dev/null
+++
b/service/src/main/java/org/apache/celeborn/server/common/service/model/ClusterInfo.java
@@ -0,0 +1,78 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.celeborn.server.common.service.model;
+
+import java.util.Date;
+
+public class ClusterInfo {
+
+ private Integer id;
+ private String name;
+ private String namespace;
+ private String endpoint;
+ private Date gmtCreate;
+ private Date gmtModify;
+
+ public Integer getId() {
+ return id;
+ }
+
+ public void setId(Integer id) {
+ this.id = id;
+ }
+
+ public String getName() {
+ return name;
+ }
+
+ public void setName(String name) {
+ this.name = name;
+ }
+
+ public String getNamespace() {
+ return namespace;
+ }
+
+ public void setNamespace(String namespace) {
+ this.namespace = namespace;
+ }
+
+ public String getEndpoint() {
+ return endpoint;
+ }
+
+ public void setEndpoint(String endpoint) {
+ this.endpoint = endpoint;
+ }
+
+ public Date getGmtCreate() {
+ return gmtCreate;
+ }
+
+ public void setGmtCreate(Date gmtCreate) {
+ this.gmtCreate = gmtCreate;
+ }
+
+ public Date getGmtModify() {
+ return gmtModify;
+ }
+
+ public void setGmtModify(Date gmtModify) {
+ this.gmtModify = gmtModify;
+ }
+}
diff --git
a/service/src/main/java/org/apache/celeborn/server/common/service/model/ClusterSystemConfig.java
b/service/src/main/java/org/apache/celeborn/server/common/service/model/ClusterSystemConfig.java
new file mode 100644
index 000000000..7f119c425
--- /dev/null
+++
b/service/src/main/java/org/apache/celeborn/server/common/service/model/ClusterSystemConfig.java
@@ -0,0 +1,86 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.celeborn.server.common.service.model;
+
+import java.util.Date;
+
+public class ClusterSystemConfig {
+ private Integer id;
+ private Integer clusterId;
+ private String configKey;
+ private String configValue;
+ private String type;
+ private Date gmtCreate;
+ private Date gmtModify;
+
+ public Integer getId() {
+ return id;
+ }
+
+ public void setId(Integer id) {
+ this.id = id;
+ }
+
+ public Integer getClusterId() {
+ return clusterId;
+ }
+
+ public void setClusterId(Integer clusterId) {
+ this.clusterId = clusterId;
+ }
+
+ public String getConfigKey() {
+ return configKey;
+ }
+
+ public void setConfigKey(String configKey) {
+ this.configKey = configKey;
+ }
+
+ public String getConfigValue() {
+ return configValue;
+ }
+
+ public void setConfigValue(String configValue) {
+ this.configValue = configValue;
+ }
+
+ public String getType() {
+ return type;
+ }
+
+ public void setType(String type) {
+ this.type = type;
+ }
+
+ public Date getGmtCreate() {
+ return gmtCreate;
+ }
+
+ public void setGmtCreate(Date gmtCreate) {
+ this.gmtCreate = gmtCreate;
+ }
+
+ public Date getGmtModify() {
+ return gmtModify;
+ }
+
+ public void setGmtModify(Date gmtModify) {
+ this.gmtModify = gmtModify;
+ }
+}
diff --git
a/service/src/main/java/org/apache/celeborn/server/common/service/model/ClusterTenantConfig.java
b/service/src/main/java/org/apache/celeborn/server/common/service/model/ClusterTenantConfig.java
new file mode 100644
index 000000000..b3c7d068d
--- /dev/null
+++
b/service/src/main/java/org/apache/celeborn/server/common/service/model/ClusterTenantConfig.java
@@ -0,0 +1,138 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.celeborn.server.common.service.model;
+
+import java.util.Date;
+
+import org.apache.commons.lang3.StringUtils;
+import org.apache.commons.lang3.tuple.Pair;
+
+public class ClusterTenantConfig {
+
+ private Integer id;
+ private Integer clusterId;
+ private String tenantId;
+ private String level;
+ private String user;
+ private String configKey;
+ private String configValue;
+ private String type;
+ private Date gmtCreate;
+ private Date gmtModify;
+
+ public Integer getId() {
+ return id;
+ }
+
+ public void setId(Integer id) {
+ this.id = id;
+ }
+
+ public Integer getClusterId() {
+ return clusterId;
+ }
+
+ public void setClusterId(Integer clusterId) {
+ this.clusterId = clusterId;
+ }
+
+ public String getTenantId() {
+ return tenantId;
+ }
+
+ public void setTenantId(String tenantId) {
+ this.tenantId = tenantId;
+ }
+
+ public String getLevel() {
+ return level;
+ }
+
+ public void setLevel(String level) {
+ this.level = level;
+ }
+
+ public String getUser() {
+ return StringUtils.isBlank(user) ? null : user;
+ }
+
+ public void setUser(String user) {
+ this.user = user;
+ }
+
+ public String getConfigKey() {
+ return configKey;
+ }
+
+ public void setConfigKey(String configKey) {
+ this.configKey = configKey;
+ }
+
+ public String getConfigValue() {
+ return configValue;
+ }
+
+ public void setConfigValue(String configValue) {
+ this.configValue = configValue;
+ }
+
+ public String getType() {
+ return type;
+ }
+
+ public void setType(String type) {
+ this.type = type;
+ }
+
+ public Date getGmtCreate() {
+ return gmtCreate;
+ }
+
+ public void setGmtCreate(Date gmtCreate) {
+ this.gmtCreate = gmtCreate;
+ }
+
+ public Date getGmtModify() {
+ return gmtModify;
+ }
+
+ public void setGmtModify(Date gmtModify) {
+ this.gmtModify = gmtModify;
+ }
+
+ public Pair getTenantInfo() {
+ return Pair.of(tenantId, user);
+ }
+
+ @Override
+ public String toString() {
+ final StringBuilder sb = new StringBuilder("ClusterTenantConfig{");
+ sb.append("id=").append(id);
+ sb.append(", clusterId=").append(clusterId);
+ sb.append(", tenantId='").append(tenantId).append('\'');
+ sb.append(", level='").append(level).append('\'');
+ sb.append(", user='").append(user).append('\'');
+ sb.append(", configKey='").append(configKey).append('\'');
+ sb.append(", configValue='").append(configValue).append('\'');
+ sb.append(", type='").append(type).append('\'');
+ sb.append(", gmtCreate=").append(gmtCreate);
+ sb.append(", gmtModify=").append(gmtModify);
+ sb.append('}');
+ return sb.toString();
+ }
+}
diff --git
a/service/src/main/java/org/apache/celeborn/server/common/service/config/ConfigService.java
b/service/src/main/java/org/apache/celeborn/server/common/service/store/IServiceManager.java
similarity index 64%
copy from
service/src/main/java/org/apache/celeborn/server/common/service/config/ConfigService.java
copy to
service/src/main/java/org/apache/celeborn/server/common/service/store/IServiceManager.java
index 80e02dc23..c2a8a9d0e 100644
---
a/service/src/main/java/org/apache/celeborn/server/common/service/config/ConfigService.java
+++
b/service/src/main/java/org/apache/celeborn/server/common/service/store/IServiceManager.java
@@ -15,24 +15,21 @@
* limitations under the License.
*/
-package org.apache.celeborn.server.common.service.config;
+package org.apache.celeborn.server.common.service.store;
-public interface ConfigService {
+import java.util.List;
- SystemConfig getSystemConfig();
+import org.apache.celeborn.server.common.service.config.SystemConfig;
+import org.apache.celeborn.server.common.service.config.TenantConfig;
+import org.apache.celeborn.server.common.service.model.ClusterInfo;
+
+public interface IServiceManager {
- TenantConfig getRawTenantConfig(String tenantId);
+ int createCluster(ClusterInfo clusterInfo);
- default DynamicConfig getTenantConfig(String tenantId) {
- TenantConfig tenantConfig = getRawTenantConfig(tenantId);
- if (tenantConfig == null || tenantConfig.getConfigs().isEmpty()) {
- return getSystemConfig();
- } else {
- return tenantConfig;
- }
- }
+ ClusterInfo getClusterInfo(String clusterName);
- void refreshAllCache();
+ List<TenantConfig> getAllTenantConfigs();
- void shutdown();
+ SystemConfig getSystemConfig();
}
diff --git
a/service/src/main/java/org/apache/celeborn/server/common/service/store/db/DBSessionFactory.java
b/service/src/main/java/org/apache/celeborn/server/common/service/store/db/DBSessionFactory.java
new file mode 100644
index 000000000..1c3457661
--- /dev/null
+++
b/service/src/main/java/org/apache/celeborn/server/common/service/store/db/DBSessionFactory.java
@@ -0,0 +1,95 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.celeborn.server.common.service.store.db;
+
+import java.io.IOException;
+import java.util.Map;
+import java.util.Properties;
+
+import javax.sql.DataSource;
+
+import com.zaxxer.hikari.HikariConfig;
+import com.zaxxer.hikari.HikariDataSource;
+import org.apache.ibatis.mapping.Environment;
+import org.apache.ibatis.session.Configuration;
+import org.apache.ibatis.session.SqlSessionFactory;
+import org.apache.ibatis.session.SqlSessionFactoryBuilder;
+import org.apache.ibatis.transaction.TransactionFactory;
+import org.apache.ibatis.transaction.jdbc.JdbcTransactionFactory;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import org.apache.celeborn.common.CelebornConf;
+import
org.apache.celeborn.server.common.service.store.db.mapper.ClusterInfoMapper;
+import
org.apache.celeborn.server.common.service.store.db.mapper.ClusterSystemConfigMapper;
+import
org.apache.celeborn.server.common.service.store.db.mapper.ClusterTenantConfigMapper;
+
+public class DBSessionFactory {
+ private static final Logger LOG =
LoggerFactory.getLogger(DBSessionFactory.class);
+ private static volatile SqlSessionFactory _instance;
+
+ public static SqlSessionFactory get(CelebornConf celebornConf) throws
IOException {
+ if (_instance == null) {
+ synchronized (DBSessionFactory.class) {
+ if (_instance == null) {
+ Properties properties = new Properties();
+ properties.setProperty(
+ "driverClassName",
celebornConf.dynamicConfigStoreDbHikariDriverClassName());
+ properties.setProperty("jdbcUrl",
celebornConf.dynamicConfigStoreDbHikariJdbcUrl());
+ properties.setProperty("username",
celebornConf.dynamicConfigStoreDbHikariUsername());
+ properties.setProperty("password",
celebornConf.dynamicConfigStoreDbHikariPassword());
+ properties.setProperty(
+ "connectionTimeout",
+
String.valueOf(celebornConf.dynamicConfigStoreDbHikariConnectionTimeout()));
+ properties.setProperty(
+ "idleTimeout",
String.valueOf(celebornConf.dynamicConfigStoreDbHikariIdleTimeout()));
+ properties.setProperty(
+ "maxLifetime",
String.valueOf(celebornConf.dynamicConfigStoreDbHikariMaxLifetime()));
+ properties.setProperty(
+ "maximumPoolSize",
+
String.valueOf(celebornConf.dynamicConfigStoreDbHikariMaximumPoolSize()));
+
+ for (Map.Entry<String, String> dbPropertiesEntry :
+
celebornConf.dynamicConfigStoreDbHikariCustomConfigs().entrySet()) {
+ properties.setProperty(
+
dbPropertiesEntry.getKey().replace("celeborn.dynamicConfig.store.db.hikari.",
""),
+ dbPropertiesEntry.getValue());
+ }
+
+ HikariConfig config = new HikariConfig(properties);
+ DataSource dataSource = new HikariDataSource(config);
+
+ TransactionFactory transactionFactory = new JdbcTransactionFactory();
+ Environment environment = new Environment("celeborn",
transactionFactory, dataSource);
+
+ Configuration configuration = new Configuration(environment);
+ configuration.setMapUnderscoreToCamelCase(true);
+ configuration.addMapper(ClusterInfoMapper.class);
+ configuration.addMapper(ClusterSystemConfigMapper.class);
+ configuration.addMapper(ClusterTenantConfigMapper.class);
+
+ SqlSessionFactoryBuilder builder = new SqlSessionFactoryBuilder();
+ _instance = builder.build(configuration);
+ LOG.info("Init sqlSessionFactory success");
+ }
+ }
+ }
+
+ return _instance;
+ }
+}
diff --git
a/service/src/main/java/org/apache/celeborn/server/common/service/store/db/DbServiceManagerImpl.java
b/service/src/main/java/org/apache/celeborn/server/common/service/store/db/DbServiceManagerImpl.java
new file mode 100644
index 000000000..e957cfed5
--- /dev/null
+++
b/service/src/main/java/org/apache/celeborn/server/common/service/store/db/DbServiceManagerImpl.java
@@ -0,0 +1,138 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.celeborn.server.common.service.store.db;
+
+import java.io.IOException;
+import java.util.ArrayList;
+import java.util.Date;
+import java.util.List;
+import java.util.Map;
+import java.util.stream.Collectors;
+
+import org.apache.ibatis.session.SqlSession;
+import org.apache.ibatis.session.SqlSessionFactory;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import org.apache.celeborn.common.CelebornConf;
+import org.apache.celeborn.server.common.service.config.ConfigLevel;
+import org.apache.celeborn.server.common.service.config.ConfigService;
+import org.apache.celeborn.server.common.service.config.SystemConfig;
+import org.apache.celeborn.server.common.service.config.TenantConfig;
+import org.apache.celeborn.server.common.service.model.ClusterInfo;
+import org.apache.celeborn.server.common.service.model.ClusterSystemConfig;
+import org.apache.celeborn.server.common.service.model.ClusterTenantConfig;
+import org.apache.celeborn.server.common.service.store.IServiceManager;
+import
org.apache.celeborn.server.common.service.store.db.mapper.ClusterInfoMapper;
+import
org.apache.celeborn.server.common.service.store.db.mapper.ClusterSystemConfigMapper;
+import
org.apache.celeborn.server.common.service.store.db.mapper.ClusterTenantConfigMapper;
+import org.apache.celeborn.server.common.service.utils.JsonUtils;
+
+public class DbServiceManagerImpl implements IServiceManager {
+ private static final Logger LOG =
LoggerFactory.getLogger(DbServiceManagerImpl.class);
+ private final CelebornConf celebornConf;
+ private final ConfigService configService;
+ private final SqlSessionFactory sqlSessionFactory;
+ private final int clusterId;
+ private final int pageSize;
+
+ public DbServiceManagerImpl(CelebornConf celebornConf, ConfigService
configServer)
+ throws IOException {
+ this.celebornConf = celebornConf;
+ this.sqlSessionFactory = DBSessionFactory.get(celebornConf);
+ this.configService = configServer;
+ this.pageSize = celebornConf.dynamicConfigStoreDbFetchPageSize();
+ this.clusterId = createCluster(getClusterInfoFromEnv());
+ }
+
+ @Override
+ public int createCluster(ClusterInfo clusterInfo) {
+ ClusterInfo clusterInfoFromDB = getClusterInfo(clusterInfo.getName());
+ if (clusterInfoFromDB == null) {
+ try (SqlSession sqlSession = sqlSessionFactory.openSession(true)) {
+ ClusterInfoMapper mapper =
sqlSession.getMapper(ClusterInfoMapper.class);
+ clusterInfo.setGmtCreate(new Date());
+ clusterInfo.setGmtModify(new Date());
+ mapper.insert(clusterInfo);
+ LOG.info("Create cluster {} successfully.",
JsonUtils.toJson(clusterInfo));
+ } catch (Exception e) {
+ LOG.warn("Create cluster {} failed: {}.",
JsonUtils.toJson(clusterInfo), e.getMessage(), e);
+ }
+ clusterInfoFromDB = getClusterInfo(clusterInfo.getName());
+ if (clusterInfoFromDB == null) {
+ throw new RuntimeException("Could not get cluster info of " +
clusterInfo.getName() + ".");
+ }
+ }
+ return clusterInfoFromDB.getId();
+ }
+
+ @Override
+ public ClusterInfo getClusterInfo(String clusterName) {
+ try (SqlSession sqlSession = sqlSessionFactory.openSession()) {
+ ClusterInfoMapper mapper = sqlSession.getMapper(ClusterInfoMapper.class);
+ return mapper.getClusterInfo(clusterName);
+ }
+ }
+
+ @Override
+ public List<TenantConfig> getAllTenantConfigs() {
+ try (SqlSession sqlSession = sqlSessionFactory.openSession()) {
+ ClusterTenantConfigMapper mapper =
sqlSession.getMapper(ClusterTenantConfigMapper.class);
+ int totalNum = mapper.getClusterTenantConfigsNum(clusterId,
ConfigLevel.TENANT.name());
+ int offset = 0;
+ List<ClusterTenantConfig> clusterAllTenantConfigs = new ArrayList<>();
+ while (offset < totalNum) {
+ List<ClusterTenantConfig> clusterTenantConfigs =
+ mapper.getClusterTenantConfigs(clusterId,
ConfigLevel.TENANT.name(), offset, pageSize);
+ clusterAllTenantConfigs.addAll(clusterTenantConfigs);
+ offset = offset + pageSize;
+ }
+
+ Map<String, List<ClusterTenantConfig>> tenantConfigMaps =
+ clusterAllTenantConfigs.stream()
+ .collect(
+ Collectors.groupingBy(clusterTenantConfig ->
clusterTenantConfig.getTenantId()));
+ return tenantConfigMaps.entrySet().stream()
+ .map(t -> new TenantConfig(configService, t.getKey(), null,
t.getValue()))
+ .collect(Collectors.toList());
+ }
+ }
+
+ @Override
+ public SystemConfig getSystemConfig() {
+ try (SqlSession sqlSession = sqlSessionFactory.openSession()) {
+ ClusterSystemConfigMapper mapper =
sqlSession.getMapper(ClusterSystemConfigMapper.class);
+ List<ClusterSystemConfig> clusterSystemConfig =
mapper.getClusterSystemConfig(clusterId);
+ return new SystemConfig(celebornConf, clusterSystemConfig);
+ }
+ }
+
+ public ClusterInfo getClusterInfoFromEnv() {
+ Map<String, String> env = System.getenv();
+ String clusterName = env.getOrDefault("CELEBORN_CLUSTER_NAME",
celebornConf.clusterName());
+ String namespace = env.getOrDefault("CELEBORN_CLUSTER_NAMESPACE", "");
+ String endpoint = env.getOrDefault("CELEBORN_CLUSTER_ENDPOINT", "");
+
+ ClusterInfo clusterInfo = new ClusterInfo();
+ clusterInfo.setName(clusterName);
+ clusterInfo.setNamespace(namespace);
+ clusterInfo.setEndpoint(endpoint);
+
+ return clusterInfo;
+ }
+}
diff --git
a/service/src/main/java/org/apache/celeborn/server/common/service/config/TenantConfig.java
b/service/src/main/java/org/apache/celeborn/server/common/service/store/db/mapper/ClusterInfoMapper.java
similarity index 54%
copy from
service/src/main/java/org/apache/celeborn/server/common/service/config/TenantConfig.java
copy to
service/src/main/java/org/apache/celeborn/server/common/service/store/db/mapper/ClusterInfoMapper.java
index c5e4f983c..cb888c122 100644
---
a/service/src/main/java/org/apache/celeborn/server/common/service/config/TenantConfig.java
+++
b/service/src/main/java/org/apache/celeborn/server/common/service/store/db/mapper/ClusterInfoMapper.java
@@ -15,26 +15,21 @@
* limitations under the License.
*/
-package org.apache.celeborn.server.common.service.config;
+package org.apache.celeborn.server.common.service.store.db.mapper;
-import java.util.Map;
+import org.apache.ibatis.annotations.Insert;
+import org.apache.ibatis.annotations.Select;
-public class TenantConfig extends DynamicConfig {
- private final String tenantId;
- private final ConfigService configService;
+import org.apache.celeborn.server.common.service.model.ClusterInfo;
- public TenantConfig(ConfigService configService, String tenantId,
Map<String, String> configs) {
- this.configService = configService;
- this.configs.putAll(configs);
- this.tenantId = tenantId;
- }
+public interface ClusterInfoMapper {
- public String getTenantId() {
- return tenantId;
- }
+ @Insert(
+ "INSERT INTO celeborn_cluster_info(name, namespace, endpoint,
gmt_create, gmt_modify) "
+ + "VALUES (#{name}, #{namespace}, #{endpoint}, #{gmtCreate},
#{gmtModify})")
+ void insert(ClusterInfo clusterInfo);
- @Override
- public DynamicConfig getParentLevelConfig() {
- return configService.getSystemConfig();
- }
+ @Select(
+ "SELECT id, name, namespace, endpoint, gmt_create, gmt_modify FROM
celeborn_cluster_info WHERE name = #{clusterName}")
+ ClusterInfo getClusterInfo(String clusterName);
}
diff --git
a/service/src/main/java/org/apache/celeborn/server/common/service/config/ConfigService.java
b/service/src/main/java/org/apache/celeborn/server/common/service/store/db/mapper/ClusterSystemConfigMapper.java
similarity index 61%
copy from
service/src/main/java/org/apache/celeborn/server/common/service/config/ConfigService.java
copy to
service/src/main/java/org/apache/celeborn/server/common/service/store/db/mapper/ClusterSystemConfigMapper.java
index 80e02dc23..a14d575b3 100644
---
a/service/src/main/java/org/apache/celeborn/server/common/service/config/ConfigService.java
+++
b/service/src/main/java/org/apache/celeborn/server/common/service/store/db/mapper/ClusterSystemConfigMapper.java
@@ -15,24 +15,18 @@
* limitations under the License.
*/
-package org.apache.celeborn.server.common.service.config;
+package org.apache.celeborn.server.common.service.store.db.mapper;
-public interface ConfigService {
+import java.util.List;
- SystemConfig getSystemConfig();
+import org.apache.ibatis.annotations.Select;
- TenantConfig getRawTenantConfig(String tenantId);
+import org.apache.celeborn.server.common.service.model.ClusterSystemConfig;
- default DynamicConfig getTenantConfig(String tenantId) {
- TenantConfig tenantConfig = getRawTenantConfig(tenantId);
- if (tenantConfig == null || tenantConfig.getConfigs().isEmpty()) {
- return getSystemConfig();
- } else {
- return tenantConfig;
- }
- }
+public interface ClusterSystemConfigMapper {
- void refreshAllCache();
-
- void shutdown();
+ @Select(
+ "SELECT id, cluster_id, config_key, config_value, type, gmt_create,
gmt_modify "
+ + "FROM celeborn_cluster_system_config WHERE cluster_id =
#{clusterId}")
+ List<ClusterSystemConfig> getClusterSystemConfig(int clusterId);
}
diff --git
a/service/src/main/java/org/apache/celeborn/server/common/service/store/db/mapper/ClusterTenantConfigMapper.java
b/service/src/main/java/org/apache/celeborn/server/common/service/store/db/mapper/ClusterTenantConfigMapper.java
new file mode 100644
index 000000000..af4539f66
--- /dev/null
+++
b/service/src/main/java/org/apache/celeborn/server/common/service/store/db/mapper/ClusterTenantConfigMapper.java
@@ -0,0 +1,42 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.celeborn.server.common.service.store.db.mapper;
+
+import java.util.List;
+
+import org.apache.ibatis.annotations.Param;
+import org.apache.ibatis.annotations.Select;
+
+import org.apache.celeborn.server.common.service.model.ClusterTenantConfig;
+
+public interface ClusterTenantConfigMapper {
+
+ @Select(
+ "SELECT id, cluster_id, tenant_id, level, user, config_key,
config_value, type, gmt_create, gmt_modify "
+ + "FROM celeborn_cluster_tenant_config WHERE cluster_id =
#{clusterId} AND level=#{level} LIMIT #{offset}, #{pageSize}")
+ List<ClusterTenantConfig> getClusterTenantConfigs(
+ @Param("clusterId") int clusterId,
+ @Param("level") String configLevel,
+ @Param("offset") int offset,
+ @Param("pageSize") int pageSize);
+
+ @Select(
+ "SELECT count(*) FROM celeborn_cluster_tenant_config WHERE cluster_id =
#{clusterId} AND level=#{level}")
+ int getClusterTenantConfigsNum(
+ @Param("clusterId") int clusterId, @Param("level") String configLevel);
+}
diff --git
a/service/src/main/java/org/apache/celeborn/server/common/service/utils/JsonUtils.java
b/service/src/main/java/org/apache/celeborn/server/common/service/utils/JsonUtils.java
new file mode 100644
index 000000000..77e99e5f5
--- /dev/null
+++
b/service/src/main/java/org/apache/celeborn/server/common/service/utils/JsonUtils.java
@@ -0,0 +1,50 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.celeborn.server.common.service.utils;
+
+import com.fasterxml.jackson.annotation.JsonAutoDetect;
+import com.fasterxml.jackson.annotation.PropertyAccessor;
+import com.fasterxml.jackson.core.JsonProcessingException;
+import com.fasterxml.jackson.databind.DeserializationFeature;
+import com.fasterxml.jackson.databind.ObjectMapper;
+import org.apache.commons.lang3.StringUtils;
+
+public class JsonUtils {
+ private static final ObjectMapper MAPPER = new ObjectMapper();
+
+ static {
+ MAPPER.disable(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES);
+ MAPPER.setVisibility(PropertyAccessor.FIELD,
JsonAutoDetect.Visibility.ANY);
+ }
+
+ public static String toJson(Object obj) {
+ try {
+ return MAPPER.writeValueAsString(obj);
+ } catch (JsonProcessingException e) {
+ throw new RuntimeException(e);
+ }
+ }
+
+ public static <T> T fromJson(String json, Class<T> clazz) {
+ try {
+ return StringUtils.isEmpty(json) ? null : MAPPER.readValue(json, clazz);
+ } catch (Exception e) {
+ throw new RuntimeException(e);
+ }
+ }
+}
diff --git a/service/src/main/resources/sql/mysql/celeborn-0.5.0-mysql.sql
b/service/src/main/resources/sql/mysql/celeborn-0.5.0-mysql.sql
new file mode 100644
index 000000000..2240c9c8d
--- /dev/null
+++ b/service/src/main/resources/sql/mysql/celeborn-0.5.0-mysql.sql
@@ -0,0 +1,56 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+CREATE TABLE IF NOT EXISTS celeborn_cluster_info
+(
+ id int NOT NULL AUTO_INCREMENT,
+ name varchar(255) NOT NULL COMMENT 'celeborn cluster name',
+ namespace varchar(255) DEFAULT NULL COMMENT 'celeborn cluster namespace',
+ endpoint varchar(255) DEFAULT NULL COMMENT 'celeborn cluster endpoint',
+ gmt_create timestamp NOT NULL,
+ gmt_modify timestamp NOT NULL,
+ PRIMARY KEY (id),
+ UNIQUE KEY `index_cluster_unique_name` (`name`)
+);
+
+CREATE TABLE IF NOT EXISTS celeborn_cluster_system_config
+(
+ id int NOT NULL AUTO_INCREMENT,
+ cluster_id int NOT NULL,
+ config_key varchar(255) NOT NULL,
+ config_value varchar(255) NOT NULL,
+ type varchar(255) DEFAULT NULL COMMENT 'conf categories, such as
quota',
+ gmt_create timestamp NOT NULL,
+ gmt_modify timestamp NOT NULL,
+ PRIMARY KEY (id),
+ UNIQUE KEY `index_unique_system_config_key` (`cluster_id`, `config_key`)
+);
+
+CREATE TABLE IF NOT EXISTS celeborn_cluster_tenant_config
+(
+ id int NOT NULL AUTO_INCREMENT,
+ cluster_id int NOT NULL,
+ tenant_id varchar(255) NOT NULL,
+ level varchar(255) NOT NULL COMMENT 'config level, valid level is
TENANT,USER',
+ user varchar(255) DEFAULT NULL COMMENT 'tenant sub user',
+ config_key varchar(255) NOT NULL,
+ config_value varchar(255) NOT NULL,
+ type varchar(255) DEFAULT NULL COMMENT 'conf categories, such as
quota',
+ gmt_create timestamp NOT NULL,
+ gmt_modify timestamp NOT NULL,
+ PRIMARY KEY (id),
+ UNIQUE KEY `index_unique_tenant_config_key` (`cluster_id`, `tenant_id`,
`user`, `config_key`)
+);
diff --git
a/service/src/test/java/org/apache/celeborn/server/common/service/config/ConfigServiceSuiteJ.java
b/service/src/test/java/org/apache/celeborn/server/common/service/config/ConfigServiceSuiteJ.java
index d94db2f02..03cdc9c30 100644
---
a/service/src/test/java/org/apache/celeborn/server/common/service/config/ConfigServiceSuiteJ.java
+++
b/service/src/test/java/org/apache/celeborn/server/common/service/config/ConfigServiceSuiteJ.java
@@ -17,52 +17,76 @@
package org.apache.celeborn.server.common.service.config;
+import java.io.IOException;
+import java.sql.SQLException;
+
+import org.apache.ibatis.session.SqlSession;
+import org.apache.ibatis.session.SqlSessionFactory;
+import org.junit.After;
import org.junit.Assert;
import org.junit.Test;
import org.apache.celeborn.common.CelebornConf;
import
org.apache.celeborn.server.common.service.config.DynamicConfig.ConfigType;
+import org.apache.celeborn.server.common.service.store.db.DBSessionFactory;
public class ConfigServiceSuiteJ {
+ private ConfigService configService;
+
+ @Test
+ public void testDbConfig() throws IOException {
+ CelebornConf celebornConf = new CelebornConf();
+ celebornConf.set(
+ CelebornConf.DYNAMIC_CONFIG_STORE_DB_HIKARI_JDBC_URL(),
+ "jdbc:h2:mem:test;MODE=MYSQL;INIT=RUNSCRIPT FROM
'classpath:celeborn-0.5.0-h2.sql'\\;"
+ + "RUNSCRIPT FROM
'classpath:celeborn-0.5.0-h2-ut-data.sql';DB_CLOSE_DELAY=-1;");
+ celebornConf.set(
+ CelebornConf.DYNAMIC_CONFIG_STORE_DB_HIKARI_DRIVER_CLASS_NAME(),
"org.h2.Driver");
+
celebornConf.set(CelebornConf.DYNAMIC_CONFIG_STORE_DB_HIKARI_MAXIMUM_POOL_SIZE(),
"1");
+ configService = new DbConfigServiceImpl(celebornConf);
+ verifyConfig(configService);
+
+ SqlSessionFactory sqlSessionFactory = DBSessionFactory.get(celebornConf);
+ try (SqlSession sqlSession = sqlSessionFactory.openSession(true)) {
+ sqlSession
+ .getConnection()
+ .createStatement()
+ .execute(
+ "UPDATE celeborn_cluster_system_config SET config_value = 100
WHERE config_key='celeborn.test.int.only'");
+ } catch (SQLException e) {
+ throw new RuntimeException(e);
+ }
+
+ configService.refreshAllCache();
+ verifyConfigChanged(configService);
+ }
@Test
- public void testFsConfig() {
+ public void testFsConfig() throws IOException {
CelebornConf celebornConf = new CelebornConf();
String file = getClass().getResource("/dynamicConfig.yaml").getFile();
celebornConf.set(CelebornConf.QUOTA_CONFIGURATION_PATH(), file);
celebornConf.set(CelebornConf.DYNAMIC_CONFIG_REFRESH_INTERVAL(), 5L);
- FsConfigServiceImpl fsConfigService = new
FsConfigServiceImpl(celebornConf);
- try {
- verifyConfig(fsConfigService);
-
- // change -> refresh config
- file = getClass().getResource("/dynamicConfig_2.yaml").getFile();
- celebornConf.set(CelebornConf.QUOTA_CONFIGURATION_PATH(), file);
-
- fsConfigService.refreshAllCache();
- SystemConfig systemConfig = fsConfigService.getSystemConfig();
-
- // verify systemConfig's intConf
- Integer intConfValue =
- systemConfig.getValue("celeborn.test.int.only", null, Integer.TYPE,
ConfigType.STRING);
- Assert.assertEquals(intConfValue.intValue(), 100);
-
- // verify systemConfig's bytesConf -- defer to celebornConf
- Long value =
- systemConfig.getValue(
- CelebornConf.SHUFFLE_PARTITION_SPLIT_THRESHOLD().key(),
- CelebornConf.SHUFFLE_PARTITION_SPLIT_THRESHOLD(),
- Long.TYPE,
- ConfigType.BYTES);
- Assert.assertEquals(value.longValue(), 1073741824);
- } finally {
- fsConfigService.shutdown();
+ configService = new FsConfigServiceImpl(celebornConf);
+ verifyConfig(configService);
+ // change -> refresh config
+ file = getClass().getResource("/dynamicConfig_2.yaml").getFile();
+ celebornConf.set(CelebornConf.QUOTA_CONFIGURATION_PATH(), file);
+ configService.refreshAllCache();
+
+ verifyConfigChanged(configService);
+ }
+
+ @After
+ public void teardown() {
+ if (configService != null) {
+ configService.shutdown();
}
}
public void verifyConfig(ConfigService configService) {
// ------------- Verify SystemConfig ----------------- //
- SystemConfig systemConfig = configService.getSystemConfig();
+ SystemConfig systemConfig = configService.getSystemConfigFromCache();
// verify systemConfig's bytesConf -- use systemConfig
Long value =
systemConfig.getValue(
@@ -113,7 +137,7 @@ public class ConfigServiceSuiteJ {
Assert.assertEquals(intConfValue.intValue(), 10);
// ------------- Verify TenantConfig ----------------- //
- DynamicConfig tenantConfig = configService.getTenantConfig("tenant_id");
+ DynamicConfig tenantConfig =
configService.getTenantConfigFromCache("tenant_id");
// verify tenantConfig's bytesConf -- use tenantConf
value =
tenantConfig.getValue(
@@ -156,7 +180,7 @@ public class ConfigServiceSuiteJ {
ConfigType.BYTES);
Assert.assertNull(value);
- DynamicConfig tenantConfigNone =
configService.getTenantConfig("tenant_id_none");
+ DynamicConfig tenantConfigNone =
configService.getTenantConfigFromCache("tenant_id_none");
// verify tenantConfig's bytesConf -- defer to systemConf
value =
tenantConfigNone.getValue(
@@ -176,4 +200,22 @@ public class ConfigServiceSuiteJ {
tenantConfigNone.getWithDefaultValue("none", 10L, Long.TYPE,
ConfigType.STRING);
Assert.assertEquals(withDefaultValue.longValue(), 10);
}
+
+ public void verifyConfigChanged(ConfigService configService) {
+
+ SystemConfig systemConfig = configService.getSystemConfigFromCache();
+ // verify systemConfig's intConf
+ Integer intConfValue =
+ systemConfig.getValue("celeborn.test.int.only", null, Integer.TYPE,
ConfigType.STRING);
+ Assert.assertEquals(intConfValue.intValue(), 100);
+
+ // verify systemConfig's bytesConf -- defer to celebornConf
+ Long value =
+ systemConfig.getValue(
+ CelebornConf.SHUFFLE_PARTITION_SPLIT_THRESHOLD().key(),
+ CelebornConf.SHUFFLE_PARTITION_SPLIT_THRESHOLD(),
+ Long.TYPE,
+ ConfigType.BYTES);
+ Assert.assertEquals(value.longValue(), 1073741824);
+ }
}
diff --git a/service/src/test/resources/celeborn-0.5.0-h2-ut-data.sql
b/service/src/test/resources/celeborn-0.5.0-h2-ut-data.sql
new file mode 100644
index 000000000..7e0b9b813
--- /dev/null
+++ b/service/src/test/resources/celeborn-0.5.0-h2-ut-data.sql
@@ -0,0 +1,37 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+INSERT INTO celeborn_cluster_info ( `id`, `name`, `namespace`, `endpoint`,
`gmt_create`, `gmt_modify` )
+VALUES
+ ( 1, 'default', 'celeborn-1', 'celeborn-namespace.endpoint.com',
'2023-08-26 22:08:30', '2023-08-26 22:08:30' );
+INSERT INTO `celeborn_cluster_system_config` ( `id`, `cluster_id`,
`config_key`, `config_value`, `type`, `gmt_create`, `gmt_modify` )
+VALUES
+ ( 1, 1, 'celeborn.client.push.buffer.initial.size', '102400', 'QUOTA',
'2023-08-26 22:08:30', '2023-08-26 22:08:30' ),
+ ( 2, 1, 'celeborn.client.push.buffer.max.size', '1024000', 'QUOTA',
'2023-08-26 22:08:30', '2023-08-26 22:08:30' ),
+ ( 3, 1, 'celeborn.worker.fetch.heartbeat.enabled', 'true', 'QUOTA',
'2023-08-26 22:08:30', '2023-08-26 22:08:30' ),
+ ( 5, 1, 'celeborn.client.push.buffer.initial.size.only', '10240', 'QUOTA',
'2023-08-26 22:08:30', '2023-08-26 22:08:30' ),
+ ( 6, 1, 'celeborn.test.timeoutMs.only', '100s', 'QUOTA', '2023-08-26
22:08:30', '2023-08-26 22:08:30' ),
+ ( 7, 1, 'celeborn.test.enabled.only', 'false', 'QUOTA', '2023-08-26
22:08:30', '2023-08-26 22:08:30' ),
+ ( 8, 1, 'celeborn.test.int.only', '10', 'QUOTA', '2023-08-26 22:08:30',
'2023-08-26 22:08:30' );
+INSERT INTO `celeborn_cluster_tenant_config` ( `id`, `cluster_id`,
`tenant_id`, `level`, `user`, `config_key`, `config_value`, `type`,
`gmt_create`, `gmt_modify` )
+VALUES
+ ( 1, 1, 'tenant_id', 'TENANT', '',
'celeborn.client.push.buffer.initial.size', '10240', 'QUOTA', '2023-08-26
22:08:30', '2023-08-26 22:08:30' ),
+ ( 2, 1, 'tenant_id', 'TENANT', '',
'celeborn.client.push.buffer.initial.size.only', '102400', 'QUOTA', '2023-08-26
22:08:30', '2023-08-26 22:08:30' ),
+ ( 3, 1, 'tenant_id', 'TENANT', '',
'celeborn.worker.fetch.heartbeat.enabled', 'false', 'QUOTA', '2023-08-26
22:08:30', '2023-08-26 22:08:30' ),
+ ( 4, 1, 'tenant_id', 'TENANT', '', 'celeborn.test.tenant.timeoutMs.only',
'100s', 'QUOTA', '2023-08-26 22:08:30', '2023-08-26 22:08:30' ),
+ ( 5, 1, 'tenant_id', 'TENANT', '', 'celeborn.test.tenant.enabled.only',
'false', 'QUOTA', '2023-08-26 22:08:30', '2023-08-26 22:08:30' ),
+ ( 6, 1, 'tenant_id', 'TENANT', '', 'celeborn.test.tenant.int.only',
'100s', 'QUOTA', '2023-08-26 22:08:30', '2023-08-26 22:08:30' );
diff --git a/service/src/test/resources/celeborn-0.5.0-h2.sql
b/service/src/test/resources/celeborn-0.5.0-h2.sql
new file mode 100644
index 000000000..b57de8028
--- /dev/null
+++ b/service/src/test/resources/celeborn-0.5.0-h2.sql
@@ -0,0 +1,64 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+CREATE SCHEMA IF NOT EXISTS TEST;
+
+SET SCHEMA TEST;
+DROP TABLE IF exists celeborn_cluster_info;
+DROP TABLE IF exists celeborn_cluster_system_config;
+DROP TABLE IF exists celeborn_cluster_tenant_config;
+
+CREATE TABLE IF NOT EXISTS celeborn_cluster_info
+(
+ id int NOT NULL AUTO_INCREMENT,
+ name varchar(255) NOT NULL COMMENT 'celeborn cluster name',
+ namespace varchar(255) DEFAULT NULL COMMENT 'celeborn cluster namespace',
+ endpoint varchar(255) DEFAULT NULL COMMENT 'celeborn cluster endpoint',
+ gmt_create timestamp NOT NULL,
+ gmt_modify timestamp NOT NULL,
+ PRIMARY KEY (id),
+ UNIQUE KEY `index_cluster_unique_name` (`name`)
+);
+
+CREATE TABLE IF NOT EXISTS celeborn_cluster_system_config
+(
+ id int NOT NULL AUTO_INCREMENT,
+ cluster_id int NOT NULL,
+ config_key varchar(255) NOT NULL,
+ config_value varchar(255) NOT NULL,
+ type varchar(255) DEFAULT NULL COMMENT 'conf categories, such as
quota',
+ gmt_create timestamp NOT NULL,
+ gmt_modify timestamp NOT NULL,
+ PRIMARY KEY (id),
+ UNIQUE KEY `index_unique_system_config_key` (`cluster_id`, `config_key`)
+);
+
+CREATE TABLE IF NOT EXISTS celeborn_cluster_tenant_config
+(
+ id int NOT NULL AUTO_INCREMENT,
+ cluster_id int NOT NULL,
+ tenant_id varchar(255) NOT NULL,
+ level varchar(255) NOT NULL COMMENT 'config level, valid level is
TENANT,USER',
+ `user` varchar(255) DEFAULT NULL COMMENT 'tenant sub user',
+ config_key varchar(255) NOT NULL,
+ config_value varchar(255) NOT NULL,
+ type varchar(255) DEFAULT NULL COMMENT 'conf categories, such as
quota',
+ gmt_create timestamp NOT NULL,
+ gmt_modify timestamp NOT NULL,
+ PRIMARY KEY (id),
+ UNIQUE KEY `index_unique_tenant_config_key` (`cluster_id`, `tenant_id`,
`user`, `config_key`)
+);