From ba923b7dc6bc59bbf07d8c92b0bca3028e785d34 Mon Sep 17 00:00:00 2001 From: Yongzao <532741407@qq.com> Date: Thu, 13 Aug 2026 18:08:47 +0800 Subject: [PATCH] [TDB-280][to dev/1.3] Harden RegionMaintainer retries and region operation atomicity --- .../org/apache/iotdb/rpc/TSStatusCode.java | 2 + .../rpc/DataNodeTSStatusRPCHandler.java | 14 +- .../manager/partition/PartitionManager.java | 366 +++++++++--------- .../impl/schema/DeleteDatabaseProcedure.java | 111 ++++-- .../rpc/DataNodeTSStatusRPCHandlerTest.java | 63 +++ .../PartitionManagerRegionMaintainTest.java | 151 ++++++++ .../persistence/PartitionInfoTest.java | 19 + .../schema/DeleteDatabaseProcedureTest.java | 118 ++++++ .../impl/DataNodeInternalRPCServiceImpl.java | 18 +- .../thrift/impl/DataNodeRegionManager.java | 18 +- .../iotdb/db/schemaengine/SchemaEngine.java | 5 +- .../db/service/RegionMigrateService.java | 29 +- .../iotdb/db/storageengine/StorageEngine.java | 104 ++--- ...aNodeInternalRPCServiceImplStatusTest.java | 65 ++++ .../DataNodeInternalRPCServiceImplTest.java | 26 ++ .../RegionMigrateServiceStatusTest.java | 44 +++ 16 files changed, 855 insertions(+), 298 deletions(-) create mode 100644 iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/client/async/handlers/rpc/DataNodeTSStatusRPCHandlerTest.java create mode 100644 iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/partition/PartitionManagerRegionMaintainTest.java create mode 100644 iotdb-core/datanode/src/test/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeInternalRPCServiceImplStatusTest.java create mode 100644 iotdb-core/datanode/src/test/java/org/apache/iotdb/db/service/RegionMigrateServiceStatusTest.java 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 af4f4727c7a2..49e75c6da904 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 a44e3781c251..848dd732212f 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 void onComplete(TSStatus response) { // 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 void onComplete(TSStatus response) { 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 7db7ac50d625..e3a50cda8d8e 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.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 void maintainRegionReplicas() { return; } - // Group tasks by region id - Map> 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 selectedRegionMaintainTask = new ArrayList<>(); - RegionMaintainType currentType = null; - for (Map.Entry> 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> tasksByRegion = + groupRegionMaintainTasks(regionMaintainTaskList); + Set deferredRegions = new HashSet<>(); + while (!tasksByRegion.isEmpty()) { + Map> headsByType = + getRegionMaintainTaskHeads(tasksByRegion, deferredRegions); + Set completedRegions = new HashSet<>(); + for (Map.Entry> entry : + headsByType.entrySet()) { + completedRegions.addAll( + submitRegionMaintainTasks(entry.getKey(), entry.getValue())); } - Set successfulTask = new HashSet<>(); - switch (currentType) { - case CREATE: - // create region - switch (selectedRegionMaintainTask.get(0).getRegionId().getType()) { - case SchemaRegion: - // create SchemaRegion - DataNodeAsyncRequestContext - 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 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 - 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 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 deleteRegionHandler = - new DataNodeAsyncRequestContext<>(CnToDnAsyncRequestType.DELETE_REGION); - Map 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 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> groupRegionMaintainTasks( + List tasks) { + Map> tasksByRegion = new LinkedHashMap<>(); + for (RegionMaintainTask task : tasks) { + tasksByRegion.computeIfAbsent(task.getRegionId(), key -> new ArrayDeque<>()).add(task); + } + return tasksByRegion; + } + + static Map> getRegionMaintainTaskHeads( + Map> tasksByRegion) { + return getRegionMaintainTaskHeads(tasksByRegion, Collections.emptySet()); + } + + static Map> getRegionMaintainTaskHeads( + Map> tasksByRegion, + Set deferredRegions) { + Map> headsByType = + new EnumMap<>(RegionMaintainType.class); + for (Map.Entry> entry : tasksByRegion.entrySet()) { + if (deferredRegions.contains(entry.getKey())) { + continue; + } + Queue 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> tasksByRegion, + Set completedRegions) { + for (TConsensusGroupId regionId : completedRegions) { + tasksByRegion.computeIfPresent( + regionId, + (key, queue) -> { + queue.poll(); + return queue.isEmpty() ? null : queue; + }); + } + } + + static void deferFailedRegionMaintainTasks( + Set deferredRegions, + Map> submittedTaskHeads, + Set completedRegions) { + submittedTaskHeads.values().stream() + .flatMap(List::stream) + .map(RegionMaintainTask::getRegionId) + .filter(regionId -> !completedRegions.contains(regionId)) + .forEach(deferredRegions::add); + } + + private Set submitRegionMaintainTasks( + RegionMaintainType taskType, List tasks) { + return taskType == RegionMaintainType.CREATE + ? submitRegionCreateTasks(tasks) + : submitRegionDeleteTasks(tasks); + } + + private Set submitRegionCreateTasks(List tasks) { + Map> tasksByRegionType = + new EnumMap<>(TConsensusGroupType.class); + for (RegionMaintainTask task : tasks) { + RegionCreateTask createTask = (RegionCreateTask) task; + tasksByRegionType + .computeIfAbsent(createTask.getRegionId().getType(), key -> new ArrayList<>()) + .add(createTask); + } + + Set completedRegions = new HashSet<>(); + for (Map.Entry> entry : + tasksByRegionType.entrySet()) { + completedRegions.addAll(submitRegionCreateTasks(entry.getKey(), entry.getValue())); + } + return completedRegions; + } + + private Set submitRegionCreateTasks( + TConsensusGroupType regionType, List tasks) { + DataNodeAsyncRequestContext requestContext = + new DataNodeAsyncRequestContext<>( + regionType == TConsensusGroupType.SchemaRegion + ? CnToDnAsyncRequestType.CREATE_SCHEMA_REGION + : CnToDnAsyncRequestType.CREATE_DATA_REGION); + Map 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 submitRegionDeleteTasks(List tasks) { + DataNodeAsyncRequestContext requestContext = + new DataNodeAsyncRequestContext<>(CnToDnAsyncRequestType.DELETE_REGION); + Map 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 collectCompletedRegionMaintainTasks( + RegionMaintainType taskType, + Map responseMap, + Map regionByRequestIndex) { + Set 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 b6ad21128af3..f5d63aa1fc06 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 @@ protected Flow executeFromState( deleteDatabaseSchema.getName()); // Submit RegionDeleteTasks - OfferRegionMaintainTasksPlan dataRegionDeleteTaskOfferPlan = - new OfferRegionMaintainTasksPlan(); List regionReplicaSets = env.getAllReplicaSets(deleteDatabaseSchema.getName()); - List schemaRegionReplicaSets = new ArrayList<>(); + OfferRegionMaintainTasksPlan regionDeleteTaskOfferPlan = + buildDataRegionDeleteTaskOfferPlan(regionReplicaSets); + List 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 asyncClientHandler = @@ -166,7 +145,7 @@ protected Flow executeFromState( .sendAsyncRequestWithRetry(asyncClientHandler); for (Map.Entry 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 @@ protected Flow executeFromState( } 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 @@ protected Flow executeFromState( } 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 @@ protected Flow executeFromState( 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 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 getSchemaRegionReplicaSets( + List regionReplicaSets) { + List schemaRegionReplicaSets = new ArrayList<>(); + regionReplicaSets.stream() + .filter( + regionReplicaSet -> + regionReplicaSet.getRegionId().getType() == TConsensusGroupType.SchemaRegion) + .forEach(schemaRegionReplicaSets::add); + return schemaRegionReplicaSets; + } + + static void appendFailedSchemaRegionDeleteTasks( + OfferRegionMaintainTasksPlan offerPlan, + Map 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 000000000000..7f06a44edffe --- /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 000000000000..4157f98495c1 --- /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> tasksByRegion = + PartitionManager.groupRegionMaintainTasks( + Arrays.asList(region0Head, region0Next, region1Head)); + + Map> firstRound = + PartitionManager.getRegionMaintainTaskHeads(tasksByRegion); + assertEquals( + Arrays.asList(region0Head, region1Head), firstRound.get(RegionMaintainType.DELETE)); + assertFalse(firstRound.containsKey(RegionMaintainType.CREATE)); + + Set 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 regionsByRequestIndex = new HashMap<>(); + regionsByRequestIndex.put(0, schemaRegion); + regionsByRequestIndex.put(1, dataRegion); + Map responses = new HashMap<>(); + responses.put(0, status(TSStatusCode.SUCCESS_STATUS)); + responses.put(1, status(TSStatusCode.CREATE_REGION_ERROR)); + + Set 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 c15cefff6579..f6f209c93564 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.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 void testGetRegionType() { 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 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 b12f49d9bd7d..2f6d365bc7a5 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 void serializeDeserializeTest() { 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 replicaSets = + Arrays.asList( + new TRegionReplicaSet(dataRegionId, Arrays.asList(dataNode0, dataNode1)), + new TRegionReplicaSet(schemaRegionId, Collections.singletonList(dataNode0))); + + OfferRegionMaintainTasksPlan offerPlan = + DeleteDatabaseProcedure.buildDataRegionDeleteTaskOfferPlan(replicaSets); + List 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 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 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 99e314affa04..9edb75c641cd 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 TSStatus updateTemplate(final TUpdateTemplateReq req) { 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 TSStatus deleteRegion(TConsensusGroupId tconsensusGroupId) { 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 TSStatus deleteRegion(TConsensusGroupId tconsensusGroupId) { 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 7d4631a0742a..b0b77f6736cb 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 TSStatus createSchemaRegion(TRegionReplicaSet regionReplicaSet, String st 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 TSStatus createDataRegion(TRegionReplicaSet regionReplicaSet, String stor 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 TSStatus createNewRegion(ConsensusGroupId regionId, String storageGroup) } 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 1e185e40cda1..c8faf2f059a4 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 @@ private ISchemaRegion createSchemaRegionWithoutExistenceCheck( 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 synchronized void deleteSchemaRegion(SchemaRegionId schemaRegionId) 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 bc4b61e763a0..8a097445a0d9 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 void run() { 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 @@ private TSStatus deletePeer() { } 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 @@ private TSStatus deleteRegion() { 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 static boolean isFailed(TSStatus status) { 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 10a3aec8395b..67d4abf59076 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 void createDataRegion(DataRegionId regionId, String databaseName) } } - 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 000000000000..6717a3a32161 --- /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 c8ac2c8880ea..74eb6cec9871 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.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 void testCreateMultiTimeSeries() throws MetadataException { 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 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 000000000000..c85caf56dd47 --- /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()))); + } +}