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 3fbbe5f9a42 [To dev/1.3] Harden RegionMaintainer retries and region
operation atomicity (#18480)
3fbbe5f9a42 is described below
commit 3fbbe5f9a420c23beb09a474e40c47d39e65a3fa
Author: Yongzao <[email protected]>
AuthorDate: Wed Aug 19 11:55:03 2026 +0800
[To dev/1.3] Harden RegionMaintainer retries and region operation atomicity
(#18480)
---
.../java/org/apache/iotdb/rpc/TSStatusCode.java | 2 +
.../handlers/rpc/DataNodeTSStatusRPCHandler.java | 14 +-
.../manager/partition/PartitionManager.java | 366 ++++++++++-----------
.../impl/schema/DeleteDatabaseProcedure.java | 111 ++++---
.../rpc/DataNodeTSStatusRPCHandlerTest.java | 63 ++++
.../PartitionManagerRegionMaintainTest.java | 151 +++++++++
.../confignode/persistence/PartitionInfoTest.java | 19 ++
.../impl/schema/DeleteDatabaseProcedureTest.java | 118 +++++++
.../impl/DataNodeInternalRPCServiceImpl.java | 18 +-
.../thrift/impl/DataNodeRegionManager.java | 18 +-
.../apache/iotdb/db/schemaengine/SchemaEngine.java | 5 +-
.../iotdb/db/service/RegionMigrateService.java | 29 +-
.../iotdb/db/storageengine/StorageEngine.java | 104 +++---
.../DataNodeInternalRPCServiceImplStatusTest.java | 65 ++++
.../DataNodeInternalRPCServiceImplTest.java | 26 ++
.../db/service/RegionMigrateServiceStatusTest.java | 44 +++
16 files changed, 855 insertions(+), 298 deletions(-)
diff --git
a/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TSStatusCode.java
b/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TSStatusCode.java
index af4f4727c7a..49e75c6da90 100644
---
a/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TSStatusCode.java
+++
b/iotdb-client/service-rpc/src/main/java/org/apache/iotdb/rpc/TSStatusCode.java
@@ -166,6 +166,8 @@ public enum TSStatusCode {
RECONSTRUCT_REGION_ERROR(908),
EXTEND_REGION_ERROR(909),
REMOVE_REGION_PEER_ERROR(910),
+ REGION_ALREADY_EXISTS(911),
+ REGION_NOT_EXIST(912),
// Cluster Manager
ADD_CONFIGNODE_ERROR(1000),
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/DataNodeTSStatusRPCHandler.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/DataNodeTSStatusRPCHandler.java
index a44e3781c25..848dd732212 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/DataNodeTSStatusRPCHandler.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/DataNodeTSStatusRPCHandler.java
@@ -52,7 +52,7 @@ public class DataNodeTSStatusRPCHandler extends
DataNodeAsyncRequestRPCHandler<T
// Put response
responseMap.put(requestId, response);
- if (response.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ if (isRequestCompleted(requestType, response)) {
// Remove only if success
nodeLocationMap.remove(requestId);
LOGGER.info("Successfully {} on DataNode: {}", requestType,
formattedTargetLocation);
@@ -68,6 +68,18 @@ public class DataNodeTSStatusRPCHandler extends
DataNodeAsyncRequestRPCHandler<T
countDownLatch.countDown();
}
+ static boolean isRequestCompleted(CnToDnAsyncRequestType requestType,
TSStatus response) {
+ if (response.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ return true;
+ }
+ if (requestType == CnToDnAsyncRequestType.DELETE_REGION) {
+ return response.getCode() ==
TSStatusCode.REGION_NOT_EXIST.getStatusCode();
+ }
+ return (requestType == CnToDnAsyncRequestType.CREATE_DATA_REGION
+ || requestType == CnToDnAsyncRequestType.CREATE_SCHEMA_REGION)
+ && response.getCode() ==
TSStatusCode.REGION_ALREADY_EXISTS.getStatusCode();
+ }
+
@Override
public void onError(Exception e) {
String errorMsg =
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 7db7ac50d62..e3a50cda8d8 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
@@ -98,11 +98,13 @@ import org.apache.tsfile.utils.Pair;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import java.util.ArrayDeque;
import java.util.ArrayList;
import java.util.Collections;
+import java.util.EnumMap;
import java.util.HashMap;
import java.util.HashSet;
-import java.util.LinkedList;
+import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
@@ -1222,204 +1224,200 @@ public class PartitionManager {
return;
}
- // Group tasks by region id
- Map<TConsensusGroupId, Queue<RegionMaintainTask>>
regionMaintainTaskMap =
- new HashMap<>();
- for (RegionMaintainTask regionMaintainTask :
regionMaintainTaskList) {
- regionMaintainTaskMap
- .computeIfAbsent(regionMaintainTask.getRegionId(), k ->
new LinkedList<>())
- .add(regionMaintainTask);
- }
-
- while (!regionMaintainTaskMap.isEmpty()) {
- // Select same type task from each region group
- List<RegionMaintainTask> selectedRegionMaintainTask = new
ArrayList<>();
- RegionMaintainType currentType = null;
- for (Map.Entry<TConsensusGroupId, Queue<RegionMaintainTask>>
entry :
- regionMaintainTaskMap.entrySet()) {
- RegionMaintainTask regionMaintainTask =
entry.getValue().peek();
- if (regionMaintainTask == null) {
- continue;
- }
-
- if (currentType == null) {
- currentType = regionMaintainTask.getType();
- selectedRegionMaintainTask.add(entry.getValue().peek());
- } else {
- if (!currentType.equals(regionMaintainTask.getType())) {
- continue;
- }
-
- if (currentType.equals(RegionMaintainType.DELETE)
- || entry
- .getKey()
- .getType()
-
.equals(selectedRegionMaintainTask.get(0).getRegionId().getType())) {
- // Delete or same create task
-
selectedRegionMaintainTask.add(entry.getValue().peek());
- }
- }
- }
-
- if (selectedRegionMaintainTask.isEmpty()) {
- break;
+ Map<TConsensusGroupId, Queue<RegionMaintainTask>>
tasksByRegion =
+ groupRegionMaintainTasks(regionMaintainTaskList);
+ Set<TConsensusGroupId> deferredRegions = new HashSet<>();
+ while (!tasksByRegion.isEmpty()) {
+ Map<RegionMaintainType, List<RegionMaintainTask>>
headsByType =
+ getRegionMaintainTaskHeads(tasksByRegion,
deferredRegions);
+ Set<TConsensusGroupId> completedRegions = new HashSet<>();
+ for (Map.Entry<RegionMaintainType, List<RegionMaintainTask>>
entry :
+ headsByType.entrySet()) {
+ completedRegions.addAll(
+ submitRegionMaintainTasks(entry.getKey(),
entry.getValue()));
}
- Set<TConsensusGroupId> successfulTask = new HashSet<>();
- switch (currentType) {
- case CREATE:
- // create region
- switch
(selectedRegionMaintainTask.get(0).getRegionId().getType()) {
- case SchemaRegion:
- // create SchemaRegion
- DataNodeAsyncRequestContext<TCreateSchemaRegionReq,
TSStatus>
- createSchemaRegionHandler =
- new DataNodeAsyncRequestContext<>(
-
CnToDnAsyncRequestType.CREATE_SCHEMA_REGION);
- for (RegionMaintainTask regionMaintainTask :
selectedRegionMaintainTask) {
- RegionCreateTask schemaRegionCreateTask =
- (RegionCreateTask) regionMaintainTask;
- LOGGER.info(
- "Start to create Region: {} on DataNode: {}",
-
schemaRegionCreateTask.getRegionReplicaSet().getRegionId(),
- schemaRegionCreateTask.getTargetDataNode());
- createSchemaRegionHandler.putRequest(
- schemaRegionCreateTask.getRegionId().getId(),
- new TCreateSchemaRegionReq(
-
schemaRegionCreateTask.getRegionReplicaSet(),
- schemaRegionCreateTask.getStorageGroup()));
- createSchemaRegionHandler.putNodeLocation(
- schemaRegionCreateTask.getRegionId().getId(),
- schemaRegionCreateTask.getTargetDataNode());
- }
-
-
CnToDnInternalServiceAsyncRequestManager.getInstance()
-
.sendAsyncRequestWithRetry(createSchemaRegionHandler);
-
- for (Map.Entry<Integer, TSStatus> entry :
-
createSchemaRegionHandler.getResponseMap().entrySet()) {
- if (entry.getValue().getCode()
- ==
TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
- successfulTask.add(
- new TConsensusGroupId(
- TConsensusGroupType.SchemaRegion,
entry.getKey()));
- }
- }
- break;
- case DataRegion:
- // Create DataRegion
- DataNodeAsyncRequestContext<TCreateDataRegionReq,
TSStatus>
- createDataRegionHandler =
- new DataNodeAsyncRequestContext<>(
-
CnToDnAsyncRequestType.CREATE_DATA_REGION);
- for (RegionMaintainTask regionMaintainTask :
selectedRegionMaintainTask) {
- RegionCreateTask dataRegionCreateTask =
- (RegionCreateTask) regionMaintainTask;
- LOGGER.info(
- "Start to create Region: {} on DataNode: {}",
-
dataRegionCreateTask.getRegionReplicaSet().getRegionId(),
- dataRegionCreateTask.getTargetDataNode());
- createDataRegionHandler.putRequest(
- dataRegionCreateTask.getRegionId().getId(),
- new TCreateDataRegionReq(
- dataRegionCreateTask.getRegionReplicaSet(),
- dataRegionCreateTask.getStorageGroup()));
- createDataRegionHandler.putNodeLocation(
- dataRegionCreateTask.getRegionId().getId(),
- dataRegionCreateTask.getTargetDataNode());
- }
-
-
CnToDnInternalServiceAsyncRequestManager.getInstance()
-
.sendAsyncRequestWithRetry(createDataRegionHandler);
-
- for (Map.Entry<Integer, TSStatus> entry :
-
createDataRegionHandler.getResponseMap().entrySet()) {
- if (entry.getValue().getCode()
- ==
TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
- successfulTask.add(
- new TConsensusGroupId(
- TConsensusGroupType.DataRegion,
entry.getKey()));
- }
- }
- break;
- }
- break;
- case DELETE:
- // delete region
- DataNodeAsyncRequestContext<TConsensusGroupId, TSStatus>
deleteRegionHandler =
- new
DataNodeAsyncRequestContext<>(CnToDnAsyncRequestType.DELETE_REGION);
- Map<Integer, TConsensusGroupId> regionIdMap = new
HashMap<>();
- for (RegionMaintainTask regionMaintainTask :
selectedRegionMaintainTask) {
- RegionDeleteTask regionDeleteTask = (RegionDeleteTask)
regionMaintainTask;
- LOGGER.info(
- "Start to delete Region: {} on DataNode: {}",
- regionDeleteTask.getRegionId(),
- regionDeleteTask.getTargetDataNode());
- deleteRegionHandler.putRequest(
- regionDeleteTask.getRegionId().getId(),
regionDeleteTask.getRegionId());
- deleteRegionHandler.putNodeLocation(
- regionDeleteTask.getRegionId().getId(),
- regionDeleteTask.getTargetDataNode());
- regionIdMap.put(
- regionDeleteTask.getRegionId().getId(),
regionDeleteTask.getRegionId());
- }
-
- long startTime = System.currentTimeMillis();
- CnToDnInternalServiceAsyncRequestManager.getInstance()
- .sendAsyncRequestWithRetry(deleteRegionHandler);
-
- LOGGER.info(
- "Deleting regions costs {}ms",
(System.currentTimeMillis() - startTime));
-
- for (Map.Entry<Integer, TSStatus> entry :
- deleteRegionHandler.getResponseMap().entrySet()) {
- if (entry.getValue().getCode()
- == TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
- successfulTask.add(regionIdMap.get(entry.getKey()));
- }
- }
- break;
- }
-
- if (successfulTask.isEmpty()) {
+ if (completedRegions.isEmpty()) {
break;
}
- for (TConsensusGroupId regionId : successfulTask) {
- regionMaintainTaskMap.compute(
- regionId,
- (k, v) -> {
- if (v == null) {
- throw new IllegalStateException();
- }
- v.poll();
- if (v.isEmpty()) {
- return null;
- } else {
- return v;
- }
- });
- }
-
- // Poll the head entry if success
try {
- getConsensusManager()
- .write(new
PollSpecificRegionMaintainTaskPlan(successfulTask));
+ TSStatus pollStatus =
+ getConsensusManager()
+ .write(new
PollSpecificRegionMaintainTaskPlan(completedRegions));
+ if (pollStatus.getCode() !=
TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ break;
+ }
} catch (ConsensusException e) {
LOGGER.warn(CONSENSUS_WRITE_ERROR, e);
- }
-
- if (successfulTask.size() <
selectedRegionMaintainTask.size()) {
- // Here we just break and wait until next schedule task
- // due to all the RegionMaintainEntry should be executed by
- // the order of they were offered
break;
}
+
+ pollCompletedRegionMaintainTaskHeads(tasksByRegion,
completedRegions);
+
+ // Failed heads remain persisted and block only their own
regions until the next
+ // scheduled retry.
+ deferFailedRegionMaintainTasks(deferredRegions, headsByType,
completedRegions);
}
}
});
}
+ static Map<TConsensusGroupId, Queue<RegionMaintainTask>>
groupRegionMaintainTasks(
+ List<RegionMaintainTask> tasks) {
+ Map<TConsensusGroupId, Queue<RegionMaintainTask>> tasksByRegion = new
LinkedHashMap<>();
+ for (RegionMaintainTask task : tasks) {
+ tasksByRegion.computeIfAbsent(task.getRegionId(), key -> new
ArrayDeque<>()).add(task);
+ }
+ return tasksByRegion;
+ }
+
+ static Map<RegionMaintainType, List<RegionMaintainTask>>
getRegionMaintainTaskHeads(
+ Map<TConsensusGroupId, Queue<RegionMaintainTask>> tasksByRegion) {
+ return getRegionMaintainTaskHeads(tasksByRegion, Collections.emptySet());
+ }
+
+ static Map<RegionMaintainType, List<RegionMaintainTask>>
getRegionMaintainTaskHeads(
+ Map<TConsensusGroupId, Queue<RegionMaintainTask>> tasksByRegion,
+ Set<TConsensusGroupId> deferredRegions) {
+ Map<RegionMaintainType, List<RegionMaintainTask>> headsByType =
+ new EnumMap<>(RegionMaintainType.class);
+ for (Map.Entry<TConsensusGroupId, Queue<RegionMaintainTask>> entry :
tasksByRegion.entrySet()) {
+ if (deferredRegions.contains(entry.getKey())) {
+ continue;
+ }
+ Queue<RegionMaintainTask> taskQueue = entry.getValue();
+ RegionMaintainTask task = taskQueue.peek();
+ if (task != null) {
+ headsByType.computeIfAbsent(task.getType(), key -> new
ArrayList<>()).add(task);
+ }
+ }
+ return headsByType;
+ }
+
+ static void pollCompletedRegionMaintainTaskHeads(
+ Map<TConsensusGroupId, Queue<RegionMaintainTask>> tasksByRegion,
+ Set<TConsensusGroupId> completedRegions) {
+ for (TConsensusGroupId regionId : completedRegions) {
+ tasksByRegion.computeIfPresent(
+ regionId,
+ (key, queue) -> {
+ queue.poll();
+ return queue.isEmpty() ? null : queue;
+ });
+ }
+ }
+
+ static void deferFailedRegionMaintainTasks(
+ Set<TConsensusGroupId> deferredRegions,
+ Map<RegionMaintainType, List<RegionMaintainTask>> submittedTaskHeads,
+ Set<TConsensusGroupId> completedRegions) {
+ submittedTaskHeads.values().stream()
+ .flatMap(List::stream)
+ .map(RegionMaintainTask::getRegionId)
+ .filter(regionId -> !completedRegions.contains(regionId))
+ .forEach(deferredRegions::add);
+ }
+
+ private Set<TConsensusGroupId> submitRegionMaintainTasks(
+ RegionMaintainType taskType, List<RegionMaintainTask> tasks) {
+ return taskType == RegionMaintainType.CREATE
+ ? submitRegionCreateTasks(tasks)
+ : submitRegionDeleteTasks(tasks);
+ }
+
+ private Set<TConsensusGroupId>
submitRegionCreateTasks(List<RegionMaintainTask> tasks) {
+ Map<TConsensusGroupType, List<RegionCreateTask>> tasksByRegionType =
+ new EnumMap<>(TConsensusGroupType.class);
+ for (RegionMaintainTask task : tasks) {
+ RegionCreateTask createTask = (RegionCreateTask) task;
+ tasksByRegionType
+ .computeIfAbsent(createTask.getRegionId().getType(), key -> new
ArrayList<>())
+ .add(createTask);
+ }
+
+ Set<TConsensusGroupId> completedRegions = new HashSet<>();
+ for (Map.Entry<TConsensusGroupType, List<RegionCreateTask>> entry :
+ tasksByRegionType.entrySet()) {
+ completedRegions.addAll(submitRegionCreateTasks(entry.getKey(),
entry.getValue()));
+ }
+ return completedRegions;
+ }
+
+ private Set<TConsensusGroupId> submitRegionCreateTasks(
+ TConsensusGroupType regionType, List<RegionCreateTask> tasks) {
+ DataNodeAsyncRequestContext<Object, TSStatus> requestContext =
+ new DataNodeAsyncRequestContext<>(
+ regionType == TConsensusGroupType.SchemaRegion
+ ? CnToDnAsyncRequestType.CREATE_SCHEMA_REGION
+ : CnToDnAsyncRequestType.CREATE_DATA_REGION);
+ Map<Integer, TConsensusGroupId> regionByRequestIndex = new HashMap<>();
+ for (int requestIndex = 0; requestIndex < tasks.size(); requestIndex++) {
+ RegionCreateTask task = tasks.get(requestIndex);
+ LOGGER.info(
+ "Start to create Region: {} on DataNode: {}",
+ task.getRegionReplicaSet().getRegionId(),
+ task.getTargetDataNode());
+ Object request =
+ regionType == TConsensusGroupType.SchemaRegion
+ ? new TCreateSchemaRegionReq(task.getRegionReplicaSet(),
task.getStorageGroup())
+ : new TCreateDataRegionReq(task.getRegionReplicaSet(),
task.getStorageGroup());
+ requestContext.putRequest(requestIndex, request);
+ requestContext.putNodeLocation(requestIndex, task.getTargetDataNode());
+ regionByRequestIndex.put(requestIndex, task.getRegionId());
+ }
+ CnToDnInternalServiceAsyncRequestManager.getInstance()
+ .sendAsyncRequestWithRetry(requestContext);
+ return collectCompletedRegionMaintainTasks(
+ RegionMaintainType.CREATE, requestContext.getResponseMap(),
regionByRequestIndex);
+ }
+
+ private Set<TConsensusGroupId>
submitRegionDeleteTasks(List<RegionMaintainTask> tasks) {
+ DataNodeAsyncRequestContext<TConsensusGroupId, TSStatus> requestContext =
+ new
DataNodeAsyncRequestContext<>(CnToDnAsyncRequestType.DELETE_REGION);
+ Map<Integer, TConsensusGroupId> regionByRequestIndex = new HashMap<>();
+ for (int requestIndex = 0; requestIndex < tasks.size(); requestIndex++) {
+ RegionDeleteTask task = (RegionDeleteTask) tasks.get(requestIndex);
+ LOGGER.info(
+ "Start to delete Region: {} on DataNode: {}",
+ task.getRegionId(),
+ task.getTargetDataNode());
+ requestContext.putRequest(requestIndex, task.getRegionId());
+ requestContext.putNodeLocation(requestIndex, task.getTargetDataNode());
+ regionByRequestIndex.put(requestIndex, task.getRegionId());
+ }
+ CnToDnInternalServiceAsyncRequestManager.getInstance()
+ .sendAsyncRequestWithRetry(requestContext);
+ return collectCompletedRegionMaintainTasks(
+ RegionMaintainType.DELETE, requestContext.getResponseMap(),
regionByRequestIndex);
+ }
+
+ static Set<TConsensusGroupId> collectCompletedRegionMaintainTasks(
+ RegionMaintainType taskType,
+ Map<Integer, TSStatus> responseMap,
+ Map<Integer, TConsensusGroupId> regionByRequestIndex) {
+ Set<TConsensusGroupId> completedRegions = new HashSet<>();
+ responseMap.forEach(
+ (requestIndex, status) -> {
+ if (isRegionMaintainTaskCompleted(taskType, status)) {
+ TConsensusGroupId regionId =
regionByRequestIndex.get(requestIndex);
+ if (regionId != null) {
+ completedRegions.add(regionId);
+ }
+ }
+ });
+ return completedRegions;
+ }
+
+ static boolean isRegionMaintainTaskCompleted(RegionMaintainType taskType,
TSStatus status) {
+ if (status.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ return true;
+ }
+ return taskType == RegionMaintainType.CREATE
+ ? status.getCode() ==
TSStatusCode.REGION_ALREADY_EXISTS.getStatusCode()
+ : status.getCode() == TSStatusCode.REGION_NOT_EXIST.getStatusCode();
+ }
+
public void startRegionCleaner() {
synchronized (scheduleMonitor) {
if (currentRegionMaintainerFuture == null) {
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteDatabaseProcedure.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteDatabaseProcedure.java
index b6ad21128af..f5d63aa1fc0 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteDatabaseProcedure.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteDatabaseProcedure.java
@@ -112,38 +112,17 @@ public class DeleteDatabaseProcedure
deleteDatabaseSchema.getName());
// Submit RegionDeleteTasks
- OfferRegionMaintainTasksPlan dataRegionDeleteTaskOfferPlan =
- new OfferRegionMaintainTasksPlan();
List<TRegionReplicaSet> regionReplicaSets =
env.getAllReplicaSets(deleteDatabaseSchema.getName());
- List<TRegionReplicaSet> schemaRegionReplicaSets = new ArrayList<>();
+ OfferRegionMaintainTasksPlan regionDeleteTaskOfferPlan =
+ buildDataRegionDeleteTaskOfferPlan(regionReplicaSets);
+ List<TRegionReplicaSet> schemaRegionReplicaSets =
+ getSchemaRegionReplicaSets(regionReplicaSets);
regionReplicaSets.forEach(
- regionReplicaSet -> {
- // Clear heartbeat cache along the way
- env.getConfigManager()
- .getLoadManager()
-
.removeRegionGroupRelatedCache(regionReplicaSet.getRegionId());
-
- if (regionReplicaSet
- .getRegionId()
- .getType()
- .equals(TConsensusGroupType.SchemaRegion)) {
- schemaRegionReplicaSets.add(regionReplicaSet);
- } else {
- regionReplicaSet
- .getDataNodeLocations()
- .forEach(
- targetDataNode ->
-
dataRegionDeleteTaskOfferPlan.appendRegionMaintainTask(
- new RegionDeleteTask(
- targetDataNode,
regionReplicaSet.getRegionId())));
- }
- });
-
- if
(!dataRegionDeleteTaskOfferPlan.getRegionMaintainTaskList().isEmpty()) {
- // submit async data region delete task
-
env.getConfigManager().getConsensusManager().write(dataRegionDeleteTaskOfferPlan);
- }
+ regionReplicaSet ->
+ env.getConfigManager()
+ .getLoadManager()
+
.removeRegionGroupRelatedCache(regionReplicaSet.getRegionId()));
// try sync delete schemaengine region
DataNodeAsyncRequestContext<TConsensusGroupId, TSStatus>
asyncClientHandler =
@@ -166,7 +145,7 @@ public class DeleteDatabaseProcedure
.sendAsyncRequestWithRetry(asyncClientHandler);
for (Map.Entry<Integer, TSStatus> entry :
asyncClientHandler.getResponseMap().entrySet()) {
- if (entry.getValue().getCode() ==
TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ if (isRegionDeleteCompleted(entry.getValue())) {
LOG.info(
"[DeleteDatabaseProcedure] Successfully delete
SchemaRegion[{}] on {}",
asyncClientHandler.getRequest(entry.getKey()),
@@ -181,16 +160,17 @@ public class DeleteDatabaseProcedure
}
if (!schemaRegionDeleteTaskMap.isEmpty()) {
- // submit async schemaengine region delete task for failed sync
execution
- OfferRegionMaintainTasksPlan schemaRegionDeleteTaskOfferPlan =
- new OfferRegionMaintainTasksPlan();
- schemaRegionDeleteTaskMap
- .values()
-
.forEach(schemaRegionDeleteTaskOfferPlan::appendRegionMaintainTask);
-
env.getConfigManager().getConsensusManager().write(schemaRegionDeleteTaskOfferPlan);
+ // submit async schemaengine region delete tasks for failed sync
executions
+ appendFailedSchemaRegionDeleteTasks(
+ regionDeleteTaskOfferPlan, schemaRegionDeleteTaskMap);
}
}
+ if (!offerRegionDeleteTasks(env, regionDeleteTaskOfferPlan)) {
+ setNextState(DeleteStorageGroupState.DELETE_DATABASE_SCHEMA);
+ return Flow.HAS_MORE_STATE;
+ }
+
env.getConfigManager()
.getLoadManager()
.clearDataPartitionPolicyTable(deleteDatabaseSchema.getName());
@@ -212,6 +192,8 @@ public class DeleteDatabaseProcedure
} else if (getCycles() > RETRY_THRESHOLD) {
setFailure(
new ProcedureException("[DeleteDatabaseProcedure] Delete
DatabaseSchema failed"));
+ } else {
+ setNextState(DeleteStorageGroupState.DELETE_DATABASE_SCHEMA);
}
}
} catch (ConsensusException | TException | IOException e) {
@@ -230,12 +212,67 @@ public class DeleteDatabaseProcedure
e);
if (getCycles() > RETRY_THRESHOLD) {
setFailure(new ProcedureException("[DeleteDatabaseProcedure] State
stuck at " + state));
+ } else {
+ setNextState(state);
}
}
}
return Flow.HAS_MORE_STATE;
}
+ static boolean isRegionDeleteCompleted(TSStatus status) {
+ return status.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()
+ || status.getCode() == TSStatusCode.REGION_NOT_EXIST.getStatusCode();
+ }
+
+ static OfferRegionMaintainTasksPlan buildDataRegionDeleteTaskOfferPlan(
+ List<TRegionReplicaSet> regionReplicaSets) {
+ OfferRegionMaintainTasksPlan offerPlan = new
OfferRegionMaintainTasksPlan();
+ regionReplicaSets.stream()
+ .filter(
+ regionReplicaSet ->
+ regionReplicaSet.getRegionId().getType() ==
TConsensusGroupType.DataRegion)
+ .forEach(
+ regionReplicaSet ->
+ regionReplicaSet
+ .getDataNodeLocations()
+ .forEach(
+ targetDataNode ->
+ offerPlan.appendRegionMaintainTask(
+ new RegionDeleteTask(
+ targetDataNode,
regionReplicaSet.getRegionId()))));
+ return offerPlan;
+ }
+
+ private static List<TRegionReplicaSet> getSchemaRegionReplicaSets(
+ List<TRegionReplicaSet> regionReplicaSets) {
+ List<TRegionReplicaSet> schemaRegionReplicaSets = new ArrayList<>();
+ regionReplicaSets.stream()
+ .filter(
+ regionReplicaSet ->
+ regionReplicaSet.getRegionId().getType() ==
TConsensusGroupType.SchemaRegion)
+ .forEach(schemaRegionReplicaSets::add);
+ return schemaRegionReplicaSets;
+ }
+
+ static void appendFailedSchemaRegionDeleteTasks(
+ OfferRegionMaintainTasksPlan offerPlan,
+ Map<Integer, RegionDeleteTask> failedSchemaRegionDeleteTasks) {
+
failedSchemaRegionDeleteTasks.values().forEach(offerPlan::appendRegionMaintainTask);
+ }
+
+ private boolean offerRegionDeleteTasks(
+ ConfigNodeProcedureEnv env, OfferRegionMaintainTasksPlan offerPlan)
+ throws ConsensusException {
+ return offerPlan.getRegionMaintainTaskList().isEmpty()
+ || isRegionDeleteTaskOfferSuccessful(
+ env.getConfigManager().getConsensusManager().write(offerPlan));
+ }
+
+ static boolean isRegionDeleteTaskOfferSuccessful(TSStatus status) {
+ return status.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode();
+ }
+
@Override
protected void rollbackState(ConfigNodeProcedureEnv env,
DeleteStorageGroupState state)
throws IOException, InterruptedException {
diff --git
a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/client/async/handlers/rpc/DataNodeTSStatusRPCHandlerTest.java
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/client/async/handlers/rpc/DataNodeTSStatusRPCHandlerTest.java
new file mode 100644
index 00000000000..7f06a44edff
--- /dev/null
+++
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/client/async/handlers/rpc/DataNodeTSStatusRPCHandlerTest.java
@@ -0,0 +1,63 @@
+/*
+ * 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.client.async.handlers.rpc;
+
+import org.apache.iotdb.common.rpc.thrift.TSStatus;
+import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType;
+import org.apache.iotdb.rpc.TSStatusCode;
+
+import org.junit.Test;
+
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertTrue;
+
+public class DataNodeTSStatusRPCHandlerTest {
+
+ @Test
+ public void testRegionOperationTerminalStatuses() {
+ assertTrue(
+ DataNodeTSStatusRPCHandler.isRequestCompleted(
+ CnToDnAsyncRequestType.CREATE_DATA_REGION,
status(TSStatusCode.SUCCESS_STATUS)));
+ assertTrue(
+ DataNodeTSStatusRPCHandler.isRequestCompleted(
+ CnToDnAsyncRequestType.CREATE_DATA_REGION,
status(TSStatusCode.REGION_ALREADY_EXISTS)));
+ assertTrue(
+ DataNodeTSStatusRPCHandler.isRequestCompleted(
+ CnToDnAsyncRequestType.CREATE_SCHEMA_REGION,
+ status(TSStatusCode.REGION_ALREADY_EXISTS)));
+ assertTrue(
+ DataNodeTSStatusRPCHandler.isRequestCompleted(
+ CnToDnAsyncRequestType.DELETE_REGION,
status(TSStatusCode.REGION_NOT_EXIST)));
+
+ assertFalse(
+ DataNodeTSStatusRPCHandler.isRequestCompleted(
+ CnToDnAsyncRequestType.CREATE_SCHEMA_REGION,
status(TSStatusCode.CREATE_REGION_ERROR)));
+ assertFalse(
+ DataNodeTSStatusRPCHandler.isRequestCompleted(
+ CnToDnAsyncRequestType.DELETE_REGION,
status(TSStatusCode.DELETE_REGION_ERROR)));
+ assertFalse(
+ DataNodeTSStatusRPCHandler.isRequestCompleted(
+ CnToDnAsyncRequestType.SET_TTL,
status(TSStatusCode.REGION_ALREADY_EXISTS)));
+ }
+
+ private static TSStatus status(TSStatusCode statusCode) {
+ return new TSStatus(statusCode.getStatusCode());
+ }
+}
diff --git
a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/partition/PartitionManagerRegionMaintainTest.java
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/partition/PartitionManagerRegionMaintainTest.java
new file mode 100644
index 00000000000..4157f98495c
--- /dev/null
+++
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/partition/PartitionManagerRegionMaintainTest.java
@@ -0,0 +1,151 @@
+/*
+ * 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.manager.partition;
+
+import org.apache.iotdb.common.rpc.thrift.TConsensusGroupId;
+import org.apache.iotdb.common.rpc.thrift.TConsensusGroupType;
+import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation;
+import org.apache.iotdb.common.rpc.thrift.TRegionReplicaSet;
+import org.apache.iotdb.common.rpc.thrift.TSStatus;
+import
org.apache.iotdb.confignode.persistence.partition.maintainer.RegionCreateTask;
+import
org.apache.iotdb.confignode.persistence.partition.maintainer.RegionDeleteTask;
+import
org.apache.iotdb.confignode.persistence.partition.maintainer.RegionMaintainTask;
+import
org.apache.iotdb.confignode.persistence.partition.maintainer.RegionMaintainType;
+import org.apache.iotdb.rpc.TSStatusCode;
+
+import org.junit.Test;
+
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Queue;
+import java.util.Set;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertSame;
+
+public class PartitionManagerRegionMaintainTest {
+
+ @Test
+ public void testPerRegionFifoAndPartialFailureRetry() {
+ TConsensusGroupId region0 = regionId(TConsensusGroupType.DataRegion, 0);
+ TConsensusGroupId region1 = regionId(TConsensusGroupType.DataRegion, 1);
+ RegionMaintainTask region0Head = deleteTask(region0);
+ RegionMaintainTask region0Next = createTask(region0);
+ RegionMaintainTask region1Head = deleteTask(region1);
+ Map<TConsensusGroupId, Queue<RegionMaintainTask>> tasksByRegion =
+ PartitionManager.groupRegionMaintainTasks(
+ Arrays.asList(region0Head, region0Next, region1Head));
+
+ Map<RegionMaintainType, List<RegionMaintainTask>> firstRound =
+ PartitionManager.getRegionMaintainTaskHeads(tasksByRegion);
+ assertEquals(
+ Arrays.asList(region0Head, region1Head),
firstRound.get(RegionMaintainType.DELETE));
+ assertFalse(firstRound.containsKey(RegionMaintainType.CREATE));
+
+ Set<TConsensusGroupId> deferredRegions = new HashSet<>();
+ PartitionManager.deferFailedRegionMaintainTasks(
+ deferredRegions, firstRound, Collections.singleton(region1));
+ PartitionManager.pollCompletedRegionMaintainTaskHeads(
+ tasksByRegion, Collections.singleton(region1));
+ assertSame(region0Head, tasksByRegion.get(region0).peek());
+ assertFalse(tasksByRegion.containsKey(region1));
+ assertFalse(
+ PartitionManager.getRegionMaintainTaskHeads(tasksByRegion,
deferredRegions)
+ .containsKey(RegionMaintainType.DELETE));
+
+ tasksByRegion =
+ PartitionManager.groupRegionMaintainTasks(
+ Arrays.asList(region0Head, region0Next, region1Head));
+ firstRound = PartitionManager.getRegionMaintainTaskHeads(tasksByRegion);
+ deferredRegions = new HashSet<>();
+ PartitionManager.deferFailedRegionMaintainTasks(
+ deferredRegions, firstRound, Collections.singleton(region0));
+ PartitionManager.pollCompletedRegionMaintainTaskHeads(
+ tasksByRegion, Collections.singleton(region0));
+ assertSame(
+ region0Next,
+ PartitionManager.getRegionMaintainTaskHeads(tasksByRegion,
deferredRegions)
+ .get(RegionMaintainType.CREATE)
+ .get(0));
+ assertSame(region1Head, tasksByRegion.get(region1).peek());
+
+ deferredRegions.clear();
+ assertSame(
+ region1Head,
+ PartitionManager.getRegionMaintainTaskHeads(tasksByRegion)
+ .get(RegionMaintainType.DELETE)
+ .get(0));
+ }
+
+ @Test
+ public void testCompletedStatusAndRequestIndexMapping() {
+ assertCompleted(RegionMaintainType.CREATE, TSStatusCode.SUCCESS_STATUS,
true);
+ assertCompleted(RegionMaintainType.CREATE,
TSStatusCode.REGION_ALREADY_EXISTS, true);
+ assertCompleted(RegionMaintainType.CREATE, TSStatusCode.REGION_NOT_EXIST,
false);
+ assertCompleted(RegionMaintainType.CREATE,
TSStatusCode.CREATE_REGION_ERROR, false);
+ assertCompleted(RegionMaintainType.DELETE, TSStatusCode.SUCCESS_STATUS,
true);
+ assertCompleted(RegionMaintainType.DELETE, TSStatusCode.REGION_NOT_EXIST,
true);
+ assertCompleted(RegionMaintainType.DELETE,
TSStatusCode.REGION_ALREADY_EXISTS, false);
+ assertCompleted(RegionMaintainType.DELETE,
TSStatusCode.DELETE_REGION_ERROR, false);
+
+ TConsensusGroupId schemaRegion =
regionId(TConsensusGroupType.SchemaRegion, 7);
+ TConsensusGroupId dataRegion = regionId(TConsensusGroupType.DataRegion, 7);
+ Map<Integer, TConsensusGroupId> regionsByRequestIndex = new HashMap<>();
+ regionsByRequestIndex.put(0, schemaRegion);
+ regionsByRequestIndex.put(1, dataRegion);
+ Map<Integer, TSStatus> responses = new HashMap<>();
+ responses.put(0, status(TSStatusCode.SUCCESS_STATUS));
+ responses.put(1, status(TSStatusCode.CREATE_REGION_ERROR));
+
+ Set<TConsensusGroupId> completed =
+ PartitionManager.collectCompletedRegionMaintainTasks(
+ RegionMaintainType.CREATE, responses, regionsByRequestIndex);
+ assertEquals(Collections.singleton(schemaRegion), completed);
+ }
+
+ private static void assertCompleted(
+ RegionMaintainType type, TSStatusCode statusCode, boolean expected) {
+ assertEquals(
+ expected, PartitionManager.isRegionMaintainTaskCompleted(type,
status(statusCode)));
+ }
+
+ private static TSStatus status(TSStatusCode statusCode) {
+ return new TSStatus(statusCode.getStatusCode());
+ }
+
+ private static RegionDeleteTask deleteTask(TConsensusGroupId regionId) {
+ return new RegionDeleteTask(new TDataNodeLocation(), regionId);
+ }
+
+ private static RegionCreateTask createTask(TConsensusGroupId regionId) {
+ TRegionReplicaSet replicaSet =
+ new TRegionReplicaSet(regionId, Collections.singletonList(new
TDataNodeLocation()));
+ return new RegionCreateTask(new TDataNodeLocation(), "root.test",
replicaSet);
+ }
+
+ private static TConsensusGroupId regionId(TConsensusGroupType type, int id) {
+ return new TConsensusGroupId(type, id);
+ }
+}
diff --git
a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/persistence/PartitionInfoTest.java
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/persistence/PartitionInfoTest.java
index c15cefff657..f6f209c9356 100644
---
a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/persistence/PartitionInfoTest.java
+++
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/persistence/PartitionInfoTest.java
@@ -36,10 +36,12 @@ import
org.apache.iotdb.confignode.consensus.request.write.partition.CreateDataP
import
org.apache.iotdb.confignode.consensus.request.write.partition.CreateSchemaPartitionPlan;
import
org.apache.iotdb.confignode.consensus.request.write.region.CreateRegionGroupsPlan;
import
org.apache.iotdb.confignode.consensus.request.write.region.OfferRegionMaintainTasksPlan;
+import
org.apache.iotdb.confignode.consensus.request.write.region.PollSpecificRegionMaintainTaskPlan;
import
org.apache.iotdb.confignode.consensus.response.partition.RegionInfoListResp;
import org.apache.iotdb.confignode.persistence.partition.PartitionInfo;
import
org.apache.iotdb.confignode.persistence.partition.maintainer.RegionCreateTask;
import
org.apache.iotdb.confignode.persistence.partition.maintainer.RegionDeleteTask;
+import
org.apache.iotdb.confignode.persistence.partition.maintainer.RegionMaintainTask;
import org.apache.iotdb.confignode.rpc.thrift.TDatabaseSchema;
import org.apache.iotdb.confignode.rpc.thrift.TShowRegionReq;
@@ -191,6 +193,23 @@ public class PartitionInfoTest {
Assert.assertEquals(Optional.empty(), partitionInfo.getRegionType(-1));
}
+ @Test
+ public void testPollSpecificRegionMaintainTaskOnlyRemovesEachRegionHead() {
+ OfferRegionMaintainTasksPlan offerPlan =
generateOfferRegionMaintainTasksPlan();
+ partitionInfo.offerRegionMaintainTasks(offerPlan);
+
+ TConsensusGroupId dataRegionId = new
TConsensusGroupId(TConsensusGroupType.DataRegion, 0);
+ partitionInfo.pollSpecificRegionMaintainTask(
+ new
PollSpecificRegionMaintainTaskPlan(Collections.singleton(dataRegionId)));
+
+ List<RegionMaintainTask> remainingTasks =
partitionInfo.getRegionMaintainEntryList();
+ Assert.assertEquals(2, remainingTasks.size());
+ Assert.assertEquals(dataRegionId, remainingTasks.get(0).getRegionId());
+ Assert.assertEquals(
+ new TConsensusGroupId(TConsensusGroupType.SchemaRegion, 2),
+ remainingTasks.get(1).getRegionId());
+ }
+
@Test
public void testShowRegion() {
for (int i = 0; i < 2; i++) {
diff --git
a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteDatabaseProcedureTest.java
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteDatabaseProcedureTest.java
index b12f49d9bd7..2f6d365bc7a 100644
---
a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteDatabaseProcedureTest.java
+++
b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteDatabaseProcedureTest.java
@@ -19,16 +19,41 @@
package org.apache.iotdb.confignode.procedure.impl.schema;
+import org.apache.iotdb.common.rpc.thrift.TConsensusGroupId;
+import org.apache.iotdb.common.rpc.thrift.TConsensusGroupType;
+import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation;
+import org.apache.iotdb.common.rpc.thrift.TRegionReplicaSet;
+import org.apache.iotdb.common.rpc.thrift.TSStatus;
+import org.apache.iotdb.confignode.consensus.request.ConfigPhysicalPlan;
+import
org.apache.iotdb.confignode.consensus.request.write.region.OfferRegionMaintainTasksPlan;
+import org.apache.iotdb.confignode.manager.ConfigManager;
+import org.apache.iotdb.confignode.manager.consensus.ConsensusManager;
+import org.apache.iotdb.confignode.manager.load.LoadManager;
+import
org.apache.iotdb.confignode.persistence.partition.maintainer.RegionDeleteTask;
+import
org.apache.iotdb.confignode.persistence.partition.maintainer.RegionMaintainTask;
+import
org.apache.iotdb.confignode.persistence.partition.maintainer.RegionMaintainType;
+import org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv;
+import
org.apache.iotdb.confignode.procedure.state.schema.DeleteStorageGroupState;
import org.apache.iotdb.confignode.procedure.store.ProcedureFactory;
import org.apache.iotdb.confignode.rpc.thrift.TDatabaseSchema;
+import org.apache.iotdb.rpc.TSStatusCode;
import org.apache.tsfile.utils.PublicBAOS;
import org.junit.Test;
+import org.mockito.ArgumentCaptor;
+import org.mockito.Mockito;
import java.io.DataOutputStream;
import java.nio.ByteBuffer;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertTrue;
import static org.junit.Assert.fail;
public class DeleteDatabaseProcedureTest {
@@ -53,4 +78,97 @@ public class DeleteDatabaseProcedureTest {
fail();
}
}
+
+ @Test
+ public void testDataRegionDeleteTasksArePreparedBeforeMetadataCleanup() {
+ TDataNodeLocation dataNode0 = new TDataNodeLocation().setDataNodeId(0);
+ TDataNodeLocation dataNode1 = new TDataNodeLocation().setDataNodeId(1);
+ TConsensusGroupId dataRegionId = new
TConsensusGroupId(TConsensusGroupType.DataRegion, 10);
+ TConsensusGroupId schemaRegionId = new
TConsensusGroupId(TConsensusGroupType.SchemaRegion, 20);
+ List<TRegionReplicaSet> replicaSets =
+ Arrays.asList(
+ new TRegionReplicaSet(dataRegionId, Arrays.asList(dataNode0,
dataNode1)),
+ new TRegionReplicaSet(schemaRegionId,
Collections.singletonList(dataNode0)));
+
+ OfferRegionMaintainTasksPlan offerPlan =
+
DeleteDatabaseProcedure.buildDataRegionDeleteTaskOfferPlan(replicaSets);
+ List<RegionMaintainTask> tasks = offerPlan.getRegionMaintainTaskList();
+
+ assertEquals(2, tasks.size());
+ assertEquals(RegionMaintainType.DELETE, tasks.get(0).getType());
+ assertEquals(dataRegionId, tasks.get(0).getRegionId());
+ assertEquals(dataNode0, tasks.get(0).getTargetDataNode());
+ assertEquals(dataNode1, tasks.get(1).getTargetDataNode());
+ }
+
+ @Test
+ public void testIdempotentDeleteStatusIsCompleted() {
+ assertTrue(
+ DeleteDatabaseProcedure.isRegionDeleteCompleted(
+ new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode())));
+ assertTrue(
+ DeleteDatabaseProcedure.isRegionDeleteCompleted(
+ new TSStatus(TSStatusCode.REGION_NOT_EXIST.getStatusCode())));
+ assertFalse(
+ DeleteDatabaseProcedure.isRegionDeleteCompleted(
+ new TSStatus(TSStatusCode.DELETE_REGION_ERROR.getStatusCode())));
+ }
+
+ @Test
+ public void testFailedSynchronousSchemaRegionDeleteIsQueued() {
+ TDataNodeLocation dataNode = new TDataNodeLocation().setDataNodeId(0);
+ TConsensusGroupId schemaRegionId = new
TConsensusGroupId(TConsensusGroupType.SchemaRegion, 20);
+ RegionDeleteTask failedTask = new RegionDeleteTask(dataNode,
schemaRegionId);
+ Map<Integer, RegionDeleteTask> failedTasks = new HashMap<>();
+ failedTasks.put(0, failedTask);
+ OfferRegionMaintainTasksPlan offerPlan = new
OfferRegionMaintainTasksPlan();
+
+ DeleteDatabaseProcedure.appendFailedSchemaRegionDeleteTasks(offerPlan,
failedTasks);
+
+ assertEquals(Collections.singletonList(failedTask),
offerPlan.getRegionMaintainTaskList());
+ }
+
+ @Test
+ public void testRegionDeleteTaskOfferMustSucceedBeforeMetadataCleanup() {
+ assertTrue(
+ DeleteDatabaseProcedure.isRegionDeleteTaskOfferSuccessful(
+ new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode())));
+ assertFalse(
+ DeleteDatabaseProcedure.isRegionDeleteTaskOfferSuccessful(
+ new
TSStatus(TSStatusCode.EXECUTE_STATEMENT_ERROR.getStatusCode())));
+ }
+
+ @Test
+ public void testFailedTaskOfferPreventsPartitionMetadataCleanup() throws
Exception {
+ TDataNodeLocation dataNode = new TDataNodeLocation().setDataNodeId(0);
+ TConsensusGroupId dataRegionId = new
TConsensusGroupId(TConsensusGroupType.DataRegion, 10);
+ TRegionReplicaSet dataRegion =
+ new TRegionReplicaSet(dataRegionId,
Collections.singletonList(dataNode));
+ ConfigNodeProcedureEnv env = Mockito.mock(ConfigNodeProcedureEnv.class);
+ ConfigManager configManager = Mockito.mock(ConfigManager.class);
+ ConsensusManager consensusManager = Mockito.mock(ConsensusManager.class);
+ LoadManager loadManager = Mockito.mock(LoadManager.class);
+ Mockito.when(env.getAllReplicaSets("root.sg"))
+ .thenReturn(Collections.singletonList(dataRegion));
+ Mockito.when(env.getConfigManager()).thenReturn(configManager);
+
Mockito.when(configManager.getConsensusManager()).thenReturn(consensusManager);
+ Mockito.when(configManager.getLoadManager()).thenReturn(loadManager);
+ Mockito.when(consensusManager.write(Mockito.any(ConfigPhysicalPlan.class)))
+ .thenReturn(new
TSStatus(TSStatusCode.EXECUTE_STATEMENT_ERROR.getStatusCode()));
+
+ DeleteDatabaseProcedure procedure =
+ new DeleteDatabaseProcedure(new TDatabaseSchema("root.sg"), false);
+ procedure.executeFromState(env,
DeleteStorageGroupState.DELETE_DATABASE_SCHEMA);
+
+ ArgumentCaptor<ConfigPhysicalPlan> planCaptor =
+ ArgumentCaptor.forClass(ConfigPhysicalPlan.class);
+ Mockito.verify(consensusManager).write(planCaptor.capture());
+ assertTrue(planCaptor.getValue() instanceof OfferRegionMaintainTasksPlan);
+ assertEquals(
+ 1,
+ ((OfferRegionMaintainTasksPlan)
planCaptor.getValue()).getRegionMaintainTaskList().size());
+ Mockito.verify(loadManager,
Mockito.never()).clearDataPartitionPolicyTable(Mockito.anyString());
+ Mockito.verify(env, Mockito.never())
+ .deleteDatabaseConfig(Mockito.anyString(), Mockito.anyBoolean());
+ }
}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeInternalRPCServiceImpl.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeInternalRPCServiceImpl.java
index 99e314affa0..9edb75c641c 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeInternalRPCServiceImpl.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeInternalRPCServiceImpl.java
@@ -2099,6 +2099,7 @@ public class DataNodeInternalRPCServiceImpl implements
IDataNodeRPCService.Iface
public TSStatus deleteRegion(TConsensusGroupId tconsensusGroupId) {
ConsensusGroupId consensusGroupId =
ConsensusGroupId.Factory.createFromTConsensusGroupId(tconsensusGroupId);
+ boolean consensusGroupDeleted = true;
if (consensusGroupId instanceof DataRegionId) {
try {
DataRegionConsensusImpl.getInstance().deleteLocalPeer(consensusGroupId);
@@ -2106,8 +2107,10 @@ public class DataNodeInternalRPCServiceImpl implements
IDataNodeRPCService.Iface
if (!(e instanceof ConsensusGroupNotExistException)) {
return RpcUtils.getStatus(TSStatusCode.DELETE_REGION_ERROR,
e.getMessage());
}
+ consensusGroupDeleted = false;
}
- return regionManager.deleteDataRegion((DataRegionId) consensusGroupId);
+ return getDeleteRegionStatus(
+ regionManager.deleteDataRegion((DataRegionId) consensusGroupId),
consensusGroupDeleted);
} else {
try {
SchemaRegionConsensusImpl.getInstance().deleteLocalPeer(consensusGroupId);
@@ -2115,11 +2118,22 @@ public class DataNodeInternalRPCServiceImpl implements
IDataNodeRPCService.Iface
if (!(e instanceof ConsensusGroupNotExistException)) {
return RpcUtils.getStatus(TSStatusCode.DELETE_REGION_ERROR,
e.getMessage());
}
+ consensusGroupDeleted = false;
}
- return regionManager.deleteSchemaRegion((SchemaRegionId)
consensusGroupId);
+ return getDeleteRegionStatus(
+ regionManager.deleteSchemaRegion((SchemaRegionId) consensusGroupId),
+ consensusGroupDeleted);
}
}
+ static TSStatus getDeleteRegionStatus(TSStatus localRegionStatus, boolean
consensusGroupDeleted) {
+ if (consensusGroupDeleted
+ && localRegionStatus.getCode() ==
TSStatusCode.REGION_NOT_EXIST.getStatusCode()) {
+ return RpcUtils.getStatus(TSStatusCode.SUCCESS_STATUS);
+ }
+ return localRegionStatus;
+ }
+
@Override
public TRegionLeaderChangeResp changeRegionLeader(TRegionLeaderChangeReq
req) {
LOGGER.info("[ChangeRegionLeader] {}", req);
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeRegionManager.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeRegionManager.java
index 7d4631a0742..b0b77f6736c 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeRegionManager.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeRegionManager.java
@@ -136,7 +136,7 @@ public class DataNodeRegionManager {
tsStatus.setMessage(
String.format("Create Schema Region failed because of %s",
e2.getMessage()));
} catch (ConsensusGroupAlreadyExistException e) {
- tsStatus = new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
+ tsStatus = new
TSStatus(TSStatusCode.REGION_ALREADY_EXISTS.getStatusCode());
tsStatus.setMessage(String.format("SchemaRegion %d already exists.",
schemaRegionId.getId()));
} catch (ConsensusException e) {
tsStatus = new
TSStatus(TSStatusCode.CREATE_REGION_ERROR.getStatusCode());
@@ -166,7 +166,7 @@ public class DataNodeRegionManager {
tsStatus = new
TSStatus(TSStatusCode.CREATE_REGION_ERROR.getStatusCode());
tsStatus.setMessage(String.format("Create Data Region failed because of
%s", e.getMessage()));
} catch (ConsensusGroupAlreadyExistException e) {
- tsStatus = new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
+ tsStatus = new
TSStatus(TSStatusCode.REGION_ALREADY_EXISTS.getStatusCode());
tsStatus.setMessage(String.format("DataRegion %d already exists.",
dataRegionId.getId()));
} catch (ConsensusException e) {
tsStatus = new
TSStatus(TSStatusCode.CREATE_REGION_ERROR.getStatusCode());
@@ -200,17 +200,21 @@ public class DataNodeRegionManager {
}
public TSStatus deleteDataRegion(DataRegionId dataRegionId) {
- storageEngine.deleteDataRegion(dataRegionId);
- dataRegionLockMap.remove(dataRegionId);
- return RpcUtils.getStatus(TSStatusCode.SUCCESS_STATUS, "Execute
successfully");
+ TSStatus status = storageEngine.deleteDataRegion(dataRegionId);
+ if (status.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ dataRegionLockMap.remove(dataRegionId);
+ }
+ return status;
}
public TSStatus deleteSchemaRegion(SchemaRegionId schemaRegionId) {
try {
- schemaEngine.deleteSchemaRegion(schemaRegionId);
+ if (!schemaEngine.deleteSchemaRegion(schemaRegionId)) {
+ return RpcUtils.getStatus(TSStatusCode.REGION_NOT_EXIST);
+ }
PipeDataNodeAgent.runtime().schemaListener(schemaRegionId).close();
schemaRegionLockMap.remove(schemaRegionId);
- return RpcUtils.getStatus(TSStatusCode.SUCCESS_STATUS, "Execute
successfully");
+ return RpcUtils.getStatus(TSStatusCode.SUCCESS_STATUS);
} catch (MetadataException e) {
LOGGER.error("{}: MetaData error: ", IoTDBConstant.GLOBAL_DB_NAME, e);
return RpcUtils.getStatus(TSStatusCode.METADATA_ERROR, e.getMessage());
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/SchemaEngine.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/SchemaEngine.java
index 1e185e40cda..c8faf2f059a 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/SchemaEngine.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/SchemaEngine.java
@@ -331,12 +331,12 @@ public class SchemaEngine {
return schemaRegion;
}
- public synchronized void deleteSchemaRegion(SchemaRegionId schemaRegionId)
+ public synchronized boolean deleteSchemaRegion(SchemaRegionId schemaRegionId)
throws MetadataException {
ISchemaRegion schemaRegion = schemaRegionMap.get(schemaRegionId);
if (schemaRegion == null) {
logger.warn("SchemaRegion(id = {}) has been deleted, skiped",
schemaRegionId);
- return;
+ return false;
}
schemaRegion.deleteSchemaRegion();
schemaMetricManager.removeSchemaRegionMetric(schemaRegionId.getId());
@@ -360,6 +360,7 @@ public class SchemaEngine {
FileUtils.deleteFileOrDirectory(sgDir);
}
}
+ return true;
}
public int getSchemaRegionNumber() {
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/RegionMigrateService.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/RegionMigrateService.java
index bc4b61e763a..8a097445a0d 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/RegionMigrateService.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/RegionMigrateService.java
@@ -502,18 +502,20 @@ public class RegionMigrateService implements IService {
originalDataNode,
TRegionMigrateFailedType.RemoveConsensusGroupFailed,
runResult);
+ return;
}
// deleteRegion: delete region data
runResult = deleteRegion();
- if (isFailed(runResult)) {
+ if (!isDeleteRegionCompleted(runResult)) {
taskFail(
taskId,
tRegionId,
originalDataNode,
TRegionMigrateFailedType.DeleteRegionFailed,
runResult);
+ return;
}
taskSucceed(taskId, tRegionId, "DeletePeer");
@@ -533,6 +535,8 @@ public class RegionMigrateService implements IService {
} else {
SchemaRegionConsensusImpl.getInstance().deleteLocalPeer(regionId);
}
+ } catch (ConsensusGroupNotExistException e) {
+ // The peer was already removed by an earlier attempt, so continue
with local cleanup.
} catch (ConsensusException e) {
String errorMsg =
String.format(
@@ -560,23 +564,10 @@ public class RegionMigrateService implements IService {
REGION_MIGRATE_PROCESS,
tRegionId,
originalDataNode);
- TSStatus status = new
TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
ConsensusGroupId regionId =
ConsensusGroupId.Factory.createFromTConsensusGroupId(tRegionId);
- try {
- if (regionId instanceof DataRegionId) {
- DataNodeRegionManager.getInstance().deleteDataRegion((DataRegionId)
regionId);
- } else {
-
DataNodeRegionManager.getInstance().deleteSchemaRegion((SchemaRegionId)
regionId);
- }
- } catch (Exception e) {
- taskLogger.error("{}, deleteRegion {} error", REGION_MIGRATE_PROCESS,
regionId, e);
- status.setCode(TSStatusCode.DELETE_REGION_ERROR.getStatusCode());
- status.setMessage("deleteRegion " + regionId + " error, " +
e.getMessage());
- return status;
- }
- status.setMessage("deleteRegion " + regionId + " succeed");
- taskLogger.info("{}, Succeed to deleteRegion {}",
REGION_MIGRATE_PROCESS, regionId);
- return status;
+ return regionId instanceof DataRegionId
+ ?
DataNodeRegionManager.getInstance().deleteDataRegion((DataRegionId) regionId)
+ :
DataNodeRegionManager.getInstance().deleteSchemaRegion((SchemaRegionId)
regionId);
}
}
@@ -630,6 +621,10 @@ public class RegionMigrateService implements IService {
return !isSucceed(status);
}
+ static boolean isDeleteRegionCompleted(TSStatus status) {
+ return isSucceed(status) || status.getCode() ==
TSStatusCode.REGION_NOT_EXIST.getStatusCode();
+ }
+
private static TEndPoint getConsensusEndPoint(
TDataNodeLocation nodeLocation, ConsensusGroupId regionId) {
if (regionId instanceof DataRegionId) {
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/StorageEngine.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/StorageEngine.java
index 10a3aec8395..67d4abf5907 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/StorageEngine.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/StorageEngine.java
@@ -767,58 +767,66 @@ public class StorageEngine implements IService {
}
}
- public void deleteDataRegion(DataRegionId regionId) {
- if (!dataRegionMap.containsKey(regionId) ||
deletingDataRegionMap.containsKey(regionId)) {
- return;
+ public TSStatus deleteDataRegion(DataRegionId regionId) {
+ DataRegion region = dataRegionMap.get(regionId);
+ if (region == null) {
+ return RpcUtils.getStatus(
+ deletingDataRegionMap.containsKey(regionId)
+ ? TSStatusCode.DELETE_REGION_ERROR
+ : TSStatusCode.REGION_NOT_EXIST);
}
- DataRegion region =
- deletingDataRegionMap.computeIfAbsent(regionId, k ->
dataRegionMap.remove(regionId));
- if (region != null) {
- region.markDeleted();
- try {
- region.abortCompaction();
- region.syncDeleteDataFiles();
- region.deleteFolder(systemDir);
- if
(CONFIG.getDataRegionConsensusProtocolClass().equals(ConsensusFactory.IOT_CONSENSUS)
- || CONFIG
- .getDataRegionConsensusProtocolClass()
- .equals(ConsensusFactory.FAST_IOT_CONSENSUS)
- || CONFIG
- .getDataRegionConsensusProtocolClass()
- .equals(ConsensusFactory.IOT_CONSENSUS_V2)) {
- // delete wal
- WALManager.getInstance()
- .deleteWALNode(
- region.getDatabaseName() + FILE_NAME_SEPARATOR +
region.getDataRegionIdString());
- // delete snapshot
- for (String dataDir : CONFIG.getLocalDataDirs()) {
- File regionSnapshotDir =
- new File(
- dataDir + File.separator +
IoTDBConstant.SNAPSHOT_FOLDER_NAME,
- region.getDatabaseName() + FILE_NAME_SEPARATOR +
regionId.getId());
- if (regionSnapshotDir.exists()) {
- try {
- FileUtils.deleteDirectory(regionSnapshotDir);
- } catch (IOException e) {
- LOGGER.error("Failed to delete snapshot dir {}",
regionSnapshotDir, e);
- }
- }
+ if (deletingDataRegionMap.putIfAbsent(regionId, region) != null) {
+ return RpcUtils.getStatus(TSStatusCode.DELETE_REGION_ERROR);
+ }
+ if (!dataRegionMap.remove(regionId, region)) {
+ deletingDataRegionMap.remove(regionId, region);
+ return RpcUtils.getStatus(TSStatusCode.DELETE_REGION_ERROR);
+ }
+ try {
+ if (!region.isDeleted()) {
+ region.markDeleted();
+ }
+ region.abortCompaction();
+ region.syncDeleteDataFiles();
+ region.deleteFolder(systemDir);
+ if
(CONFIG.getDataRegionConsensusProtocolClass().equals(ConsensusFactory.IOT_CONSENSUS)
+ || CONFIG
+ .getDataRegionConsensusProtocolClass()
+ .equals(ConsensusFactory.FAST_IOT_CONSENSUS)
+ || CONFIG
+ .getDataRegionConsensusProtocolClass()
+ .equals(ConsensusFactory.IOT_CONSENSUS_V2)) {
+ // delete wal
+ WALManager.getInstance()
+ .deleteWALNode(
+ region.getDatabaseName() + FILE_NAME_SEPARATOR +
region.getDataRegionIdString());
+ // delete snapshot
+ for (String dataDir : CONFIG.getLocalDataDirs()) {
+ File regionSnapshotDir =
+ new File(
+ dataDir + File.separator +
IoTDBConstant.SNAPSHOT_FOLDER_NAME,
+ region.getDatabaseName() + FILE_NAME_SEPARATOR +
regionId.getId());
+ if (regionSnapshotDir.exists()) {
+ FileUtils.deleteDirectory(regionSnapshotDir);
}
}
- WRITING_METRICS.removeDataRegionMemoryCostMetrics(regionId);
- WRITING_METRICS.removeFlushingMemTableStatusMetrics(regionId);
- WRITING_METRICS.removeActiveMemtableCounterMetrics(regionId);
- FileMetrics.getInstance()
- .deleteRegion(region.getDatabaseName(),
region.getDataRegionIdString());
- } catch (Exception e) {
- LOGGER.error(
- "Error occurs when deleting data region {}-{}",
- region.getDatabaseName(),
- region.getDataRegionIdString(),
- e);
- } finally {
- deletingDataRegionMap.remove(regionId);
}
+ WRITING_METRICS.removeDataRegionMemoryCostMetrics(regionId);
+ WRITING_METRICS.removeFlushingMemTableStatusMetrics(regionId);
+ WRITING_METRICS.removeActiveMemtableCounterMetrics(regionId);
+ FileMetrics.getInstance()
+ .deleteRegion(region.getDatabaseName(),
region.getDataRegionIdString());
+ return RpcUtils.getStatus(TSStatusCode.SUCCESS_STATUS);
+ } catch (Exception e) {
+ dataRegionMap.putIfAbsent(regionId, region);
+ LOGGER.error(
+ "Error occurs when deleting data region {}-{}",
+ region.getDatabaseName(),
+ region.getDataRegionIdString(),
+ e);
+ return RpcUtils.getStatus(TSStatusCode.DELETE_REGION_ERROR,
e.getMessage());
+ } finally {
+ deletingDataRegionMap.remove(regionId, region);
}
}
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeInternalRPCServiceImplStatusTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeInternalRPCServiceImplStatusTest.java
new file mode 100644
index 00000000000..6717a3a3216
--- /dev/null
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeInternalRPCServiceImplStatusTest.java
@@ -0,0 +1,65 @@
+/*
+ * 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.db.protocol.thrift.impl;
+
+import org.apache.iotdb.common.rpc.thrift.TSStatus;
+import org.apache.iotdb.db.conf.IoTDBDescriptor;
+import org.apache.iotdb.rpc.TSStatusCode;
+
+import org.junit.Assert;
+import org.junit.BeforeClass;
+import org.junit.Test;
+
+public class DataNodeInternalRPCServiceImplStatusTest {
+
+ @BeforeClass
+ public static void setUpDataNodeId() {
+ IoTDBDescriptor.getInstance().getConfig().setDataNodeId(0);
+ }
+
+ @Test
+ public void testDeleteRegionStatusCombinesConsensusAndLocalResults() {
+ Assert.assertEquals(
+ TSStatusCode.SUCCESS_STATUS.getStatusCode(),
+ DataNodeInternalRPCServiceImpl.getDeleteRegionStatus(
+ new TSStatus(TSStatusCode.REGION_NOT_EXIST.getStatusCode()),
true)
+ .getCode());
+ Assert.assertEquals(
+ TSStatusCode.REGION_NOT_EXIST.getStatusCode(),
+ DataNodeInternalRPCServiceImpl.getDeleteRegionStatus(
+ new TSStatus(TSStatusCode.REGION_NOT_EXIST.getStatusCode()),
false)
+ .getCode());
+ Assert.assertEquals(
+ TSStatusCode.SUCCESS_STATUS.getStatusCode(),
+ DataNodeInternalRPCServiceImpl.getDeleteRegionStatus(
+ new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode()),
false)
+ .getCode());
+ Assert.assertEquals(
+ TSStatusCode.DELETE_REGION_ERROR.getStatusCode(),
+ DataNodeInternalRPCServiceImpl.getDeleteRegionStatus(
+ new
TSStatus(TSStatusCode.DELETE_REGION_ERROR.getStatusCode()), true)
+ .getCode());
+ Assert.assertEquals(
+ TSStatusCode.DELETE_REGION_ERROR.getStatusCode(),
+ DataNodeInternalRPCServiceImpl.getDeleteRegionStatus(
+ new
TSStatus(TSStatusCode.DELETE_REGION_ERROR.getStatusCode()), false)
+ .getCode());
+ }
+}
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/service/DataNodeInternalRPCServiceImplTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/service/DataNodeInternalRPCServiceImplTest.java
index c8ac2c8880e..74eb6cec987 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/service/DataNodeInternalRPCServiceImplTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/service/DataNodeInternalRPCServiceImplTest.java
@@ -52,10 +52,12 @@ import org.apache.iotdb.db.schemaengine.SchemaEngine;
import org.apache.iotdb.db.storageengine.dataregion.DataRegion;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
import org.apache.iotdb.db.utils.EnvironmentUtils;
+import org.apache.iotdb.mpp.rpc.thrift.TCreateSchemaRegionReq;
import org.apache.iotdb.mpp.rpc.thrift.TPlanNode;
import org.apache.iotdb.mpp.rpc.thrift.TSendBatchPlanNodeReq;
import org.apache.iotdb.mpp.rpc.thrift.TSendBatchPlanNodeResp;
import org.apache.iotdb.mpp.rpc.thrift.TSendSinglePlanNodeReq;
+import org.apache.iotdb.rpc.TSStatusCode;
import org.apache.ratis.util.FileUtils;
import org.apache.tsfile.enums.TSDataType;
@@ -382,6 +384,30 @@ public class DataNodeInternalRPCServiceImplTest {
Assert.assertTrue(response.getResponses().get(0).accepted);
}
+ @Test
+ public void testRegionOperationRetryReturnsAlreadyCompletedStatus() {
+ TRegionReplicaSet regionReplicaSet = genRegionReplicaSet();
+ regionReplicaSet.setRegionId(new
TConsensusGroupId(TConsensusGroupType.SchemaRegion, 2));
+ TCreateSchemaRegionReq createReq =
+ new TCreateSchemaRegionReq()
+ .setRegionReplicaSet(regionReplicaSet)
+ .setStorageGroup("root.retry_test");
+
+ Assert.assertEquals(
+ TSStatusCode.SUCCESS_STATUS.getStatusCode(),
+
dataNodeInternalRPCServiceImpl.createSchemaRegion(createReq).getCode());
+ Assert.assertEquals(
+ TSStatusCode.REGION_ALREADY_EXISTS.getStatusCode(),
+
dataNodeInternalRPCServiceImpl.createSchemaRegion(createReq).getCode());
+
+ Assert.assertEquals(
+ TSStatusCode.SUCCESS_STATUS.getStatusCode(),
+
dataNodeInternalRPCServiceImpl.deleteRegion(regionReplicaSet.getRegionId()).getCode());
+ Assert.assertEquals(
+ TSStatusCode.REGION_NOT_EXIST.getStatusCode(),
+
dataNodeInternalRPCServiceImpl.deleteRegion(regionReplicaSet.getRegionId()).getCode());
+ }
+
private TRegionReplicaSet genRegionReplicaSet() {
List<TDataNodeLocation> dataNodeList = new ArrayList<>();
dataNodeList.add(
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/service/RegionMigrateServiceStatusTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/service/RegionMigrateServiceStatusTest.java
new file mode 100644
index 00000000000..c85caf56dd4
--- /dev/null
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/service/RegionMigrateServiceStatusTest.java
@@ -0,0 +1,44 @@
+/*
+ * 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.db.service;
+
+import org.apache.iotdb.common.rpc.thrift.TSStatus;
+import org.apache.iotdb.rpc.TSStatusCode;
+
+import org.junit.Test;
+
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertTrue;
+
+public class RegionMigrateServiceStatusTest {
+
+ @Test
+ public void testIdempotentDeleteCompletionStatus() {
+ assertTrue(
+ RegionMigrateService.isDeleteRegionCompleted(
+ new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode())));
+ assertTrue(
+ RegionMigrateService.isDeleteRegionCompleted(
+ new TSStatus(TSStatusCode.REGION_NOT_EXIST.getStatusCode())));
+ assertFalse(
+ RegionMigrateService.isDeleteRegionCompleted(
+ new TSStatus(TSStatusCode.DELETE_REGION_ERROR.getStatusCode())));
+ }
+}