From 46580d9f56bc92153c4ef3172e02d4aa6805aa29 Mon Sep 17 00:00:00 2001 From: Yongzao Date: Wed, 5 Aug 2026 19:51:36 +0800 Subject: [PATCH] [ConfigNode] Revert procedure-based RegionGroup cleanup and fix RegionMaintainer retries (#18404) --- .../org/apache/iotdb/rpc/TSStatusCode.java | 2 + .../confignode/i18n/ConfigNodeMessages.java | 4 - .../confignode/i18n/ManagerMessages.java | 5 +- .../confignode/i18n/ProcedureMessages.java | 14 - .../confignode/i18n/ConfigNodeMessages.java | 6 - .../confignode/i18n/ManagerMessages.java | 5 +- .../confignode/i18n/ProcedureMessages.java | 16 - .../client/async/CnToDnAsyncRequestType.java | 1 + ...oDnInternalServiceAsyncRequestManager.java | 5 + .../confignode/manager/ProcedureManager.java | 4 + .../manager/partition/PartitionManager.java | 334 +++++++++------ .../persistence/partition/PartitionInfo.java | 29 +- .../region/CreateRegionGroupsProcedure.java | 44 +- .../region/RemoveRegionGroupProcedure.java | 389 ------------------ .../impl/schema/DeleteDatabaseProcedure.java | 113 +++-- .../state/RemoveRegionGroupState.java | 24 -- .../state/schema/DeleteDatabaseState.java | 5 +- .../procedure/store/ProcedureFactory.java | 6 - .../procedure/store/ProcedureType.java | 1 - .../PartitionManagerRegionMaintainTest.java | 70 ++++ .../persistence/PartitionInfoTest.java | 27 -- .../RemoveRegionGroupProcedureTest.java | 100 ----- .../impl/DataNodeInternalRPCServiceImpl.java | 19 +- .../thrift/impl/DataNodeRegionManager.java | 18 +- .../iotdb/db/schemaengine/SchemaEngine.java | 5 +- .../db/service/RegionMigrateService.java | 12 +- .../iotdb/db/storageengine/StorageEngine.java | 7 +- .../DataNodeInternalRPCServiceImplTest.java | 26 ++ 28 files changed, 458 insertions(+), 833 deletions(-) delete mode 100644 iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/region/RemoveRegionGroupProcedure.java delete mode 100644 iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/state/RemoveRegionGroupState.java create mode 100644 iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/partition/PartitionManagerRegionMaintainTest.java delete mode 100644 iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/region/RemoveRegionGroupProcedureTest.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 64c289bdf795a..ae988b8d1c73b 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 @@ -197,6 +197,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/i18n/en/org/apache/iotdb/confignode/i18n/ConfigNodeMessages.java b/iotdb-core/confignode/src/main/i18n/en/org/apache/iotdb/confignode/i18n/ConfigNodeMessages.java index 6b21e33317da9..3a5f279f222ce 100644 --- a/iotdb-core/confignode/src/main/i18n/en/org/apache/iotdb/confignode/i18n/ConfigNodeMessages.java +++ b/iotdb-core/confignode/src/main/i18n/en/org/apache/iotdb/confignode/i18n/ConfigNodeMessages.java @@ -634,10 +634,6 @@ private ConfigNodeMessages() {} "Unexpected interruption while waiting for ConfigNode leader load warm-up."; public static final String EXCEPTION_PROCEDURE_FILE_ARG_EXCEEDS_THE_LOAD_BUFFER_LIMIT_ARG_ACTUAL_SIZE_ARG_62375B4C = "Procedure file %s exceeds the load buffer limit %s, actual size %s"; - public static final String MESSAGE_DROPPING_LEGACY_REGION_DELETE_TASK_FOR_ARG_WHILE_REPLAYING_OFFER_PLAN_REGION_DELETION_IS_NOW_HANDLED_BY_REMOVEREGIONGROUPPROCEDURE_2A81A649 = - "Dropping legacy region-delete task for {} while replaying offer plan; region deletion is now handled by RemoveRegionGroupProcedure."; - public static final String MESSAGE_DROPPING_LEGACY_REGION_DELETE_TASK_FOR_ARG_WHILE_LOADING_SNAPSHOT_REGION_DELETION_IS_NOW_HANDLED_BY_REMOVEREGIONGROUPPROCEDURE_A9D409A0 = - "Dropping legacy region-delete task for {} while loading snapshot; region deletion is now handled by RemoveRegionGroupProcedure."; public static final String MESSAGE_CONFIGNODE_LEADER_IS_WARMING_UP_BEFORE_SERVING_THE_REGISTERING_CONFIGNODE_WILL_WAIT_2E051639 = "ConfigNode leader is warming up before serving the registering ConfigNode, will wait"; public static final String MESSAGE_CONFIGNODE_LEADER_IS_WARMING_UP_BEFORE_SERVING_THE_REGISTERING_CONFIGNODE_WILL_WAIT_AND_RETRY_STATUS_ARG_RETRY_ARG_3C924873 = diff --git a/iotdb-core/confignode/src/main/i18n/en/org/apache/iotdb/confignode/i18n/ManagerMessages.java b/iotdb-core/confignode/src/main/i18n/en/org/apache/iotdb/confignode/i18n/ManagerMessages.java index 20316edc39713..22efb8ad9e576 100644 --- a/iotdb-core/confignode/src/main/i18n/en/org/apache/iotdb/confignode/i18n/ManagerMessages.java +++ b/iotdb-core/confignode/src/main/i18n/en/org/apache/iotdb/confignode/i18n/ManagerMessages.java @@ -73,6 +73,7 @@ public final class ManagerMessages { "DataRegionGroupExtensionPolicy %s doesn't exist."; public static final String DECREASE_REFERENCE_COUNT_FOR_SNAPSHOT_ERROR = "Decrease reference count for snapshot {} error."; + public static final String DELETING_REGIONS_COSTS_MS = "Deleting regions costs {}ms"; public static final String DETECTED_HISTORICAL_PIPE_COMPLETION_REPORT_FROM_DATANODE = "Detected historical pipe completion report from DataNode {} for pipe {}. remainingEventCount: {}, remainingTime: {}, completedDataNodes: {}"; public static final String DETECTED_COMPLETION_OF_PIPE_STATIC_META_REMOVE_IT = @@ -440,6 +441,8 @@ public final class ManagerMessages { "Start to create Region: {} on DataNode: {}"; public static final String START_TO_CREATE_UDF_ON_DATA_NODES_NEEDTOSAVEJAR = "Start to create UDF [{}] on Data Nodes, needToSaveJar[{}]"; + public static final String START_TO_DELETE_REGION_ON_DATANODE = + "Start to delete Region: {} on DataNode: {}"; public static final String START_TRANSFER_OF = "Start transfer of {}"; public static final String STOP_SUBMITTING_CQ_BECAUSE = "Stop submitting CQ {} because {}"; public static final String STOP_SUBMITTING_CQ_BECAUSE_CURRENT_NODE_IS_NOT_LEADER_OR = @@ -522,8 +525,6 @@ public final class ManagerMessages { "Unexpected interruption during waiting for configNode leader ready."; public static final String UNEXPECTED_INTERRUPTION_DURING_WAITING_FOR_GET_CLUSTER_ID = "Unexpected interruption during waiting for get cluster id."; - public static final String UNEXPECTED_NON_CREATE_REGION_MAINTAIN_TASK_SKIPPED = - "Unexpected non-create task in the RegionMaintainer queue; skipping it (the queue only recreates region replicas now, and region deletion is handled by RemoveRegionGroupProcedure)."; public static final String UNEXPECTED_NULL_PROCEDURE_PARAMETERS_FOR_WAITINGPROCEDUREFINISHED = "Unexpected null procedure parameters for waitingProcedureFinished"; public static final String UNKNOWN_DATAPARTITION_ALLOCATION_STRATEGY_USING_INHERIT_STRATEGY_BY_DEFAULT = diff --git a/iotdb-core/confignode/src/main/i18n/en/org/apache/iotdb/confignode/i18n/ProcedureMessages.java b/iotdb-core/confignode/src/main/i18n/en/org/apache/iotdb/confignode/i18n/ProcedureMessages.java index f531d9ee8cc9d..44fad775c3511 100644 --- a/iotdb-core/confignode/src/main/i18n/en/org/apache/iotdb/confignode/i18n/ProcedureMessages.java +++ b/iotdb-core/confignode/src/main/i18n/en/org/apache/iotdb/confignode/i18n/ProcedureMessages.java @@ -632,16 +632,6 @@ public final class ProcedureMessages { public static final String PID_ADDREGION_STATE_FAILED = "[pid{}][AddRegion] state {} failed"; public static final String PID_ADDREGION_SUCCESS_HAS_BEEN_ADDED_TO_DATANODE_PROCEDURE_TOOK = "[pid{}][AddRegion] success, {} has been added to DataNode {}. Procedure took {} (start at {})."; - public static final String PID_REMOVEREGIONGROUP_STARTED_WILL_BE_DELETED = - "[pid{}][RemoveRegionGroup] started, region group {} will be deleted from DataNodes {}."; - public static final String PID_REMOVEREGIONGROUP_STARTED_REPLICA_WILL_BE_DELETED_FROM_DATANODE = - "[pid{}][RemoveRegionGroup] region {} will be deleted from DataNode {}."; - public static final String PID_REMOVEREGIONGROUP_STATE_FAILED = - "[pid{}][RemoveRegionGroup] state {} failed"; - public static final String PID_REMOVEREGIONGROUP_DELETE_REPLICA_FAILED = - "[pid{}][RemoveRegionGroup] failed to delete a replica of region {} (attempt {}), will keep retrying until it is deleted. reason: {}"; - public static final String PID_REMOVEREGIONGROUP_SUCCESS_PROCEDURE_TOOK = - "[pid{}][RemoveRegionGroup] success, region group {} has been deleted. Procedure took {} (started at {})."; public static final String PID_MIGRATEREGION_STARTED_WILL_BE_MIGRATED_FROM_DATANODE_TO = "[pid{}][MigrateRegion] started, {} will be migrated from DataNode {} to {}."; public static final String PID_MIGRATEREGION_STATE_COMPLETE = @@ -1399,10 +1389,6 @@ private ProcedureMessages() {} public static final String LOG_ARG_8393DD4A = "{}"; public static final String MESSAGE_HALT_PID_ARG_ACTIVECOUNT_ARG_411F3EBF = "Halt pid={}, activeCount={}"; public static final String MESSAGE_EXCEPTION_HAPPENED_WHEN_WORKER_ARG_EXECUTE_PROCEDURE_ARG_6E3AD27D = "Exception happened when worker {} execute procedure {}"; - public static final String EXCEPTION_CANNOT_DERIVE_A_COLLISION_FREE_DELETE_TASKID_PROCID_ARG_DELETETASKSEQ_ARG_EXCEED_THE_71B7046A = - "cannot derive a collision-free delete taskId: procId=%d, deleteTaskSeq=%d exceed the "; - public static final String EXCEPTION_CANNOT_DERIVE_A_COLLISION_FREE_DELETE_TASKID_PROCID_ARG_DELETETASKSEQ_ARG_EXCEED_THE_ARG_ARG_BIT_BUDGET_015C598D = - "cannot derive a collision-free delete taskId: procId=%d, deleteTaskSeq=%d exceed the %d/%d-bit budget"; public static final String MESSAGE_FAILED_TO_SHOW_DATAPARTITIONTABLE_INTEGRITY_CHECK_PROGRESS_5EE98694 = "Failed to show DataPartitionTable integrity check progress"; public static final String MESSAGE_ENCOUNTERED_UNEXPECTED_DATAPARTITIONTABLEINTEGRITYCHECKPROCEDURESTATE_ARG_WHEN_SHOWING_PROGRESS_5FA2739F = "Encountered unexpected DataPartitionTableIntegrityCheckProcedureState {} when showing progress"; diff --git a/iotdb-core/confignode/src/main/i18n/zh/org/apache/iotdb/confignode/i18n/ConfigNodeMessages.java b/iotdb-core/confignode/src/main/i18n/zh/org/apache/iotdb/confignode/i18n/ConfigNodeMessages.java index afd1b24329b66..2956f59e67481 100644 --- a/iotdb-core/confignode/src/main/i18n/zh/org/apache/iotdb/confignode/i18n/ConfigNodeMessages.java +++ b/iotdb-core/confignode/src/main/i18n/zh/org/apache/iotdb/confignode/i18n/ConfigNodeMessages.java @@ -671,12 +671,6 @@ private ConfigNodeMessages() {} public static final String EXCEPTION_PROCEDURE_FILE_ARG_EXCEEDS_THE_LOAD_BUFFER_LIMIT_ARG_ACTUAL_SIZE_ARG_62375B4C = "Procedure 文件 %s 超过了加载缓冲区限制 %s,实际大小为 %s"; - public static final String - MESSAGE_DROPPING_LEGACY_REGION_DELETE_TASK_FOR_ARG_WHILE_REPLAYING_OFFER_PLAN_REGION_DELETION_IS_NOW_HANDLED_BY_REMOVEREGIONGROUPPROCEDURE_2A81A649 = - "重放 offer plan 时丢弃 {} 的遗留 region-delete 任务;region 删除现已由 RemoveRegionGroupProcedure 处理。"; - public static final String - MESSAGE_DROPPING_LEGACY_REGION_DELETE_TASK_FOR_ARG_WHILE_LOADING_SNAPSHOT_REGION_DELETION_IS_NOW_HANDLED_BY_REMOVEREGIONGROUPPROCEDURE_A9D409A0 = - "加载快照时丢弃 {} 的遗留 region-delete 任务;region 删除现已由 RemoveRegionGroupProcedure 处理。"; public static final String MESSAGE_CONFIGNODE_LEADER_IS_WARMING_UP_BEFORE_SERVING_THE_REGISTERING_CONFIGNODE_WILL_WAIT_2E051639 = "ConfigNode leader 在服务注册中的 ConfigNode 前正在进行预热,将等待"; diff --git a/iotdb-core/confignode/src/main/i18n/zh/org/apache/iotdb/confignode/i18n/ManagerMessages.java b/iotdb-core/confignode/src/main/i18n/zh/org/apache/iotdb/confignode/i18n/ManagerMessages.java index 8dc9afd1e8bf0..a14c90e5d9baf 100644 --- a/iotdb-core/confignode/src/main/i18n/zh/org/apache/iotdb/confignode/i18n/ManagerMessages.java +++ b/iotdb-core/confignode/src/main/i18n/zh/org/apache/iotdb/confignode/i18n/ManagerMessages.java @@ -73,6 +73,7 @@ public final class ManagerMessages { "DataRegionGroupExtensionPolicy %s 不存在。"; public static final String DECREASE_REFERENCE_COUNT_FOR_SNAPSHOT_ERROR = "减少快照 {} 的引用计数失败。"; + public static final String DELETING_REGIONS_COSTS_MS = "删除 region 耗时 {}ms"; public static final String DETECTED_HISTORICAL_PIPE_COMPLETION_REPORT_FROM_DATANODE = "检测到来自 DataNode {} 的历史 pipe 完成上报,pipe {}。remainingEventCount: {}, remainingTime: {}, completedDataNodes: {}"; public static final String DETECTED_COMPLETION_OF_PIPE_STATIC_META_REMOVE_IT = @@ -435,6 +436,8 @@ public final class ManagerMessages { "开始在 DataNode: {} 上创建 Region:{}"; public static final String START_TO_CREATE_UDF_ON_DATA_NODES_NEEDTOSAVEJAR = "开始在 Data Nodes 上创建 UDF [{}],needToSaveJar[{}]"; + public static final String START_TO_DELETE_REGION_ON_DATANODE = + "开始在 DataNode: {} 上删除 Region:{}"; public static final String START_TRANSFER_OF = "开始传输 {}"; public static final String STOP_SUBMITTING_CQ_BECAUSE = "停止提交 CQ {},原因:{}"; public static final String STOP_SUBMITTING_CQ_BECAUSE_CURRENT_NODE_IS_NOT_LEADER_OR = @@ -517,8 +520,6 @@ public final class ManagerMessages { "等待 configNode leader 就绪过程中发生意外中断。"; public static final String UNEXPECTED_INTERRUPTION_DURING_WAITING_FOR_GET_CLUSTER_ID = "等待获取 cluster id 过程中发生意外中断。"; - public static final String UNEXPECTED_NON_CREATE_REGION_MAINTAIN_TASK_SKIPPED = - "RegionMaintainer 队列中出现意外的非 create 任务;跳过处理(该队列目前仅用于重建 region 副本,region 删除由 RemoveRegionGroupProcedure 处理)。"; public static final String UNEXPECTED_NULL_PROCEDURE_PARAMETERS_FOR_WAITINGPROCEDUREFINISHED = "waitingProcedureFinished 的 procedure 参数为空"; public static final String UNKNOWN_DATAPARTITION_ALLOCATION_STRATEGY_USING_INHERIT_STRATEGY_BY_DEFAULT = diff --git a/iotdb-core/confignode/src/main/i18n/zh/org/apache/iotdb/confignode/i18n/ProcedureMessages.java b/iotdb-core/confignode/src/main/i18n/zh/org/apache/iotdb/confignode/i18n/ProcedureMessages.java index f7de1c3fb3ad7..eb418935ddaff 100644 --- a/iotdb-core/confignode/src/main/i18n/zh/org/apache/iotdb/confignode/i18n/ProcedureMessages.java +++ b/iotdb-core/confignode/src/main/i18n/zh/org/apache/iotdb/confignode/i18n/ProcedureMessages.java @@ -598,16 +598,6 @@ public final class ProcedureMessages { public static final String PID_ADDREGION_STATE_FAILED = "[pid{}][AddRegion] 状态 {} 失败"; public static final String PID_ADDREGION_SUCCESS_HAS_BEEN_ADDED_TO_DATANODE_PROCEDURE_TOOK = "[pid{}][AddRegion] 成功,{} 已添加到 DataNode {}。Procedure 耗时 {}(开始于 {})。"; - public static final String PID_REMOVEREGIONGROUP_STARTED_WILL_BE_DELETED = - "[pid{}][RemoveRegionGroup] 开始,region group {} 将从 DataNode {} 上删除。"; - public static final String PID_REMOVEREGIONGROUP_STARTED_REPLICA_WILL_BE_DELETED_FROM_DATANODE = - "[pid{}][RemoveRegionGroup] region {} 将从 DataNode {} 上删除。"; - public static final String PID_REMOVEREGIONGROUP_STATE_FAILED = - "[pid{}][RemoveRegionGroup] 状态 {} 失败"; - public static final String PID_REMOVEREGIONGROUP_DELETE_REPLICA_FAILED = - "[pid{}][RemoveRegionGroup] 删除 region {} 的一个副本失败(第 {} 次尝试),将持续重试直到删除成功。原因:{}"; - public static final String PID_REMOVEREGIONGROUP_SUCCESS_PROCEDURE_TOOK = - "[pid{}][RemoveRegionGroup] 成功,region group {} 已删除。过程耗时 {}(开始于 {})。"; public static final String PID_MIGRATEREGION_STARTED_WILL_BE_MIGRATED_FROM_DATANODE_TO = "[pid{}][MigrateRegion] 开始,{} 将从 DataNode {} 迁移到 {}。"; public static final String PID_MIGRATEREGION_STATE_COMPLETE = "[pid{}][MigrateRegion] 状态 {} 完成"; @@ -1466,12 +1456,6 @@ private ProcedureMessages() {} public static final String MESSAGE_EXCEPTION_HAPPENED_WHEN_WORKER_ARG_EXECUTE_PROCEDURE_ARG_6E3AD27D = "worker {} 执行 procedure {} 时发生异常"; - public static final String - EXCEPTION_CANNOT_DERIVE_A_COLLISION_FREE_DELETE_TASKID_PROCID_ARG_DELETETASKSEQ_ARG_EXCEED_THE_71B7046A = - "无法推导出无冲突的 delete taskId:procId=%d,deleteTaskSeq=%d 超出了 "; - public static final String - EXCEPTION_CANNOT_DERIVE_A_COLLISION_FREE_DELETE_TASKID_PROCID_ARG_DELETETASKSEQ_ARG_EXCEED_THE_ARG_ARG_BIT_BUDGET_015C598D = - "无法推导出无冲突的 delete taskId:procId=%d,deleteTaskSeq=%d 超出了 %d/%d 位的预算"; public static final String MESSAGE_FAILED_TO_SHOW_DATAPARTITIONTABLE_INTEGRITY_CHECK_PROGRESS_5EE98694 = "显示 DataPartitionTable 完整性检查进度失败"; diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/CnToDnAsyncRequestType.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/CnToDnAsyncRequestType.java index cc28fed0ff56e..2946884e37346 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/CnToDnAsyncRequestType.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/CnToDnAsyncRequestType.java @@ -39,6 +39,7 @@ public enum CnToDnAsyncRequestType { // Region Maintenance CREATE_DATA_REGION, CREATE_SCHEMA_REGION, + DELETE_REGION, RESET_PEER_LIST, NOTIFY_REGION_MIGRATION, UPDATE_REGION_ROUTE_MAP, diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/CnToDnInternalServiceAsyncRequestManager.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/CnToDnInternalServiceAsyncRequestManager.java index 9c90c469137de..4f66330f6ec1d 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/CnToDnInternalServiceAsyncRequestManager.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/client/async/CnToDnInternalServiceAsyncRequestManager.java @@ -19,6 +19,7 @@ package org.apache.iotdb.confignode.client.async; +import org.apache.iotdb.common.rpc.thrift.TConsensusGroupId; import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation; import org.apache.iotdb.common.rpc.thrift.TFlushReq; import org.apache.iotdb.common.rpc.thrift.TNodeLocations; @@ -142,6 +143,10 @@ protected void initActionMapBuilder() { (req, client, handler) -> client.createDataRegion( (TCreateDataRegionReq) req, (DataNodeTSStatusRPCHandler) handler)); + actionMapBuilder.put( + CnToDnAsyncRequestType.DELETE_REGION, + (req, client, handler) -> + client.deleteRegion((TConsensusGroupId) req, (DataNodeTSStatusRPCHandler) handler)); actionMapBuilder.put( CnToDnAsyncRequestType.CREATE_SCHEMA_REGION, (req, client, handler) -> diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java index 8f5b27b7e2fa4..0fe19854f7b4d 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java @@ -59,6 +59,7 @@ import org.apache.iotdb.confignode.consensus.request.write.procedure.UpdateProcedurePlan; import org.apache.iotdb.confignode.consensus.request.write.region.CreateRegionGroupsPlan; import org.apache.iotdb.confignode.i18n.ManagerMessages; +import org.apache.iotdb.confignode.manager.partition.PartitionManager; import org.apache.iotdb.confignode.persistence.ProcedureInfo; import org.apache.iotdb.confignode.procedure.PartitionTableAutoCleaner; import org.apache.iotdb.confignode.procedure.Procedure; @@ -337,6 +338,9 @@ public TSStatus deleteDatabases( } List results = new ArrayList<>(procedures.size()); procedures.forEach(procedure -> results.add(waitingProcedureFinished(procedure))); + // Clear the previously deleted regions + final PartitionManager partitionManager = getConfigManager().getPartitionManager(); + partitionManager.getRegionMaintainer().submit(partitionManager::maintainRegionReplicas); if (results.stream() .allMatch(result -> result.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode())) { return StatusUtils.OK; 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 7efdae278db3c..28655923cca92 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 @@ -84,7 +84,9 @@ import org.apache.iotdb.confignode.manager.schema.ClusterSchemaManager; 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.persistence.partition.maintainer.RegionMaintainType; import org.apache.iotdb.confignode.procedure.impl.partition.DataPartitionTableIntegrityCheckProcedure; import org.apache.iotdb.confignode.rpc.thrift.TCountTimeSlotListReq; import org.apache.iotdb.confignode.rpc.thrift.TGetRegionGroupsByTimeReq; @@ -108,7 +110,6 @@ import java.nio.ByteBuffer; import java.util.ArrayList; import java.util.Collections; -import java.util.EnumMap; import java.util.HashMap; import java.util.HashSet; import java.util.LinkedHashMap; @@ -1356,139 +1357,228 @@ public String getRegionDatabase(TConsensusGroupId regionId) { /** * Called by {@link PartitionManager#regionMaintainer}. * - *

Periodically recreate the failed RegionReplicas offered to the RegionMaintainer queue. - * Region deletion is owned by {@code RemoveRegionGroupProcedure} and is no longer handled here. + *

Periodically maintain the RegionReplicas to be created or deleted */ public void maintainRegionReplicas() { // The consensusManager of configManager may not be fully initialized at this time - if (getConsensusManager() == null || !getConsensusManager().isLeader()) { - return; - } - - // Group the queued tasks into one FIFO sub-queue per region. The queue only ever holds - // RegionCreateTasks now (delete tasks are filtered out at the PartitionInfo ingestion points), - // and a region may carry several of them when more than one of its replicas failed to create. - final Map> tasksByRegion = new HashMap<>(); - for (RegionMaintainTask task : partitionInfo.getRegionMaintainEntryList()) { - if (!(task instanceof RegionCreateTask)) { - // Unreachable: the queue only holds create tasks now (legacy delete tasks are dropped at - // the - // PartitionInfo ingestion points). Guard against a regression so an unexpected task type - // cannot silently stall the loop. - LOGGER.warn(ManagerMessages.UNEXPECTED_NON_CREATE_REGION_MAINTAIN_TASK_SKIPPED); - continue; - } - tasksByRegion - .computeIfAbsent(task.getRegionId(), k -> new LinkedList<>()) - .add((RegionCreateTask) task); - } - - // Drain the sub-queues head-by-head. Each round takes the head of every region, batches those - // heads by region type into a single create RPC per type, then durably polls the tasks that - // succeeded. Tasks of the same region are advanced one at a time to preserve their offer order. - while (!tasksByRegion.isEmpty()) { - final Map> headsByType = - new EnumMap<>(TConsensusGroupType.class); - for (Queue queue : tasksByRegion.values()) { - final RegionCreateTask head = queue.peek(); - headsByType.computeIfAbsent(head.getRegionId().getType(), k -> new ArrayList<>()).add(head); - } + Optional.ofNullable(getConsensusManager()) + .ifPresent( + consensusManager -> { + if (getConsensusManager().isLeader()) { + List regionMaintainTaskList = + partitionInfo.getRegionMaintainEntryList(); + + if (regionMaintainTaskList.isEmpty()) { + 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()); + } + } + } - final Set successfulRegions = new HashSet<>(); - int selectedCount = 0; - for (Map.Entry> entry : headsByType.entrySet()) { - selectedCount += entry.getValue().size(); - successfulRegions.addAll(submitRegionCreateTasks(entry.getKey(), entry.getValue())); - } + if (selectedRegionMaintainTask.isEmpty()) { + break; + } - if (successfulRegions.isEmpty()) { - break; - } + 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( + ManagerMessages.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 (isRegionMaintainTaskCompleted( + RegionMaintainType.CREATE, entry.getValue())) { + 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( + ManagerMessages.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 (isRegionMaintainTaskCompleted( + RegionMaintainType.CREATE, entry.getValue())) { + 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( + ManagerMessages.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( + ManagerMessages.DELETING_REGIONS_COSTS_MS, + (System.currentTimeMillis() - startTime)); + + for (Map.Entry entry : + deleteRegionHandler.getResponseMap().entrySet()) { + if (isRegionMaintainTaskCompleted( + RegionMaintainType.DELETE, entry.getValue())) { + successfulTask.add(regionIdMap.get(entry.getKey())); + } + } + break; + } - // Advance the in-memory sub-queues so the next round picks the following task of each region. - for (TConsensusGroupId regionId : successfulRegions) { - tasksByRegion.computeIfPresent( - regionId, - (k, queue) -> { - queue.poll(); - return queue.isEmpty() ? null : queue; - }); - } + if (successfulTask.isEmpty()) { + break; + } - // Durably remove the head of every successfully created region from the persisted queue. - try { - getConsensusManager().write(new PollSpecificRegionMaintainTaskPlan(successfulRegions)); - } catch (ConsensusException e) { - LOGGER.warn(CONSENSUS_WRITE_ERROR, e); - } + 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; + } + }); + } - if (successfulRegions.size() < selectedCount) { - // Some tasks failed this round; stop and retry on the next schedule so that the tasks of - // each region keep being executed in the order they were offered. - break; - } - } - } + // Poll the head entry if success + try { + getConsensusManager() + .write(new PollSpecificRegionMaintainTaskPlan(successfulTask)); + } catch (ConsensusException e) { + LOGGER.warn(CONSENSUS_WRITE_ERROR, e); + } - /** - * Send a batched create RPC for the given heads, all of which share the given region type, and - * return the ids of the regions whose replica was created successfully. - */ - private Set submitRegionCreateTasks( - TConsensusGroupType regionType, List createTasks) { - final Set successfulRegions = new HashSet<>(); - switch (regionType) { - case SchemaRegion: - final DataNodeAsyncRequestContext schemaHandler = - new DataNodeAsyncRequestContext<>(CnToDnAsyncRequestType.CREATE_SCHEMA_REGION); - for (RegionCreateTask task : createTasks) { - LOGGER.info( - ManagerMessages.START_TO_CREATE_REGION_ON_DATANODE, - task.getRegionReplicaSet().getRegionId(), - task.getTargetDataNode()); - schemaHandler.putRequest( - task.getRegionId().getId(), - new TCreateSchemaRegionReq(task.getRegionReplicaSet(), task.getStorageGroup())); - schemaHandler.putNodeLocation(task.getRegionId().getId(), task.getTargetDataNode()); - } - CnToDnInternalServiceAsyncRequestManager.getInstance() - .sendAsyncRequestWithRetry(schemaHandler); - collectSuccessfulRegions( - schemaHandler.getResponseMap(), TConsensusGroupType.SchemaRegion, successfulRegions); - break; - case DataRegion: - final DataNodeAsyncRequestContext dataHandler = - new DataNodeAsyncRequestContext<>(CnToDnAsyncRequestType.CREATE_DATA_REGION); - for (RegionCreateTask task : createTasks) { - LOGGER.info( - ManagerMessages.START_TO_CREATE_REGION_ON_DATANODE, - task.getRegionReplicaSet().getRegionId(), - task.getTargetDataNode()); - dataHandler.putRequest( - task.getRegionId().getId(), - new TCreateDataRegionReq(task.getRegionReplicaSet(), task.getStorageGroup())); - dataHandler.putNodeLocation(task.getRegionId().getId(), task.getTargetDataNode()); - } - CnToDnInternalServiceAsyncRequestManager.getInstance() - .sendAsyncRequestWithRetry(dataHandler); - collectSuccessfulRegions( - dataHandler.getResponseMap(), TConsensusGroupType.DataRegion, successfulRegions); - break; - default: - break; - } - return successfulRegions; + 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; + } + } + } + }); } - private void collectSuccessfulRegions( - Map responseMap, - TConsensusGroupType regionType, - Set successfulRegions) { - for (Map.Entry entry : responseMap.entrySet()) { - if (entry.getValue().getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()) { - successfulRegions.add(new TConsensusGroupId(regionType, entry.getKey())); - } + static boolean isRegionMaintainTaskCompleted( + RegionMaintainType regionMaintainType, TSStatus status) { + if (status.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()) { + return true; } + return regionMaintainType == RegionMaintainType.CREATE + ? status.getCode() == TSStatusCode.REGION_ALREADY_EXISTS.getStatusCode() + : status.getCode() == TSStatusCode.REGION_NOT_EXIST.getStatusCode(); } public void startRegionCleaner() { diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/partition/PartitionInfo.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/partition/PartitionInfo.java index e71b61d7b338c..af9429e59532c 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/partition/PartitionInfo.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/partition/PartitionInfo.java @@ -64,7 +64,6 @@ import org.apache.iotdb.confignode.exception.DatabaseNotExistsException; import org.apache.iotdb.confignode.i18n.ConfigNodeMessages; import org.apache.iotdb.confignode.persistence.partition.maintainer.RegionMaintainTask; -import org.apache.iotdb.confignode.persistence.partition.maintainer.RegionMaintainType; import org.apache.iotdb.confignode.rpc.thrift.TRegionInfo; import org.apache.iotdb.confignode.rpc.thrift.TShowRegionReq; import org.apache.iotdb.confignode.rpc.thrift.TTimeSlotList; @@ -233,20 +232,7 @@ public TSStatus createRegionGroups(CreateRegionGroupsPlan plan) { public TSStatus offerRegionMaintainTasks( OfferRegionMaintainTasksPlan offerRegionMaintainTasksPlan) { synchronized (regionMaintainTaskList) { - // The RegionMaintainer queue only recreates failed region replicas now; region deletion is - // owned by RemoveRegionGroupProcedure. Drop any legacy DELETE task that an upgraded node may - // replay from an old consensus log, so it cannot get stuck in the queue and block the - // recreation of that region's other replicas. - for (RegionMaintainTask task : offerRegionMaintainTasksPlan.getRegionMaintainTaskList()) { - if (RegionMaintainType.DELETE.equals(task.getType())) { - LOGGER.info( - ConfigNodeMessages - .MESSAGE_DROPPING_LEGACY_REGION_DELETE_TASK_FOR_ARG_WHILE_REPLAYING_OFFER_PLAN_REGION_DELETION_IS_NOW_HANDLED_BY_REMOVEREGIONGROUPPROCEDURE_2A81A649, - task.getRegionId()); - continue; - } - regionMaintainTaskList.add(task); - } + regionMaintainTaskList.addAll(offerRegionMaintainTasksPlan.getRegionMaintainTaskList()); return RpcUtils.SUCCESS_STATUS; } } @@ -1087,22 +1073,11 @@ public void processLoadSnapshot(final File snapshotDir) throws TException, IOExc databasePartitionTables.put(database, databasePartitionTable); } - // restore the RegionMaintainer queue + // restore deletedRegionSet length = ReadWriteIOUtils.readInt(fileInputStream); for (int i = 0; i < length; i++) { final RegionMaintainTask task = RegionMaintainTask.Factory.create(fileInputStream, protocol); - // The RegionMaintainer queue only recreates failed region replicas now; region deletion is - // owned by RemoveRegionGroupProcedure. Drop any legacy DELETE task carried over from an - // upgraded snapshot so it cannot get stuck at the head of a region's queue and block the - // recreation of that region's other replicas. - if (RegionMaintainType.DELETE.equals(task.getType())) { - LOGGER.info( - ConfigNodeMessages - .MESSAGE_DROPPING_LEGACY_REGION_DELETE_TASK_FOR_ARG_WHILE_LOADING_SNAPSHOT_REGION_DELETION_IS_NOW_HANDLED_BY_REMOVEREGIONGROUPPROCEDURE_A9D409A0, - task.getRegionId()); - continue; - } regionMaintainTaskList.add(task); } } diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/region/CreateRegionGroupsProcedure.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/region/CreateRegionGroupsProcedure.java index 276cdf432d9e6..e9cce807e77fc 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/region/CreateRegionGroupsProcedure.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/region/CreateRegionGroupsProcedure.java @@ -36,6 +36,7 @@ import org.apache.iotdb.confignode.i18n.ProcedureMessages; import org.apache.iotdb.confignode.manager.load.cache.region.RegionHeartbeatSample; import org.apache.iotdb.confignode.persistence.partition.maintainer.RegionCreateTask; +import org.apache.iotdb.confignode.persistence.partition.maintainer.RegionDeleteTask; import org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv; import org.apache.iotdb.confignode.procedure.exception.ProcedureException; import org.apache.iotdb.confignode.procedure.impl.StateMachineProcedure; @@ -50,9 +51,7 @@ import java.io.DataOutputStream; import java.io.IOException; import java.nio.ByteBuffer; -import java.util.ArrayList; import java.util.HashMap; -import java.util.List; import java.util.Map; import java.util.Objects; import java.util.Set; @@ -108,12 +107,6 @@ protected Flow executeFromState( case SHUNT_REGION_REPLICAS: persistPlan = new CreateRegionGroupsPlan(); final OfferRegionMaintainTasksPlan offerPlan = new OfferRegionMaintainTasksPlan(); - // RegionGroups that failed to reach a serving quorum have their redundant (already-created) - // replicas removed via an independent root RemoveRegionGroupProcedure. Submitting them as - // root procedures (instead of children) keeps this procedure from waiting for or being - // failed by the cleanup: each one retries forever until those replicas are deleted, while - // this procedure proceeds to activate the region groups that did form a quorum. - final List removeRegionGroupProcedures = new ArrayList<>(); // Filter those RegionGroups that created successfully createRegionGroupsPlan .getRegionGroupMap() @@ -161,11 +154,7 @@ protected Flow executeFromState( .CREATEREGIONGROUPS_FAILED_TO_CREATE_SOME_REPLICAS_OF_REGIONGROUP_BUT_THIS, regionReplicaSet.getRegionId()); } else { - // The redundant RegionReplicas (the ones that did get created) should - // be deleted otherwise - final TRegionReplicaSet redundantReplicas = - new TRegionReplicaSet() - .setRegionId(regionReplicaSet.getRegionId()); + // The redundant RegionReplicas should be deleted otherwise regionReplicaSet .getDataNodeLocations() .forEach( @@ -173,13 +162,12 @@ protected Flow executeFromState( if (!failedRegionReplicas .getDataNodeLocations() .contains(targetDataNode)) { - redundantReplicas.addToDataNodeLocations(targetDataNode); + RegionDeleteTask deleteTask = + new RegionDeleteTask( + targetDataNode, regionReplicaSet.getRegionId()); + offerPlan.appendRegionMaintainTask(deleteTask); } }); - if (redundantReplicas.getDataNodeLocationsSize() > 0) { - removeRegionGroupProcedures.add( - new RemoveRegionGroupProcedure(redundantReplicas)); - } LOGGER.info( ProcedureMessages @@ -200,26 +188,6 @@ protected Flow executeFromState( LOGGER.warn( ConfigNodeMessages.FAILED_IN_THE_WRITE_API_EXECUTING_THE_CONSENSUS_LAYER_DUE, e); } - // Submit the redundant-replica cleanups as independent root procedures. This is - // intentionally NOT guarded by isStateDeserialized(): the executor persists a procedure at - // a state BEFORE that state's body has run (it advances the state on the previous cycle, - // then may stop at the inter-state boundary on a leader switch — see - // ProcedureExecutor#executeProcedure), so a recovery that lands on SHUNT_REGION_REPLICAS - // means the submissions have NOT happened yet. Skipping them would leave the - // already-created - // replicas of sub-quorum region groups on disk with no cleanup and no partition-table - // record - // (the else branch above never persisted them). Re-submitting on recovery is safe instead: - // the cleanups are recomputed from the serialized failedRegionReplicaSets, each gets a - // fresh - // procId and performs an idempotent delete, so a duplicate is harmless whereas a skip - // leaks. - removeRegionGroupProcedures.forEach( - removeRegionGroupProcedure -> - env.getConfigManager() - .getProcedureManager() - .getExecutor() - .submitProcedure(removeRegionGroupProcedure)); setNextState(CreateRegionGroupsState.REBALANCE_DATA_PARTITION_POLICY); break; case REBALANCE_DATA_PARTITION_POLICY: diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/region/RemoveRegionGroupProcedure.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/region/RemoveRegionGroupProcedure.java deleted file mode 100644 index b2b49a688d204..0000000000000 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/region/RemoveRegionGroupProcedure.java +++ /dev/null @@ -1,389 +0,0 @@ -/* - * 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.procedure.impl.region; - -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.commons.exception.runtime.ThriftSerDeException; -import org.apache.iotdb.commons.queryengine.utils.DateTimeUtils; -import org.apache.iotdb.commons.utils.CommonDateTimeUtils; -import org.apache.iotdb.commons.utils.TestOnly; -import org.apache.iotdb.commons.utils.ThriftCommonsSerDeUtils; -import org.apache.iotdb.confignode.i18n.ProcedureMessages; -import org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv; -import org.apache.iotdb.confignode.procedure.env.RegionMaintainHandler; -import org.apache.iotdb.confignode.procedure.exception.ProcedureException; -import org.apache.iotdb.confignode.procedure.state.RemoveRegionGroupState; -import org.apache.iotdb.confignode.procedure.store.ProcedureType; -import org.apache.iotdb.mpp.rpc.thrift.TRegionMigrateResult; - -import org.apache.tsfile.utils.ReadWriteIOUtils; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -import java.io.DataOutputStream; -import java.io.IOException; -import java.nio.ByteBuffer; -import java.util.List; -import java.util.Objects; -import java.util.stream.Collectors; - -import static org.apache.iotdb.commons.utils.KillPoint.KillPoint.setKillPoint; -import static org.apache.iotdb.confignode.procedure.env.RegionMaintainHandler.simplifiedLocation; -import static org.apache.iotdb.rpc.TSStatusCode.SUCCESS_STATUS; - -/** - * Delete a whole region group: every replica's consensus peer and all of its data on every DataNode - * that hosts it. - * - *

Each replica is removed with a local {@code deleteLocalPeer} (the {@code - * submitDeleteOldRegionPeerTask} path), which needs no consensus quorum and tolerates an - * already-absent peer, so it works for a group of any size — including a sub-quorum group that - * never finished forming. The DataNode runs the deletion asynchronously and this procedure polls - * for the result, so a slow deletion is never wrongly reported as finished. - * - *

This procedure is submitted as an independent root procedure (not a child) by its callers, - * which only enqueue the deletion and return immediately. It therefore owns the deletion end to - * end: on any failure it retries the current replica forever (backing off between attempts) instead - * of giving up, because there is no parent left to fall back to and the region's peer/data must not - * be left on disk. Each genuine re-attempt uses a FRESH DataNode-side taskId (the DataNode dedups - * by taskId and caches a terminal result forever, so reusing one taskId would make every retry a - * no-op that never re-runs the delete); the in-flight taskId is persisted so a leader change - * re-polls the same task rather than double-submitting. It carries its own {@link - * TRegionReplicaSet} copy, so it can finish even after the caller has dropped the partition table, - * and it survives ConfigNode leader change / restart. - */ -public class RemoveRegionGroupProcedure extends RegionOperationProcedure { - private static final Logger LOGGER = LoggerFactory.getLogger(RemoveRegionGroupProcedure.class); - - private static final long DELETE_REPLICA_RETRY_INTERVAL_MS = 5_000; - - private TRegionReplicaSet regionReplicaSet; - - // The index of the replica currently being deleted. Persisted and advanced only after that - // replica - // is deleted, so after a ConfigNode leader change the procedure resumes on the first replica it - // has - // not finished deleting. - private int currentReplicaIndex; - - // Number of failed attempts on the replica at currentReplicaIndex, used only for logging. Retries - // are unbounded, so this is not a budget. Transient: a leader change restarts the counter for the - // current replica. - private transient int attemptedForCurrentReplica; - - // Monotonic count of delete tasks this procedure has submitted, across all replicas. Persisted - // and - // only ever incremented. It is the low half of the DataNode-side taskId (see deleteTaskId): a - // fresh value per genuine re-attempt makes the DataNode re-run the delete instead of replaying a - // cached terminal result for a reused taskId (the DataNode dedups by taskId and never clears the - // cache), which is the bug this fixes. It never resets, so every taskId this procedure emits is - // distinct even across replicas and retries. - private long deleteTaskSeq; - - // Whether a delete task for the replica at currentReplicaIndex has already been submitted (and - // thus - // deleteTaskSeq already identifies an in-flight task to re-poll) rather than needing a fresh one. - // Persisted so a leader change mid-attempt re-polls the SAME in-flight task instead of submitting - // a - // duplicate; cleared on success or when a terminal failure forces a fresh re-attempt. - private boolean deleteTaskSubmitted; - - // Bit budget for deleteTaskId(): sign bit (=> negative) + PROC_ID_BITS + SEQ_BITS must be <= 64. - private static final int SEQ_BITS = 20; - private static final int PROC_ID_BITS = 43; - - public RemoveRegionGroupProcedure() { - super(); - } - - public RemoveRegionGroupProcedure(TRegionReplicaSet regionReplicaSet) { - super(regionReplicaSet.getRegionId()); - this.regionReplicaSet = regionReplicaSet; - } - - @TestOnly - void setCurrentReplicaIndex(int currentReplicaIndex) { - this.currentReplicaIndex = currentReplicaIndex; - } - - @TestOnly - void setDeleteTaskState(long deleteTaskSeq, boolean deleteTaskSubmitted) { - this.deleteTaskSeq = deleteTaskSeq; - this.deleteTaskSubmitted = deleteTaskSubmitted; - } - - @TestOnly - long deleteTaskIdForTest() { - return deleteTaskId(); - } - - @Override - protected Flow executeFromState(ConfigNodeProcedureEnv env, RemoveRegionGroupState state) - throws InterruptedException { - final List dataNodeLocations = - regionReplicaSet == null ? null : regionReplicaSet.getDataNodeLocations(); - if (dataNodeLocations == null) { - // A null replica set means deserialization failed. Retrying cannot recover the lost - // locations, - // so fail loudly instead of silently reporting the group as deleted (which would leave the - // region's peer/data on disk with no record of where it lives). - setFailure( - new ProcedureException(ProcedureMessages.UNSUPPORTED_STATE + "missing regionReplicaSet")); - return Flow.NO_MORE_STATE; - } - final RegionMaintainHandler handler = env.getRegionMaintainHandler(); - switch (state) { - case DELETE_REGION_REPLICAS: - if (currentReplicaIndex == 0 && attemptedForCurrentReplica == 0) { - LOGGER.info( - ProcedureMessages.PID_REMOVEREGIONGROUP_STARTED_WILL_BE_DELETED, - getProcId(), - regionId, - dataNodeLocations.stream() - .map(RegionMaintainHandler::simplifiedLocation) - .collect(Collectors.toList())); - } - if (currentReplicaIndex >= dataNodeLocations.size()) { - // Requirement: every successfully completed maintain task must be logged. - LOGGER.info( - ProcedureMessages.PID_REMOVEREGIONGROUP_SUCCESS_PROCEDURE_TOOK, - getProcId(), - regionId, - CommonDateTimeUtils.convertMillisecondToDurationStr( - System.currentTimeMillis() - getSubmittedTime()), - DateTimeUtils.convertLongToDate(getSubmittedTime(), "ms")); - return Flow.NO_MORE_STATE; - } - - final TDataNodeLocation targetDataNode = dataNodeLocations.get(currentReplicaIndex); - LOGGER.info( - ProcedureMessages.PID_REMOVEREGIONGROUP_STARTED_REPLICA_WILL_BE_DELETED_FROM_DATANODE, - getProcId(), - regionId, - simplifiedLocation(targetDataNode)); - - // Start a fresh attempt (fresh taskId) unless we are resuming an already-submitted one - // after - // a leader change, in which case we re-poll the SAME task rather than submitting a - // duplicate. - if (!deleteTaskSubmitted) { - deleteTaskSeq++; - deleteTaskSubmitted = true; - } - final long deleteTaskId = deleteTaskId(); - - // deleteLocalPeer is idempotent (it tolerates an already-absent peer), and re-submitting - // the - // same taskId re-polls the same DataNode task, so resuming after a leader change is safe. - final TSStatus submitStatus; - final TRegionMigrateResult result; - try { - submitStatus = - handler.submitDeleteOldRegionPeerTask(deleteTaskId, targetDataNode, regionId); - setKillPoint(state); - if (submitStatus.getCode() != SUCCESS_STATUS.getStatusCode()) { - return retryCurrentReplica( - String.format( - "submit delete task for region %s to DataNode %s failed: %s", - regionId, simplifiedLocation(targetDataNode), submitStatus)); - } - result = handler.waitTaskFinish(deleteTaskId, targetDataNode); - } catch (InterruptedException e) { - throw e; - } catch (Exception e) { - LOGGER.error(ProcedureMessages.PID_REMOVEREGIONGROUP_STATE_FAILED, getProcId(), state, e); - return retryCurrentReplica( - String.format( - "delete region %s from DataNode %s threw %s", - regionId, simplifiedLocation(targetDataNode), e)); - } - - switch (result.getTaskStatus()) { - case SUCCESS: - // Advance to the next replica with a fresh retry counter and a fresh delete task. - currentReplicaIndex++; - attemptedForCurrentReplica = 0; - deleteTaskSubmitted = false; - setNextState(RemoveRegionGroupState.DELETE_REGION_REPLICAS); - return Flow.HAS_MORE_STATE; - case PROCESSING: - // waitTaskFinish() only returns PROCESSING when its polling loop was interrupted, i.e. - // this ConfigNode is shutting down / losing leadership. The delete task is still - // running on the DataNode, so persist and re-poll after recovery: stay on this replica - // without advancing it, without consuming a retry attempt, and keeping deleteTaskSeq / - // deleteTaskSubmitted so the re-poll targets the same in-flight task. - setNextState(RemoveRegionGroupState.DELETE_REGION_REPLICAS); - return Flow.HAS_MORE_STATE; - case TASK_NOT_EXIST: - case FAIL: - default: - return retryCurrentReplica( - String.format( - "delete region %s from DataNode %s, task status is %s", - regionId, simplifiedLocation(targetDataNode), result.getTaskStatus())); - } - default: - setFailure(new ProcedureException(ProcedureMessages.UNSUPPORTED_STATE + state.name())); - return Flow.NO_MORE_STATE; - } - } - - /** - * Retry the replica at {@link #currentReplicaIndex} after a backoff. This procedure never gives - * up on a replica: because it is submitted as an independent root procedure, there is no parent - * to fall back to, and skipping or failing would leave the region's peer/data on disk. So it - * backs off and re-runs the same state until the replica is deleted, which eventually succeeds - * once the target DataNode is reachable: the delete is idempotent, and clearing {@link - * #deleteTaskSubmitted} here makes the next attempt use a FRESH DataNode-side taskId (a new - * {@link #deleteTaskSeq}), so the DataNode actually re-executes the delete instead of returning a - * cached terminal result for the previous taskId. - */ - private Flow retryCurrentReplica(String reason) throws InterruptedException { - attemptedForCurrentReplica++; - LOGGER.warn( - ProcedureMessages.PID_REMOVEREGIONGROUP_DELETE_REPLICA_FAILED, - getProcId(), - regionId, - attemptedForCurrentReplica, - reason); - // Force a fresh delete task on the next attempt so the DataNode re-runs the delete rather than - // replaying a cached FAIL/SUCCESS for this taskId. - deleteTaskSubmitted = false; - Thread.sleep(DELETE_REPLICA_RETRY_INTERVAL_MS); - setNextState(RemoveRegionGroupState.DELETE_REGION_REPLICAS); - return Flow.HAS_MORE_STATE; - } - - /** - * The DataNode-side taskId for the current attempt, derived from this procedure's (globally - * unique, consensus-replicated) procId and its monotonic {@link #deleteTaskSeq}. It is packed - * into the NEGATIVE i64 space, which is disjoint from every real procId (all {@code >= 0}); other - * region-maintain procedures (add/remove peer) use {@code getProcId()} directly as the taskId - * against the same DataNode task map, so a negative id can never collide with theirs. Unlike - * minting from the procedure-store id allocator, this needs nothing extra replicated: procId is - * already replicated and deleteTaskSeq is persisted with this procedure, so the taskId is stable - * across a leader change and never regresses. - * - *

Layout: sign bit set (=> negative) | {@value PROC_ID_BITS} bits of procId | {@value - * SEQ_BITS} bits of deleteTaskSeq. The bounds are astronomically beyond any real cluster (a - * procId needs 2^43 procedures; a single group delete needs 2^20 retries), and are asserted - * rather than silently wrapped so a violation fails the procedure loudly instead of emitting a - * colliding id. - */ - private long deleteTaskId() { - final long procId = getProcId(); - if (procId < 0 || procId >= (1L << PROC_ID_BITS) || deleteTaskSeq >= (1L << SEQ_BITS)) { - throw new IllegalStateException( - String.format( - ProcedureMessages - .EXCEPTION_CANNOT_DERIVE_A_COLLISION_FREE_DELETE_TASKID_PROCID_ARG_DELETETASKSEQ_ARG_EXCEED_THE_ARG_ARG_BIT_BUDGET_015C598D, - procId, - deleteTaskSeq, - PROC_ID_BITS, - SEQ_BITS)); - } - return Long.MIN_VALUE | (procId << SEQ_BITS) | deleteTaskSeq; - } - - @Override - protected void rollbackState(ConfigNodeProcedureEnv env, RemoveRegionGroupState state) - throws IOException, InterruptedException, ProcedureException {} - - @Override - protected RemoveRegionGroupState getState(int stateId) { - return RemoveRegionGroupState.values()[stateId]; - } - - @Override - protected int getStateId(RemoveRegionGroupState removeRegionGroupState) { - return removeRegionGroupState.ordinal(); - } - - @Override - protected RemoveRegionGroupState getInitialState() { - return RemoveRegionGroupState.DELETE_REGION_REPLICAS; - } - - @Override - public void serialize(DataOutputStream stream) throws IOException { - stream.writeShort(ProcedureType.REMOVE_REGION_GROUP_PROCEDURE.getTypeCode()); - super.serialize(stream); - ThriftCommonsSerDeUtils.serializeTRegionReplicaSet(regionReplicaSet, stream); - ReadWriteIOUtils.write(currentReplicaIndex, stream); - // Persist the delete-task cursor so a leader change re-derives the SAME in-flight taskId and - // re-polls it (deleteTaskSubmitted == true) instead of submitting a duplicate, and so the - // monotonic deleteTaskSeq never regresses. - ReadWriteIOUtils.write(deleteTaskSeq, stream); - ReadWriteIOUtils.write(deleteTaskSubmitted, stream); - } - - @Override - public void deserialize(ByteBuffer byteBuffer) { - super.deserialize(byteBuffer); - try { - regionReplicaSet = ThriftCommonsSerDeUtils.deserializeTRegionReplicaSet(byteBuffer); - regionId = regionReplicaSet.getRegionId(); - currentReplicaIndex = ReadWriteIOUtils.readInt(byteBuffer); - // deleteTaskSeq/deleteTaskSubmitted were appended after the first version of this procedure. - // That first version only ever existed on the unreleased branch that added this procedure - // (never in a release), but a dev/CI cluster could persist a blob without these trailing - // fields; tolerate it by defaulting to "no in-flight task" instead of reading past the end. - if (byteBuffer.hasRemaining()) { - deleteTaskSeq = ReadWriteIOUtils.readLong(byteBuffer); - deleteTaskSubmitted = ReadWriteIOUtils.readBool(byteBuffer); - } - } catch (ThriftSerDeException e) { - LOGGER.error(ProcedureMessages.ERROR_IN_DESERIALIZE, this.getClass(), e); - } - } - - @Override - public boolean equals(Object obj) { - if (!(obj instanceof RemoveRegionGroupProcedure)) { - return false; - } - RemoveRegionGroupProcedure procedure = (RemoveRegionGroupProcedure) obj; - return this.currentReplicaIndex == procedure.currentReplicaIndex - && this.deleteTaskSeq == procedure.deleteTaskSeq - && this.deleteTaskSubmitted == procedure.deleteTaskSubmitted - && Objects.equals(this.regionReplicaSet, procedure.regionReplicaSet); - } - - @Override - public int hashCode() { - return Objects.hash(regionReplicaSet, currentReplicaIndex, deleteTaskSeq, deleteTaskSubmitted); - } - - @Override - public String toString() { - return "RemoveRegionGroupProcedure{" - + "regionReplicaSet=" - + regionReplicaSet - + ", currentReplicaIndex=" - + currentReplicaIndex - + ", deleteTaskSeq=" - + deleteTaskSeq - + ", deleteTaskSubmitted=" - + deleteTaskSubmitted - + '}'; - } -} 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 b547e8ebc7d39..53d56da67bc86 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 @@ -19,21 +19,29 @@ 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.commons.exception.runtime.ThriftSerDeException; import org.apache.iotdb.commons.service.metric.MetricService; import org.apache.iotdb.commons.utils.ThriftConfigNodeSerDeUtils; +import org.apache.iotdb.confignode.client.async.CnToDnAsyncRequestType; +import org.apache.iotdb.confignode.client.async.CnToDnInternalServiceAsyncRequestManager; +import org.apache.iotdb.confignode.client.async.handlers.DataNodeAsyncRequestContext; import org.apache.iotdb.confignode.consensus.request.write.database.PreDeleteDatabasePlan; +import org.apache.iotdb.confignode.consensus.request.write.region.OfferRegionMaintainTasksPlan; import org.apache.iotdb.confignode.i18n.ProcedureMessages; import org.apache.iotdb.confignode.manager.partition.PartitionMetrics; +import org.apache.iotdb.confignode.persistence.partition.maintainer.RegionDeleteTask; import org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv; import org.apache.iotdb.confignode.procedure.exception.ProcedureException; import org.apache.iotdb.confignode.procedure.impl.StateMachineProcedure; -import org.apache.iotdb.confignode.procedure.impl.region.RemoveRegionGroupProcedure; import org.apache.iotdb.confignode.procedure.state.schema.DeleteDatabaseState; import org.apache.iotdb.confignode.procedure.store.ProcedureType; import org.apache.iotdb.confignode.rpc.thrift.TDatabaseSchema; +import org.apache.iotdb.consensus.exception.ConsensusException; import org.apache.iotdb.rpc.TSStatusCode; import org.apache.thrift.TException; @@ -43,7 +51,10 @@ import java.io.DataOutputStream; import java.io.IOException; import java.nio.ByteBuffer; +import java.util.ArrayList; +import java.util.HashMap; import java.util.List; +import java.util.Map; import java.util.Objects; public class DeleteDatabaseProcedure @@ -104,40 +115,88 @@ protected Flow executeFromState(final ConfigNodeProcedureEnv env, final DeleteDa ProcedureMessages.LOG_DELETEDATABASEPROCEDURE_DELETE_DATABASESCHEMA_ARG_A49A47AC, deleteDatabaseSchema.getName()); - // Enqueue deletion of every region group (both schema and data regions) of this database. - // Each is submitted as an INDEPENDENT root RemoveRegionGroupProcedure rather than a - // child: - // this procedure only submits the deletions and then returns, so it can neither wait for - // nor be failed/rolled-back by a slow or failing region deletion. Each carries its own - // copy of the replica set, so the deletion still completes (and survives leader change / - // restart) even after the next state drops the partition table. - // - // Submission is intentionally NOT guarded by isStateDeserialized(): the executor persists - // a procedure at a state BEFORE that state's body has run (it advances the state on the - // previous cycle, then may stop at the inter-state boundary on a leader switch — see - // ProcedureExecutor#executeProcedure). So a recovery that lands on this state means the - // submission has NOT happened yet; skipping it would drop every region group's cleanup - // while the next state still drops the partition table, orphaning the region peers/data - // on - // disk with no record of where they live. Re-submitting on recovery is safe instead: - // every RemoveRegionGroupProcedure gets a fresh procId and performs an idempotent delete, - // so a duplicate is harmless whereas a skip leaks data. + // Submit RegionDeleteTasks + final OfferRegionMaintainTasksPlan dataRegionDeleteTaskOfferPlan = + new OfferRegionMaintainTasksPlan(); final List regionReplicaSets = env.getAllReplicaSets(deleteDatabaseSchema.getName()); + final List schemaRegionReplicaSets = new ArrayList<>(); regionReplicaSets.forEach( regionReplicaSet -> { // Clear heartbeat cache along the way env.getConfigManager() .getLoadManager() .removeRegionGroupRelatedCache(regionReplicaSet.getRegionId()); - env.getConfigManager() - .getProcedureManager() - .getExecutor() - .submitProcedure(new RemoveRegionGroupProcedure(regionReplicaSet)); + + if (regionReplicaSet + .getRegionId() + .getType() + .equals(TConsensusGroupType.SchemaRegion)) { + schemaRegionReplicaSets.add(regionReplicaSet); + } else { + regionReplicaSet + .getDataNodeLocations() + .forEach( + targetDataNode -> + dataRegionDeleteTaskOfferPlan.appendRegionMaintainTask( + new RegionDeleteTask( + targetDataNode, regionReplicaSet.getRegionId()))); + } }); - setNextState(DeleteDatabaseState.DELETE_DATABASE_CONFIG); - break; - case DELETE_DATABASE_CONFIG: + + if (!dataRegionDeleteTaskOfferPlan.getRegionMaintainTaskList().isEmpty()) { + // submit async data region delete task + env.getConfigManager().getConsensusManager().write(dataRegionDeleteTaskOfferPlan); + } + + // try sync delete schemaengine region + final DataNodeAsyncRequestContext asyncClientHandler = + new DataNodeAsyncRequestContext<>(CnToDnAsyncRequestType.DELETE_REGION); + final Map schemaRegionDeleteTaskMap = new HashMap<>(); + int requestIndex = 0; + for (final TRegionReplicaSet schemaRegionReplicaSet : schemaRegionReplicaSets) { + for (final TDataNodeLocation dataNodeLocation : + schemaRegionReplicaSet.getDataNodeLocations()) { + asyncClientHandler.putRequest(requestIndex, schemaRegionReplicaSet.getRegionId()); + asyncClientHandler.putNodeLocation(requestIndex, dataNodeLocation); + schemaRegionDeleteTaskMap.put( + requestIndex, + new RegionDeleteTask(dataNodeLocation, schemaRegionReplicaSet.getRegionId())); + requestIndex++; + } + } + if (!schemaRegionDeleteTaskMap.isEmpty()) { + CnToDnInternalServiceAsyncRequestManager.getInstance() + .sendAsyncRequestWithRetry(asyncClientHandler); + for (final Map.Entry entry : + asyncClientHandler.getResponseMap().entrySet()) { + if (entry.getValue().getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()) { + LOG.info( + ProcedureMessages + .LOG_DELETEDATABASEPROCEDURE_SUCCESSFULLY_DELETE_SCHEMAREGION_ARG_ARG_BA0535DA, + asyncClientHandler.getRequest(entry.getKey()), + schemaRegionDeleteTaskMap.get(entry.getKey()).getTargetDataNode()); + schemaRegionDeleteTaskMap.remove(entry.getKey()); + } else { + LOG.warn( + ProcedureMessages + .LOG_DELETEDATABASEPROCEDURE_FAILED_DELETE_SCHEMAREGION_ARG_ARG_SUBMIT_ASYNC_DELETION_8C3E6DE3, + asyncClientHandler.getRequest(entry.getKey()), + schemaRegionDeleteTaskMap.get(entry.getKey()).getTargetDataNode()); + } + } + + if (!schemaRegionDeleteTaskMap.isEmpty()) { + // submit async schemaengine region delete task for failed sync execution + final OfferRegionMaintainTasksPlan schemaRegionDeleteTaskOfferPlan = + new OfferRegionMaintainTasksPlan(); + schemaRegionDeleteTaskMap + .values() + .forEach(schemaRegionDeleteTaskOfferPlan::appendRegionMaintainTask); + env.getConfigManager().getConsensusManager().write(schemaRegionDeleteTaskOfferPlan); + } + } + env.getConfigManager() .getLoadManager() .clearDataPartitionPolicyTable(deleteDatabaseSchema.getName()); @@ -166,7 +225,7 @@ protected Flow executeFromState(final ConfigNodeProcedureEnv env, final DeleteDa ProcedureMessages.DELETEDATABASEPROCEDURE_DELETE_DATABASESCHEMA_FAILED)); } } - } catch (final TException | IOException e) { + } catch (final ConsensusException | TException | IOException e) { if (isRollbackSupported(state)) { setFailure( new ProcedureException( diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/state/RemoveRegionGroupState.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/state/RemoveRegionGroupState.java deleted file mode 100644 index 9b7d99c8a2323..0000000000000 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/state/RemoveRegionGroupState.java +++ /dev/null @@ -1,24 +0,0 @@ -/* - * 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.procedure.state; - -public enum RemoveRegionGroupState { - DELETE_REGION_REPLICAS, -} diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/state/schema/DeleteDatabaseState.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/state/schema/DeleteDatabaseState.java index c0ad67fe81190..cde6b2bdd8111 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/state/schema/DeleteDatabaseState.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/state/schema/DeleteDatabaseState.java @@ -22,8 +22,5 @@ public enum DeleteDatabaseState { PRE_DELETE_DATABASE, INVALIDATE_CACHE, - DELETE_DATABASE_SCHEMA, - // Delete the DatabasePartitionTable and related config after all region groups have been deleted - // by the RemoveRegionGroupProcedure children spawned in DELETE_DATABASE_SCHEMA. - DELETE_DATABASE_CONFIG + DELETE_DATABASE_SCHEMA } diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/store/ProcedureFactory.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/store/ProcedureFactory.java index 4cc3af5480b12..26392e7c897fd 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/store/ProcedureFactory.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/store/ProcedureFactory.java @@ -43,7 +43,6 @@ import org.apache.iotdb.confignode.procedure.impl.region.NotifyRegionMigrationProcedure; import org.apache.iotdb.confignode.procedure.impl.region.ReconstructRegionProcedure; import org.apache.iotdb.confignode.procedure.impl.region.RegionMigrateProcedure; -import org.apache.iotdb.confignode.procedure.impl.region.RemoveRegionGroupProcedure; import org.apache.iotdb.confignode.procedure.impl.region.RemoveRegionPeerProcedure; import org.apache.iotdb.confignode.procedure.impl.schema.AlterEncodingCompressorProcedure; import org.apache.iotdb.confignode.procedure.impl.schema.AlterLogicalViewProcedure; @@ -146,9 +145,6 @@ public Procedure create(ByteBuffer buffer) throws IOException { case NOTIFY_REGION_MIGRATION_PROCEDURE: procedure = new NotifyRegionMigrationProcedure(); break; - case REMOVE_REGION_GROUP_PROCEDURE: - procedure = new RemoveRegionGroupProcedure(); - break; case ALTER_ENCODING_COMPRESSOR_PROCEDURE: procedure = new AlterEncodingCompressorProcedure(false); break; @@ -467,8 +463,6 @@ public static ProcedureType getProcedureType(final Procedure procedure) { return ProcedureType.RECONSTRUCT_REGION_PROCEDURE; } else if (procedure instanceof NotifyRegionMigrationProcedure) { return ProcedureType.NOTIFY_REGION_MIGRATION_PROCEDURE; - } else if (procedure instanceof RemoveRegionGroupProcedure) { - return ProcedureType.REMOVE_REGION_GROUP_PROCEDURE; } else if (procedure instanceof CreateTriggerProcedure) { return ProcedureType.CREATE_TRIGGER_PROCEDURE; } else if (procedure instanceof DropTriggerProcedure) { diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/store/ProcedureType.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/store/ProcedureType.java index c6b885bfac2b3..1cd6a46a4dcd1 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/store/ProcedureType.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/store/ProcedureType.java @@ -41,7 +41,6 @@ public enum ProcedureType { ADD_REGION_PEER_PROCEDURE((short) 204), REMOVE_REGION_PEER_PROCEDURE((short) 205), NOTIFY_REGION_MIGRATION_PROCEDURE((short) 206), - REMOVE_REGION_GROUP_PROCEDURE((short) 207), @TestOnly CREATE_MANY_DATABASES_PROCEDURE((short) 250), 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 0000000000000..ae920e43d9c0f --- /dev/null +++ b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/partition/PartitionManagerRegionMaintainTest.java @@ -0,0 +1,70 @@ +/* + * 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.TSStatus; +import org.apache.iotdb.confignode.persistence.partition.maintainer.RegionMaintainType; +import org.apache.iotdb.rpc.TSStatusCode; + +import org.junit.Test; + +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; + +public class PartitionManagerRegionMaintainTest { + + @Test + public void testCreateRegionCompletedStatus() { + assertTrue( + PartitionManager.isRegionMaintainTaskCompleted( + RegionMaintainType.CREATE, new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode()))); + assertTrue( + PartitionManager.isRegionMaintainTaskCompleted( + RegionMaintainType.CREATE, + new TSStatus(TSStatusCode.REGION_ALREADY_EXISTS.getStatusCode()))); + assertFalse( + PartitionManager.isRegionMaintainTaskCompleted( + RegionMaintainType.CREATE, + new TSStatus(TSStatusCode.REGION_NOT_EXIST.getStatusCode()))); + assertFalse( + PartitionManager.isRegionMaintainTaskCompleted( + RegionMaintainType.CREATE, + new TSStatus(TSStatusCode.CREATE_REGION_ERROR.getStatusCode()))); + } + + @Test + public void testDeleteRegionCompletedStatus() { + assertTrue( + PartitionManager.isRegionMaintainTaskCompleted( + RegionMaintainType.DELETE, new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode()))); + assertTrue( + PartitionManager.isRegionMaintainTaskCompleted( + RegionMaintainType.DELETE, + new TSStatus(TSStatusCode.REGION_NOT_EXIST.getStatusCode()))); + assertFalse( + PartitionManager.isRegionMaintainTaskCompleted( + RegionMaintainType.DELETE, + new TSStatus(TSStatusCode.REGION_ALREADY_EXISTS.getStatusCode()))); + assertFalse( + PartitionManager.isRegionMaintainTaskCompleted( + RegionMaintainType.DELETE, + new TSStatus(TSStatusCode.DELETE_REGION_ERROR.getStatusCode()))); + } +} 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 8c32378cfcb1a..afccb0c0eba12 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 @@ -40,8 +40,6 @@ 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.persistence.partition.maintainer.RegionMaintainType; import org.apache.iotdb.confignode.rpc.thrift.TDatabaseSchema; import org.apache.iotdb.confignode.rpc.thrift.TShowRegionReq; @@ -154,31 +152,6 @@ public void testSnapshot() throws TException, IOException { Assert.assertEquals(partitionInfo, partitionInfo1); } - @Test - public void testLegacyRegionDeleteTasksAreFiltered() throws TException, IOException { - // Region deletion is owned by RemoveRegionGroupProcedure; the RegionMaintainer queue only - // recreates failed replicas. A legacy RegionDeleteTask (offered by an old version and replayed - // from a consensus log, or carried over in a snapshot) must be dropped rather than queued, so - // it cannot block the recreation of that region's other replicas. - - // The offer plan mixes two RegionCreateTasks with one legacy RegionDeleteTask. - partitionInfo.offerRegionMaintainTasks(generateOfferRegionMaintainTasksPlan()); - - // The DELETE task is filtered out at offer time; only the two CREATE tasks remain queued. - List queuedTasks = partitionInfo.getRegionMaintainEntryList(); - Assert.assertEquals(2, queuedTasks.size()); - for (RegionMaintainTask task : queuedTasks) { - Assert.assertEquals(RegionMaintainType.CREATE, task.getType()); - } - - // A snapshot round-trip keeps the CREATE tasks and never resurrects a DELETE task. - Assert.assertTrue(partitionInfo.processTakeSnapshot(snapshotDir)); - PartitionInfo loaded = new PartitionInfo(); - loaded.processLoadSnapshot(snapshotDir); - Assert.assertEquals(partitionInfo, loaded); - Assert.assertEquals(2, loaded.getRegionMaintainEntryList().size()); - } - @Test public void testGetRegionType() { diff --git a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/region/RemoveRegionGroupProcedureTest.java b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/region/RemoveRegionGroupProcedureTest.java deleted file mode 100644 index e239d64c12aa4..0000000000000 --- a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/region/RemoveRegionGroupProcedureTest.java +++ /dev/null @@ -1,100 +0,0 @@ -/* - * 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.procedure.impl.region; - -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.TEndPoint; -import org.apache.iotdb.common.rpc.thrift.TRegionReplicaSet; -import org.apache.iotdb.confignode.procedure.store.ProcedureFactory; - -import org.apache.tsfile.utils.PublicBAOS; -import org.junit.Assert; -import org.junit.Test; - -import java.io.DataOutputStream; -import java.nio.ByteBuffer; -import java.util.Arrays; - -public class RemoveRegionGroupProcedureTest { - @Test - public void serDeTest() throws Exception { - final TRegionReplicaSet regionReplicaSet = - new TRegionReplicaSet( - new TConsensusGroupId(TConsensusGroupType.DataRegion, 10), - Arrays.asList( - new TDataNodeLocation( - 1, - new TEndPoint("127.0.0.1", 0), - new TEndPoint("127.0.0.1", 1), - new TEndPoint("127.0.0.1", 2), - new TEndPoint("127.0.0.1", 3), - new TEndPoint("127.0.0.1", 4)), - new TDataNodeLocation( - 2, - new TEndPoint("127.0.0.1", 10), - new TEndPoint("127.0.0.1", 11), - new TEndPoint("127.0.0.1", 12), - new TEndPoint("127.0.0.1", 13), - new TEndPoint("127.0.0.1", 14)))); - final RemoveRegionGroupProcedure procedure = new RemoveRegionGroupProcedure(regionReplicaSet); - // A non-zero cursor so the round-trip actually exercises currentReplicaIndex (de)serialization; - // equals/hashCode include it, so a dropped/garbled cursor would fail the assertion. - procedure.setCurrentReplicaIndex(1); - // Non-default delete-task cursor so the round-trip exercises deleteTaskSeq/deleteTaskSubmitted - // too; equals/hashCode include them, so a dropped/garbled value would fail the assertion. - procedure.setDeleteTaskState(42L, true); - try (PublicBAOS byteArrayOutputStream = new PublicBAOS(); - DataOutputStream outputStream = new DataOutputStream(byteArrayOutputStream)) { - procedure.serialize(outputStream); - ByteBuffer buffer = - ByteBuffer.wrap(byteArrayOutputStream.getBuf(), 0, byteArrayOutputStream.size()); - // Exercises ProcedureType.REMOVE_REGION_GROUP_PROCEDURE + ProcedureFactory registration as - // well as the procedure's own serialize/deserialize. - Assert.assertEquals(procedure, ProcedureFactory.getInstance().create(buffer)); - } - } - - @Test - public void deleteTaskIdIsNegativeAndUnique() { - // The DataNode taskResultMap is keyed only by taskId and is shared with add/remove-peer tasks, - // which use a procedure's (non-negative) procId directly as the taskId. So a delete taskId must - // be strictly negative (disjoint from every procId) and distinct per (procId, deleteTaskSeq), - // otherwise a later peer op could be silently deduped against a lingering delete-task entry. - final TRegionReplicaSet regionReplicaSet = - new TRegionReplicaSet( - new TConsensusGroupId(TConsensusGroupType.DataRegion, 1), - Arrays.asList(new TDataNodeLocation())); - final java.util.Set seen = new java.util.HashSet<>(); - for (long procId : new long[] {0L, 1L, 100L, 1L << 20, (1L << 43) - 1}) { - for (long seq : new long[] {1L, 2L, 100L, (1L << 20) - 1}) { - final RemoveRegionGroupProcedure procedure = - new RemoveRegionGroupProcedure(regionReplicaSet); - procedure.setProcId(procId); - procedure.setDeleteTaskState(seq, true); - final long taskId = procedure.deleteTaskIdForTest(); - Assert.assertTrue("taskId must be negative: " + taskId, taskId < 0); - Assert.assertTrue( - "taskId must be unique for (" + procId + "," + seq + ")", seen.add(taskId)); - } - } - } -} 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 9f62a18c73a0b..28d75ba4f4940 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 @@ -2824,6 +2824,7 @@ public TSStatus deleteRegion(TConsensusGroupId tconsensusGroupId) { } ConsensusGroupId consensusGroupId = ConsensusGroupId.Factory.createFromTConsensusGroupId(tconsensusGroupId); + boolean consensusGroupDeleted = true; if (consensusGroupId instanceof DataRegionId) { try { DataRegionConsensusImpl.getInstance().deleteLocalPeer(consensusGroupId); @@ -2831,8 +2832,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); @@ -2840,11 +2843,23 @@ 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); } } + private 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(DataNodeMiscMessages.CHANGE_REGION_LEADER, 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 9704ebc7c051a..9501c8e677da2 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 @@ -138,7 +138,7 @@ public TSStatus createSchemaRegion( tsStatus.setMessage( String.format(DataNodeMiscMessages.CREATE_SCHEMA_REGION_FAILED_FMT, e2.getMessage())); } catch (final ConsensusGroupAlreadyExistException e) { - tsStatus = new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode()); + tsStatus = new TSStatus(TSStatusCode.REGION_ALREADY_EXISTS.getStatusCode()); tsStatus.setMessage( String.format( DataNodeMiscMessages.SCHEMA_REGION_ALREADY_EXISTS_FMT, schemaRegionId.getId())); @@ -171,7 +171,7 @@ public TSStatus createDataRegion(TRegionReplicaSet regionReplicaSet, String stor tsStatus.setMessage( String.format(DataNodeMiscMessages.CREATE_DATA_REGION_FAILED_FMT, e.getMessage())); } catch (ConsensusGroupAlreadyExistException e) { - tsStatus = new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode()); + tsStatus = new TSStatus(TSStatusCode.REGION_ALREADY_EXISTS.getStatusCode()); tsStatus.setMessage( String.format(DataNodeMiscMessages.DATA_REGION_ALREADY_EXISTS_FMT, dataRegionId.getId())); } catch (ConsensusException e) { @@ -208,17 +208,21 @@ public TSStatus createNewRegion(final ConsensusGroupId regionId, final String st } 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(DataNodeMiscMessages.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 ce93319a68658..9ae48b3c02149 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 @@ -335,12 +335,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(DataNodeSchemaMessages.SCHEMA_REGION_ALREADY_DELETED, schemaRegionId); - return; + return false; } schemaRegion.deleteSchemaRegion(); schemaMetricManager.removeSchemaRegionMetric(schemaRegionId.getId()); @@ -364,6 +364,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 69200271ee054..3b22d97880ff8 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 @@ -511,19 +511,17 @@ public void run() { // deletePeer: remove the peer from the consensus group TSStatus runResult = deletePeer(); if (isFailed(runResult)) { - // A failed delete must NOT fall through to taskSucceed, otherwise the ConfigNode would - // forget the task while the peer/data is still present. taskFail( taskId, tRegionId, originalDataNode, TRegionMigrateFailedType.RemoveConsensusGroupFailed, runResult); - return; } // deleteRegion: delete region data runResult = deleteRegion(); + if (isFailed(runResult)) { taskFail( taskId, @@ -531,7 +529,6 @@ public void run() { originalDataNode, TRegionMigrateFailedType.DeleteRegionFailed, runResult); - return; } taskSucceed(taskId, tRegionId, "DeletePeer"); @@ -551,13 +548,6 @@ private TSStatus deletePeer() { } else { SchemaRegionConsensusImpl.getInstance().deleteLocalPeer(regionId); } - } catch (ConsensusGroupNotExistException e) { - // The peer is already absent (e.g. a retry after a previous attempt removed it, or the - // region group is being deleted outright). Treat it as success and continue to delete data. - taskLogger.info( - "{}, The local peer of region {} does not exist, skip deleting it", - REGION_MIGRATE_PROCESS, - regionId); } catch (ConsensusException e) { String errorMsg = String.format( 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 2062ba709784c..1a0c828fd711e 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 @@ -819,9 +819,9 @@ public void createDataRegion(DataRegionId regionId, String databaseName) } } - public void deleteDataRegion(DataRegionId regionId) { + public TSStatus deleteDataRegion(DataRegionId regionId) { if (!dataRegionMap.containsKey(regionId) || deletingDataRegionMap.containsKey(regionId)) { - return; + return RpcUtils.getStatus(TSStatusCode.REGION_NOT_EXIST); } DataRegion region = deletingDataRegionMap.computeIfAbsent(regionId, k -> dataRegionMap.remove(regionId)); @@ -878,10 +878,13 @@ public void deleteDataRegion(DataRegionId regionId) { region.getDatabaseName(), region.getDataRegionIdString(), e); + return RpcUtils.getStatus(TSStatusCode.DELETE_REGION_ERROR, e.getMessage()); } finally { deletingDataRegionMap.remove(regionId); } + return RpcUtils.getStatus(TSStatusCode.SUCCESS_STATUS); } + return RpcUtils.getStatus(TSStatusCode.REGION_NOT_EXIST); } /** 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 c6059d39d049e..ff337f3454c97 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 @@ -57,12 +57,14 @@ 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.TDataNodeHeartbeatReq; import org.apache.iotdb.mpp.rpc.thrift.TDataNodeHeartbeatResp; 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; @@ -469,6 +471,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(