This is an automated email from the ASF dual-hosted git repository. liyuheng pushed a commit to branch cp/16196 in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit feac280f3b04a63494393d8ba6dddcd253082708 Author: Li Yu Heng <[email protected]> AuthorDate: Wed Aug 20 10:58:02 2025 +0800 Region operation "extend" and "remove" support multi regions in one SQL (#16196) (cherry picked from commit 0aed7f81eddbad0a7c0d029e46c2bb91d68dea13) --- .../IoTDBRegionGroupExpandAndShrinkForIoTV1IT.java | 329 ++++++++++++++++++++- .../org/apache/iotdb/db/qp/sql/IoTDBSqlParser.g4 | 4 +- .../iotdb/confignode/manager/ConfigManager.java | 4 +- .../iotdb/confignode/manager/ProcedureManager.java | 64 +++- .../config/executor/ClusterConfigTaskExecutor.java | 6 +- .../db/queryengine/plan/parser/ASTVisitor.java | 12 +- .../metadata/region/ExtendRegionStatement.java | 10 +- .../metadata/region/RemoveRegionStatement.java | 10 +- .../src/main/thrift/confignode.thrift | 4 +- 9 files changed, 411 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..aab710912eb 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); + + // expect one data region, one schema region + // plus one system data region, one system schema region + Assert.assertEquals(4, regionMap.size()); + + // 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, regionIds.get(0), 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..197012447e3 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,8 @@ import java.util.Map; import java.util.Optional; import java.util.Set; import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.locks.ReentrantLock; +import java.util.function.BiFunction; import java.util.stream.Collectors; import java.util.stream.Stream; @@ -766,7 +768,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 +1050,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 +1141,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..d399cf2ea5e 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,8 @@ public class ClusterConfigTaskExecutor implements IConfigTaskExecutor { CONFIG_NODE_CLIENT_MANAGER.borrowClient(ConfigNodeInfo.CONFIG_REGION_ID)) { final TExtendRegionReq req = new TExtendRegionReq( - extendRegionStatement.getRegionId(), extendRegionStatement.getDataNodeId()); + extendRegionTask.getStatement().getRegionIds(), + extendRegionTask.getStatement().getDataNodeId() final TSStatus status = configNodeClient.extendRegion(req); if (status.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) { future.setException(new IoTDBException(status.message, status.code)); @@ -2895,7 +2896,8 @@ public class ClusterConfigTaskExecutor implements IConfigTaskExecutor { CONFIG_NODE_CLIENT_MANAGER.borrowClient(ConfigNodeInfo.CONFIG_REGION_ID)) { final TRemoveRegionReq req = new TRemoveRegionReq( - removeRegionStatement.getRegionId(), removeRegionStatement.getDataNodeId()); + removeRegionTask.getStatement().getRegionIds(), + removeRegionTask.getStatement().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..3ccf328f86a 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,8 @@ import java.util.function.Consumer; import java.util.regex.Pattern; import java.util.stream.Collectors; +import static com.google.common.collect.ImmutableList.toImmutableList; +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 +4190,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 }
