This is an automated email from the ASF dual-hosted git repository.
CRZbulabula pushed a commit to branch dev/1.3
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/dev/1.3 by this push:
new dcb31e7842b [To dev/1.3] Supply the max_schema/data_region_group_num
param to modify schema when create or alter database (#18391)
dcb31e7842b is described below
commit dcb31e7842ba698601a2cde969431b9ab0596833
Author: libo <[email protected]>
AuthorDate: Tue Aug 4 20:51:22 2026 +0800
[To dev/1.3] Supply the max_schema/data_region_group_num param to modify
schema when create or alter database (#18391)
* Supply the max_schema/data_region_group_num param to modify schema when
create or alter database (#17988)
* The new value cannot be less than the current maximum quota. The quota
uses the default value if unspecified in the ALTER DATABASE statement (#18154)
* Fix max data region group quota calculation (#18263)
* Add system database RegionGroup quota IT coverage
* Fix system database RegionGroup quota validation
* Align system database RegionGroup quota initialization
---
.../IoTDBDatabaseAutoDataRegionGroupQuotaIT.java | 122 ++++++++++
.../IoTDBDatabaseMixedRegionGroupPolicyIT.java | 78 +++++++
.../it/database/IoTDBDatabaseRegionControlIT.java | 135 ++++++++++-
.../database/IoTDBSystemDatabaseRegionGroupIT.java | 90 ++++++++
.../pipe/it/autocreate/IoTDBPipeIdempotentIT.java | 6 +-
.../org/apache/iotdb/db/qp/sql/IdentifierParser.g4 | 6 +-
.../org/apache/iotdb/db/qp/sql/IoTDBSqlParser.g4 | 4 +-
.../antlr4/org/apache/iotdb/db/qp/sql/SqlLexer.g4 | 8 +-
.../manager/partition/PartitionManager.java | 10 +-
.../manager/schema/ClusterSchemaManager.java | 252 +++++++++++++++------
.../persistence/schema/ClusterSchemaInfo.java | 24 +-
.../thrift/ConfigNodeRPCServiceProcessor.java | 32 ++-
.../thrift/ConfigNodeRPCServiceProcessorTest.java | 34 +++
.../config/metadata/DatabaseSchemaTask.java | 9 +-
.../db/queryengine/plan/parser/ASTVisitor.java | 12 +-
.../metadata/DatabaseSchemaStatement.java | 28 +--
.../reporter/iotdb/IoTDBSessionReporter.java | 4 +-
17 files changed, 715 insertions(+), 139 deletions(-)
diff --git
a/integration-test/src/test/java/org/apache/iotdb/confignode/it/database/IoTDBDatabaseAutoDataRegionGroupQuotaIT.java
b/integration-test/src/test/java/org/apache/iotdb/confignode/it/database/IoTDBDatabaseAutoDataRegionGroupQuotaIT.java
new file mode 100644
index 00000000000..1cb9d5e6690
--- /dev/null
+++
b/integration-test/src/test/java/org/apache/iotdb/confignode/it/database/IoTDBDatabaseAutoDataRegionGroupQuotaIT.java
@@ -0,0 +1,122 @@
+/*
+ * 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.iotdb.confignode.it.database;
+
+import org.apache.iotdb.consensus.ConsensusFactory;
+import org.apache.iotdb.it.env.EnvFactory;
+import org.apache.iotdb.it.framework.IoTDBTestRunner;
+import org.apache.iotdb.itbase.category.ClusterIT;
+
+import org.junit.After;
+import org.junit.Assert;
+import org.junit.Before;
+import org.junit.Test;
+import org.junit.experimental.categories.Category;
+import org.junit.runner.RunWith;
+
+import java.sql.Connection;
+import java.sql.ResultSet;
+import java.sql.SQLException;
+import java.sql.Statement;
+
+@RunWith(IoTDBTestRunner.class)
+@Category({ClusterIT.class})
+public class IoTDBDatabaseAutoDataRegionGroupQuotaIT {
+
+ private static final int CONFIG_NODE_NUM = 1;
+ private static final int DATA_NODE_NUM = 3;
+ private static final int SCHEMA_REPLICATION_FACTOR = 3;
+ private static final int DATA_REPLICATION_FACTOR = 1;
+ private static final int DATA_REGION_PER_DATA_NODE = 2;
+
+ @Before
+ public void setUp() throws Exception {
+ EnvFactory.getEnv()
+ .getConfig()
+ .getCommonConfig()
+ .setSchemaRegionGroupExtensionPolicy("CUSTOM")
+ .setDataRegionGroupExtensionPolicy("AUTO")
+
.setSchemaRegionConsensusProtocolClass(ConsensusFactory.RATIS_CONSENSUS)
+ .setDefaultSchemaRegionGroupNumPerDatabase(1)
+ .setDefaultDataRegionGroupNumPerDatabase(1)
+ .setSchemaReplicationFactor(SCHEMA_REPLICATION_FACTOR)
+ .setDataReplicationFactor(DATA_REPLICATION_FACTOR)
+ .setDataRegionPerDataNode(DATA_REGION_PER_DATA_NODE);
+ EnvFactory.getEnv().initClusterEnvironment(CONFIG_NODE_NUM, DATA_NODE_NUM);
+ }
+
+ @After
+ public void tearDown() {
+ EnvFactory.getEnv().cleanClusterEnvironment();
+ }
+
+ @Test
+ public void testMaxDataRegionGroupNumUsesDataReplicationFactor() throws
SQLException {
+ int expectedMaxDataRegionGroupNum =
+ (int)
+ Math.ceil((double) DATA_REGION_PER_DATA_NODE * DATA_NODE_NUM /
DATA_REPLICATION_FACTOR);
+
+ try (Connection connection = EnvFactory.getEnv().getConnection();
+ Statement statement = connection.createStatement()) {
+ statement.execute("CREATE DATABASE root.data_rf WITH
MAX_SCHEMA_REGION_GROUP_NUM=2");
+
+ try (ResultSet resultSet = statement.executeQuery("SHOW DATABASES
DETAILS root.data_rf")) {
+ Assert.assertTrue(resultSet.next());
+ Assert.assertEquals(2, resultSet.getInt("MaxSchemaRegionGroupNum"));
+ Assert.assertEquals(
+ expectedMaxDataRegionGroupNum,
resultSet.getInt("MaxDataRegionGroupNum"));
+ Assert.assertFalse(resultSet.next());
+ }
+ }
+ }
+
+ @Test
+ public void testMaxDataRegionGroupNumRejectedUnderAutoPolicy() throws
SQLException {
+ try (Connection connection = EnvFactory.getEnv().getConnection();
+ Statement statement = connection.createStatement()) {
+ SQLException createException =
+ Assert.assertThrows(
+ SQLException.class,
+ () ->
+ statement.execute(
+ "CREATE DATABASE root.auto_create WITH
MAX_DATA_REGION_GROUP_NUM=4"));
+ Assert.assertTrue(
+ createException
+ .getMessage()
+ .contains(
+ "max_data_region_group_num can only be set when "
+ + "data_region_group_extension_policy is CUSTOM"));
+
+ statement.execute("CREATE DATABASE root.auto_alter");
+ SQLException alterException =
+ Assert.assertThrows(
+ SQLException.class,
+ () ->
+ statement.execute(
+ "ALTER DATABASE root.auto_alter WITH
MAX_DATA_REGION_GROUP_NUM=4"));
+ Assert.assertTrue(
+ alterException
+ .getMessage()
+ .contains(
+ "max_data_region_group_num can only be set when "
+ + "data_region_group_extension_policy is CUSTOM"));
+ }
+ }
+}
diff --git
a/integration-test/src/test/java/org/apache/iotdb/confignode/it/database/IoTDBDatabaseMixedRegionGroupPolicyIT.java
b/integration-test/src/test/java/org/apache/iotdb/confignode/it/database/IoTDBDatabaseMixedRegionGroupPolicyIT.java
new file mode 100644
index 00000000000..f55f76dbe78
--- /dev/null
+++
b/integration-test/src/test/java/org/apache/iotdb/confignode/it/database/IoTDBDatabaseMixedRegionGroupPolicyIT.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.iotdb.confignode.it.database;
+
+import org.apache.iotdb.it.env.EnvFactory;
+import org.apache.iotdb.it.framework.IoTDBTestRunner;
+import org.apache.iotdb.itbase.category.ClusterIT;
+import org.apache.iotdb.itbase.category.LocalStandaloneIT;
+
+import org.junit.After;
+import org.junit.Assert;
+import org.junit.Before;
+import org.junit.Test;
+import org.junit.experimental.categories.Category;
+import org.junit.runner.RunWith;
+
+import java.sql.Connection;
+import java.sql.ResultSet;
+import java.sql.SQLException;
+import java.sql.Statement;
+
+@RunWith(IoTDBTestRunner.class)
+@Category({LocalStandaloneIT.class, ClusterIT.class})
+public class IoTDBDatabaseMixedRegionGroupPolicyIT {
+
+ private static final int DATA_REGION_PER_DATA_NODE = 4;
+
+ @Before
+ public void setUp() throws Exception {
+ EnvFactory.getEnv()
+ .getConfig()
+ .getCommonConfig()
+ .setSchemaRegionGroupExtensionPolicy("CUSTOM")
+ .setDataRegionGroupExtensionPolicy("AUTO")
+ .setDefaultSchemaRegionGroupNumPerDatabase(1)
+ .setDefaultDataRegionGroupNumPerDatabase(1)
+ .setDataReplicationFactor(1)
+ .setDataRegionPerDataNode(DATA_REGION_PER_DATA_NODE);
+ EnvFactory.getEnv().initClusterEnvironment(1, 1);
+ }
+
+ @After
+ public void tearDown() {
+ EnvFactory.getEnv().cleanClusterEnvironment();
+ }
+
+ @Test
+ public void testAutoPolicyStillAdjustsWhenTheOtherPolicyIsCustom() throws
SQLException {
+ try (Connection connection = EnvFactory.getEnv().getConnection();
+ Statement statement = connection.createStatement()) {
+ statement.execute("CREATE DATABASE root.mixed WITH
MAX_SCHEMA_REGION_GROUP_NUM=2");
+
+ try (ResultSet resultSet = statement.executeQuery("SHOW DATABASES
DETAILS root.mixed")) {
+ Assert.assertTrue(resultSet.next());
+ Assert.assertEquals(2, resultSet.getInt("MaxSchemaRegionGroupNum"));
+ Assert.assertEquals(DATA_REGION_PER_DATA_NODE,
resultSet.getInt("MaxDataRegionGroupNum"));
+ Assert.assertFalse(resultSet.next());
+ }
+ }
+ }
+}
diff --git
a/integration-test/src/test/java/org/apache/iotdb/confignode/it/database/IoTDBDatabaseRegionControlIT.java
b/integration-test/src/test/java/org/apache/iotdb/confignode/it/database/IoTDBDatabaseRegionControlIT.java
index b66fadeb35f..e0b87e98d1f 100644
---
a/integration-test/src/test/java/org/apache/iotdb/confignode/it/database/IoTDBDatabaseRegionControlIT.java
+++
b/integration-test/src/test/java/org/apache/iotdb/confignode/it/database/IoTDBDatabaseRegionControlIT.java
@@ -39,6 +39,7 @@ import org.junit.runner.RunWith;
import java.io.IOException;
import java.sql.Connection;
+import java.sql.ResultSet;
import java.sql.SQLException;
import java.sql.Statement;
import java.util.Collections;
@@ -58,6 +59,8 @@ public class IoTDBDatabaseRegionControlIT {
EnvFactory.getEnv()
.getConfig()
.getCommonConfig()
+ .setSchemaRegionGroupExtensionPolicy("CUSTOM")
+ .setDataRegionGroupExtensionPolicy("CUSTOM")
.setDefaultSchemaRegionGroupNumPerDatabase(testDefaultSchemaRegionGroupNumPerDatabase)
.setDefaultDataRegionGroupNumPerDatabase(testDefaultDataRegionGroupNumPerDatabase);
@@ -132,7 +135,7 @@ public class IoTDBDatabaseRegionControlIT {
final int testDataRegionGroupNum = 3;
String createDatabaseSQL =
String.format(
- "CREATE DATABASE %s WITH SCHEMA_REGION_GROUP_NUM=%d,
DATA_REGION_GROUP_NUM=%d;",
+ "CREATE DATABASE %s WITH MAX_SCHEMA_REGION_GROUP_NUM=%d,
MAX_DATA_REGION_GROUP_NUM=%d;",
database, testSchemaRegionGroupNum, testDataRegionGroupNum);
statement.execute(createDatabaseSQL);
insertBatchData(statement, database, 0);
@@ -162,6 +165,53 @@ public class IoTDBDatabaseRegionControlIT {
}
}
+ @Test
+ public void testRecreateWithPartialMaxRegionGroupNumUsesDefaultCounterpart()
throws SQLException {
+ try (final Connection connection = EnvFactory.getEnv().getConnection();
+ final Statement statement = connection.createStatement()) {
+ final String database = "root.rg_tree_create_schema";
+ statement.execute(String.format("CREATE DATABASE %s;", database));
+
+ try (final ResultSet resultSet =
+ statement.executeQuery("SHOW DATABASES DETAILS " + database)) {
+ Assert.assertTrue(resultSet.next());
+ Assert.assertEquals(
+ testDefaultSchemaRegionGroupNumPerDatabase,
+ resultSet.getInt("MaxSchemaRegionGroupNum"));
+ Assert.assertEquals(
+ testDefaultDataRegionGroupNumPerDatabase,
resultSet.getInt("MaxDataRegionGroupNum"));
+ Assert.assertFalse(resultSet.next());
+ }
+
+ statement.execute(String.format("DROP DATABASE %s;", database));
+ statement.execute(
+ String.format("CREATE DATABASE %s WITH
MAX_SCHEMA_REGION_GROUP_NUM=3;", database));
+
+ try (final ResultSet resultSet =
+ statement.executeQuery("SHOW DATABASES DETAILS " + database)) {
+ Assert.assertTrue(resultSet.next());
+ Assert.assertEquals(3, resultSet.getInt("MaxSchemaRegionGroupNum"));
+ Assert.assertEquals(
+ testDefaultDataRegionGroupNumPerDatabase,
resultSet.getInt("MaxDataRegionGroupNum"));
+ Assert.assertFalse(resultSet.next());
+ }
+
+ final String dataDatabase = "root.rg_tree_create_data";
+ statement.execute(
+ String.format("CREATE DATABASE %s WITH
MAX_DATA_REGION_GROUP_NUM=4;", dataDatabase));
+
+ try (final ResultSet resultSet =
+ statement.executeQuery("SHOW DATABASES DETAILS " + dataDatabase)) {
+ Assert.assertTrue(resultSet.next());
+ Assert.assertEquals(
+ testDefaultSchemaRegionGroupNumPerDatabase,
+ resultSet.getInt("MaxSchemaRegionGroupNum"));
+ Assert.assertEquals(4, resultSet.getInt("MaxDataRegionGroupNum"));
+ Assert.assertFalse(resultSet.next());
+ }
+ }
+ }
+
@Test
public void testRegionGroupNumControlThroughAlter()
throws SQLException, ClientManagerException, IOException,
InterruptedException, TException {
@@ -204,7 +254,7 @@ public class IoTDBDatabaseRegionControlIT {
final int testDataRegionGroupNum = 3;
String alterDatabaseSQL =
String.format(
- "ALTER DATABASE %s WITH SCHEMA_REGION_GROUP_NUM=%d,
DATA_REGION_GROUP_NUM=%d;",
+ "ALTER DATABASE %s WITH MAX_SCHEMA_REGION_GROUP_NUM=%d,
MAX_DATA_REGION_GROUP_NUM=%d;",
database, testSchemaRegionGroupNum, testDataRegionGroupNum);
statement.execute(alterDatabaseSQL);
insertBatchData(statement, database, batchSize);
@@ -233,4 +283,85 @@ public class IoTDBDatabaseRegionControlIT {
Assert.assertEquals(testDataRegionGroupNum, dataRegionGroupNum.get());
}
}
+
+ @Test
+ public void testAlterMaxDataRegionGroupNumCannotDecrease() throws
SQLException {
+ try (final Connection connection = EnvFactory.getEnv().getConnection();
+ final Statement statement = connection.createStatement()) {
+ final String database = "root.rg_tree_decrease";
+ statement.execute(
+ String.format("CREATE DATABASE %s WITH
MAX_DATA_REGION_GROUP_NUM=8;", database));
+ statement.execute(
+ String.format("ALTER DATABASE %s WITH
MAX_DATA_REGION_GROUP_NUM=16;", database));
+
+ final SQLException exception =
+ Assert.assertThrows(
+ SQLException.class,
+ () ->
+ statement.execute(
+ String.format(
+ "ALTER DATABASE %s WITH
MAX_DATA_REGION_GROUP_NUM=12;", database)));
+ Assert.assertTrue(
+ exception.getMessage(),
+ exception
+ .getMessage()
+ .contains(
+ "MaxDataRegionGroupNum should be greater than or equal to
current max "
+ + "DataRegionGroupNum: 16."));
+
+ try (final ResultSet resultSet =
+ statement.executeQuery("SHOW DATABASES DETAILS " + database)) {
+ Assert.assertTrue(resultSet.next());
+ Assert.assertEquals(16, resultSet.getInt("MaxDataRegionGroupNum"));
+ Assert.assertFalse(resultSet.next());
+ }
+ }
+ }
+
+ @Test
+ public void testAlterMaxSchemaRegionGroupNumCannotDecrease() throws
SQLException {
+ try (final Connection connection = EnvFactory.getEnv().getConnection();
+ final Statement statement = connection.createStatement()) {
+ final String schemaDatabase = "root.rg_tree_decrease_schema";
+ statement.execute(
+ String.format("CREATE DATABASE %s WITH
MAX_SCHEMA_REGION_GROUP_NUM=4;", schemaDatabase));
+ statement.execute(
+ String.format("ALTER DATABASE %s WITH
MAX_SCHEMA_REGION_GROUP_NUM=5;", schemaDatabase));
+
+ final SQLException schemaException =
+ Assert.assertThrows(
+ SQLException.class,
+ () ->
+ statement.execute(
+ String.format(
+ "ALTER DATABASE %s WITH
MAX_SCHEMA_REGION_GROUP_NUM=3;",
+ schemaDatabase)));
+ Assert.assertTrue(
+ schemaException.getMessage(),
+ schemaException
+ .getMessage()
+ .contains(
+ "MaxSchemaRegionGroupNum should be greater than or equal to
current max "
+ + "SchemaRegionGroupNum: 5."));
+
+ try (final ResultSet resultSet =
+ statement.executeQuery("SHOW DATABASES DETAILS " + schemaDatabase)) {
+ Assert.assertTrue(resultSet.next());
+ Assert.assertEquals(5, resultSet.getInt("MaxSchemaRegionGroupNum"));
+ Assert.assertFalse(resultSet.next());
+ }
+ }
+ }
+
+ @Test
+ public void testDeprecatedRegionGroupNumSqlRejected() throws SQLException {
+ try (final Connection connection = EnvFactory.getEnv().getConnection();
+ final Statement statement = connection.createStatement()) {
+ Assert.assertThrows(
+ SQLException.class,
+ () ->
+ statement.execute(
+ "CREATE DATABASE root.paradise3 WITH
SCHEMA_REGION_GROUP_NUM=3, DATA_REGION_GROUP_NUM=4;"));
+ }
+ }
}
diff --git
a/integration-test/src/test/java/org/apache/iotdb/confignode/it/database/IoTDBSystemDatabaseRegionGroupIT.java
b/integration-test/src/test/java/org/apache/iotdb/confignode/it/database/IoTDBSystemDatabaseRegionGroupIT.java
new file mode 100644
index 00000000000..b72b06f246e
--- /dev/null
+++
b/integration-test/src/test/java/org/apache/iotdb/confignode/it/database/IoTDBSystemDatabaseRegionGroupIT.java
@@ -0,0 +1,90 @@
+/*
+ * 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.iotdb.confignode.it.database;
+
+import org.apache.iotdb.it.env.EnvFactory;
+import org.apache.iotdb.it.framework.IoTDBTestRunner;
+import org.apache.iotdb.itbase.category.ClusterIT;
+import org.apache.iotdb.itbase.category.LocalStandaloneIT;
+
+import org.junit.After;
+import org.junit.Assert;
+import org.junit.Test;
+import org.junit.experimental.categories.Category;
+import org.junit.runner.RunWith;
+
+import java.sql.Connection;
+import java.sql.ResultSet;
+import java.sql.SQLException;
+import java.sql.Statement;
+
+/** Verifies the single-RegionGroup invariant of the system database. */
+@RunWith(IoTDBTestRunner.class)
+@Category({LocalStandaloneIT.class, ClusterIT.class})
+public class IoTDBSystemDatabaseRegionGroupIT {
+
+ private static final String SYSTEM_DATABASE = "root.__system";
+ private static final int DEFAULT_REGION_GROUP_NUM = 3;
+
+ @After
+ public void tearDown() {
+ EnvFactory.getEnv().cleanClusterEnvironment();
+ }
+
+ @Test
+ public void testSystemDatabaseRegionGroupNumUnderAutoPolicy() throws
Exception {
+ assertSystemDatabaseUsesSingleRegionGroup("AUTO");
+ }
+
+ @Test
+ public void testSystemDatabaseRegionGroupNumUnderCustomPolicy() throws
Exception {
+ assertSystemDatabaseUsesSingleRegionGroup("CUSTOM");
+ }
+
+ private static void assertSystemDatabaseUsesSingleRegionGroup(String policy)
throws Exception {
+ EnvFactory.getEnv()
+ .getConfig()
+ .getCommonConfig()
+ .setSchemaRegionGroupExtensionPolicy(policy)
+ .setDataRegionGroupExtensionPolicy(policy)
+ .setDefaultSchemaRegionGroupNumPerDatabase(DEFAULT_REGION_GROUP_NUM)
+ .setDefaultDataRegionGroupNumPerDatabase(DEFAULT_REGION_GROUP_NUM);
+ EnvFactory.getEnv().initClusterEnvironment(1, 1);
+
+ try (Connection connection = EnvFactory.getEnv().getConnection();
+ Statement statement = connection.createStatement()) {
+ statement.execute("INSERT INTO root.__system.audit.d1(timestamp, s)
VALUES(1, 1)");
+ assertSingleRegionGroup(statement);
+ }
+ }
+
+ private static void assertSingleRegionGroup(Statement statement) throws
SQLException {
+ try (ResultSet resultSet =
+ statement.executeQuery("SHOW DATABASES DETAILS " + SYSTEM_DATABASE)) {
+ Assert.assertTrue("Expected system database to be created",
resultSet.next());
+ Assert.assertEquals(SYSTEM_DATABASE, resultSet.getString("Database"));
+ Assert.assertEquals(1, resultSet.getInt("SchemaRegionGroupNum"));
+ Assert.assertEquals(1, resultSet.getInt("MaxSchemaRegionGroupNum"));
+ Assert.assertEquals(1, resultSet.getInt("DataRegionGroupNum"));
+ Assert.assertEquals(1, resultSet.getInt("MaxDataRegionGroupNum"));
+ Assert.assertFalse(resultSet.next());
+ }
+ }
+}
diff --git
a/integration-test/src/test/java/org/apache/iotdb/pipe/it/autocreate/IoTDBPipeIdempotentIT.java
b/integration-test/src/test/java/org/apache/iotdb/pipe/it/autocreate/IoTDBPipeIdempotentIT.java
index c5001f9e469..5470abfd528 100644
---
a/integration-test/src/test/java/org/apache/iotdb/pipe/it/autocreate/IoTDBPipeIdempotentIT.java
+++
b/integration-test/src/test/java/org/apache/iotdb/pipe/it/autocreate/IoTDBPipeIdempotentIT.java
@@ -63,6 +63,8 @@ public class IoTDBPipeIdempotentIT extends
AbstractPipeDualAutoIT {
// Limit the schemaRegion number to 1 to guarantee the after sql
executed on the same region
// of the tested idempotent sql.
.setDefaultSchemaRegionGroupNumPerDatabase(1)
+ .setSchemaRegionGroupExtensionPolicy("CUSTOM")
+ .setDataRegionGroupExtensionPolicy("CUSTOM")
.setConfigNodeConsensusProtocolClass(ConsensusFactory.RATIS_CONSENSUS)
.setSchemaRegionConsensusProtocolClass(ConsensusFactory.RATIS_CONSENSUS)
.setPipeMemoryManagementEnabled(false)
@@ -71,6 +73,8 @@ public class IoTDBPipeIdempotentIT extends
AbstractPipeDualAutoIT {
.getConfig()
.getCommonConfig()
.setAutoCreateSchemaEnabled(true)
+ .setSchemaRegionGroupExtensionPolicy("CUSTOM")
+ .setDataRegionGroupExtensionPolicy("CUSTOM")
.setConfigNodeConsensusProtocolClass(ConsensusFactory.RATIS_CONSENSUS)
.setSchemaRegionConsensusProtocolClass(ConsensusFactory.RATIS_CONSENSUS)
.setPipeMemoryManagementEnabled(false)
@@ -284,7 +288,7 @@ public class IoTDBPipeIdempotentIT extends
AbstractPipeDualAutoIT {
public void testAlterDatabaseIdempotent() throws Exception {
testIdempotent(
Collections.singletonList("create database root.sg1"),
- "ALTER DATABASE root.sg1 WITH SCHEMA_REGION_GROUP_NUM=2,
DATA_REGION_GROUP_NUM=3;",
+ "ALTER DATABASE root.sg1 WITH MAX_SCHEMA_REGION_GROUP_NUM=2,
MAX_DATA_REGION_GROUP_NUM=3;",
"create database root.sg2",
"count databases",
"count,",
diff --git
a/iotdb-core/antlr/src/main/antlr4/org/apache/iotdb/db/qp/sql/IdentifierParser.g4
b/iotdb-core/antlr/src/main/antlr4/org/apache/iotdb/db/qp/sql/IdentifierParser.g4
index 632911fb8af..0087a5335db 100644
---
a/iotdb-core/antlr/src/main/antlr4/org/apache/iotdb/db/qp/sql/IdentifierParser.g4
+++
b/iotdb-core/antlr/src/main/antlr4/org/apache/iotdb/db/qp/sql/IdentifierParser.g4
@@ -76,7 +76,7 @@ keyWords
| CREATE
| DATA
| DATA_REPLICATION_FACTOR
- | DATA_REGION_GROUP_NUM
+ | MAX_DATA_REGION_GROUP_NUM
| DATABASE
| DATABASES
| DATANODE
@@ -201,7 +201,7 @@ keyWords
| RUNNING
| SCHEMA
| SCHEMA_REPLICATION_FACTOR
- | SCHEMA_REGION_GROUP_NUM
+ | MAX_SCHEMA_REGION_GROUP_NUM
| SELECT
| SERIESSLOTID
| SESSION
@@ -273,4 +273,4 @@ keyWords
| OPTION
| INF
| CURRENT_TIMESTAMP
- ;
\ No newline at end of file
+ ;
diff --git
a/iotdb-core/antlr/src/main/antlr4/org/apache/iotdb/db/qp/sql/IoTDBSqlParser.g4
b/iotdb-core/antlr/src/main/antlr4/org/apache/iotdb/db/qp/sql/IoTDBSqlParser.g4
index f37ef736851..a8897f99224 100644
---
a/iotdb-core/antlr/src/main/antlr4/org/apache/iotdb/db/qp/sql/IoTDBSqlParser.g4
+++
b/iotdb-core/antlr/src/main/antlr4/org/apache/iotdb/db/qp/sql/IoTDBSqlParser.g4
@@ -115,8 +115,8 @@ databaseAttributeKey
| SCHEMA_REPLICATION_FACTOR
| DATA_REPLICATION_FACTOR
| TIME_PARTITION_INTERVAL
- | SCHEMA_REGION_GROUP_NUM
- | DATA_REGION_GROUP_NUM
+ | MAX_SCHEMA_REGION_GROUP_NUM
+ | MAX_DATA_REGION_GROUP_NUM
;
// ---- Drop Database
diff --git
a/iotdb-core/antlr/src/main/antlr4/org/apache/iotdb/db/qp/sql/SqlLexer.g4
b/iotdb-core/antlr/src/main/antlr4/org/apache/iotdb/db/qp/sql/SqlLexer.g4
index f90a0d3dbfd..d5847b2d47d 100644
--- a/iotdb-core/antlr/src/main/antlr4/org/apache/iotdb/db/qp/sql/SqlLexer.g4
+++ b/iotdb-core/antlr/src/main/antlr4/org/apache/iotdb/db/qp/sql/SqlLexer.g4
@@ -1094,12 +1094,12 @@ TIME_PARTITION_INTERVAL
: T I M E '_' P A R T I T I O N '_' I N T E R V A L
;
-SCHEMA_REGION_GROUP_NUM
- : S C H E M A '_' R E G I O N '_' G R O U P '_' N U M
+MAX_SCHEMA_REGION_GROUP_NUM
+ : M A X '_' S C H E M A '_' R E G I O N '_' G R O U P '_' N U M
;
-DATA_REGION_GROUP_NUM
- : D A T A '_' R E G I O N '_' G R O U P '_' N U M
+MAX_DATA_REGION_GROUP_NUM
+ : M A X '_' D A T A '_' R E G I O N '_' G R O U P '_' N U M
;
CURRENT_TIMESTAMP
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/partition/PartitionManager.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/partition/PartitionManager.java
index 056b819d5bc..7db7ac50d62 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/partition/PartitionManager.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/partition/PartitionManager.java
@@ -548,14 +548,14 @@ public class PartitionManager {
for (Map.Entry<String, Integer> entry :
unassignedPartitionSlotsCountMap.entrySet()) {
final String database = entry.getKey();
- int minRegionGroupNum =
- getClusterSchemaManager().getMinRegionGroupNum(database,
consensusGroupType);
+ int maxRegionGroupNum =
+ getClusterSchemaManager().getMaxRegionGroupNum(database,
consensusGroupType);
int allocatedRegionGroupCount =
partitionInfo.getRegionGroupCount(database, consensusGroupType);
- // Extend RegionGroups until allocatedRegionGroupCount ==
minRegionGroupNum
- if (allocatedRegionGroupCount < minRegionGroupNum) {
- allotmentMap.put(database, minRegionGroupNum -
allocatedRegionGroupCount);
+ // Extend RegionGroups until allocatedRegionGroupCount ==
maxRegionGroupNum
+ if (allocatedRegionGroupCount < maxRegionGroupNum) {
+ allotmentMap.put(database, maxRegionGroupNum -
allocatedRegionGroupCount);
}
}
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/schema/ClusterSchemaManager.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/schema/ClusterSchemaManager.java
index 56fabd6aa53..e267cfd9250 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/schema/ClusterSchemaManager.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/schema/ClusterSchemaManager.java
@@ -67,6 +67,7 @@ import
org.apache.iotdb.confignode.manager.consensus.ConsensusManager;
import org.apache.iotdb.confignode.manager.node.NodeManager;
import org.apache.iotdb.confignode.manager.partition.PartitionManager;
import org.apache.iotdb.confignode.manager.partition.PartitionMetrics;
+import
org.apache.iotdb.confignode.manager.partition.RegionGroupExtensionPolicy;
import org.apache.iotdb.confignode.persistence.schema.ClusterSchemaInfo;
import org.apache.iotdb.confignode.rpc.thrift.TDatabaseInfo;
import org.apache.iotdb.confignode.rpc.thrift.TDatabaseSchema;
@@ -215,31 +216,23 @@ public class ClusterSchemaManager {
return result;
}
- if (databaseSchema.isSetMinSchemaRegionGroupNum()) {
- // Validate alter SchemaRegionGroupNum
- int minSchemaRegionGroupNum =
- getMinRegionGroupNum(databaseSchema.getName(),
TConsensusGroupType.SchemaRegion);
- if (databaseSchema.getMinSchemaRegionGroupNum() <=
minSchemaRegionGroupNum) {
- result = new
TSStatus(TSStatusCode.DATABASE_CONFIG_ERROR.getStatusCode());
- result.setMessage(
- String.format(
- "Failed to alter database. The SchemaRegionGroupNum could only
be increased. "
- + "Current SchemaRegionGroupNum: %d, Alter
SchemaRegionGroupNum: %d",
- minSchemaRegionGroupNum,
databaseSchema.getMinSchemaRegionGroupNum()));
+ if (databaseSchema.isSetMaxSchemaRegionGroupNum()) {
+ result =
+ validateMaxRegionGroupNumOnAlter(
+ databaseSchema.getName(),
+ TConsensusGroupType.SchemaRegion,
+ databaseSchema.getMaxSchemaRegionGroupNum());
+ if (result.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
return result;
}
}
- if (databaseSchema.isSetMinDataRegionGroupNum()) {
- // Validate alter DataRegionGroupNum
- int minDataRegionGroupNum =
- getMinRegionGroupNum(databaseSchema.getName(),
TConsensusGroupType.DataRegion);
- if (databaseSchema.getMinDataRegionGroupNum() <= minDataRegionGroupNum) {
- result = new
TSStatus(TSStatusCode.DATABASE_CONFIG_ERROR.getStatusCode());
- result.setMessage(
- String.format(
- "Failed to alter database. The DataRegionGroupNum could only
be increased. "
- + "Current DataRegionGroupNum: %d, Alter
DataRegionGroupNum: %d",
- minDataRegionGroupNum,
databaseSchema.getMinDataRegionGroupNum()));
+ if (databaseSchema.isSetMaxDataRegionGroupNum()) {
+ result =
+ validateMaxRegionGroupNumOnAlter(
+ databaseSchema.getName(),
+ TConsensusGroupType.DataRegion,
+ databaseSchema.getMaxDataRegionGroupNum());
+ if (result.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
return result;
}
}
@@ -453,6 +446,14 @@ public class ClusterSchemaManager {
* each Database based on existing cluster resources
*/
public synchronized void adjustMaxRegionGroupNum() {
+ final boolean isAdjustSchemaRegionGroupNum =
+
!CONF.getSchemaRegionGroupExtensionPolicy().equals(RegionGroupExtensionPolicy.CUSTOM);
+ final boolean isAdjustDataRegionGroupNum =
+
!CONF.getDataRegionGroupExtensionPolicy().equals(RegionGroupExtensionPolicy.CUSTOM);
+ if (!isAdjustSchemaRegionGroupNum && !isAdjustDataRegionGroupNum) {
+ return;
+ }
+
// Get all DatabaseSchemas
Map<String, TDatabaseSchema> databaseSchemaMap =
getMatchedDatabaseSchemasByName(getDatabaseNames());
@@ -480,60 +481,38 @@ public class ClusterSchemaManager {
continue;
}
- // Adjust maxSchemaRegionGroupNum for each Database.
- // All Databases share the DataNodes equally.
- // The allocated SchemaRegionGroups will not be shrunk.
- final int allocatedSchemaRegionGroupCount;
- try {
- allocatedSchemaRegionGroupCount =
- getPartitionManager()
- .getRegionGroupCount(databaseSchema.getName(),
TConsensusGroupType.SchemaRegion);
- } catch (final DatabaseNotExistsException e) {
- // ignore the pre deleted database
- continue;
+ int maxSchemaRegionGroupNum =
databaseSchema.getMaxSchemaRegionGroupNum();
+ if (isAdjustSchemaRegionGroupNum) {
+ try {
+ maxSchemaRegionGroupNum =
+ adjustRegionGroupNum(
+ TConsensusGroupType.SchemaRegion,
+ databaseSchema,
+ dataNodeNum,
+ databaseNum,
+ totalCpuCoreNum);
+ } catch (final DatabaseNotExistsException e) {
+ // ignore the pre deleted database
+ continue;
+ }
}
- final int maxSchemaRegionGroupNum =
- calcMaxRegionGroupNum(
- databaseSchema.getMinSchemaRegionGroupNum(),
- SCHEMA_REGION_PER_DATA_NODE,
- dataNodeNum,
- databaseNum,
- databaseSchema.getSchemaReplicationFactor(),
- allocatedSchemaRegionGroupCount);
- LOGGER.info(
- "[AdjustRegionGroupNum] The maximum number of SchemaRegionGroups for
Database: {} is adjusted to: {}",
- databaseSchema.getName(),
- maxSchemaRegionGroupNum);
-
- // Adjust maxDataRegionGroupNum for each Database.
- // All Databases share the DataNodes equally.
- // The allocated DataRegionGroups will not be shrunk.
- final int allocatedDataRegionGroupCount;
- try {
- allocatedDataRegionGroupCount =
- getPartitionManager()
- .getRegionGroupCount(databaseSchema.getName(),
TConsensusGroupType.DataRegion);
- } catch (final DatabaseNotExistsException e) {
- // ignore the pre deleted database
- continue;
+ int maxDataRegionGroupNum = databaseSchema.getMaxDataRegionGroupNum();
+ if (isAdjustDataRegionGroupNum) {
+ try {
+ maxDataRegionGroupNum =
+ adjustRegionGroupNum(
+ TConsensusGroupType.DataRegion,
+ databaseSchema,
+ dataNodeNum,
+ databaseNum,
+ totalCpuCoreNum);
+ } catch (final DatabaseNotExistsException e) {
+ // ignore the pre deleted database
+ continue;
+ }
}
- final int maxDataRegionGroupNum =
- calcMaxRegionGroupNum(
- databaseSchema.getMinDataRegionGroupNum(),
- DATA_REGION_PER_DATA_NODE == 0
- ? CONF.getDataRegionPerDataNodeProportion()
- : DATA_REGION_PER_DATA_NODE,
- DATA_REGION_PER_DATA_NODE == 0 ? totalCpuCoreNum : dataNodeNum,
- databaseNum,
- databaseSchema.getDataReplicationFactor(),
- allocatedDataRegionGroupCount);
- LOGGER.info(
- "[AdjustRegionGroupNum] The maximum number of DataRegionGroups for
Database: {} is adjusted to: {}",
- databaseSchema.getName(),
- maxDataRegionGroupNum);
-
adjustMaxRegionGroupNumPlan.putEntry(
databaseSchema.getName(), new Pair<>(maxSchemaRegionGroupNum,
maxDataRegionGroupNum));
}
@@ -566,6 +545,135 @@ public class ClusterSchemaManager {
allocatedRegionGroupCount));
}
+ /** Adjust the automatic max quota of a schema or data RegionGroup. */
+ public int adjustRegionGroupNum(
+ TConsensusGroupType consensusGroupType,
+ TDatabaseSchema databaseSchema,
+ int dataNodeNum,
+ int databaseNum,
+ int totalCpuCoreNum)
+ throws DatabaseNotExistsException {
+ final boolean isSchemaRegion = consensusGroupType ==
TConsensusGroupType.SchemaRegion;
+ final int allocatedRegionGroupCount =
+ getPartitionManager().getRegionGroupCount(databaseSchema.getName(),
consensusGroupType);
+
+ int maxRegionGroupNum =
+ calcMaxRegionGroupNum(
+ isSchemaRegion
+ ? databaseSchema.getMinSchemaRegionGroupNum()
+ : databaseSchema.getMinDataRegionGroupNum(),
+ isSchemaRegion
+ ? SCHEMA_REGION_PER_DATA_NODE
+ : (DATA_REGION_PER_DATA_NODE == 0
+ ? CONF.getDataRegionPerDataNodeProportion()
+ : DATA_REGION_PER_DATA_NODE),
+ isSchemaRegion
+ ? dataNodeNum
+ : (DATA_REGION_PER_DATA_NODE == 0 ? totalCpuCoreNum :
dataNodeNum),
+ databaseNum,
+ isSchemaRegion
+ ? databaseSchema.getSchemaReplicationFactor()
+ : databaseSchema.getDataReplicationFactor(),
+ allocatedRegionGroupCount);
+ LOGGER.info(
+ "[AdjustRegionGroupNum] The maximum number of {}RegionGroups for
Database: {} is adjusted to: {}",
+ isSchemaRegion ? "Schema" : "Data",
+ databaseSchema.getName(),
+ maxRegionGroupNum);
+
+ return maxRegionGroupNum;
+ }
+
+ public static TSStatus validateMaxRegionGroupNumOnCreation(
+ TDatabaseSchema databaseSchema, TConsensusGroupType consensusGroupType) {
+ return validateMaxRegionGroupNum(
+ consensusGroupType,
+ TConsensusGroupType.SchemaRegion.equals(consensusGroupType)
+ ? databaseSchema.getMaxSchemaRegionGroupNum()
+ : databaseSchema.getMaxDataRegionGroupNum(),
+ true);
+ }
+
+ private static TSStatus validateMaxRegionGroupNum(
+ TConsensusGroupType consensusGroupType, int maxRegionGroupNum, boolean
isCreate) {
+ final boolean isSchemaRegion =
TConsensusGroupType.SchemaRegion.equals(consensusGroupType);
+ final RegionGroupExtensionPolicy policy =
+ isSchemaRegion
+ ? CONF.getSchemaRegionGroupExtensionPolicy()
+ : CONF.getDataRegionGroupExtensionPolicy();
+ final String configKey =
+ isSchemaRegion ? "max_schema_region_group_num" :
"max_data_region_group_num";
+ final String fieldName = isSchemaRegion ? "MaxSchemaRegionGroupNum" :
"MaxDataRegionGroupNum";
+
+ if (!policy.equals(RegionGroupExtensionPolicy.CUSTOM)) {
+ return new TSStatus(TSStatusCode.DATABASE_CONFIG_ERROR.getStatusCode())
+ .setMessage(
+ String.format(
+ "Failed to %s database. The %s can only be set when
%s_region_group_extension_policy is CUSTOM.",
+ isCreate ? "create" : "alter", configKey, isSchemaRegion ?
"schema" : "data"));
+ }
+
+ final int defaultRegionGroupNum =
+ isSchemaRegion
+ ? CONF.getDefaultSchemaRegionGroupNumPerDatabase()
+ : CONF.getDefaultDataRegionGroupNumPerDatabase();
+ if (maxRegionGroupNum < defaultRegionGroupNum) {
+ return new TSStatus(TSStatusCode.DATABASE_CONFIG_ERROR.getStatusCode())
+ .setMessage(
+ String.format(
+ "%s should be greater than or equal to default
%sRegionGroupNum: %d.",
+ fieldName, isSchemaRegion ? "Schema" : "Data",
defaultRegionGroupNum));
+ }
+
+ return StatusUtils.OK;
+ }
+
+ private TSStatus validateMaxRegionGroupNumOnAlter(
+ String database, TConsensusGroupType consensusGroupType, int
maxRegionGroupNum) {
+ TSStatus status = validateMaxRegionGroupNum(consensusGroupType,
maxRegionGroupNum, false);
+ if (status.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ return status;
+ }
+
+ final boolean isSchemaRegion =
TConsensusGroupType.SchemaRegion.equals(consensusGroupType);
+ final String fieldName = isSchemaRegion ? "MaxSchemaRegionGroupNum" :
"MaxDataRegionGroupNum";
+ final int minRegionGroupNum = getMinRegionGroupNum(database,
consensusGroupType);
+ if (maxRegionGroupNum < minRegionGroupNum) {
+ return new TSStatus(TSStatusCode.DATABASE_CONFIG_ERROR.getStatusCode())
+ .setMessage(
+ String.format(
+ "%s should be greater than or equal to current min
%sRegionGroupNum: %d.",
+ fieldName, isSchemaRegion ? "Schema" : "Data",
minRegionGroupNum));
+ }
+
+ final int currentMaxRegionGroupNum = getMaxRegionGroupNum(database,
consensusGroupType);
+ if (maxRegionGroupNum < currentMaxRegionGroupNum) {
+ return new TSStatus(TSStatusCode.DATABASE_CONFIG_ERROR.getStatusCode())
+ .setMessage(
+ String.format(
+ "%s should be greater than or equal to current max
%sRegionGroupNum: %d.",
+ fieldName, isSchemaRegion ? "Schema" : "Data",
currentMaxRegionGroupNum));
+ }
+
+ final int allocatedRegionGroupCount;
+ try {
+ allocatedRegionGroupCount =
+ getPartitionManager().getRegionGroupCount(database,
consensusGroupType);
+ } catch (DatabaseNotExistsException e) {
+ return new TSStatus(TSStatusCode.DATABASE_NOT_EXIST.getStatusCode())
+ .setMessage(e.getMessage());
+ }
+ if (maxRegionGroupNum < allocatedRegionGroupCount) {
+ return new TSStatus(TSStatusCode.DATABASE_CONFIG_ERROR.getStatusCode())
+ .setMessage(
+ String.format(
+ "%s should be greater than or equal to allocated
%sRegionGroupNum: %d.",
+ fieldName, isSchemaRegion ? "Schema" : "Data",
allocatedRegionGroupCount));
+ }
+
+ return StatusUtils.OK;
+ }
+
// ======================================================
// Leader scheduling interfaces
// ======================================================
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/schema/ClusterSchemaInfo.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/schema/ClusterSchemaInfo.java
index 4d6f2522052..5a659907402 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/schema/ClusterSchemaInfo.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/schema/ClusterSchemaInfo.java
@@ -178,31 +178,15 @@ public class ClusterSchemaInfo implements
SnapshotProcessor {
TDatabaseSchema currentSchema =
mTree.getDatabaseNodeByDatabasePath(partialPathName).getAsMNode().getDatabaseSchema();
// TODO: Support alter other fields
- if (alterSchema.isSetMinSchemaRegionGroupNum()) {
-
currentSchema.setMinSchemaRegionGroupNum(alterSchema.getMinSchemaRegionGroupNum());
- currentSchema.setMaxSchemaRegionGroupNum(
- Math.max(
- currentSchema.getMinSchemaRegionGroupNum(),
- currentSchema.getMaxSchemaRegionGroupNum()));
- LOGGER.info(
- "[AdjustRegionGroupNum] The minimum number of SchemaRegionGroups
for Database: {} is adjusted to: {}",
- currentSchema.getName(),
- currentSchema.getMinSchemaRegionGroupNum());
+ if (alterSchema.isSetMaxSchemaRegionGroupNum()) {
+
currentSchema.setMaxSchemaRegionGroupNum(alterSchema.getMaxSchemaRegionGroupNum());
LOGGER.info(
"[AdjustRegionGroupNum] The maximum number of SchemaRegionGroups
for Database: {} is adjusted to: {}",
currentSchema.getName(),
currentSchema.getMaxSchemaRegionGroupNum());
}
- if (alterSchema.isSetMinDataRegionGroupNum()) {
-
currentSchema.setMinDataRegionGroupNum(alterSchema.getMinDataRegionGroupNum());
- currentSchema.setMaxDataRegionGroupNum(
- Math.max(
- currentSchema.getMinDataRegionGroupNum(),
- currentSchema.getMaxDataRegionGroupNum()));
- LOGGER.info(
- "[AdjustRegionGroupNum] The minimum number of DataRegionGroups for
Database: {} is adjusted to: {}",
- currentSchema.getName(),
- currentSchema.getMinDataRegionGroupNum());
+ if (alterSchema.isSetMaxDataRegionGroupNum()) {
+
currentSchema.setMaxDataRegionGroupNum(alterSchema.getMaxDataRegionGroupNum());
LOGGER.info(
"[AdjustRegionGroupNum] The maximum number of DataRegionGroups for
Database: {} is adjusted to: {}",
currentSchema.getName(),
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/service/thrift/ConfigNodeRPCServiceProcessor.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/service/thrift/ConfigNodeRPCServiceProcessor.java
index 289e40ecec3..9fbb6b2e657 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/service/thrift/ConfigNodeRPCServiceProcessor.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/service/thrift/ConfigNodeRPCServiceProcessor.java
@@ -20,6 +20,7 @@
package org.apache.iotdb.confignode.service.thrift;
import org.apache.iotdb.common.rpc.thrift.TConfigNodeLocation;
+import org.apache.iotdb.common.rpc.thrift.TConsensusGroupType;
import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation;
import org.apache.iotdb.common.rpc.thrift.TFlushReq;
import org.apache.iotdb.common.rpc.thrift.TNodeLocations;
@@ -78,6 +79,7 @@ import
org.apache.iotdb.confignode.consensus.response.partition.RegionInfoListRe
import org.apache.iotdb.confignode.consensus.response.ttl.ShowTTLResp;
import org.apache.iotdb.confignode.manager.ConfigManager;
import org.apache.iotdb.confignode.manager.consensus.ConsensusManager;
+import org.apache.iotdb.confignode.manager.schema.ClusterSchemaManager;
import org.apache.iotdb.confignode.rpc.thrift.IConfigNodeRPCService;
import org.apache.iotdb.confignode.rpc.thrift.TAINodeConfigurationResp;
import org.apache.iotdb.confignode.rpc.thrift.TAINodeRegisterReq;
@@ -432,6 +434,7 @@ public class ConfigNodeRPCServiceProcessor implements
IConfigNodeRPCService.Ifac
if (isSystemDatabase) {
databaseSchema.setMinSchemaRegionGroupNum(1);
+ databaseSchema.setMaxSchemaRegionGroupNum(1);
} else if (!databaseSchema.isSetMinSchemaRegionGroupNum()) {
databaseSchema.setMinSchemaRegionGroupNum(
configNodeConfig.getDefaultSchemaRegionGroupNumPerDatabase());
@@ -444,6 +447,7 @@ public class ConfigNodeRPCServiceProcessor implements
IConfigNodeRPCService.Ifac
if (isSystemDatabase) {
databaseSchema.setMinDataRegionGroupNum(1);
+ databaseSchema.setMaxDataRegionGroupNum(1);
} else if (!databaseSchema.isSetMinDataRegionGroupNum()) {
databaseSchema.setMinDataRegionGroupNum(
configNodeConfig.getDefaultDataRegionGroupNumPerDatabase());
@@ -453,14 +457,36 @@ public class ConfigNodeRPCServiceProcessor implements
IConfigNodeRPCService.Ifac
.setMessage("Failed to create database. The dataRegionGroupNum
should be positive.");
}
+ if (!isSystemDatabase) {
+ if (databaseSchema.isSetMaxSchemaRegionGroupNum()) {
+ TSStatus status =
+ ClusterSchemaManager.validateMaxRegionGroupNumOnCreation(
+ databaseSchema, TConsensusGroupType.SchemaRegion);
+ if (status.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ errorResp = status;
+ }
+ }
+ if (databaseSchema.isSetMaxDataRegionGroupNum()) {
+ TSStatus status =
+ ClusterSchemaManager.validateMaxRegionGroupNumOnCreation(
+ databaseSchema, TConsensusGroupType.DataRegion);
+ if (status.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ errorResp = status;
+ }
+ }
+ }
+
if (errorResp != null) {
LOGGER.warn("Execute SetDatabase: {} with result: {}", databaseSchema,
errorResp);
return errorResp;
}
- // The maxRegionGroupNum is equal to the minRegionGroupNum when initialize
-
databaseSchema.setMaxSchemaRegionGroupNum(databaseSchema.getMinSchemaRegionGroupNum());
-
databaseSchema.setMaxDataRegionGroupNum(databaseSchema.getMinDataRegionGroupNum());
+ if (!databaseSchema.isSetMaxSchemaRegionGroupNum()) {
+
databaseSchema.setMaxSchemaRegionGroupNum(databaseSchema.getMinSchemaRegionGroupNum());
+ }
+ if (!databaseSchema.isSetMaxDataRegionGroupNum()) {
+
databaseSchema.setMaxDataRegionGroupNum(databaseSchema.getMinDataRegionGroupNum());
+ }
DatabaseSchemaPlan setPlan =
new DatabaseSchemaPlan(ConfigPhysicalPlanType.CreateDatabase,
databaseSchema);
diff --git
a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/service/thrift/ConfigNodeRPCServiceProcessorTest.java
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/service/thrift/ConfigNodeRPCServiceProcessorTest.java
index c4d993b5a79..75ef34c9975 100644
---
a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/service/thrift/ConfigNodeRPCServiceProcessorTest.java
+++
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/service/thrift/ConfigNodeRPCServiceProcessorTest.java
@@ -24,13 +24,16 @@ import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation;
import org.apache.iotdb.common.rpc.thrift.TEndPoint;
import org.apache.iotdb.common.rpc.thrift.TSStatus;
import org.apache.iotdb.commons.conf.CommonConfig;
+import org.apache.iotdb.commons.schema.SchemaConstant;
import org.apache.iotdb.confignode.conf.ConfigNodeConfig;
+import
org.apache.iotdb.confignode.consensus.request.write.database.DatabaseSchemaPlan;
import
org.apache.iotdb.confignode.consensus.response.datanode.DataNodeRegisterResp;
import org.apache.iotdb.confignode.manager.ConfigManager;
import org.apache.iotdb.confignode.rpc.thrift.TDataNodeRegisterReq;
import org.apache.iotdb.confignode.rpc.thrift.TDataNodeRegisterResp;
import org.apache.iotdb.confignode.rpc.thrift.TDataNodeRestartReq;
import org.apache.iotdb.confignode.rpc.thrift.TDataNodeRestartResp;
+import org.apache.iotdb.confignode.rpc.thrift.TDatabaseSchema;
import org.apache.iotdb.confignode.rpc.thrift.TRuntimeConfiguration;
import org.apache.iotdb.confignode.service.ConfigNode;
import org.apache.iotdb.rpc.TSStatusCode;
@@ -48,6 +51,37 @@ import java.util.Collections;
public class ConfigNodeRPCServiceProcessorTest extends TestCase {
+ public void testSetSystemDatabaseRegionGroupQuota() {
+ CommonConfig commonConfig = Mockito.mock(CommonConfig.class);
+ ConfigNodeConfig configNodeConfig = Mockito.mock(ConfigNodeConfig.class);
+ ConfigNode configNode = Mockito.mock(ConfigNode.class);
+ ConfigManager configManager = Mockito.mock(ConfigManager.class);
+
Mockito.when(configNodeConfig.getDefaultSchemaRegionGroupNumPerDatabase()).thenReturn(3);
+
Mockito.when(configNodeConfig.getDefaultDataRegionGroupNumPerDatabase()).thenReturn(3);
+
Mockito.when(configManager.setDatabase(Mockito.any(DatabaseSchemaPlan.class)))
+ .thenReturn(new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode()));
+ ConfigNodeRPCServiceProcessor sut =
+ new ConfigNodeRPCServiceProcessor(
+ commonConfig, configNodeConfig, configNode, configManager);
+
+ TDatabaseSchema systemDatabase =
+ new TDatabaseSchema(SchemaConstant.SYSTEM_DATABASE)
+ .setMaxSchemaRegionGroupNum(1)
+ .setMaxDataRegionGroupNum(1);
+
+ TSStatus result = sut.setDatabase(systemDatabase);
+
+ Assert.assertEquals(TSStatusCode.SUCCESS_STATUS.getStatusCode(),
result.getCode());
+ ArgumentCaptor<DatabaseSchemaPlan> planCaptor =
+ ArgumentCaptor.forClass(DatabaseSchemaPlan.class);
+ Mockito.verify(configManager).setDatabase(planCaptor.capture());
+ TDatabaseSchema actualSchema = planCaptor.getValue().getSchema();
+ Assert.assertEquals(1, actualSchema.getMinSchemaRegionGroupNum());
+ Assert.assertEquals(1, actualSchema.getMaxSchemaRegionGroupNum());
+ Assert.assertEquals(1, actualSchema.getMinDataRegionGroupNum());
+ Assert.assertEquals(1, actualSchema.getMaxDataRegionGroupNum());
+ }
+
/**
* This test should be a normal data-node registration where a valid ip is
used as address of the
* rpc-service. Nothing special should happen here.
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/metadata/DatabaseSchemaTask.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/metadata/DatabaseSchemaTask.java
index 48eb140f094..4af0e64e269 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/metadata/DatabaseSchemaTask.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/metadata/DatabaseSchemaTask.java
@@ -66,11 +66,12 @@ public class DatabaseSchemaTask implements IConfigTask {
if (databaseSchemaStatement.getTimePartitionInterval() != null) {
databaseSchema.setTimePartitionInterval(databaseSchemaStatement.getTimePartitionInterval());
}
- if (databaseSchemaStatement.getSchemaRegionGroupNum() != null) {
-
databaseSchema.setMinSchemaRegionGroupNum(databaseSchemaStatement.getSchemaRegionGroupNum());
+ if (databaseSchemaStatement.getMaxSchemaRegionGroupNum() != null) {
+ databaseSchema.setMaxSchemaRegionGroupNum(
+ databaseSchemaStatement.getMaxSchemaRegionGroupNum());
}
- if (databaseSchemaStatement.getDataRegionGroupNum() != null) {
-
databaseSchema.setMinDataRegionGroupNum(databaseSchemaStatement.getDataRegionGroupNum());
+ if (databaseSchemaStatement.getMaxDataRegionGroupNum() != null) {
+
databaseSchema.setMaxDataRegionGroupNum(databaseSchemaStatement.getMaxDataRegionGroupNum());
}
return databaseSchema;
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/parser/ASTVisitor.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/parser/ASTVisitor.java
index 77af818e124..75c841d4140 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/parser/ASTVisitor.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/parser/ASTVisitor.java
@@ -2630,12 +2630,12 @@ public class ASTVisitor extends
IoTDBSqlParserBaseVisitor<Statement> {
} else if (attributeKey.TIME_PARTITION_INTERVAL() != null) {
long timePartitionInterval =
Long.parseLong(attribute.INTEGER_LITERAL().getText());
databaseSchemaStatement.setTimePartitionInterval(timePartitionInterval);
- } else if (attributeKey.SCHEMA_REGION_GROUP_NUM() != null) {
- int schemaRegionGroupNum =
Integer.parseInt(attribute.INTEGER_LITERAL().getText());
- databaseSchemaStatement.setSchemaRegionGroupNum(schemaRegionGroupNum);
- } else if (attributeKey.DATA_REGION_GROUP_NUM() != null) {
- int dataRegionGroupNum =
Integer.parseInt(attribute.INTEGER_LITERAL().getText());
- databaseSchemaStatement.setDataRegionGroupNum(dataRegionGroupNum);
+ } else if (attributeKey.MAX_SCHEMA_REGION_GROUP_NUM() != null) {
+ int maxSchemaRegionGroupNum =
Integer.parseInt(attribute.INTEGER_LITERAL().getText());
+
databaseSchemaStatement.setMaxSchemaRegionGroupNum(maxSchemaRegionGroupNum);
+ } else if (attributeKey.MAX_DATA_REGION_GROUP_NUM() != null) {
+ int maxDataRegionGroupNum =
Integer.parseInt(attribute.INTEGER_LITERAL().getText());
+
databaseSchemaStatement.setMaxDataRegionGroupNum(maxDataRegionGroupNum);
}
}
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/metadata/DatabaseSchemaStatement.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/metadata/DatabaseSchemaStatement.java
index 24fd17bd80a..53c6331e6e2 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/metadata/DatabaseSchemaStatement.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/metadata/DatabaseSchemaStatement.java
@@ -42,8 +42,8 @@ public class DatabaseSchemaStatement extends Statement
implements IConfigStateme
private Integer schemaReplicationFactor = null;
private Integer dataReplicationFactor = null;
private Long timePartitionInterval = null;
- private Integer schemaRegionGroupNum = null;
- private Integer dataRegionGroupNum = null;
+ private Integer maxSchemaRegionGroupNum = null;
+ private Integer maxDataRegionGroupNum = null;
private boolean enablePrintExceptionLog = true;
public DatabaseSchemaStatement(DatabaseSchemaStatementType subType) {
@@ -96,20 +96,20 @@ public class DatabaseSchemaStatement extends Statement
implements IConfigStateme
this.timePartitionInterval = timePartitionInterval;
}
- public Integer getSchemaRegionGroupNum() {
- return schemaRegionGroupNum;
+ public Integer getMaxSchemaRegionGroupNum() {
+ return maxSchemaRegionGroupNum;
}
- public void setSchemaRegionGroupNum(Integer schemaRegionGroupNum) {
- this.schemaRegionGroupNum = schemaRegionGroupNum;
+ public void setMaxSchemaRegionGroupNum(Integer maxSchemaRegionGroupNum) {
+ this.maxSchemaRegionGroupNum = maxSchemaRegionGroupNum;
}
- public Integer getDataRegionGroupNum() {
- return dataRegionGroupNum;
+ public Integer getMaxDataRegionGroupNum() {
+ return maxDataRegionGroupNum;
}
- public void setDataRegionGroupNum(Integer dataRegionGroupNum) {
- this.dataRegionGroupNum = dataRegionGroupNum;
+ public void setMaxDataRegionGroupNum(Integer maxDataRegionGroupNum) {
+ this.maxDataRegionGroupNum = maxDataRegionGroupNum;
}
public boolean getEnablePrintExceptionLog() {
@@ -164,10 +164,10 @@ public class DatabaseSchemaStatement extends Statement
implements IConfigStateme
+ dataReplicationFactor
+ ", timePartitionInterval="
+ timePartitionInterval
- + ", schemaRegionGroupNum="
- + schemaRegionGroupNum
- + ", dataRegionGroupNum="
- + dataRegionGroupNum
+ + ", maxSchemaRegionGroupNum="
+ + maxSchemaRegionGroupNum
+ + ", maxDataRegionGroupNum="
+ + maxDataRegionGroupNum
+ '}';
}
diff --git
a/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/reporter/iotdb/IoTDBSessionReporter.java
b/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/reporter/iotdb/IoTDBSessionReporter.java
index b70ecebfff2..3603bac1a81 100644
---
a/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/reporter/iotdb/IoTDBSessionReporter.java
+++
b/iotdb-core/metrics/interface/src/main/java/org/apache/iotdb/metrics/reporter/iotdb/IoTDBSessionReporter.java
@@ -75,9 +75,7 @@ public class IoTDBSessionReporter extends IoTDBReporter {
if (!result.hasNext()) {
try (SessionDataSetWrapper result2 =
this.sessionPool.executeQueryStatement(
- "CREATE DATABASE "
- + metricConfig.getInternalDatabase()
- + " WITH SCHEMA_REGION_GROUP_NUM=1,
DATA_REGION_GROUP_NUM=1")) {
+ "CREATE DATABASE " + metricConfig.getInternalDatabase())) {
if (!result2.hasNext()) {
LOGGER.error("IoTDBSessionReporter checkOrCreateDatabase failed.");
}