This is an automated email from the ASF dual-hosted git repository.
yongzao 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 ca4539ff4e1 [To dev/1.3] Region operation "extend" and "remove"
support multi regions in one SQL (#16234)
ca4539ff4e1 is described below
commit ca4539ff4e19e57c0fa82368d066e6f3c53f4020
Author: Li Yu Heng <[email protected]>
AuthorDate: Fri Aug 22 18:00:59 2025 +0800
[To dev/1.3] Region operation "extend" and "remove" support multi regions
in one SQL (#16234)
---
.../IoTDBRegionGroupExpandAndShrinkForIoTV1IT.java | 329 ++++++++++++++++++++-
.../org/apache/iotdb/db/qp/sql/IoTDBSqlParser.g4 | 4 +-
.../iotdb/confignode/manager/ConfigManager.java | 4 +-
.../iotdb/confignode/manager/ProcedureManager.java | 63 +++-
.../config/executor/ClusterConfigTaskExecutor.java | 4 +-
.../db/queryengine/plan/parser/ASTVisitor.java | 11 +-
.../metadata/region/ExtendRegionStatement.java | 10 +-
.../metadata/region/RemoveRegionStatement.java | 10 +-
.../src/main/thrift/confignode.thrift | 4 +-
9 files changed, 407 insertions(+), 32 deletions(-)
diff --git
a/integration-test/src/test/java/org/apache/iotdb/confignode/it/regionmigration/pass/commit/IoTDBRegionGroupExpandAndShrinkForIoTV1IT.java
b/integration-test/src/test/java/org/apache/iotdb/confignode/it/regionmigration/pass/commit/IoTDBRegionGroupExpandAndShrinkForIoTV1IT.java
index 3c4aa16a401..afc1163fef0 100644
---
a/integration-test/src/test/java/org/apache/iotdb/confignode/it/regionmigration/pass/commit/IoTDBRegionGroupExpandAndShrinkForIoTV1IT.java
+++
b/integration-test/src/test/java/org/apache/iotdb/confignode/it/regionmigration/pass/commit/IoTDBRegionGroupExpandAndShrinkForIoTV1IT.java
@@ -37,11 +37,14 @@ import org.slf4j.LoggerFactory;
import java.sql.Connection;
import java.sql.Statement;
+import java.util.ArrayList;
+import java.util.List;
import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.TimeUnit;
import java.util.function.Predicate;
+import java.util.stream.Collectors;
import static org.apache.iotdb.util.MagicUtils.makeItCloseQuietly;
@@ -51,6 +54,8 @@ public class IoTDBRegionGroupExpandAndShrinkForIoTV1IT
extends IoTDBRegionOperationReliabilityITFramework {
private static final String EXPAND_FORMAT = "extend region %d to %d";
private static final String SHRINK_FORMAT = "remove region %d from %d";
+ private static final String MULTI_EXPAND_FORMAT = "extend region %s to %d";
+ private static final String MULTI_SHRINK_FORMAT = "remove region %s from %d";
private static Logger LOGGER =
LoggerFactory.getLogger(IoTDBRegionGroupExpandAndShrinkForIoTV1IT.class);
@@ -65,7 +70,7 @@ public class IoTDBRegionGroupExpandAndShrinkForIoTV1IT
* <p>4. Check
*/
@Test
- public void normal1C5DTest() throws Exception {
+ public void singleRegionTest() throws Exception {
EnvFactory.getEnv()
.getConfig()
.getCommonConfig()
@@ -170,4 +175,326 @@ public class IoTDBRegionGroupExpandAndShrinkForIoTV1IT
LOGGER.info("Region {} has shrunk from DataNode {}", selectedRegion,
targetDataNode);
}
+
+ /**
+ * Test multi-region expand and shrink operations with normal flow: 1.
Multi-expand: expand
+ * multiple regions to target DataNode 2. Multi-shrink: shrink multiple
regions from target
+ * DataNode
+ */
+ @Test
+ public void multiRegionNormalTest() throws Exception {
+ EnvFactory.getEnv()
+ .getConfig()
+ .getCommonConfig()
+ .setDataRegionConsensusProtocolClass(ConsensusFactory.IOT_CONSENSUS)
+
.setSchemaRegionConsensusProtocolClass(ConsensusFactory.RATIS_CONSENSUS)
+ .setDataReplicationFactor(1)
+ .setSchemaReplicationFactor(1);
+
+ EnvFactory.getEnv().initClusterEnvironment(1, 5);
+
+ try (final Connection connection =
makeItCloseQuietly(EnvFactory.getEnv().getConnection());
+ final Statement statement =
makeItCloseQuietly(connection.createStatement());
+ SyncConfigNodeIServiceClient client =
+ (SyncConfigNodeIServiceClient)
EnvFactory.getEnv().getLeaderConfigNodeConnection()) {
+ // prepare data
+ statement.execute(INSERTION1);
+ statement.execute(FLUSH_COMMAND);
+
+ // collect necessary information
+ Map<Integer, Set<Integer>> regionMap = getAllRegionMap(statement);
+ Set<Integer> allDataNodeId = getAllDataNodes(statement);
+
+ // at least one data region, one schema region
+ // maybe more, one system data region, one system schema region
+ Assert.assertTrue(regionMap.size() >= 2);
+
+ // select multiple regions for testing
+ List<Integer> selectedRegions = new ArrayList<>(regionMap.keySet());
+ selectedRegions = selectedRegions.subList(0, Math.min(3,
selectedRegions.size()));
+
+ // find target DataNode that doesn't contain any of the selected regions
+ int targetDataNode =
+ findDataNodeNotContainsAnyRegion(allDataNodeId, regionMap,
selectedRegions);
+
+ LOGGER.info("Selected regions for multi-region test: {}",
selectedRegions);
+ LOGGER.info("Target DataNode: {}", targetDataNode);
+
+ // multi-expand: expand all selected regions to target DataNode
+ multiRegionGroupExpand(statement, client, selectedRegions,
targetDataNode);
+
+ // verify expand result
+ regionMap = getAllRegionMap(statement);
+ for (int regionId : selectedRegions) {
+ Assert.assertTrue(
+ "Region " + regionId + " should contain target DataNode " +
targetDataNode,
+ regionMap.get(regionId).contains(targetDataNode));
+ }
+ LOGGER.info("Multi-region expand test passed");
+
+ // multi-shrink: shrink all selected regions from target DataNode
+ multiRegionGroupShrink(statement, client, selectedRegions,
targetDataNode);
+
+ // verify shrink result
+ regionMap = getAllRegionMap(statement);
+ for (int regionId : selectedRegions) {
+ Assert.assertFalse(
+ "Region " + regionId + " should not contain target DataNode " +
targetDataNode,
+ regionMap.get(regionId).contains(targetDataNode));
+ }
+ LOGGER.info("Multi-region shrink test passed");
+ }
+ }
+
+ /** Test multi-region expand with partial regions already in target DataNode
*/
+ @Test
+ public void multiRegionExpandPartialExistTest() throws Exception {
+ EnvFactory.getEnv()
+ .getConfig()
+ .getCommonConfig()
+ .setDataRegionConsensusProtocolClass(ConsensusFactory.IOT_CONSENSUS)
+
.setSchemaRegionConsensusProtocolClass(ConsensusFactory.RATIS_CONSENSUS)
+ .setDataReplicationFactor(1)
+ .setSchemaReplicationFactor(1);
+
+ EnvFactory.getEnv().initClusterEnvironment(1, 5);
+
+ try (final Connection connection =
makeItCloseQuietly(EnvFactory.getEnv().getConnection());
+ final Statement statement =
makeItCloseQuietly(connection.createStatement());
+ SyncConfigNodeIServiceClient client =
+ (SyncConfigNodeIServiceClient)
EnvFactory.getEnv().getLeaderConfigNodeConnection()) {
+ // prepare data
+ statement.execute(INSERTION1);
+ statement.execute(FLUSH_COMMAND);
+
+ Map<Integer, Set<Integer>> regionMap = getAllRegionMap(statement);
+ Set<Integer> allDataNodeId = getAllDataNodes(statement);
+
+ List<Integer> allRegions = new ArrayList<>(regionMap.keySet());
+ List<Integer> selectedRegions = allRegions.subList(0, Math.min(3,
allRegions.size()));
+
+ int targetDataNode =
+ findDataNodeNotContainsAnyRegion(allDataNodeId, regionMap,
selectedRegions);
+
+ // first expand some regions individually
+ List<Integer> preExpandRegions =
+ selectedRegions.subList(0, Math.min(2, selectedRegions.size()));
+ for (int regionId : preExpandRegions) {
+ regionGroupExpand(statement, client, regionId, targetDataNode);
+ }
+
+ // now try to expand all regions (including already expanded ones)
+ LOGGER.info(
+ "Testing multi-expand with regions {} to DataNode {}, where {}
already exist",
+ selectedRegions,
+ targetDataNode,
+ preExpandRegions);
+
+ multiRegionGroupExpand(statement, client, selectedRegions,
targetDataNode);
+
+ // verify all regions are in target DataNode
+ regionMap = getAllRegionMap(statement);
+ for (int regionId : selectedRegions) {
+ Assert.assertTrue(
+ "Region " + regionId + " should contain target DataNode " +
targetDataNode,
+ regionMap.get(regionId).contains(targetDataNode));
+ }
+ LOGGER.info("Multi-region expand partial exist test passed");
+ }
+ }
+
+ /** Test multi-region shrink with partial regions not in target DataNode */
+ @Test
+ public void multiRegionShrinkPartialNotExistTest() throws Exception {
+ EnvFactory.getEnv()
+ .getConfig()
+ .getCommonConfig()
+ .setDataRegionConsensusProtocolClass(ConsensusFactory.IOT_CONSENSUS)
+
.setSchemaRegionConsensusProtocolClass(ConsensusFactory.RATIS_CONSENSUS)
+ .setDataReplicationFactor(1)
+ .setSchemaReplicationFactor(1);
+
+ EnvFactory.getEnv().initClusterEnvironment(1, 5);
+
+ try (final Connection connection =
makeItCloseQuietly(EnvFactory.getEnv().getConnection());
+ final Statement statement =
makeItCloseQuietly(connection.createStatement());
+ SyncConfigNodeIServiceClient client =
+ (SyncConfigNodeIServiceClient)
EnvFactory.getEnv().getLeaderConfigNodeConnection()) {
+ // prepare data
+ statement.execute(INSERTION1);
+ statement.execute(FLUSH_COMMAND);
+
+ Map<Integer, Set<Integer>> regionMap = getAllRegionMap(statement);
+ Set<Integer> allDataNodeId = getAllDataNodes(statement);
+
+ List<Integer> allRegions = new ArrayList<>(regionMap.keySet());
+ List<Integer> selectedRegions = allRegions.subList(0, Math.min(3,
allRegions.size()));
+
+ int targetDataNode =
+ findDataNodeNotContainsAnyRegion(allDataNodeId, regionMap,
selectedRegions);
+
+ // first expand all regions to target DataNode
+ multiRegionGroupExpand(statement, client, selectedRegions,
targetDataNode);
+
+ // then shrink some regions individually
+ List<Integer> preShrinkRegions =
+ selectedRegions.subList(0, Math.min(2, selectedRegions.size()));
+ for (int regionId : preShrinkRegions) {
+ regionGroupShrink(statement, client, regionId, targetDataNode);
+ }
+
+ // now try to shrink all regions (including already shrunk ones)
+ LOGGER.info(
+ "Testing multi-shrink with regions {} from DataNode {}, where {}
already removed",
+ selectedRegions,
+ targetDataNode,
+ preShrinkRegions);
+
+ multiRegionGroupShrink(statement, client, selectedRegions,
targetDataNode);
+
+ // verify all regions are not in target DataNode
+ regionMap = getAllRegionMap(statement);
+ for (int regionId : selectedRegions) {
+ Assert.assertFalse(
+ "Region " + regionId + " should not contain target DataNode " +
targetDataNode,
+ regionMap.get(regionId).contains(targetDataNode));
+ }
+ LOGGER.info("Multi-region shrink partial not exist test passed");
+ }
+ }
+
+ private void multiRegionGroupExpand(
+ Statement statement,
+ SyncConfigNodeIServiceClient client,
+ List<Integer> regionIds,
+ int targetDataNode)
+ throws Exception {
+ String command = buildMultiRegionCommand(MULTI_EXPAND_FORMAT, regionIds,
targetDataNode);
+
+ Predicate<TShowRegionResp> expandPredicate =
+ tShowRegionResp -> {
+ Map<Integer, Set<Integer>> newRegionMap =
+ getRunningRegionMap(tShowRegionResp.getRegionInfoList());
+ return regionIds.stream()
+ .allMatch(
+ regionId -> {
+ Set<Integer> dataNodes = newRegionMap.get(regionId);
+ return dataNodes != null &&
dataNodes.contains(targetDataNode);
+ });
+ };
+
+ executeMultiRegionOperation(
+ statement,
+ client,
+ command,
+ regionIds,
+ expandPredicate,
+ Optional.of(targetDataNode),
+ Optional.empty(),
+ "expand");
+ }
+
+ private void multiRegionGroupShrink(
+ Statement statement,
+ SyncConfigNodeIServiceClient client,
+ List<Integer> regionIds,
+ int targetDataNode)
+ throws Exception {
+ String command = buildMultiRegionCommand(MULTI_SHRINK_FORMAT, regionIds,
targetDataNode);
+
+ Predicate<TShowRegionResp> shrinkPredicate =
+ tShowRegionResp -> {
+ Map<Integer, Set<Integer>> newRegionMap =
+ getRegionMap(tShowRegionResp.getRegionInfoList());
+ return regionIds.stream()
+ .allMatch(
+ regionId -> {
+ Set<Integer> dataNodes = newRegionMap.get(regionId);
+ return dataNodes == null ||
!dataNodes.contains(targetDataNode);
+ });
+ };
+
+ executeMultiRegionOperation(
+ statement,
+ client,
+ command,
+ regionIds,
+ shrinkPredicate,
+ Optional.empty(),
+ Optional.of(targetDataNode),
+ "shrink");
+ }
+
+ private String buildMultiRegionCommand(
+ String format, List<Integer> regionIds, int targetDataNode) {
+ String regionIdStr =
regionIds.stream().map(String::valueOf).collect(Collectors.joining(","));
+ return String.format(format, regionIdStr, targetDataNode);
+ }
+
+ private void executeMultiRegionOperation(
+ Statement statement,
+ SyncConfigNodeIServiceClient client,
+ String command,
+ List<Integer> regionIds,
+ Predicate<TShowRegionResp> predicate,
+ Optional<Integer> expectedDataNode,
+ Optional<Integer> notExpectedDataNode,
+ String operationType) {
+
+ LOGGER.info("Executing multi-region {} command: {}", operationType,
command);
+
+ Awaitility.await()
+ .atMost(30, TimeUnit.SECONDS)
+ .pollInterval(2, TimeUnit.SECONDS)
+ .until(
+ () -> {
+ try {
+ statement.execute(command);
+ return true;
+ } catch (Exception e) {
+ String errorMessage = e.getMessage();
+ // If error message contains both "successfully submitted" and
"failed to submit",
+ // consider it as partial success and continue
+ if (errorMessage != null
+ && errorMessage.contains("successfully submitted")
+ && errorMessage.contains("failed to submit")) {
+ LOGGER.warn(
+ "Multi-region {} partially succeeded: {}",
operationType, errorMessage);
+ return true;
+ }
+ LOGGER.warn(
+ "Multi-region {} command execution failed, retrying: {}",
+ operationType,
+ errorMessage);
+ return false;
+ }
+ });
+
+ // Use the first region for awaitUntilSuccess (framework limitation)
+ awaitUntilSuccess(client, predicate, expectedDataNode,
notExpectedDataNode);
+
+ String targetDescription =
+ expectedDataNode.isPresent()
+ ? "to DataNode " + expectedDataNode.get()
+ : "from DataNode " + notExpectedDataNode.get();
+ LOGGER.info(
+ "Regions {} have {} {}",
+ regionIds,
+ operationType.equals("expand") ? "expanded" : "shrunk",
+ targetDescription);
+ }
+
+ private int findDataNodeNotContainsAnyRegion(
+ Set<Integer> allDataNodeId, Map<Integer, Set<Integer>> regionMap,
List<Integer> regionIds) {
+ return allDataNodeId.stream()
+ .filter(
+ dataNodeId ->
+ regionIds.stream()
+ .noneMatch(regionId ->
regionMap.get(regionId).contains(dataNodeId)))
+ .findFirst()
+ .orElseThrow(
+ () ->
+ new RuntimeException(
+ "Cannot find DataNode that doesn't contain any of the
regions"));
+ }
}
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 02b5062cb02..f28a2e0a063 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
@@ -541,11 +541,11 @@ reconstructRegion
;
extendRegion
- : EXTEND REGION regionId=INTEGER_LITERAL TO
targetDataNodeId=INTEGER_LITERAL
+ : EXTEND REGION regionIds+=INTEGER_LITERAL (COMMA
regionIds+=INTEGER_LITERAL)* TO targetDataNodeId=INTEGER_LITERAL
;
removeRegion
- : REMOVE REGION regionId=INTEGER_LITERAL FROM
targetDataNodeId=INTEGER_LITERAL
+ : REMOVE REGION regionIds+=INTEGER_LITERAL (COMMA
regionIds+=INTEGER_LITERAL)* FROM targetDataNodeId=INTEGER_LITERAL
;
verifyConnection
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConfigManager.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConfigManager.java
index 13e734d0bc7..b9613de8ed7 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConfigManager.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConfigManager.java
@@ -2348,7 +2348,7 @@ public class ConfigManager implements IManager {
public TSStatus extendRegion(TExtendRegionReq req) {
TSStatus status = confirmLeader();
return status.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()
- ? procedureManager.extendRegion(req)
+ ? procedureManager.extendRegions(req)
: status;
}
@@ -2356,7 +2356,7 @@ public class ConfigManager implements IManager {
public TSStatus removeRegion(TRemoveRegionReq req) {
TSStatus status = confirmLeader();
return status.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()
- ? procedureManager.removeRegion(req)
+ ? procedureManager.removeRegions(req)
: status;
}
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java
index 460ed5276ca..5f6450769f5 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java
@@ -149,6 +149,7 @@ import java.util.Map;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.ScheduledExecutorService;
+import java.util.function.BiFunction;
import java.util.stream.Collectors;
import java.util.stream.Stream;
@@ -766,7 +767,7 @@ public class ProcedureManager {
failMessage =
String.format(
"Target DataNode %s already contains region %s",
- targetDataNode.getDataNodeId(), req.getRegionId());
+ targetDataNode.getDataNodeId(), regionId);
}
if (failMessage != null) {
@@ -1048,14 +1049,60 @@ public class ProcedureManager {
return RpcUtils.SUCCESS_STATUS;
}
- public TSStatus extendRegion(TExtendRegionReq req) {
+ public TSStatus extendRegions(TExtendRegionReq req) {
+ return processExtendOrRemoveRegions(
+ req.getRegionId(), req, this::extendOneRegion,
TSStatusCode.EXTEND_REGION_ERROR);
+ }
+
+ public TSStatus removeRegions(TRemoveRegionReq req) {
+ return processExtendOrRemoveRegions(
+ req.getRegionId(), req, this::removeOneRegion,
TSStatusCode.REMOVE_REGION_PEER_ERROR);
+ }
+
+ private <R> TSStatus processExtendOrRemoveRegions(
+ Iterable<Integer> regionIds,
+ R req,
+ BiFunction<Integer, R, TSStatus> regionAction,
+ TSStatusCode errorCode) {
+ TSStatus resp = new TSStatus();
+ StringBuilder messageBuilder = new StringBuilder();
+
+ int total = 0, success = 0;
+ for (int regionId : regionIds) {
+ total++;
+ TSStatus subStatus = regionAction.apply(regionId, req);
+ if (subStatus.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ messageBuilder.append("region ").append(regionId).append(":
Successfully submitted\n");
+ success++;
+ } else {
+ messageBuilder
+ .append("region ")
+ .append(regionId)
+ .append(": ")
+ .append(subStatus.getMessage())
+ .append('\n');
+ }
+ resp.addToSubStatus(subStatus);
+ }
+
+ messageBuilder.insert(
+ 0,
+ String.format(
+ "Total regions: %d, successfully submitted: %d, failed to submit:
%d\n",
+ total, success, total - success));
+
+ resp.setCode(
+ total == success ? TSStatusCode.SUCCESS_STATUS.getStatusCode() :
errorCode.getStatusCode());
+ resp.setMessage(messageBuilder.toString());
+ return resp;
+ }
+
+ private TSStatus extendOneRegion(int theRegionId, TExtendRegionReq req) {
try (AutoCloseableLock ignoredLock =
AutoCloseableLock.acquire(env.getSubmitRegionMigrateLock())) {
TConsensusGroupId regionId;
Optional<TConsensusGroupId> optional =
- configManager
- .getPartitionManager()
- .generateTConsensusGroupIdByRegionId(req.getRegionId());
+
configManager.getPartitionManager().generateTConsensusGroupIdByRegionId(theRegionId);
if (optional.isPresent()) {
regionId = optional.get();
} else {
@@ -1093,14 +1140,12 @@ public class ProcedureManager {
}
}
- public TSStatus removeRegion(TRemoveRegionReq req) {
+ private TSStatus removeOneRegion(int theRegionId, TRemoveRegionReq req) {
try (AutoCloseableLock ignoredLock =
AutoCloseableLock.acquire(env.getSubmitRegionMigrateLock())) {
TConsensusGroupId regionId;
Optional<TConsensusGroupId> optional =
- configManager
- .getPartitionManager()
- .generateTConsensusGroupIdByRegionId(req.getRegionId());
+
configManager.getPartitionManager().generateTConsensusGroupIdByRegionId(theRegionId);
if (optional.isPresent()) {
regionId = optional.get();
} else {
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/executor/ClusterConfigTaskExecutor.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/executor/ClusterConfigTaskExecutor.java
index 3a3261276f2..b128f3ba361 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/executor/ClusterConfigTaskExecutor.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/execution/config/executor/ClusterConfigTaskExecutor.java
@@ -2873,7 +2873,7 @@ public class ClusterConfigTaskExecutor implements
IConfigTaskExecutor {
CONFIG_NODE_CLIENT_MANAGER.borrowClient(ConfigNodeInfo.CONFIG_REGION_ID)) {
final TExtendRegionReq req =
new TExtendRegionReq(
- extendRegionStatement.getRegionId(),
extendRegionStatement.getDataNodeId());
+ extendRegionStatement.getRegionIds(),
extendRegionStatement.getDataNodeId());
final TSStatus status = configNodeClient.extendRegion(req);
if (status.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
future.setException(new IoTDBException(status.message, status.code));
@@ -2895,7 +2895,7 @@ public class ClusterConfigTaskExecutor implements
IConfigTaskExecutor {
CONFIG_NODE_CLIENT_MANAGER.borrowClient(ConfigNodeInfo.CONFIG_REGION_ID)) {
final TRemoveRegionReq req =
new TRemoveRegionReq(
- removeRegionStatement.getRegionId(),
removeRegionStatement.getDataNodeId());
+ removeRegionStatement.getRegionIds(),
removeRegionStatement.getDataNodeId());
final TSStatus status = configNodeClient.removeRegion(req);
if (status.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
future.setException(new IoTDBException(status.message, status.code));
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 aa59bc78659..83304b0c6c2 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
@@ -262,6 +262,7 @@ import java.util.function.Consumer;
import java.util.regex.Pattern;
import java.util.stream.Collectors;
+import static java.util.stream.Collectors.toList;
import static org.apache.iotdb.commons.schema.SchemaConstant.ALL_RESULT_NODES;
import static
org.apache.iotdb.db.queryengine.plan.optimization.LimitOffsetPushDown.canPushDownLimitOffsetToGroupByTime;
import static
org.apache.iotdb.db.queryengine.plan.optimization.LimitOffsetPushDown.pushDownLimitOffsetToTimeParameter;
@@ -4188,14 +4189,16 @@ public class ASTVisitor extends
IoTDBSqlParserBaseVisitor<Statement> {
@Override
public Statement visitExtendRegion(IoTDBSqlParser.ExtendRegionContext ctx) {
- return new ExtendRegionStatement(
- Integer.parseInt(ctx.regionId.getText()),
Integer.parseInt(ctx.targetDataNodeId.getText()));
+ List<Integer> regionIds =
+ ctx.regionIds.stream().map(token ->
Integer.parseInt(token.getText())).collect(toList());
+ return new ExtendRegionStatement(regionIds,
Integer.parseInt(ctx.targetDataNodeId.getText()));
}
@Override
public Statement visitRemoveRegion(IoTDBSqlParser.RemoveRegionContext ctx) {
- return new RemoveRegionStatement(
- Integer.parseInt(ctx.regionId.getText()),
Integer.parseInt(ctx.targetDataNodeId.getText()));
+ List<Integer> regionIds =
+ ctx.regionIds.stream().map(token ->
Integer.parseInt(token.getText())).collect(toList());
+ return new RemoveRegionStatement(regionIds,
Integer.parseInt(ctx.targetDataNodeId.getText()));
}
@Override
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/metadata/region/ExtendRegionStatement.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/metadata/region/ExtendRegionStatement.java
index 0048a789f95..591c62c4b6b 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/metadata/region/ExtendRegionStatement.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/metadata/region/ExtendRegionStatement.java
@@ -32,17 +32,17 @@ import java.util.List;
public class ExtendRegionStatement extends Statement implements
IConfigStatement {
- private final int regionId;
+ private final List<Integer> regionIds;
private final int dataNodeId;
- public ExtendRegionStatement(int regionId, int dataNodeId) {
+ public ExtendRegionStatement(List<Integer> regionIds, int dataNodeId) {
super();
- this.regionId = regionId;
+ this.regionIds = regionIds;
this.dataNodeId = dataNodeId;
}
- public int getRegionId() {
- return regionId;
+ public List<Integer> getRegionIds() {
+ return regionIds;
}
public int getDataNodeId() {
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/metadata/region/RemoveRegionStatement.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/metadata/region/RemoveRegionStatement.java
index aa185ad627e..f4d3b94c682 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/metadata/region/RemoveRegionStatement.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/statement/metadata/region/RemoveRegionStatement.java
@@ -32,17 +32,17 @@ import java.util.List;
public class RemoveRegionStatement extends Statement implements
IConfigStatement {
- private final int regionId;
+ private final List<Integer> regionIds;
private final int dataNodeId;
- public RemoveRegionStatement(int regionId, int dataNodeId) {
+ public RemoveRegionStatement(List<Integer> regionIds, int dataNodeId) {
super();
- this.regionId = regionId;
+ this.regionIds = regionIds;
this.dataNodeId = dataNodeId;
}
- public int getRegionId() {
- return regionId;
+ public List<Integer> getRegionIds() {
+ return regionIds;
}
public int getDataNodeId() {
diff --git a/iotdb-protocol/thrift-confignode/src/main/thrift/confignode.thrift
b/iotdb-protocol/thrift-confignode/src/main/thrift/confignode.thrift
index 19e8e2a7d1f..87a154e235e 100644
--- a/iotdb-protocol/thrift-confignode/src/main/thrift/confignode.thrift
+++ b/iotdb-protocol/thrift-confignode/src/main/thrift/confignode.thrift
@@ -315,12 +315,12 @@ struct TReconstructRegionReq {
}
struct TExtendRegionReq {
- 1: required i32 regionId
+ 1: required list<i32> regionId
2: required i32 dataNodeId
}
struct TRemoveRegionReq {
- 1: required i32 regionId
+ 1: required list<i32> regionId
2: required i32 dataNodeId
}