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 494201b1e9fec..0b7ff7eb0edb5 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
@@ -665,4 +665,16 @@ private ConfigNodeMessages() {}
public static final String
EXCEPTION_FAILED_TO_CREATE_OR_ALTER_TOPIC_MODE_CONSENSUS_DOES_NOT_SUPPORT_TOPIC_ATTRIBUTES_ARG_3C2D0BDA =
"Failed to create or alter topic, mode=consensus does not support topic attributes %s";
+ public static final String
+ LOG_REJECT_CREATEREGIONGROUPSPLAN_BECAUSE_DATABASE_ARG_DOES_NOT_EXIST_616E0CDE =
+ "Reject CreateRegionGroupsPlan because database {} does not exist";
+ public static final String
+ LOG_REJECT_CREATEREGIONGROUPSPLAN_BECAUSE_DATABASE_ARG_IS_BEING_DELETED_C085AC01 =
+ "Reject CreateRegionGroupsPlan because database {} is being deleted";
+ public static final String
+ MESSAGE_CREATE_REGIONGROUPS_FAILED_BECAUSE_DATABASE_ARG_DOES_NOT_EXIST_AF0F2440 =
+ "Create RegionGroups failed because database %s does not exist";
+ public static final String
+ MESSAGE_CREATE_REGIONGROUPS_FAILED_BECAUSE_DATABASE_ARG_IS_BEING_DELETED_651DB780 =
+ "Create RegionGroups failed because database %s is being deleted";
}
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..45361a98d8053 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
@@ -587,6 +587,11 @@ private ManagerMessages() {}
public static final String MESSAGE_SCHEMA_ENGINE_MODE_E37ED98C = "schema_engine_mode";
public static final String MESSAGE_TAG_ATTRIBUTE_TOTAL_SIZE_AF658CFE = "tag_attribute_total_size";
public static final String MESSAGE_DATABASE_LIMIT_THRESHOLD_45C23274 = "database_limit_threshold";
+ public static final String
+ MESSAGE_DATABASE_ARG_STILL_HAS_UNFINISHED_LIFECYCLE_PROCEDURES_67573924 =
+ "Database %s still has unfinished lifecycle procedures";
+ public static final String MESSAGE_SOME_OTHER_TASK_IS_DELETING_DATABASE_ARG_7BDB2C0F =
+ "Some other task is deleting database %s";
public static final String LOG_UNEXPECTED_ERROR_HAPPENED_SETTING_SPACE_QUOTA_DATABASE_ARG_F6ED7586 = "Unexpected error happened while setting space quota on database: %s ";
public static final String LOG_UNEXPECTED_ERROR_HAPPENED_SETTING_THROTTLE_QUOTA_USER_ARG_C111BE81 = "Unexpected error happened while setting throttle quota on user: %s ";
public static final String LOG_SCHEMA_TEMPLATE_NEED_TWO_FILES_1E57542A = "schema_template need two files";
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 0193e5a04854e..d5228705755b5 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
@@ -710,4 +710,16 @@ private ConfigNodeMessages() {}
public static final String
EXCEPTION_FAILED_TO_CREATE_OR_ALTER_TOPIC_MODE_CONSENSUS_DOES_NOT_SUPPORT_TOPIC_ATTRIBUTES_ARG_3C2D0BDA =
"创建或修改 topic 失败,mode=consensus 不支持 topic 属性 %s";
+ public static final String
+ LOG_REJECT_CREATEREGIONGROUPSPLAN_BECAUSE_DATABASE_ARG_DOES_NOT_EXIST_616E0CDE =
+ "拒绝 CreateRegionGroupsPlan,因为数据库 {} 不存在";
+ public static final String
+ LOG_REJECT_CREATEREGIONGROUPSPLAN_BECAUSE_DATABASE_ARG_IS_BEING_DELETED_C085AC01 =
+ "拒绝 CreateRegionGroupsPlan,因为数据库 {} 正在删除";
+ public static final String
+ MESSAGE_CREATE_REGIONGROUPS_FAILED_BECAUSE_DATABASE_ARG_DOES_NOT_EXIST_AF0F2440 =
+ "创建 RegionGroups 失败,因为数据库 %s 不存在";
+ public static final String
+ MESSAGE_CREATE_REGIONGROUPS_FAILED_BECAUSE_DATABASE_ARG_IS_BEING_DELETED_651DB780 =
+ "创建 RegionGroups 失败,因为数据库 %s 正在删除";
}
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..5ed644c0f6dab 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
@@ -579,6 +579,11 @@ private ManagerMessages() {}
public static final String MESSAGE_SCHEMA_ENGINE_MODE_E37ED98C = "schema_engine_mode";
public static final String MESSAGE_TAG_ATTRIBUTE_TOTAL_SIZE_AF658CFE = "tag_attribute_total_size";
public static final String MESSAGE_DATABASE_LIMIT_THRESHOLD_45C23274 = "database_limit_threshold";
+ public static final String
+ MESSAGE_DATABASE_ARG_STILL_HAS_UNFINISHED_LIFECYCLE_PROCEDURES_67573924 =
+ "数据库 %s 仍有未完成的生命周期流程";
+ public static final String MESSAGE_SOME_OTHER_TASK_IS_DELETING_DATABASE_ARG_7BDB2C0F =
+ "其他任务正在删除数据库 %s";
public static final String LOG_UNEXPECTED_ERROR_HAPPENED_SETTING_SPACE_QUOTA_DATABASE_ARG_F6ED7586 = "设置数据库 %s 的空间配额时发生意外错误 ";
public static final String LOG_UNEXPECTED_ERROR_HAPPENED_SETTING_THROTTLE_QUOTA_USER_ARG_C111BE81 = "设置用户 %s 的限流配额时发生意外错误 ";
public static final String LOG_SCHEMA_TEMPLATE_NEED_TWO_FILES_1E57542A = "schema_template 需要两个文件";
diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/ConfigPhysicalPlan.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/ConfigPhysicalPlan.java
index 0d9ca912571a4..8b6546f662584 100644
--- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/ConfigPhysicalPlan.java
+++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/ConfigPhysicalPlan.java
@@ -85,6 +85,7 @@
import org.apache.iotdb.confignode.consensus.request.write.procedure.UpdateProcedurePlan;
import org.apache.iotdb.confignode.consensus.request.write.quota.SetSpaceQuotaPlan;
import org.apache.iotdb.confignode.consensus.request.write.quota.SetThrottleQuotaPlan;
+import org.apache.iotdb.confignode.consensus.request.write.region.BatchRemoveRegionCreateTasksPlan;
import org.apache.iotdb.confignode.consensus.request.write.region.CreateRegionGroupsPlan;
import org.apache.iotdb.confignode.consensus.request.write.region.OfferRegionMaintainTasksPlan;
import org.apache.iotdb.confignode.consensus.request.write.region.PollRegionMaintainTaskPlan;
@@ -266,6 +267,9 @@ public static ConfigPhysicalPlan create(final ByteBuffer buffer) throws IOExcept
case PollSpecificRegionMaintainTask:
plan = new PollSpecificRegionMaintainTaskPlan();
break;
+ case BatchRemoveRegionCreateTasks:
+ plan = new BatchRemoveRegionCreateTasksPlan();
+ break;
case CreateSchemaPartition:
plan = new CreateSchemaPartitionPlan();
break;
diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/ConfigPhysicalPlanType.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/ConfigPhysicalPlanType.java
index 1be9518141483..e6f9c9847d148 100644
--- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/ConfigPhysicalPlanType.java
+++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/ConfigPhysicalPlanType.java
@@ -72,6 +72,7 @@ public enum ConfigPhysicalPlanType {
AddRegionLocation((short) 311),
RemoveRegionLocation((short) 312),
GetRegionGroupsByTime((short) 313),
+ BatchRemoveRegionCreateTasks((short) 314),
/** Partition. */
GetSchemaPartition((short) 400),
diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/write/region/BatchRemoveRegionCreateTasksPlan.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/write/region/BatchRemoveRegionCreateTasksPlan.java
new file mode 100644
index 0000000000000..256b290eb6250
--- /dev/null
+++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/consensus/request/write/region/BatchRemoveRegionCreateTasksPlan.java
@@ -0,0 +1,82 @@
+/*
+ * 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.consensus.request.write.region;
+
+import org.apache.iotdb.confignode.consensus.request.ConfigPhysicalPlan;
+import org.apache.iotdb.confignode.consensus.request.ConfigPhysicalPlanType;
+
+import org.apache.tsfile.utils.ReadWriteIOUtils;
+
+import java.io.DataOutputStream;
+import java.io.IOException;
+import java.nio.ByteBuffer;
+import java.util.Objects;
+
+/**
+ * Removes every queued RegionCreateTask that belongs to the specified pre-deleted database.
+ *
+ *
The state machine ignores this plan when the database is missing or active, so replaying a
+ * cancellation cannot affect a later database incarnation that reuses the same name.
+ */
+public class BatchRemoveRegionCreateTasksPlan extends ConfigPhysicalPlan {
+
+ private String database;
+
+ public BatchRemoveRegionCreateTasksPlan() {
+ super(ConfigPhysicalPlanType.BatchRemoveRegionCreateTasks);
+ }
+
+ public BatchRemoveRegionCreateTasksPlan(final String database) {
+ super(ConfigPhysicalPlanType.BatchRemoveRegionCreateTasks);
+ this.database = database;
+ }
+
+ public String getDatabase() {
+ return database;
+ }
+
+ @Override
+ protected void serializeImpl(DataOutputStream stream) throws IOException {
+ stream.writeShort(getType().getPlanType());
+ ReadWriteIOUtils.write(database, stream);
+ }
+
+ @Override
+ protected void deserializeImpl(ByteBuffer buffer) throws IOException {
+ database = ReadWriteIOUtils.readString(buffer);
+ }
+
+ @Override
+ public boolean equals(Object o) {
+ if (this == o) {
+ return true;
+ }
+ if (!(o instanceof BatchRemoveRegionCreateTasksPlan)) {
+ return false;
+ }
+ final BatchRemoveRegionCreateTasksPlan that = (BatchRemoveRegionCreateTasksPlan) o;
+ return Objects.equals(database, that.database);
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(database);
+ }
+}
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..5962473bef20b 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
@@ -220,6 +220,7 @@ public class ProcedureManager {
private final PartitionTableAutoCleaner partitionTableCleaner;
private final ReentrantLock tableLock = new ReentrantLock();
+ private final ReentrantLock databaseLifecycleAdmissionLock = new ReentrantLock();
public ProcedureManager(ConfigManager configManager, ProcedureInfo procedureInfo) {
this.configManager = configManager;
@@ -1501,18 +1502,27 @@ private static TDataNodeLocation buildFakeDataNodeLocation(int dataNodeId, Strin
// endregion
+ public CreateRegionGroupsProcedure submitCreateRegionGroups(
+ final TConsensusGroupType consensusGroupType,
+ final CreateRegionGroupsPlan createRegionGroupsPlan) {
+ final CreateRegionGroupsProcedure procedure =
+ new CreateRegionGroupsProcedure(consensusGroupType, createRegionGroupsPlan);
+ // Reentrant for PartitionManager, which holds the same admission lock while allocating the
+ // plan. Keeping this guard here also makes a future direct submission visible atomically to
+ // same-name database creation.
+ try (final AutoCloseableLock ignored = acquireDatabaseLifecycleAdmissionLock()) {
+ executor.submitProcedure(procedure);
+ }
+ return procedure;
+ }
+
/**
- * Generate {@link CreateRegionGroupsProcedure} and wait until it finished.
+ * Wait until a {@link CreateRegionGroupsProcedure} finishes.
*
* @return {@link TSStatusCode#SUCCESS_STATUS} if all RegionGroups have been created successfully,
* {@link TSStatusCode#CREATE_REGION_ERROR} otherwise
*/
- public TSStatus createRegionGroups(
- final TConsensusGroupType consensusGroupType,
- final CreateRegionGroupsPlan createRegionGroupsPlan) {
- final CreateRegionGroupsProcedure procedure =
- new CreateRegionGroupsProcedure(consensusGroupType, createRegionGroupsPlan);
- executor.submitProcedure(procedure);
+ public TSStatus waitCreateRegionGroups(final CreateRegionGroupsProcedure procedure) {
final TSStatus status = waitingProcedureFinished(procedure);
if (status.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
return status;
@@ -1522,6 +1532,21 @@ public TSStatus createRegionGroups(
}
}
+ public AutoCloseableLock acquireDatabaseLifecycleAdmissionLock() {
+ return AutoCloseableLock.acquire(databaseLifecycleAdmissionLock);
+ }
+
+ public boolean hasUnfinishedDatabaseLifecycleProcedure(final String database) {
+ return executor.getProcedures().values().stream()
+ .filter(procedure -> !procedure.isFinished())
+ .anyMatch(
+ procedure ->
+ (procedure instanceof DeleteDatabaseProcedure
+ && database.equals(((DeleteDatabaseProcedure) procedure).getDatabase()))
+ || (procedure instanceof CreateRegionGroupsProcedure
+ && ((CreateRegionGroupsProcedure) procedure).containsDatabase(database)));
+ }
+
/**
* Generate {@link CreateTriggerProcedure} and wait until it finished.
*
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..f86684a114dca 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
@@ -27,6 +27,7 @@
import org.apache.iotdb.common.rpc.thrift.TSeriesPartitionSlot;
import org.apache.iotdb.common.rpc.thrift.TTimePartitionSlot;
import org.apache.iotdb.commons.cluster.RegionRoleType;
+import org.apache.iotdb.commons.cluster.RegionStatus;
import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory;
import org.apache.iotdb.commons.concurrent.ThreadName;
import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil;
@@ -59,6 +60,7 @@
import org.apache.iotdb.confignode.consensus.request.write.partition.CreateDataPartitionPlan;
import org.apache.iotdb.confignode.consensus.request.write.partition.CreateSchemaPartitionPlan;
import org.apache.iotdb.confignode.consensus.request.write.partition.RemoveRegionLocationPlan;
+import org.apache.iotdb.confignode.consensus.request.write.region.BatchRemoveRegionCreateTasksPlan;
import org.apache.iotdb.confignode.consensus.request.write.region.CreateRegionGroupsPlan;
import org.apache.iotdb.confignode.consensus.request.write.region.PollSpecificRegionMaintainTaskPlan;
import org.apache.iotdb.confignode.consensus.response.partition.CountTimeSlotListResp;
@@ -86,6 +88,7 @@
import org.apache.iotdb.confignode.persistence.partition.maintainer.RegionCreateTask;
import org.apache.iotdb.confignode.persistence.partition.maintainer.RegionMaintainTask;
import org.apache.iotdb.confignode.procedure.impl.partition.DataPartitionTableIntegrityCheckProcedure;
+import org.apache.iotdb.confignode.procedure.impl.region.CreateRegionGroupsProcedure;
import org.apache.iotdb.confignode.rpc.thrift.TCountTimeSlotListReq;
import org.apache.iotdb.confignode.rpc.thrift.TGetRegionGroupsByTimeReq;
import org.apache.iotdb.confignode.rpc.thrift.TGetRegionIdReq;
@@ -99,6 +102,7 @@
import org.apache.iotdb.rpc.RpcUtils;
import org.apache.iotdb.rpc.TSStatusCode;
+import org.apache.ratis.util.AutoCloseableLock;
import org.apache.tsfile.file.metadata.IDeviceID;
import org.apache.tsfile.file.metadata.IDeviceID.Deserializer;
import org.apache.tsfile.utils.Pair;
@@ -123,9 +127,11 @@
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.Future;
import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.ThreadLocalRandom;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.locks.ReentrantLock;
import java.util.stream.Collectors;
/** The {@link PartitionManager} manages cluster PartitionTable read and write requests. */
@@ -154,8 +160,17 @@ public class PartitionManager {
// Try to delete Regions in every 10s
private static final int REGION_MAINTAINER_WORK_INTERVAL = 10;
+ private static final int SCHEMA_REGION_CREATE_BATCH_SIZE_PER_DATA_NODE = 32;
+ private static final int DATA_REGION_CREATE_BATCH_SIZE_PER_DATA_NODE = 64;
+ private static final long REGION_CREATE_BACKOFF_BASE_NANOS =
+ TimeUnit.SECONDS.toNanos(REGION_MAINTAINER_WORK_INTERVAL);
+ private static final long REGION_CREATE_BACKOFF_MAX_NANOS = TimeUnit.MINUTES.toNanos(5);
+
private final ScheduledExecutorService regionMaintainer;
private Future> currentRegionMaintainerFuture;
+ private final ReentrantLock regionCreateTaskLock = new ReentrantLock();
+ private final Map> regionCreateBackoffMap =
+ new EnumMap<>(TConsensusGroupType.class);
private final AtomicBoolean dataPartitionTableIntegrityCheckProcedureRunning =
new AtomicBoolean(false);
@@ -741,11 +756,22 @@ private TSStatus generateAndAllocateRegionGroups(
final Map allotmentMap, final TConsensusGroupType consensusGroupType)
throws NotEnoughDataNodeException, DatabaseNotExistsException {
if (!allotmentMap.isEmpty()) {
- final CreateRegionGroupsPlan createRegionGroupsPlan =
- getLoadManager().allocateRegionGroups(allotmentMap, consensusGroupType);
- LOGGER.info(ManagerMessages.CREATEREGIONGROUPS_STARTING_TO_CREATE_THE_FOLLOWING_REGIONGROUPS);
- createRegionGroupsPlan.planLog(LOGGER);
- return getProcedureManager().createRegionGroups(consensusGroupType, createRegionGroupsPlan);
+ final CreateRegionGroupsProcedure procedure;
+ // Database creation uses the same admission lock when checking unfinished lifecycle
+ // procedures. Cover both ID allocation and submission so a delayed plan can never become
+ // invisible between same-name database deletion and recreation.
+ try (final AutoCloseableLock ignored =
+ getProcedureManager().acquireDatabaseLifecycleAdmissionLock()) {
+ final CreateRegionGroupsPlan createRegionGroupsPlan =
+ getLoadManager().allocateRegionGroups(allotmentMap, consensusGroupType);
+ LOGGER.info(
+ ManagerMessages.CREATEREGIONGROUPS_STARTING_TO_CREATE_THE_FOLLOWING_REGIONGROUPS);
+ createRegionGroupsPlan.planLog(LOGGER);
+ procedure =
+ getProcedureManager()
+ .submitCreateRegionGroups(consensusGroupType, createRegionGroupsPlan);
+ }
+ return getProcedureManager().waitCreateRegionGroups(procedure);
} else {
return RpcUtils.SUCCESS_STATUS;
}
@@ -1117,14 +1143,35 @@ public SchemaNodeManagementResp getNodePathsPartition(GetNodePathsPartitionPlan
}
}
- public void preDeleteDatabase(
+ public TSStatus preDeleteDatabase(
final String database, final PreDeleteDatabasePlan.PreDeleteType preDeleteType) {
final PreDeleteDatabasePlan preDeleteDatabasePlan =
new PreDeleteDatabasePlan(database, preDeleteType);
+ regionCreateTaskLock.lock();
+ try {
+ return getConsensusManager().write(preDeleteDatabasePlan);
+ } catch (final ConsensusException e) {
+ LOGGER.warn(CONSENSUS_WRITE_ERROR, e);
+ final TSStatus status = new TSStatus(TSStatusCode.EXECUTE_STATEMENT_ERROR.getStatusCode());
+ status.setMessage(e.getMessage());
+ return status;
+ } finally {
+ regionCreateTaskLock.unlock();
+ }
+ }
+
+ /** Durably removes all queued RegionCreateTasks of the specified database. */
+ public TSStatus batchRemoveRegionCreateTasks(final String database) {
+ regionCreateTaskLock.lock();
try {
- getConsensusManager().write(preDeleteDatabasePlan);
+ return getConsensusManager().write(new BatchRemoveRegionCreateTasksPlan(database));
} catch (final ConsensusException e) {
LOGGER.warn(CONSENSUS_WRITE_ERROR, e);
+ final TSStatus status = new TSStatus(TSStatusCode.EXECUTE_STATEMENT_ERROR.getStatusCode());
+ status.setMessage(e.getMessage());
+ return status;
+ } finally {
+ regionCreateTaskLock.unlock();
}
}
@@ -1132,6 +1179,10 @@ public boolean isDatabasePreDeleted(final String database) {
return partitionInfo.isDatabasePreDeleted(database);
}
+ public TSStatus validateCreateRegionGroups(final CreateRegionGroupsPlan plan) {
+ return partitionInfo.validateCreateRegionGroups(plan);
+ }
+
/**
* Get TSeriesPartitionSlot.
*
@@ -1365,11 +1416,46 @@ public void maintainRegionReplicas() {
return;
}
+ regionCreateTaskLock.lock();
+ try {
+ // Leadership may have changed while this invocation was waiting for an in-flight DROP.
+ if (!getConsensusManager().isLeader()) {
+ return;
+ }
+ maintainRegionReplicasUnderLock();
+ } finally {
+ regionCreateTaskLock.unlock();
+ }
+ }
+
+ private void maintainRegionReplicasUnderLock() {
+ final List persistedTasks = partitionInfo.getRegionMaintainEntryList();
+ final Map invalidTaskCountByRegion = new HashMap<>();
+ for (RegionMaintainTask task : persistedTasks) {
+ if (!(task instanceof RegionCreateTask)
+ || !isRegionCreateTaskRegionValid((RegionCreateTask) task)) {
+ invalidTaskCountByRegion.merge(task.getRegionId(), 1, Integer::sum);
+ }
+ }
+ if (!removeInvalidRegionCreateTasks(invalidTaskCountByRegion)) {
+ return;
+ }
+ final Set invalidRegionIds = invalidTaskCountByRegion.keySet();
+
+ // Do not infer that an Unknown cache entry means a missing replica until the new leader has
+ // collected enough heartbeats. Orphan/pre-deleted tasks above are still cleaned immediately.
+ if (!getLoadManager().isLoadReady()) {
+ 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()) {
+ final Map> tasksByRegion = new LinkedHashMap<>();
+ for (RegionMaintainTask task : persistedTasks) {
+ if (invalidRegionIds.contains(task.getRegionId())) {
+ continue;
+ }
if (!(task instanceof RegionCreateTask)) {
// Unreachable: the queue only holds create tasks now (legacy delete tasks are dropped at
// the
@@ -1383,51 +1469,105 @@ public void maintainRegionReplicas() {
.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);
+ final Set invalidHeadRegionIds = new HashSet<>();
+ final Map> headsByType =
+ new EnumMap<>(TConsensusGroupType.class);
+ final Map> selectedCountByTypeAndDataNode =
+ new EnumMap<>(TConsensusGroupType.class);
+ for (Queue queue : tasksByRegion.values()) {
+ final RegionCreateTask head = queue.peek();
+ if (!isRegionCreateTaskTargetValid(head)) {
+ invalidHeadRegionIds.add(head.getRegionId());
+ continue;
}
-
- 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 (isRegionCreateTargetInBackoff(head)) {
+ continue;
}
-
- if (successfulRegions.isEmpty()) {
- break;
+ final TConsensusGroupType type = head.getRegionId().getType();
+ final int dataNodeId = head.getTargetDataNode().getDataNodeId();
+ final Map selectedCountByDataNode =
+ selectedCountByTypeAndDataNode.computeIfAbsent(type, ignored -> new HashMap<>());
+ final int selectedCount = selectedCountByDataNode.getOrDefault(dataNodeId, 0);
+ if (selectedCount >= getRegionCreateBatchSize(type)) {
+ continue;
}
+ selectedCountByDataNode.put(dataNodeId, selectedCount + 1);
+ headsByType.computeIfAbsent(type, ignored -> new ArrayList<>()).add(head);
+ }
- // 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;
- });
- }
+ // A target-specific stale task only removes the head of its Region queue. A following task of
+ // the same Region may still point to another replica that is genuinely missing.
+ if (!invalidHeadRegionIds.isEmpty()
+ && !writeRegionCreateTaskPlan(
+ new PollSpecificRegionMaintainTaskPlan(invalidHeadRegionIds))) {
+ return;
+ }
- // 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);
- }
+ final Set successfulRegions = new HashSet<>();
+ for (Map.Entry> entry : headsByType.entrySet()) {
+ successfulRegions.addAll(submitRegionCreateTasks(entry.getKey(), entry.getValue()));
+ }
+ if (!successfulRegions.isEmpty()) {
+ writeRegionCreateTaskPlan(new PollSpecificRegionMaintainTaskPlan(successfulRegions));
+ }
+ }
- 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;
+ private boolean writeRegionCreateTaskPlan(ConfigPhysicalPlan plan) {
+ try {
+ return getConsensusManager().write(plan).getCode()
+ == TSStatusCode.SUCCESS_STATUS.getStatusCode();
+ } catch (ConsensusException e) {
+ LOGGER.warn(CONSENSUS_WRITE_ERROR, e);
+ return false;
+ }
+ }
+
+ private boolean removeInvalidRegionCreateTasks(
+ final Map invalidTaskCountByRegion) {
+ if (invalidTaskCountByRegion.isEmpty()) {
+ return true;
+ }
+ final int maxTaskCount =
+ invalidTaskCountByRegion.values().stream().mapToInt(Integer::intValue).max().orElse(0);
+ for (int taskIndex = 0; taskIndex < maxTaskCount; taskIndex++) {
+ final int currentTaskIndex = taskIndex;
+ final Set regionIds =
+ invalidTaskCountByRegion.entrySet().stream()
+ .filter(entry -> entry.getValue() > currentTaskIndex)
+ .map(Map.Entry::getKey)
+ .collect(Collectors.toSet());
+ if (!writeRegionCreateTaskPlan(new PollSpecificRegionMaintainTaskPlan(regionIds))) {
+ return false;
}
}
+ return true;
+ }
+
+ private boolean isRegionCreateTaskRegionValid(RegionCreateTask task) {
+ return partitionInfo.isDatabaseExisted(task.getStorageGroup())
+ && Objects.equals(
+ task.getStorageGroup(), partitionInfo.getRegionDatabase(task.getRegionId()));
+ }
+
+ private boolean isRegionCreateTaskTargetValid(RegionCreateTask task) {
+ final List currentReplicaSets =
+ partitionInfo.getReplicaSets(
+ task.getStorageGroup(), Collections.singletonList(task.getRegionId()));
+ if (currentReplicaSets.size() != 1
+ || currentReplicaSets.get(0).getDataNodeLocations().stream()
+ .noneMatch(
+ location -> location.getDataNodeId() == task.getTargetDataNode().getDataNodeId())) {
+ return false;
+ }
+ return RegionStatus.Unknown.equals(
+ getLoadManager()
+ .getRegionStatus(task.getRegionId(), task.getTargetDataNode().getDataNodeId()));
+ }
+
+ private int getRegionCreateBatchSize(TConsensusGroupType regionType) {
+ return TConsensusGroupType.SchemaRegion.equals(regionType)
+ ? SCHEMA_REGION_CREATE_BATCH_SIZE_PER_DATA_NODE
+ : DATA_REGION_CREATE_BATCH_SIZE_PER_DATA_NODE;
}
/**
@@ -1451,10 +1591,8 @@ private Set submitRegionCreateTasks(
new TCreateSchemaRegionReq(task.getRegionReplicaSet(), task.getStorageGroup()));
schemaHandler.putNodeLocation(task.getRegionId().getId(), task.getTargetDataNode());
}
- CnToDnInternalServiceAsyncRequestManager.getInstance()
- .sendAsyncRequestWithRetry(schemaHandler);
- collectSuccessfulRegions(
- schemaHandler.getResponseMap(), TConsensusGroupType.SchemaRegion, successfulRegions);
+ CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequest(schemaHandler);
+ collectSuccessfulRegions(schemaHandler.getResponseMap(), createTasks, successfulRegions);
break;
case DataRegion:
final DataNodeAsyncRequestContext dataHandler =
@@ -1469,10 +1607,8 @@ private Set submitRegionCreateTasks(
new TCreateDataRegionReq(task.getRegionReplicaSet(), task.getStorageGroup()));
dataHandler.putNodeLocation(task.getRegionId().getId(), task.getTargetDataNode());
}
- CnToDnInternalServiceAsyncRequestManager.getInstance()
- .sendAsyncRequestWithRetry(dataHandler);
- collectSuccessfulRegions(
- dataHandler.getResponseMap(), TConsensusGroupType.DataRegion, successfulRegions);
+ CnToDnInternalServiceAsyncRequestManager.getInstance().sendAsyncRequest(dataHandler);
+ collectSuccessfulRegions(dataHandler.getResponseMap(), createTasks, successfulRegions);
break;
default:
break;
@@ -1482,13 +1618,92 @@ private Set submitRegionCreateTasks(
private void collectSuccessfulRegions(
Map responseMap,
- TConsensusGroupType regionType,
+ List createTasks,
Set successfulRegions) {
- for (Map.Entry entry : responseMap.entrySet()) {
- if (entry.getValue().getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
- successfulRegions.add(new TConsensusGroupId(regionType, entry.getKey()));
+ final Map> tasksByDataNode = new HashMap<>();
+ for (RegionCreateTask task : createTasks) {
+ tasksByDataNode
+ .computeIfAbsent(task.getTargetDataNode().getDataNodeId(), ignored -> new ArrayList<>())
+ .add(task);
+ final TSStatus status = responseMap.get(task.getRegionId().getId());
+ if (status != null && status.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ successfulRegions.add(task.getRegionId());
+ }
+ }
+ for (List dataNodeTasks : tasksByDataNode.values()) {
+ RegionCreateTask failedTask = null;
+ TSStatus failedStatus = null;
+ for (RegionCreateTask task : dataNodeTasks) {
+ final TSStatus status = responseMap.get(task.getRegionId().getId());
+ if (status == null || status.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ if (failedTask == null || isDirectMemoryFailure(status)) {
+ failedTask = task;
+ failedStatus = status;
+ }
+ }
}
+ if (failedTask != null) {
+ recordRegionCreateFailure(failedTask, failedStatus);
+ } else {
+ clearRegionCreateBackoff(dataNodeTasks.get(0));
+ }
+ }
+ }
+
+ private boolean isRegionCreateTargetInBackoff(RegionCreateTask task) {
+ return Optional.ofNullable(regionCreateBackoffMap.get(task.getRegionId().getType()))
+ .map(backoffByDataNode -> backoffByDataNode.get(task.getTargetDataNode().getDataNodeId()))
+ .map(backoff -> backoff.nextAttemptNanos > System.nanoTime())
+ .orElse(false);
+ }
+
+ private void clearRegionCreateBackoff(RegionCreateTask task) {
+ Optional.ofNullable(regionCreateBackoffMap.get(task.getRegionId().getType()))
+ .ifPresent(
+ backoffByDataNode ->
+ backoffByDataNode.remove(task.getTargetDataNode().getDataNodeId()));
+ }
+
+ private void recordRegionCreateFailure(RegionCreateTask task, TSStatus status) {
+ final RegionCreateBackoff backoff =
+ regionCreateBackoffMap
+ .computeIfAbsent(task.getRegionId().getType(), ignored -> new HashMap<>())
+ .computeIfAbsent(
+ task.getTargetDataNode().getDataNodeId(), ignored -> new RegionCreateBackoff());
+ backoff.failureCount++;
+ long delayNanos =
+ Math.min(
+ REGION_CREATE_BACKOFF_MAX_NANOS,
+ REGION_CREATE_BACKOFF_BASE_NANOS << Math.min(backoff.failureCount - 1, 5));
+ if (isDirectMemoryFailure(status)) {
+ delayNanos = REGION_CREATE_BACKOFF_MAX_NANOS;
+ } else {
+ final long jitterBound = Math.max(1, delayNanos / 4);
+ delayNanos =
+ Math.min(
+ REGION_CREATE_BACKOFF_MAX_NANOS,
+ delayNanos + ThreadLocalRandom.current().nextLong(jitterBound));
+ }
+ backoff.nextAttemptNanos = System.nanoTime() + delayNanos;
+ }
+
+ private boolean isDirectMemoryFailure(TSStatus status) {
+ if (status == null || status.getMessage() == null) {
+ return false;
}
+ final String normalizedMessage = status.getMessage().toLowerCase(java.util.Locale.ROOT);
+ return normalizedMessage.contains("direct memory")
+ || normalizedMessage.contains("direct buffer")
+ || normalizedMessage.contains("directbuffer")
+ || normalizedMessage.contains("outofmemory")
+ || normalizedMessage.contains("out of memory")
+ || normalizedMessage.contains("oom")
+ || normalizedMessage.contains("内存");
+ }
+
+ private static class RegionCreateBackoff {
+ private int failureCount;
+ private long nextAttemptNanos;
}
public void startRegionCleaner() {
diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/schema/ClusterSchemaManager.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/schema/ClusterSchemaManager.java
index 357c14eca4aa0..431639f5aef08 100644
--- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/schema/ClusterSchemaManager.java
+++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/schema/ClusterSchemaManager.java
@@ -92,6 +92,7 @@
import org.apache.iotdb.confignode.i18n.ManagerMessages;
import org.apache.iotdb.confignode.i18n.ProcedureMessages;
import org.apache.iotdb.confignode.manager.IManager;
+import org.apache.iotdb.confignode.manager.ProcedureManager;
import org.apache.iotdb.confignode.manager.consensus.ConsensusManager;
import org.apache.iotdb.confignode.manager.node.NodeManager;
import org.apache.iotdb.confignode.manager.partition.PartitionManager;
@@ -118,6 +119,7 @@
import org.apache.iotdb.rpc.RpcUtils;
import org.apache.iotdb.rpc.TSStatusCode;
+import org.apache.ratis.util.AutoCloseableLock;
import org.apache.tsfile.annotations.TableModel;
import org.apache.tsfile.enums.TSDataType;
import org.apache.tsfile.file.metadata.IDeviceID;
@@ -171,57 +173,70 @@ public ClusterSchemaManager(
/** Set Database */
public TSStatus setDatabase(
final DatabaseSchemaPlan databaseSchemaPlan, final boolean isGeneratedByPipe) {
- TSStatus result;
-
final TDatabaseSchema schema = databaseSchemaPlan.getSchema();
- if (getPartitionManager().isDatabasePreDeleted(schema.getName())) {
- return RpcUtils.getStatus(
- TSStatusCode.METADATA_ERROR,
- String.format("Some other task is deleting database %s", schema.getName()));
- }
-
- createDatabaseLock.lock();
- try {
- clusterSchemaInfo.isDatabaseNameValid(
- schema.getName(), schema.isSetIsTableModel() && schema.isIsTableModel());
- if (!schema.getName().equals(SchemaConstant.SYSTEM_DATABASE)
- && !schema.getName().equals(SchemaConstant.AUDIT_DATABASE)
- && !schema.getName().equals(Audit.TABLE_MODEL_AUDIT_DATABASE)) {
- clusterSchemaInfo.checkDatabaseLimit();
+ final ProcedureManager procedureManager = configManager.getProcedureManager();
+ try (final AutoCloseableLock ignored =
+ procedureManager.acquireDatabaseLifecycleAdmissionLock()) {
+ if (procedureManager.hasUnfinishedDatabaseLifecycleProcedure(schema.getName())) {
+ return RpcUtils.getStatus(
+ TSStatusCode.METADATA_ERROR,
+ String.format(
+ ManagerMessages
+ .MESSAGE_DATABASE_ARG_STILL_HAS_UNFINISHED_LIFECYCLE_PROCEDURES_67573924,
+ schema.getName()));
}
- // Cache DatabaseSchema
- result =
- getConsensusManager()
- .write(
- isGeneratedByPipe
- ? new PipeEnrichedPlan(databaseSchemaPlan)
- : databaseSchemaPlan);
- // set ttl
- if (schema.isSetTTL()) {
- result = configManager.getTTLManager().setTTL(databaseSchemaPlan, isGeneratedByPipe);
+ if (getPartitionManager().isDatabasePreDeleted(schema.getName())) {
+ return RpcUtils.getStatus(
+ TSStatusCode.METADATA_ERROR,
+ String.format(
+ ManagerMessages.MESSAGE_SOME_OTHER_TASK_IS_DELETING_DATABASE_ARG_7BDB2C0F,
+ schema.getName()));
}
- // Bind Database metrics
- PartitionMetrics.bindDatabaseRelatedMetricsWhenUpdate(
- MetricService.getInstance(),
- configManager,
- schema.getName(),
- schema.getDataReplicationFactor(),
- schema.getSchemaReplicationFactor());
- // Adjust the maximum RegionGroup number of each Database
- adjustMaxRegionGroupNum();
- } catch (final ConsensusException e) {
- LOGGER.warn(CONSENSUS_WRITE_ERROR, e);
- result = new TSStatus(TSStatusCode.EXECUTE_STATEMENT_ERROR.getStatusCode());
- result.setMessage(e.getMessage());
- } catch (final MetadataException metadataException) {
- // Reject if Database already set
- result = new TSStatus(metadataException.getErrorCode());
- result.setMessage(metadataException.getMessage());
- } finally {
- createDatabaseLock.unlock();
- }
- return result;
+ TSStatus result;
+ createDatabaseLock.lock();
+ try {
+ clusterSchemaInfo.isDatabaseNameValid(
+ schema.getName(), schema.isSetIsTableModel() && schema.isIsTableModel());
+ if (!schema.getName().equals(SchemaConstant.SYSTEM_DATABASE)
+ && !schema.getName().equals(SchemaConstant.AUDIT_DATABASE)
+ && !schema.getName().equals(Audit.TABLE_MODEL_AUDIT_DATABASE)) {
+ clusterSchemaInfo.checkDatabaseLimit();
+ }
+ // Cache DatabaseSchema
+ result =
+ getConsensusManager()
+ .write(
+ isGeneratedByPipe
+ ? new PipeEnrichedPlan(databaseSchemaPlan)
+ : databaseSchemaPlan);
+ // set ttl
+ if (schema.isSetTTL()) {
+ result = configManager.getTTLManager().setTTL(databaseSchemaPlan, isGeneratedByPipe);
+ }
+ // Bind Database metrics
+ PartitionMetrics.bindDatabaseRelatedMetricsWhenUpdate(
+ MetricService.getInstance(),
+ configManager,
+ schema.getName(),
+ schema.getDataReplicationFactor(),
+ schema.getSchemaReplicationFactor());
+ // Adjust the maximum RegionGroup number of each Database
+ adjustMaxRegionGroupNum();
+ } catch (final ConsensusException e) {
+ LOGGER.warn(CONSENSUS_WRITE_ERROR, e);
+ result = new TSStatus(TSStatusCode.EXECUTE_STATEMENT_ERROR.getStatusCode());
+ result.setMessage(e.getMessage());
+ } catch (final MetadataException metadataException) {
+ // Reject if Database already set
+ result = new TSStatus(metadataException.getErrorCode());
+ result.setMessage(metadataException.getMessage());
+ } finally {
+ createDatabaseLock.unlock();
+ }
+
+ return result;
+ }
}
/** Alter Database */
diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/executor/ConfigPlanExecutor.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/executor/ConfigPlanExecutor.java
index 772f46baa3162..7fe8923855ac5 100644
--- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/executor/ConfigPlanExecutor.java
+++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/persistence/executor/ConfigPlanExecutor.java
@@ -111,6 +111,7 @@
import org.apache.iotdb.confignode.consensus.request.write.procedure.UpdateProcedurePlan;
import org.apache.iotdb.confignode.consensus.request.write.quota.SetSpaceQuotaPlan;
import org.apache.iotdb.confignode.consensus.request.write.quota.SetThrottleQuotaPlan;
+import org.apache.iotdb.confignode.consensus.request.write.region.BatchRemoveRegionCreateTasksPlan;
import org.apache.iotdb.confignode.consensus.request.write.region.CreateRegionGroupsPlan;
import org.apache.iotdb.confignode.consensus.request.write.region.OfferRegionMaintainTasksPlan;
import org.apache.iotdb.confignode.consensus.request.write.region.PollSpecificRegionMaintainTaskPlan;
@@ -454,6 +455,9 @@ public TSStatus executeNonQueryPlan(ConfigPhysicalPlan physicalPlan)
case PollSpecificRegionMaintainTask:
return partitionInfo.pollSpecificRegionMaintainTask(
(PollSpecificRegionMaintainTaskPlan) physicalPlan);
+ case BatchRemoveRegionCreateTasks:
+ return partitionInfo.batchRemoveRegionCreateTasks(
+ (BatchRemoveRegionCreateTasksPlan) physicalPlan);
case CreateSchemaPartition:
return partitionInfo.createSchemaPartition((CreateSchemaPartitionPlan) physicalPlan);
case CreateDataPartition:
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..40f402b99ae7f 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
@@ -49,6 +49,7 @@
import org.apache.iotdb.confignode.consensus.request.write.partition.CreateSchemaPartitionPlan;
import org.apache.iotdb.confignode.consensus.request.write.partition.RemoveRegionLocationPlan;
import org.apache.iotdb.confignode.consensus.request.write.partition.UpdateRegionLocationPlan;
+import org.apache.iotdb.confignode.consensus.request.write.region.BatchRemoveRegionCreateTasksPlan;
import org.apache.iotdb.confignode.consensus.request.write.region.CreateRegionGroupsPlan;
import org.apache.iotdb.confignode.consensus.request.write.region.OfferRegionMaintainTasksPlan;
import org.apache.iotdb.confignode.consensus.request.write.region.PollSpecificRegionMaintainTaskPlan;
@@ -63,6 +64,7 @@
import org.apache.iotdb.confignode.consensus.response.partition.SchemaPartitionResp;
import org.apache.iotdb.confignode.exception.DatabaseNotExistsException;
import org.apache.iotdb.confignode.i18n.ConfigNodeMessages;
+import org.apache.iotdb.confignode.persistence.partition.maintainer.RegionCreateTask;
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;
@@ -192,37 +194,71 @@ public TSStatus createDatabase(final DatabaseSchemaPlan plan) {
* @return {@link TSStatusCode#SUCCESS_STATUS}
*/
public TSStatus createRegionGroups(CreateRegionGroupsPlan plan) {
- TSStatus result;
- AtomicInteger maxRegionId = new AtomicInteger(Integer.MIN_VALUE);
+ updateNextRegionGroupId(plan);
+
+ final TSStatus validationStatus = validateCreateRegionGroups(plan);
+ if (validationStatus.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ return validationStatus;
+ }
plan.getRegionGroupMap()
.forEach(
(database, regionReplicaSets) -> {
- if (isDatabasePreDeleted(database)) {
- LOGGER.warn(
- ConfigNodeMessages
- .CREATEREGIONGROUPS_DATABASE_HAS_BEEN_DELETED_CORRESPONDING_REGIONGROUPS,
- database);
- return;
- }
databasePartitionTables.get(database).createRegionGroups(regionReplicaSets);
- regionReplicaSets.forEach(
- regionReplicaSet ->
- maxRegionId.set(
- Math.max(maxRegionId.get(), regionReplicaSet.getRegionId().getId())));
});
+ return new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
+ }
+
+ /** Validates all databases before any RegionGroup in a potentially batched plan is persisted. */
+ public TSStatus validateCreateRegionGroups(final CreateRegionGroupsPlan plan) {
+ for (final String database : plan.getRegionGroupMap().keySet()) {
+ final DatabasePartitionTable databasePartitionTable = databasePartitionTables.get(database);
+ if (databasePartitionTable == null) {
+ LOGGER.warn(
+ ConfigNodeMessages
+ .LOG_REJECT_CREATEREGIONGROUPSPLAN_BECAUSE_DATABASE_ARG_DOES_NOT_EXIST_616E0CDE,
+ database);
+ return new TSStatus(TSStatusCode.DATABASE_NOT_EXIST.getStatusCode())
+ .setMessage(
+ String.format(
+ ConfigNodeMessages
+ .MESSAGE_CREATE_REGIONGROUPS_FAILED_BECAUSE_DATABASE_ARG_DOES_NOT_EXIST_AF0F2440,
+ database));
+ }
+ if (!databasePartitionTable.isNotPreDeleted()) {
+ LOGGER.warn(
+ ConfigNodeMessages
+ .LOG_REJECT_CREATEREGIONGROUPSPLAN_BECAUSE_DATABASE_ARG_IS_BEING_DELETED_C085AC01,
+ database);
+ return new TSStatus(TSStatusCode.DATABASE_NOT_EXIST.getStatusCode())
+ .setMessage(
+ String.format(
+ ConfigNodeMessages
+ .MESSAGE_CREATE_REGIONGROUPS_FAILED_BECAUSE_DATABASE_ARG_IS_BEING_DELETED_651DB780,
+ database));
+ }
+ }
+
+ return new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
+ }
+
+ private void updateNextRegionGroupId(final CreateRegionGroupsPlan plan) {
+ final int maxRegionId =
+ plan.getRegionGroupMap().values().stream()
+ .flatMap(List::stream)
+ .mapToInt(regionReplicaSet -> regionReplicaSet.getRegionId().getId())
+ .max()
+ .orElse(Integer.MIN_VALUE);
+
// To ensure that the nextRegionGroupId is updated correctly when
// the ConfigNode-followers concurrently processes CreateRegionsPlan,
// we need to add a synchronization lock here
synchronized (nextRegionGroupId) {
- if (nextRegionGroupId.get() < maxRegionId.get()) {
- nextRegionGroupId.set(maxRegionId.get());
+ if (nextRegionGroupId.get() < maxRegionId) {
+ nextRegionGroupId.set(maxRegionId);
}
}
-
- result = new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
- return result;
}
/**
@@ -245,6 +281,10 @@ public TSStatus offerRegionMaintainTasks(
task.getRegionId());
continue;
}
+ final RegionCreateTask createTask = (RegionCreateTask) task;
+ if (!isRegionCreateTaskOwnedByCurrentPartitionTable(createTask)) {
+ continue;
+ }
regionMaintainTaskList.add(task);
}
return RpcUtils.SUCCESS_STATUS;
@@ -290,6 +330,37 @@ public TSStatus pollSpecificRegionMaintainTask(PollSpecificRegionMaintainTaskPla
}
}
+ /** Idempotently remove all RegionCreateTasks that belong to the specified database. */
+ public TSStatus batchRemoveRegionCreateTasks(BatchRemoveRegionCreateTasksPlan plan) {
+ synchronized (regionMaintainTaskList) {
+ final DatabasePartitionTable databasePartitionTable =
+ databasePartitionTables.get(plan.getDatabase());
+ if (databasePartitionTable == null || databasePartitionTable.isNotPreDeleted()) {
+ return RpcUtils.SUCCESS_STATUS;
+ }
+ regionMaintainTaskList.removeIf(
+ task ->
+ task instanceof RegionCreateTask
+ && Objects.equals(
+ plan.getDatabase(), ((RegionCreateTask) task).getStorageGroup()));
+ return RpcUtils.SUCCESS_STATUS;
+ }
+ }
+
+ private boolean isRegionCreateTaskOwnedByCurrentPartitionTable(RegionCreateTask task) {
+ if (!isDatabaseExisted(task.getStorageGroup())) {
+ return false;
+ }
+ return getReplicaSets(task.getStorageGroup(), Collections.singletonList(task.getRegionId()))
+ .stream()
+ .anyMatch(
+ replicaSet ->
+ replicaSet.getDataNodeLocations().stream()
+ .anyMatch(
+ location ->
+ location.getDataNodeId() == task.getTargetDataNode().getDataNodeId()));
+ }
+
/**
* Get a deep copy of RegionCleanList for RegionCleaner to maintain cluster RegionReplicas.
*
@@ -1024,10 +1095,13 @@ public boolean processTakeSnapshot(File snapshotDir) throws TException, IOExcept
databasePartitionTableEntry.getValue().serialize(bufferedOutputStream, protocol);
}
- // serialize regionCleanList
- ReadWriteIOUtils.write(regionMaintainTaskList.size(), bufferedOutputStream);
- for (RegionMaintainTask task : regionMaintainTaskList) {
- task.serialize(bufferedOutputStream, protocol);
+ // Serialize the queue under the same monitor used by every consensus mutation so the count
+ // and entries belong to one atomic snapshot.
+ synchronized (regionMaintainTaskList) {
+ ReadWriteIOUtils.write(regionMaintainTaskList.size(), bufferedOutputStream);
+ for (RegionMaintainTask task : regionMaintainTaskList) {
+ task.serialize(bufferedOutputStream, protocol);
+ }
}
// write to file
diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/env/ConfigNodeProcedureEnv.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/env/ConfigNodeProcedureEnv.java
index 76bf575681c87..e0f8fe1af330b 100644
--- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/env/ConfigNodeProcedureEnv.java
+++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/env/ConfigNodeProcedureEnv.java
@@ -53,6 +53,7 @@
import org.apache.iotdb.confignode.manager.schema.ClusterSchemaManager;
import org.apache.iotdb.confignode.persistence.partition.PartitionInfo;
import org.apache.iotdb.confignode.persistence.schema.ClusterSchemaInfo;
+import org.apache.iotdb.confignode.procedure.Procedure;
import org.apache.iotdb.confignode.procedure.exception.ProcedureException;
import org.apache.iotdb.confignode.procedure.impl.schema.SchemaUtils;
import org.apache.iotdb.confignode.procedure.scheduler.LockQueue;
@@ -102,6 +103,7 @@
import java.util.List;
import java.util.Map;
import java.util.Set;
+import java.util.TreeSet;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.locks.ReentrantLock;
import java.util.stream.Collectors;
@@ -115,6 +117,9 @@ public class ConfigNodeProcedureEnv {
/** Add or remove node lock. */
private final LockQueue nodeLock = new LockQueue();
+ /** Serializes procedures that mutate the lifecycle of the same database. */
+ private final Map databaseLockMap = new HashMap<>();
+
private final ReentrantLock schedulerLock = new ReentrantLock(true);
private final ReentrantLock submitRegionMigrateLock = new ReentrantLock(true);
@@ -159,9 +164,13 @@ public TSStatus deleteDatabaseConfig(final String name, final boolean isGenerate
* @param preDeleteType execute/rollback
* @param deleteSgName database name
*/
- public void preDeleteDatabase(
+ public TSStatus preDeleteDatabase(
final PreDeleteDatabasePlan.PreDeleteType preDeleteType, final String deleteSgName) {
- getPartitionManager().preDeleteDatabase(deleteSgName, preDeleteType);
+ return getPartitionManager().preDeleteDatabase(deleteSgName, preDeleteType);
+ }
+
+ public TSStatus batchRemoveRegionCreateTasks(final String database) {
+ return getPartitionManager().batchRemoveRegionCreateTasks(database);
}
public boolean invalidateCache(final String databaseName) throws IOException, TException {
@@ -490,6 +499,10 @@ public TSStatus persistRegionGroup(CreateRegionGroupsPlan createRegionGroupsPlan
}
}
+ public TSStatus validateCreateRegionGroups(final CreateRegionGroupsPlan createRegionGroupsPlan) {
+ return getPartitionManager().validateCreateRegionGroups(createRegionGroupsPlan);
+ }
+
/**
* Force activating RegionGroup by setting status to Running, therefore the ConfigNode-leader can
* select leader for it and use it to allocate new Partitions
@@ -1135,6 +1148,67 @@ public LockQueue getNodeLock() {
return nodeLock;
}
+ /**
+ * Atomically tries to lock all databases in lexical order.
+ *
+ * @return the first database whose lock is unavailable, or null when all locks are acquired
+ */
+ public String tryLockDatabases(final Procedure> procedure, final Set databaseNames) {
+ schedulerLock.lock();
+ try {
+ final List acquiredDatabases = new ArrayList<>();
+ for (final String database : new TreeSet<>(databaseNames)) {
+ final LockQueue lockQueue =
+ databaseLockMap.computeIfAbsent(database, key -> new LockQueue());
+ if (!lockQueue.tryLock(procedure)) {
+ acquiredDatabases.forEach(
+ acquiredDatabase -> {
+ final LockQueue acquiredLock = databaseLockMap.get(acquiredDatabase);
+ if (acquiredLock != null && acquiredLock.releaseLock(procedure)) {
+ acquiredLock.wakeWaitingProcedures(scheduler);
+ if (acquiredLock.isIdle()) {
+ databaseLockMap.remove(acquiredDatabase, acquiredLock);
+ }
+ }
+ });
+ return database;
+ }
+ acquiredDatabases.add(database);
+ }
+ return null;
+ } finally {
+ schedulerLock.unlock();
+ }
+ }
+
+ public void waitDatabaseLock(final Procedure> procedure, final String databaseName) {
+ schedulerLock.lock();
+ try {
+ databaseLockMap
+ .computeIfAbsent(databaseName, key -> new LockQueue())
+ .waitProcedure(procedure, scheduler);
+ } finally {
+ schedulerLock.unlock();
+ }
+ }
+
+ public void releaseDatabaseLocks(final Procedure> procedure, final Set databaseNames) {
+ schedulerLock.lock();
+ try {
+ for (final String database : databaseNames) {
+ final LockQueue lockQueue = databaseLockMap.get(database);
+ if (lockQueue != null && lockQueue.releaseLock(procedure)) {
+ lockQueue.wakeWaitingProcedures(scheduler);
+ if (lockQueue.isIdle()) {
+ databaseLockMap.remove(database, lockQueue);
+ }
+ }
+ }
+ } finally {
+ schedulerLock.unlock();
+ }
+ }
+
public ProcedureScheduler getScheduler() {
return scheduler;
}
diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/AbstractDatabaseProcedure.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/AbstractDatabaseProcedure.java
new file mode 100644
index 0000000000000..fc41b9afbc5a3
--- /dev/null
+++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/AbstractDatabaseProcedure.java
@@ -0,0 +1,67 @@
+/*
+ * 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;
+
+import org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv;
+import org.apache.iotdb.confignode.procedure.state.ProcedureLockState;
+
+import java.util.Set;
+
+/** A procedure that holds exclusive lifecycle locks for its databases until it finishes. */
+public abstract class AbstractDatabaseProcedure
+ extends StateMachineProcedure {
+
+ private transient String waitingDatabase;
+
+ protected AbstractDatabaseProcedure() {
+ super();
+ }
+
+ protected AbstractDatabaseProcedure(final boolean isGeneratedByPipe) {
+ super(isGeneratedByPipe);
+ }
+
+ protected abstract Set getDatabaseNames();
+
+ @Override
+ protected ProcedureLockState acquireLock(final ConfigNodeProcedureEnv env) {
+ waitingDatabase = env.tryLockDatabases(this, getDatabaseNames());
+ return waitingDatabase == null
+ ? ProcedureLockState.LOCK_ACQUIRED
+ : ProcedureLockState.LOCK_EVENT_WAIT;
+ }
+
+ @Override
+ protected void waitForLock(final ConfigNodeProcedureEnv env) {
+ if (waitingDatabase != null) {
+ env.waitDatabaseLock(this, waitingDatabase);
+ }
+ }
+
+ @Override
+ protected void releaseLock(final ConfigNodeProcedureEnv env) {
+ env.releaseDatabaseLocks(this, getDatabaseNames());
+ }
+
+ @Override
+ protected boolean holdLock(final ConfigNodeProcedureEnv env) {
+ return true;
+ }
+}
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..f963f6f9005b6 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
@@ -38,7 +38,7 @@
import org.apache.iotdb.confignode.persistence.partition.maintainer.RegionCreateTask;
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.AbstractDatabaseProcedure;
import org.apache.iotdb.confignode.procedure.state.CreateRegionGroupsState;
import org.apache.iotdb.confignode.procedure.store.ProcedureType;
import org.apache.iotdb.consensus.exception.ConsensusException;
@@ -61,7 +61,7 @@
import java.util.stream.Collectors;
public class CreateRegionGroupsProcedure
- extends StateMachineProcedure {
+ extends AbstractDatabaseProcedure {
private static final Logger LOGGER = LoggerFactory.getLogger(CreateRegionGroupsProcedure.class);
@@ -100,6 +100,18 @@ public CreateRegionGroupsProcedure(
@Override
protected Flow executeFromState(
final ConfigNodeProcedureEnv env, final CreateRegionGroupsState state) {
+ final TSStatus validationStatus = env.validateCreateRegionGroups(createRegionGroupsPlan);
+ if (validationStatus.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ // Only SHUNT_REGION_REPLICAS can have sent create RPCs without transferring ownership of the
+ // planned RegionGroups to PartitionInfo. Delete every planned replica idempotently; all later
+ // states are already owned and cleaned by DeleteDatabaseProcedure.
+ if (state == CreateRegionGroupsState.SHUNT_REGION_REPLICAS) {
+ submitPlannedRegionReplicaCleanup(env);
+ }
+ setFailure(new ProcedureException(new IoTDBException(validationStatus)));
+ return Flow.NO_MORE_STATE;
+ }
+
switch (state) {
case CREATE_REGION_GROUPS:
failedRegionReplicaSets = env.doRegionCreation(consensusGroupType, createRegionGroupsPlan);
@@ -191,6 +203,7 @@ protected Flow executeFromState(
final TSStatus persistStatus = env.persistRegionGroup(persistPlan);
if (persistStatus.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ submitPlannedRegionReplicaCleanup(env);
setFailure(new ProcedureException(new IoTDBException(persistStatus)));
return Flow.NO_MORE_STATE;
}
@@ -323,6 +336,26 @@ protected CreateRegionGroupsState getInitialState() {
return CreateRegionGroupsState.CREATE_REGION_GROUPS;
}
+ @Override
+ protected Set getDatabaseNames() {
+ return createRegionGroupsPlan.getRegionGroupMap().keySet();
+ }
+
+ public boolean containsDatabase(final String database) {
+ return createRegionGroupsPlan.getRegionGroupMap().containsKey(database);
+ }
+
+ private void submitPlannedRegionReplicaCleanup(final ConfigNodeProcedureEnv env) {
+ createRegionGroupsPlan.getRegionGroupMap().values().stream()
+ .flatMap(List::stream)
+ .forEach(
+ replicaSet ->
+ env.getConfigManager()
+ .getProcedureManager()
+ .getExecutor()
+ .submitProcedure(new RemoveRegionGroupProcedure(replicaSet)));
+ }
+
@Override
public void serialize(final DataOutputStream stream) throws IOException {
// Must serialize CREATE_REGION_GROUPS.getTypeCode() firstly
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..58e230a706ccd 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
@@ -23,13 +23,15 @@
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.TestOnly;
+import org.apache.iotdb.commons.utils.ThriftCommonsSerDeUtils;
import org.apache.iotdb.commons.utils.ThriftConfigNodeSerDeUtils;
import org.apache.iotdb.confignode.consensus.request.write.database.PreDeleteDatabasePlan;
import org.apache.iotdb.confignode.i18n.ProcedureMessages;
import org.apache.iotdb.confignode.manager.partition.PartitionMetrics;
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.AbstractDatabaseProcedure;
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;
@@ -37,22 +39,30 @@
import org.apache.iotdb.rpc.TSStatusCode;
import org.apache.thrift.TException;
+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.ArrayList;
+import java.util.Collections;
import java.util.List;
import java.util.Objects;
+import java.util.Set;
-public class DeleteDatabaseProcedure
- extends StateMachineProcedure {
+public class DeleteDatabaseProcedure extends AbstractDatabaseProcedure {
private static final Logger LOG = LoggerFactory.getLogger(DeleteDatabaseProcedure.class);
private static final int RETRY_THRESHOLD = 5;
private TDatabaseSchema deleteDatabaseSchema;
+ // Captured after PRE_DELETE is committed and persisted with this procedure. RegionIds are never
+ // reused, so cancellation and deletion cannot accidentally target a newly created database with
+ // the same name after a leader change or restart.
+ private List targetRegionReplicaSets;
+
public DeleteDatabaseProcedure(final boolean isGeneratedByPipe) {
super(isGeneratedByPipe);
}
@@ -63,6 +73,15 @@ public DeleteDatabaseProcedure(
this.deleteDatabaseSchema = deleteDatabaseSchema;
}
+ @TestOnly
+ DeleteDatabaseProcedure(
+ final TDatabaseSchema deleteDatabaseSchema,
+ final boolean isGeneratedByPipe,
+ final List targetRegionReplicaSets) {
+ this(deleteDatabaseSchema, isGeneratedByPipe);
+ this.targetRegionReplicaSets = new ArrayList<>(targetRegionReplicaSets);
+ }
+
public TDatabaseSchema getDeleteDatabaseSchema() {
return deleteDatabaseSchema;
}
@@ -83,22 +102,42 @@ protected Flow executeFromState(final ConfigNodeProcedureEnv env, final DeleteDa
LOG.info(
ProcedureMessages.LOG_DELETEDATABASEPROCEDURE_PRE_DELETE_DATABASE_ARG_6A1FEACC,
deleteDatabaseSchema.getName());
- env.preDeleteDatabase(
- PreDeleteDatabasePlan.PreDeleteType.EXECUTE, deleteDatabaseSchema.getName());
- setNextState(DeleteDatabaseState.INVALIDATE_CACHE);
+ final TSStatus preDeleteStatus =
+ env.preDeleteDatabase(
+ PreDeleteDatabasePlan.PreDeleteType.EXECUTE, deleteDatabaseSchema.getName());
+ if (preDeleteStatus.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ captureTargetRegionReplicaSets(env);
+ setNextState(DeleteDatabaseState.INVALIDATE_CACHE);
+ } else if (getCycles() > RETRY_THRESHOLD) {
+ setFailure(
+ new ProcedureException(
+ ProcedureMessages.DELETEDATABASEPROCEDURE_DELETE_DATABASESCHEMA_FAILED));
+ }
break;
case INVALIDATE_CACHE:
LOG.info(
ProcedureMessages.LOG_DELETEDATABASEPROCEDURE_INVALIDATE_CACHE_DATABASE_ARG_299FC9BC,
deleteDatabaseSchema.getName());
if (env.invalidateCache(deleteDatabaseSchema.getName())) {
- setNextState(DeleteDatabaseState.DELETE_DATABASE_SCHEMA);
+ setNextState(DeleteDatabaseState.BATCH_REMOVE_REGION_CREATE_TASKS);
} else {
setFailure(
new ProcedureException(
ProcedureMessages.DELETEDATABASEPROCEDURE_INVALIDATE_CACHE_FAILED));
}
break;
+ case BATCH_REMOVE_REGION_CREATE_TASKS:
+ captureTargetRegionReplicaSets(env);
+ final TSStatus removeTasksStatus =
+ env.batchRemoveRegionCreateTasks(deleteDatabaseSchema.getName());
+ if (removeTasksStatus.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ setNextState(DeleteDatabaseState.DELETE_DATABASE_SCHEMA);
+ } else if (getCycles() > RETRY_THRESHOLD) {
+ setFailure(
+ new ProcedureException(
+ ProcedureMessages.DELETEDATABASEPROCEDURE_DELETE_DATABASESCHEMA_FAILED));
+ }
+ break;
case DELETE_DATABASE_SCHEMA:
LOG.info(
ProcedureMessages.LOG_DELETEDATABASEPROCEDURE_DELETE_DATABASESCHEMA_ARG_A49A47AC,
@@ -122,9 +161,8 @@ protected Flow executeFromState(final ConfigNodeProcedureEnv env, final DeleteDa
// 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.
- final List regionReplicaSets =
- env.getAllReplicaSets(deleteDatabaseSchema.getName());
- regionReplicaSets.forEach(
+ captureTargetRegionReplicaSets(env);
+ targetRegionReplicaSets.forEach(
regionReplicaSet -> {
// Clear heartbeat cache along the way
env.getConfigManager()
@@ -191,6 +229,13 @@ protected Flow executeFromState(final ConfigNodeProcedureEnv env, final DeleteDa
return Flow.HAS_MORE_STATE;
}
+ private void captureTargetRegionReplicaSets(final ConfigNodeProcedureEnv env) {
+ if (targetRegionReplicaSets == null) {
+ targetRegionReplicaSets =
+ new ArrayList<>(env.getAllReplicaSets(deleteDatabaseSchema.getName()));
+ }
+ }
+
@Override
protected void rollbackState(final ConfigNodeProcedureEnv env, final DeleteDatabaseState state)
throws IOException, InterruptedException {
@@ -234,6 +279,13 @@ protected DeleteDatabaseState getInitialState() {
return DeleteDatabaseState.PRE_DELETE_DATABASE;
}
+ @Override
+ protected Set getDatabaseNames() {
+ return deleteDatabaseSchema == null
+ ? Collections.emptySet()
+ : Collections.singleton(deleteDatabaseSchema.getName());
+ }
+
public String getDatabase() {
return deleteDatabaseSchema.getName();
}
@@ -246,6 +298,14 @@ public void serialize(final DataOutputStream stream) throws IOException {
: ProcedureType.DELETE_DATABASE_PROCEDURE.getTypeCode());
super.serialize(stream);
ThriftConfigNodeSerDeUtils.serializeTDatabaseSchema(deleteDatabaseSchema, stream);
+ if (targetRegionReplicaSets == null) {
+ ReadWriteIOUtils.write(-1, stream);
+ } else {
+ ReadWriteIOUtils.write(targetRegionReplicaSets.size(), stream);
+ for (TRegionReplicaSet regionReplicaSet : targetRegionReplicaSets) {
+ ThriftCommonsSerDeUtils.serializeTRegionReplicaSet(regionReplicaSet, stream);
+ }
+ }
}
@Override
@@ -253,6 +313,16 @@ public void deserialize(final ByteBuffer byteBuffer) {
super.deserialize(byteBuffer);
try {
deleteDatabaseSchema = ThriftConfigNodeSerDeUtils.deserializeTDatabaseSchema(byteBuffer);
+ if (byteBuffer.hasRemaining()) {
+ final int size = ReadWriteIOUtils.readInt(byteBuffer);
+ if (size >= 0) {
+ targetRegionReplicaSets = new ArrayList<>(size);
+ for (int i = 0; i < size; i++) {
+ targetRegionReplicaSets.add(
+ ThriftCommonsSerDeUtils.deserializeTRegionReplicaSet(byteBuffer));
+ }
+ }
+ }
} catch (final ThriftSerDeException e) {
LOG.error(ProcedureMessages.ERROR_IN_DESERIALIZE_DELETEDATABASEPROCEDURE, e);
}
@@ -266,7 +336,8 @@ public boolean equals(final Object that) {
&& Objects.equals(thatProc.getCurrentState(), this.getCurrentState())
&& thatProc.getCycles() == this.getCycles()
&& thatProc.isGeneratedByPipe == this.isGeneratedByPipe
- && thatProc.deleteDatabaseSchema.equals(this.getDeleteDatabaseSchema());
+ && thatProc.deleteDatabaseSchema.equals(this.getDeleteDatabaseSchema())
+ && Objects.equals(thatProc.targetRegionReplicaSets, this.targetRegionReplicaSets);
}
return false;
}
@@ -274,6 +345,11 @@ public boolean equals(final Object that) {
@Override
public int hashCode() {
return Objects.hash(
- getProcId(), getCurrentState(), getCycles(), isGeneratedByPipe, deleteDatabaseSchema);
+ getProcId(),
+ getCurrentState(),
+ getCycles(),
+ isGeneratedByPipe,
+ deleteDatabaseSchema,
+ targetRegionReplicaSets);
}
}
diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/scheduler/LockQueue.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/scheduler/LockQueue.java
index e2f5935a909ab..164f778c1f22c 100644
--- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/scheduler/LockQueue.java
+++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/scheduler/LockQueue.java
@@ -64,4 +64,8 @@ public int wakeWaitingProcedures(ProcedureScheduler procedureScheduler) {
}
return count;
}
+
+ public boolean isIdle() {
+ return lockOwnerProcedure == null && deque.isEmpty();
+ }
}
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..6298d7ac02748 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
@@ -25,5 +25,7 @@ public enum DeleteDatabaseState {
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_CONFIG,
+ // Appended to preserve the serialized ordinals of procedures written by older versions.
+ BATCH_REMOVE_REGION_CREATE_TASKS
}
diff --git a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/consensus/request/ConfigPhysicalPlanSerDeTest.java b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/consensus/request/ConfigPhysicalPlanSerDeTest.java
index e844dcf6910e9..4fe56aa3f7d43 100644
--- a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/consensus/request/ConfigPhysicalPlanSerDeTest.java
+++ b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/consensus/request/ConfigPhysicalPlanSerDeTest.java
@@ -118,6 +118,7 @@
import org.apache.iotdb.confignode.consensus.request.write.procedure.UpdateProcedurePlan;
import org.apache.iotdb.confignode.consensus.request.write.quota.SetSpaceQuotaPlan;
import org.apache.iotdb.confignode.consensus.request.write.quota.SetThrottleQuotaPlan;
+import org.apache.iotdb.confignode.consensus.request.write.region.BatchRemoveRegionCreateTasksPlan;
import org.apache.iotdb.confignode.consensus.request.write.region.CreateRegionGroupsPlan;
import org.apache.iotdb.confignode.consensus.request.write.region.OfferRegionMaintainTasksPlan;
import org.apache.iotdb.confignode.consensus.request.write.region.PollRegionMaintainTaskPlan;
@@ -273,6 +274,19 @@ public void CreateDatabasePlanTest() throws IOException {
Assert.assertEquals(req0, req1);
}
+ @Test
+ public void CreateRegionGroupsPlanTest() throws IOException {
+ final CreateRegionGroupsPlan plan = new CreateRegionGroupsPlan();
+ plan.addRegionGroup(
+ "root.sg",
+ new TRegionReplicaSet(
+ new TConsensusGroupId(TConsensusGroupType.DataRegion, 1), Collections.emptyList()));
+
+ final CreateRegionGroupsPlan deserializedPlan =
+ (CreateRegionGroupsPlan) ConfigPhysicalPlan.Factory.create(plan.serializeToByteBuffer());
+ Assert.assertEquals(plan, deserializedPlan);
+ }
+
@Test
public void AlterDatabasePlanTest() throws IOException {
DatabaseSchemaPlan req0 =
@@ -423,6 +437,15 @@ public void PollRegionMaintainTaskPlan() throws IOException {
Assert.assertEquals(plan0, plan1);
}
+ @Test
+ public void BatchRemoveRegionCreateTasksPlanTest() throws IOException {
+ final BatchRemoveRegionCreateTasksPlan plan0 = new BatchRemoveRegionCreateTasksPlan("root.sg");
+ final BatchRemoveRegionCreateTasksPlan plan1 =
+ (BatchRemoveRegionCreateTasksPlan)
+ ConfigPhysicalPlan.Factory.create(plan0.serializeToByteBuffer());
+ Assert.assertEquals(plan0, plan1);
+ }
+
@Test
public void CreateSchemaPartitionPlanTest() throws IOException {
TDataNodeLocation dataNodeLocation = new TDataNodeLocation();
diff --git a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ClusterSchemaManagerTest.java b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ClusterSchemaManagerTest.java
index 261f60ed65a1d..67837d3d8d094 100644
--- a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ClusterSchemaManagerTest.java
+++ b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ClusterSchemaManagerTest.java
@@ -18,10 +18,27 @@
*/
package org.apache.iotdb.confignode.manager;
+import org.apache.iotdb.common.rpc.thrift.TSStatus;
+import org.apache.iotdb.confignode.consensus.request.ConfigPhysicalPlanType;
+import org.apache.iotdb.confignode.consensus.request.write.database.DatabaseSchemaPlan;
import org.apache.iotdb.confignode.manager.schema.ClusterSchemaManager;
+import org.apache.iotdb.confignode.manager.schema.ClusterSchemaQuotaStatistics;
+import org.apache.iotdb.confignode.persistence.schema.ClusterSchemaInfo;
+import org.apache.iotdb.confignode.rpc.thrift.TDatabaseSchema;
+import org.apache.iotdb.rpc.TSStatusCode;
+import org.apache.ratis.util.AutoCloseableLock;
import org.junit.Assert;
import org.junit.Test;
+import org.mockito.Mockito;
+
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.locks.ReentrantLock;
public class ClusterSchemaManagerTest {
@@ -37,4 +54,51 @@ public void testCalcMaxRegionGroupNum() {
// (resourceWeight * resource) / (createdStorageGroupNum * replicationFactor)
Assert.assertEquals(20, ClusterSchemaManager.calcMaxRegionGroupNum(3, 1.0, 120, 2, 3, 5));
}
+
+ @Test
+ public void testSetDatabaseWaitsForLifecycleAdmissionBeforeCheckingProcedures() throws Exception {
+ final String database = "root.sg";
+ final IManager configManager = Mockito.mock(IManager.class);
+ final ProcedureManager procedureManager = Mockito.mock(ProcedureManager.class);
+ final ReentrantLock admissionLock = new ReentrantLock();
+ final CountDownLatch admissionAttempted = new CountDownLatch(1);
+ final AtomicBoolean unfinishedProcedure = new AtomicBoolean(false);
+ Mockito.when(configManager.getProcedureManager()).thenReturn(procedureManager);
+ Mockito.when(procedureManager.acquireDatabaseLifecycleAdmissionLock())
+ .thenAnswer(
+ ignored -> {
+ admissionAttempted.countDown();
+ return AutoCloseableLock.acquire(admissionLock);
+ });
+ Mockito.when(procedureManager.hasUnfinishedDatabaseLifecycleProcedure(database))
+ .thenAnswer(ignored -> unfinishedProcedure.get());
+
+ final ClusterSchemaManager schemaManager =
+ new ClusterSchemaManager(
+ configManager,
+ Mockito.mock(ClusterSchemaInfo.class),
+ Mockito.mock(ClusterSchemaQuotaStatistics.class));
+ final DatabaseSchemaPlan plan =
+ new DatabaseSchemaPlan(
+ ConfigPhysicalPlanType.CreateDatabase, new TDatabaseSchema(database));
+ final ExecutorService executor = Executors.newSingleThreadExecutor();
+
+ admissionLock.lock();
+ try {
+ final Future statusFuture =
+ executor.submit(() -> schemaManager.setDatabase(plan, false));
+ Assert.assertTrue(admissionAttempted.await(10, TimeUnit.SECONDS));
+ unfinishedProcedure.set(true);
+ admissionLock.unlock();
+
+ final TSStatus status = statusFuture.get(10, TimeUnit.SECONDS);
+ Assert.assertEquals(TSStatusCode.METADATA_ERROR.getStatusCode(), status.getCode());
+ Mockito.verify(procedureManager).hasUnfinishedDatabaseLifecycleProcedure(database);
+ } finally {
+ if (admissionLock.isHeldByCurrentThread()) {
+ admissionLock.unlock();
+ }
+ executor.shutdownNow();
+ }
+ }
}
diff --git a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java
index 4de883b9c158f..193abc74cfa5e 100644
--- a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java
+++ b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/manager/ProcedureManagerTest.java
@@ -23,19 +23,24 @@
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.common.rpc.thrift.TSStatus;
import org.apache.iotdb.commons.cluster.NodeStatus;
import org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant;
import org.apache.iotdb.commons.schema.table.TreeViewSchema;
import org.apache.iotdb.commons.schema.table.TsTable;
+import org.apache.iotdb.confignode.consensus.request.write.region.CreateRegionGroupsPlan;
import org.apache.iotdb.confignode.manager.load.LoadManager;
import org.apache.iotdb.confignode.procedure.Procedure;
import org.apache.iotdb.confignode.procedure.ProcedureExecutor;
import org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv;
import org.apache.iotdb.confignode.procedure.env.RemoveDataNodeHandler;
import org.apache.iotdb.confignode.procedure.impl.node.RemoveDataNodesProcedure;
+import org.apache.iotdb.confignode.procedure.impl.region.CreateRegionGroupsProcedure;
import org.apache.iotdb.confignode.procedure.impl.region.RegionMigrateProcedure;
import org.apache.iotdb.confignode.procedure.impl.region.RegionMigrationPlan;
+import org.apache.iotdb.confignode.procedure.impl.schema.DeleteDatabaseProcedure;
+import org.apache.iotdb.confignode.rpc.thrift.TDatabaseSchema;
import org.junit.Assert;
import org.junit.BeforeClass;
@@ -128,6 +133,7 @@ public static void setUp() throws IOException {
when(PROCEDURE_MANAGER.getExecutor()).thenReturn(PROCEDURE_EXECUTOR);
when(PROCEDURE_EXECUTOR.getProcedures()).thenReturn(procedureMap);
+ PROCEDURE_MANAGER.setExecutor(PROCEDURE_EXECUTOR);
when(PROCEDURE_MANAGER.getEnv()).thenReturn(ENV);
when(ENV.getRemoveDataNodeHandler()).thenReturn(REMOVE_DATA_NODE_HANDLER);
}
@@ -233,4 +239,31 @@ public void testInjectTreeViewSourceAttributesPreservesUserTreePattern() {
Assert.assertFalse(
topicAttributes.containsKey(PipeSourceConstant.SOURCE_PATTERN_INCLUSION_KEY));
}
+
+ @Test
+ public void testDetectUnfinishedDatabaseLifecycleProcedures() {
+ final CreateRegionGroupsPlan createPlan = new CreateRegionGroupsPlan();
+ createPlan.addRegionGroup(
+ "root.create",
+ new TRegionReplicaSet(
+ new TConsensusGroupId(TConsensusGroupType.DataRegion, 10), List.of()));
+ final CreateRegionGroupsProcedure createProcedure =
+ new CreateRegionGroupsProcedure(TConsensusGroupType.DataRegion, createPlan);
+ createProcedure.setProcId(100);
+ final DeleteDatabaseProcedure deleteProcedure =
+ new DeleteDatabaseProcedure(new TDatabaseSchema("root.delete"), false);
+ deleteProcedure.setProcId(101);
+
+ procedureMap.clear();
+ try {
+ procedureMap.put(createProcedure.getProcId(), createProcedure);
+ procedureMap.put(deleteProcedure.getProcId(), deleteProcedure);
+
+ Assert.assertTrue(PROCEDURE_MANAGER.hasUnfinishedDatabaseLifecycleProcedure("root.create"));
+ Assert.assertTrue(PROCEDURE_MANAGER.hasUnfinishedDatabaseLifecycleProcedure("root.delete"));
+ Assert.assertFalse(PROCEDURE_MANAGER.hasUnfinishedDatabaseLifecycleProcedure("root.other"));
+ } finally {
+ procedureMap.clear();
+ }
+ }
}
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..3a41a88e50560 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
@@ -24,6 +24,7 @@
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.common.rpc.thrift.TSStatus;
import org.apache.iotdb.common.rpc.thrift.TSeriesPartitionSlot;
import org.apache.iotdb.common.rpc.thrift.TTimePartitionSlot;
import org.apache.iotdb.commons.partition.DataPartitionTable;
@@ -32,11 +33,15 @@
import org.apache.iotdb.confignode.consensus.request.ConfigPhysicalPlanType;
import org.apache.iotdb.confignode.consensus.request.read.region.GetRegionInfoListPlan;
import org.apache.iotdb.confignode.consensus.request.write.database.DatabaseSchemaPlan;
+import org.apache.iotdb.confignode.consensus.request.write.database.DeleteDatabasePlan;
+import org.apache.iotdb.confignode.consensus.request.write.database.PreDeleteDatabasePlan;
import org.apache.iotdb.confignode.consensus.request.write.partition.CreateDataPartitionPlan;
import org.apache.iotdb.confignode.consensus.request.write.partition.CreateSchemaPartitionPlan;
+import org.apache.iotdb.confignode.consensus.request.write.region.BatchRemoveRegionCreateTasksPlan;
import org.apache.iotdb.confignode.consensus.request.write.region.CreateRegionGroupsPlan;
import org.apache.iotdb.confignode.consensus.request.write.region.OfferRegionMaintainTasksPlan;
import org.apache.iotdb.confignode.consensus.response.partition.RegionInfoListResp;
+import org.apache.iotdb.confignode.exception.DatabaseNotExistsException;
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;
@@ -44,6 +49,7 @@
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;
+import org.apache.iotdb.rpc.TSStatusCode;
import org.apache.thrift.TException;
import org.apache.tsfile.external.commons.io.FileUtils;
@@ -52,8 +58,10 @@
import org.junit.Before;
import org.junit.Test;
+import java.io.DataInputStream;
import java.io.File;
import java.io.IOException;
+import java.nio.file.Files;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
@@ -148,6 +156,17 @@ public void testSnapshot() throws TException, IOException {
partitionInfo.offerRegionMaintainTasks(generateOfferRegionMaintainTasksPlan());
Assert.assertTrue(partitionInfo.processTakeSnapshot(snapshotDir));
+ try (final DataInputStream inputStream =
+ new DataInputStream(
+ Files.newInputStream(new File(snapshotDir, "partition_info.bin").toPath()))) {
+ // Keep the historical snapshot framing: next RegionGroupId followed by database count.
+ Assert.assertEquals(
+ Math.max(
+ schemaRegionReplicaSet.getRegionId().getId(),
+ dataRegionReplicaSet.getRegionId().getId()),
+ inputStream.readInt());
+ Assert.assertEquals(1, inputStream.readInt());
+ }
PartitionInfo partitionInfo1 = new PartitionInfo();
partitionInfo1.processLoadSnapshot(snapshotDir);
@@ -162,7 +181,16 @@ public void testLegacyRegionDeleteTasksAreFiltered() throws TException, IOExcept
// it cannot block the recreation of that region's other replicas.
// The offer plan mixes two RegionCreateTasks with one legacy RegionDeleteTask.
- partitionInfo.offerRegionMaintainTasks(generateOfferRegionMaintainTasksPlan());
+ final OfferRegionMaintainTasksPlan offerPlan = generateOfferRegionMaintainTasksPlan();
+ final RegionCreateTask createTask =
+ (RegionCreateTask) offerPlan.getRegionMaintainTaskList().get(0);
+ partitionInfo.createDatabase(
+ new DatabaseSchemaPlan(
+ ConfigPhysicalPlanType.CreateDatabase, new TDatabaseSchema("root.sg")));
+ final CreateRegionGroupsPlan createRegionGroupsPlan = new CreateRegionGroupsPlan();
+ createRegionGroupsPlan.addRegionGroup("root.sg", createTask.getRegionReplicaSet());
+ partitionInfo.createRegionGroups(createRegionGroupsPlan);
+ partitionInfo.offerRegionMaintainTasks(offerPlan);
// The DELETE task is filtered out at offer time; only the two CREATE tasks remain queued.
List queuedTasks = partitionInfo.getRegionMaintainEntryList();
@@ -179,6 +207,162 @@ public void testLegacyRegionDeleteTasksAreFiltered() throws TException, IOExcept
Assert.assertEquals(2, loaded.getRegionMaintainEntryList().size());
}
+ @Test
+ public void testBatchRemoveAllRegionCreateTasksAndSnapshot() throws TException, IOException {
+ final String database = "root.sg";
+ final String otherDatabase = "root.other";
+ partitionInfo.createDatabase(
+ new DatabaseSchemaPlan(
+ ConfigPhysicalPlanType.CreateDatabase, new TDatabaseSchema(database)));
+ partitionInfo.createDatabase(
+ new DatabaseSchemaPlan(
+ ConfigPhysicalPlanType.CreateDatabase, new TDatabaseSchema(otherDatabase)));
+
+ final TRegionReplicaSet region0 =
+ generateTRegionReplicaSet(0, new TConsensusGroupId(TConsensusGroupType.DataRegion, 0));
+ final TRegionReplicaSet region1 =
+ generateTRegionReplicaSet(10, new TConsensusGroupId(TConsensusGroupType.DataRegion, 1));
+ final TRegionReplicaSet otherRegion =
+ generateTRegionReplicaSet(20, new TConsensusGroupId(TConsensusGroupType.DataRegion, 2));
+ final CreateRegionGroupsPlan createRegionGroupsPlan = new CreateRegionGroupsPlan();
+ createRegionGroupsPlan.addRegionGroup(database, region0);
+ createRegionGroupsPlan.addRegionGroup(database, region1);
+ createRegionGroupsPlan.addRegionGroup(otherDatabase, otherRegion);
+ partitionInfo.createRegionGroups(createRegionGroupsPlan);
+
+ final OfferRegionMaintainTasksPlan offerPlan = new OfferRegionMaintainTasksPlan();
+ offerPlan.appendRegionMaintainTask(
+ new RegionCreateTask(region0.getDataNodeLocations().get(0), database, region0));
+ offerPlan.appendRegionMaintainTask(
+ new RegionCreateTask(region0.getDataNodeLocations().get(1), database, region0));
+ offerPlan.appendRegionMaintainTask(
+ new RegionCreateTask(region1.getDataNodeLocations().get(0), database, region1));
+ offerPlan.appendRegionMaintainTask(
+ new RegionCreateTask(
+ otherRegion.getDataNodeLocations().get(0), otherDatabase, otherRegion));
+ partitionInfo.offerRegionMaintainTasks(offerPlan);
+ Assert.assertEquals(4, partitionInfo.getRegionMaintainEntryList().size());
+
+ partitionInfo.preDeleteDatabase(
+ new PreDeleteDatabasePlan(database, PreDeleteDatabasePlan.PreDeleteType.EXECUTE));
+ partitionInfo.batchRemoveRegionCreateTasks(new BatchRemoveRegionCreateTasksPlan(database));
+ Assert.assertEquals(1, partitionInfo.getRegionMaintainEntryList().size());
+ Assert.assertEquals(
+ otherRegion.getRegionId(), partitionInfo.getRegionMaintainEntryList().get(0).getRegionId());
+
+ // Replaying the same consensus plan is idempotent, and a snapshot cannot revive removed tasks.
+ partitionInfo.batchRemoveRegionCreateTasks(new BatchRemoveRegionCreateTasksPlan(database));
+ Assert.assertTrue(partitionInfo.processTakeSnapshot(snapshotDir));
+ final PartitionInfo loaded = new PartitionInfo();
+ loaded.processLoadSnapshot(snapshotDir);
+ Assert.assertEquals(1, loaded.getRegionMaintainEntryList().size());
+ Assert.assertEquals(
+ otherRegion.getRegionId(), loaded.getRegionMaintainEntryList().get(0).getRegionId());
+ }
+
+ @Test
+ public void testCancelledTasksCannotAffectRecreatedDatabase() {
+ final String database = "root.sg";
+ partitionInfo.createDatabase(
+ new DatabaseSchemaPlan(
+ ConfigPhysicalPlanType.CreateDatabase, new TDatabaseSchema(database)));
+ final TRegionReplicaSet oldRegion =
+ generateTRegionReplicaSet(0, new TConsensusGroupId(TConsensusGroupType.DataRegion, 0));
+ final CreateRegionGroupsPlan oldCreatePlan = new CreateRegionGroupsPlan();
+ oldCreatePlan.addRegionGroup(database, oldRegion);
+ partitionInfo.createRegionGroups(oldCreatePlan);
+
+ final OfferRegionMaintainTasksPlan oldOfferPlan = new OfferRegionMaintainTasksPlan();
+ oldOfferPlan.appendRegionMaintainTask(
+ new RegionCreateTask(oldRegion.getDataNodeLocations().get(0), database, oldRegion));
+ partitionInfo.offerRegionMaintainTasks(oldOfferPlan);
+ Assert.assertEquals(1, partitionInfo.getRegionMaintainEntryList().size());
+
+ partitionInfo.preDeleteDatabase(
+ new PreDeleteDatabasePlan(database, PreDeleteDatabasePlan.PreDeleteType.EXECUTE));
+ final BatchRemoveRegionCreateTasksPlan oldCancellation =
+ new BatchRemoveRegionCreateTasksPlan(database);
+ partitionInfo.batchRemoveRegionCreateTasks(oldCancellation);
+ Assert.assertTrue(partitionInfo.getRegionMaintainEntryList().isEmpty());
+
+ // A late offer from the old create procedure is rejected after PRE_DELETE.
+ partitionInfo.offerRegionMaintainTasks(oldOfferPlan);
+ Assert.assertTrue(partitionInfo.getRegionMaintainEntryList().isEmpty());
+
+ partitionInfo.deleteDatabase(new DeleteDatabasePlan(database));
+ partitionInfo.createDatabase(
+ new DatabaseSchemaPlan(
+ ConfigPhysicalPlanType.CreateDatabase, new TDatabaseSchema(database)));
+ final TRegionReplicaSet newRegion =
+ generateTRegionReplicaSet(10, new TConsensusGroupId(TConsensusGroupType.DataRegion, 1));
+ final CreateRegionGroupsPlan newCreatePlan = new CreateRegionGroupsPlan();
+ newCreatePlan.addRegionGroup(database, newRegion);
+ partitionInfo.createRegionGroups(newCreatePlan);
+ final OfferRegionMaintainTasksPlan newOfferPlan = new OfferRegionMaintainTasksPlan();
+ newOfferPlan.appendRegionMaintainTask(
+ new RegionCreateTask(newRegion.getDataNodeLocations().get(0), database, newRegion));
+ partitionInfo.offerRegionMaintainTasks(newOfferPlan);
+
+ // Replaying the old cancellation is a no-op because the same-name database is not pre-deleted.
+ // A late task offer from the old RegionGroup is rejected by Region ownership validation.
+ partitionInfo.batchRemoveRegionCreateTasks(oldCancellation);
+ partitionInfo.offerRegionMaintainTasks(oldOfferPlan);
+ Assert.assertEquals(1, partitionInfo.getRegionMaintainEntryList().size());
+ Assert.assertEquals(
+ newRegion.getRegionId(), partitionInfo.getRegionMaintainEntryList().get(0).getRegionId());
+ }
+
+ @Test
+ public void testCreateRegionGroupsRejectsPreDeletedAndMissingDatabase()
+ throws DatabaseNotExistsException {
+ final String database = "root.lifecycle";
+ partitionInfo.createDatabase(
+ new DatabaseSchemaPlan(
+ ConfigPhysicalPlanType.CreateDatabase, new TDatabaseSchema(database)));
+
+ final CreateRegionGroupsPlan preDeletedPlan = new CreateRegionGroupsPlan();
+ preDeletedPlan.addRegionGroup(
+ database,
+ generateTRegionReplicaSet(0, new TConsensusGroupId(TConsensusGroupType.DataRegion, 1)));
+ partitionInfo.preDeleteDatabase(
+ new PreDeleteDatabasePlan(database, PreDeleteDatabasePlan.PreDeleteType.EXECUTE));
+
+ TSStatus status = partitionInfo.createRegionGroups(preDeletedPlan);
+ Assert.assertEquals(TSStatusCode.DATABASE_NOT_EXIST.getStatusCode(), status.getCode());
+ Assert.assertTrue(
+ partitionInfo.getAllReplicaSets(database, TConsensusGroupType.DataRegion).isEmpty());
+
+ final CreateRegionGroupsPlan missingPlan = new CreateRegionGroupsPlan();
+ missingPlan.addRegionGroup(
+ "root.missing",
+ generateTRegionReplicaSet(10, new TConsensusGroupId(TConsensusGroupType.SchemaRegion, 2)));
+ status = partitionInfo.createRegionGroups(missingPlan);
+ Assert.assertEquals(TSStatusCode.DATABASE_NOT_EXIST.getStatusCode(), status.getCode());
+ }
+
+ @Test
+ public void testBatchedCreateRegionGroupsPlanIsValidatedAtomically()
+ throws DatabaseNotExistsException {
+ final String existingDatabase = "root.existing";
+ partitionInfo.createDatabase(
+ new DatabaseSchemaPlan(
+ ConfigPhysicalPlanType.CreateDatabase, new TDatabaseSchema(existingDatabase)));
+
+ final CreateRegionGroupsPlan batchedPlan = new CreateRegionGroupsPlan();
+ batchedPlan.addRegionGroup(
+ existingDatabase,
+ generateTRegionReplicaSet(0, new TConsensusGroupId(TConsensusGroupType.DataRegion, 40)));
+ batchedPlan.addRegionGroup(
+ "root.missing",
+ generateTRegionReplicaSet(10, new TConsensusGroupId(TConsensusGroupType.DataRegion, 41)));
+
+ final TSStatus status = partitionInfo.createRegionGroups(batchedPlan);
+ Assert.assertEquals(TSStatusCode.DATABASE_NOT_EXIST.getStatusCode(), status.getCode());
+ Assert.assertEquals(
+ 0, partitionInfo.getRegionGroupCount(existingDatabase, TConsensusGroupType.DataRegion));
+ Assert.assertEquals(42, partitionInfo.generateNextRegionGroupId());
+ }
+
@Test
public void testGetRegionType() {
diff --git a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/CreateRegionGroupsProcedureTest.java b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/CreateRegionGroupsProcedureTest.java
index 6269953736100..8e58d5ff85a8e 100644
--- a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/CreateRegionGroupsProcedureTest.java
+++ b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/CreateRegionGroupsProcedureTest.java
@@ -24,20 +24,36 @@
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.common.rpc.thrift.TSStatus;
import org.apache.iotdb.confignode.consensus.request.write.region.CreateRegionGroupsPlan;
+import org.apache.iotdb.confignode.manager.ConfigManager;
+import org.apache.iotdb.confignode.manager.ProcedureManager;
+import org.apache.iotdb.confignode.procedure.Procedure;
+import org.apache.iotdb.confignode.procedure.ProcedureExecutor;
+import org.apache.iotdb.confignode.procedure.env.ConfigNodeProcedureEnv;
import org.apache.iotdb.confignode.procedure.impl.region.CreateRegionGroupsProcedure;
+import org.apache.iotdb.confignode.procedure.impl.region.RemoveRegionGroupProcedure;
+import org.apache.iotdb.confignode.procedure.impl.schema.DeleteDatabaseProcedure;
+import org.apache.iotdb.confignode.procedure.scheduler.ProcedureScheduler;
+import org.apache.iotdb.confignode.procedure.state.CreateRegionGroupsState;
+import org.apache.iotdb.confignode.procedure.state.ProcedureLockState;
import org.apache.iotdb.confignode.procedure.store.ProcedureFactory;
import org.apache.iotdb.confignode.procedure.store.ProcedureType;
+import org.apache.iotdb.confignode.rpc.thrift.TDatabaseSchema;
+import org.apache.iotdb.rpc.TSStatusCode;
import org.apache.tsfile.utils.PublicBAOS;
import org.junit.Assert;
import org.junit.Test;
+import org.mockito.ArgumentCaptor;
+import org.mockito.Mockito;
import java.io.DataOutputStream;
import java.io.IOException;
import java.nio.ByteBuffer;
import java.util.Collections;
import java.util.HashMap;
+import java.util.List;
import java.util.Map;
import static org.apache.iotdb.common.rpc.thrift.TConsensusGroupType.DataRegion;
@@ -47,6 +63,56 @@
public class CreateRegionGroupsProcedureTest {
+ private static class TestCreateRegionGroupsProcedure extends CreateRegionGroupsProcedure {
+
+ private TestCreateRegionGroupsProcedure() {
+ super();
+ }
+
+ private TestCreateRegionGroupsProcedure(
+ final TConsensusGroupType consensusGroupType,
+ final CreateRegionGroupsPlan createRegionGroupsPlan,
+ final CreateRegionGroupsPlan persistPlan,
+ final Map failedRegionReplicaSets) {
+ super(consensusGroupType, createRegionGroupsPlan, persistPlan, failedRegionReplicaSets);
+ }
+
+ private void executeShunt(final ConfigNodeProcedureEnv env) {
+ executeFromState(env, CreateRegionGroupsState.SHUNT_REGION_REPLICAS);
+ }
+
+ private void executeCreate(final ConfigNodeProcedureEnv env) {
+ executeFromState(env, CreateRegionGroupsState.CREATE_REGION_GROUPS);
+ }
+
+ private void executePostPersist(final ConfigNodeProcedureEnv env) {
+ executeFromState(env, CreateRegionGroupsState.REBALANCE_DATA_PARTITION_POLICY);
+ }
+
+ private ProcedureLockState acquireDatabaseLock(final ConfigNodeProcedureEnv env) {
+ return acquireLock(env);
+ }
+
+ private void releaseDatabaseLock(final ConfigNodeProcedureEnv env) {
+ releaseLock(env);
+ }
+ }
+
+ private static class TestDeleteDatabaseProcedure extends DeleteDatabaseProcedure {
+
+ private TestDeleteDatabaseProcedure(final TDatabaseSchema databaseSchema) {
+ super(databaseSchema, false);
+ }
+
+ private ProcedureLockState acquireDatabaseLock(final ConfigNodeProcedureEnv env) {
+ return acquireLock(env);
+ }
+
+ private void releaseDatabaseLock(final ConfigNodeProcedureEnv env) {
+ releaseLock(env);
+ }
+ }
+
@Test
public void serializeDeserializeTest() {
TDataNodeLocation dataNodeLocation0 = new TDataNodeLocation();
@@ -111,6 +177,7 @@ public void serializeDeserializeTest() {
ByteBuffer.wrap(byteArrayOutputStream.getBuf(), 0, byteArrayOutputStream.size());
Assert.assertEquals(ProcedureType.CREATE_REGION_GROUPS.getTypeCode(), buffer.getShort());
procedure1.deserialize(buffer);
+ Assert.assertFalse(buffer.hasRemaining());
assertEquals(procedure0, procedure1);
assertEquals(procedure0.hashCode(), procedure1.hashCode());
@@ -124,4 +191,161 @@ public void serializeDeserializeTest() {
fail();
}
}
+
+ @Test
+ public void testPersistRejectionCleansEveryPlannedRegionReplica() {
+ final TDataNodeLocation createdDataNode =
+ new TDataNodeLocation().setDataNodeId(1).setInternalEndPoint(new TEndPoint("0.0.0.1", 1));
+ final TDataNodeLocation failedDataNode =
+ new TDataNodeLocation().setDataNodeId(2).setInternalEndPoint(new TEndPoint("0.0.0.2", 2));
+ final TDataNodeLocation otherFailedDataNode =
+ new TDataNodeLocation().setDataNodeId(3).setInternalEndPoint(new TEndPoint("0.0.0.3", 3));
+ final TConsensusGroupId regionId = new TConsensusGroupId(DataRegion, 10);
+ final TConsensusGroupId otherRegionId = new TConsensusGroupId(DataRegion, 11);
+ final TRegionReplicaSet allocatedReplicaSet =
+ new TRegionReplicaSet(regionId, List.of(createdDataNode, failedDataNode));
+ final TRegionReplicaSet failedReplicaSet =
+ new TRegionReplicaSet(regionId, Collections.singletonList(failedDataNode));
+ final TRegionReplicaSet otherAllocatedReplicaSet =
+ new TRegionReplicaSet(otherRegionId, Collections.singletonList(otherFailedDataNode));
+
+ final CreateRegionGroupsPlan createPlan = new CreateRegionGroupsPlan();
+ createPlan.addRegionGroup("root.sg", allocatedReplicaSet);
+ createPlan.addRegionGroup("root.sg", otherAllocatedReplicaSet);
+ final Map failedReplicaSets = new HashMap<>();
+ failedReplicaSets.put(regionId, failedReplicaSet);
+ failedReplicaSets.put(otherRegionId, otherAllocatedReplicaSet);
+ final TestCreateRegionGroupsProcedure procedure =
+ new TestCreateRegionGroupsProcedure(
+ DataRegion, createPlan, new CreateRegionGroupsPlan(), failedReplicaSets);
+
+ final ConfigNodeProcedureEnv env = Mockito.mock(ConfigNodeProcedureEnv.class);
+ final ConfigManager configManager = Mockito.mock(ConfigManager.class);
+ final ProcedureManager procedureManager = Mockito.mock(ProcedureManager.class);
+ @SuppressWarnings("unchecked")
+ final ProcedureExecutor executor =
+ Mockito.mock(ProcedureExecutor.class);
+ Mockito.when(env.validateCreateRegionGroups(createPlan))
+ .thenReturn(new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode()));
+ Mockito.when(env.persistRegionGroup(Mockito.any()))
+ .thenReturn(new TSStatus(TSStatusCode.DATABASE_NOT_EXIST.getStatusCode()));
+ Mockito.when(env.getConfigManager()).thenReturn(configManager);
+ Mockito.when(configManager.getProcedureManager()).thenReturn(procedureManager);
+ Mockito.when(procedureManager.getExecutor()).thenReturn(executor);
+
+ procedure.executeShunt(env);
+
+ final ArgumentCaptor> cleanupCaptor =
+ ArgumentCaptor.forClass(Procedure.class);
+ Mockito.verify(executor, Mockito.times(2)).submitProcedure(cleanupCaptor.capture());
+ Assert.assertTrue(
+ cleanupCaptor.getAllValues().contains(new RemoveRegionGroupProcedure(allocatedReplicaSet)));
+ Assert.assertTrue(
+ cleanupCaptor
+ .getAllValues()
+ .contains(new RemoveRegionGroupProcedure(otherAllocatedReplicaSet)));
+ }
+
+ @Test
+ public void testFencedBeforeCreateRpcDoesNotSubmitCleanup() {
+ final String database = "root.sg";
+ final CreateRegionGroupsPlan createPlan = new CreateRegionGroupsPlan();
+ createPlan.addRegionGroup(
+ database,
+ new TRegionReplicaSet(new TConsensusGroupId(DataRegion, 10), Collections.emptyList()));
+ final TestCreateRegionGroupsProcedure procedure =
+ new TestCreateRegionGroupsProcedure(
+ DataRegion, createPlan, new CreateRegionGroupsPlan(), Collections.emptyMap());
+
+ final ConfigNodeProcedureEnv env = Mockito.mock(ConfigNodeProcedureEnv.class);
+ Mockito.when(env.validateCreateRegionGroups(createPlan))
+ .thenReturn(new TSStatus(TSStatusCode.DATABASE_CONFIG_ERROR.getStatusCode()));
+ procedure.executeCreate(env);
+
+ Mockito.verify(env, Mockito.never())
+ .doRegionCreation(Mockito.any(), Mockito.any(CreateRegionGroupsPlan.class));
+ Mockito.verify(env, Mockito.never()).getConfigManager();
+ Assert.assertTrue(procedure.isFailed());
+ }
+
+ @Test
+ public void testFencedAfterCreateRpcCleansEveryPlannedRegionReplica() {
+ final TRegionReplicaSet replicaSet =
+ new TRegionReplicaSet(
+ new TConsensusGroupId(DataRegion, 10),
+ Collections.singletonList(new TDataNodeLocation().setDataNodeId(1)));
+ final CreateRegionGroupsPlan createPlan = new CreateRegionGroupsPlan();
+ createPlan.addRegionGroup("root.sg", replicaSet);
+ final TestCreateRegionGroupsProcedure procedure =
+ new TestCreateRegionGroupsProcedure(
+ DataRegion, createPlan, new CreateRegionGroupsPlan(), Collections.emptyMap());
+
+ final ConfigNodeProcedureEnv env = Mockito.mock(ConfigNodeProcedureEnv.class);
+ final ConfigManager configManager = Mockito.mock(ConfigManager.class);
+ final ProcedureManager procedureManager = Mockito.mock(ProcedureManager.class);
+ @SuppressWarnings("unchecked")
+ final ProcedureExecutor executor =
+ Mockito.mock(ProcedureExecutor.class);
+ Mockito.when(env.validateCreateRegionGroups(createPlan))
+ .thenReturn(new TSStatus(TSStatusCode.DATABASE_NOT_EXIST.getStatusCode()));
+ Mockito.when(env.getConfigManager()).thenReturn(configManager);
+ Mockito.when(configManager.getProcedureManager()).thenReturn(procedureManager);
+ Mockito.when(procedureManager.getExecutor()).thenReturn(executor);
+
+ procedure.executeShunt(env);
+
+ final ArgumentCaptor> cleanupCaptor =
+ ArgumentCaptor.forClass(Procedure.class);
+ Mockito.verify(executor).submitProcedure(cleanupCaptor.capture());
+ Assert.assertEquals(new RemoveRegionGroupProcedure(replicaSet), cleanupCaptor.getValue());
+ Mockito.verify(env, Mockito.never()).persistRegionGroup(Mockito.any());
+ Assert.assertTrue(procedure.isFailed());
+ }
+
+ @Test
+ public void testFencedAfterPersistenceDoesNotSubmitCleanup() {
+ final CreateRegionGroupsPlan createPlan = new CreateRegionGroupsPlan();
+ createPlan.addRegionGroup(
+ "root.sg",
+ new TRegionReplicaSet(new TConsensusGroupId(DataRegion, 10), Collections.emptyList()));
+ final TestCreateRegionGroupsProcedure procedure =
+ new TestCreateRegionGroupsProcedure(
+ DataRegion, createPlan, createPlan, Collections.emptyMap());
+
+ final ConfigNodeProcedureEnv env = Mockito.mock(ConfigNodeProcedureEnv.class);
+ Mockito.when(env.validateCreateRegionGroups(createPlan))
+ .thenReturn(new TSStatus(TSStatusCode.DATABASE_NOT_EXIST.getStatusCode()));
+
+ procedure.executePostPersist(env);
+
+ Mockito.verify(env, Mockito.never()).getConfigManager();
+ Assert.assertTrue(procedure.isFailed());
+ }
+
+ @Test
+ public void testCreateAndDeleteDatabaseLifecycleAreMutuallyExclusive() {
+ final String database = "root.sg";
+ final CreateRegionGroupsPlan createPlan = new CreateRegionGroupsPlan();
+ createPlan.addRegionGroup(
+ database,
+ new TRegionReplicaSet(new TConsensusGroupId(DataRegion, 1), Collections.emptyList()));
+ final TestCreateRegionGroupsProcedure createProcedure =
+ new TestCreateRegionGroupsProcedure(
+ DataRegion, createPlan, new CreateRegionGroupsPlan(), Collections.emptyMap());
+ createProcedure.setProcId(1);
+ final TestDeleteDatabaseProcedure deleteProcedure =
+ new TestDeleteDatabaseProcedure(new TDatabaseSchema(database));
+ deleteProcedure.setProcId(2);
+
+ final ConfigNodeProcedureEnv env =
+ new ConfigNodeProcedureEnv(
+ Mockito.mock(ConfigManager.class), Mockito.mock(ProcedureScheduler.class));
+ Assert.assertEquals(ProcedureLockState.LOCK_ACQUIRED, createProcedure.acquireDatabaseLock(env));
+ Assert.assertEquals(
+ ProcedureLockState.LOCK_EVENT_WAIT, deleteProcedure.acquireDatabaseLock(env));
+
+ createProcedure.releaseDatabaseLock(env);
+ Assert.assertEquals(ProcedureLockState.LOCK_ACQUIRED, deleteProcedure.acquireDatabaseLock(env));
+ deleteProcedure.releaseDatabaseLock(env);
+ }
}
diff --git a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteDatabaseProcedureTest.java b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteDatabaseProcedureTest.java
index b12f49d9bd7da..0d599dea05f53 100644
--- a/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteDatabaseProcedureTest.java
+++ b/iotdb-core/confignode/src/test/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteDatabaseProcedureTest.java
@@ -19,6 +19,11 @@
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.TEndPoint;
+import org.apache.iotdb.common.rpc.thrift.TRegionReplicaSet;
import org.apache.iotdb.confignode.procedure.store.ProcedureFactory;
import org.apache.iotdb.confignode.rpc.thrift.TDatabaseSchema;
@@ -27,6 +32,7 @@
import java.io.DataOutputStream;
import java.nio.ByteBuffer;
+import java.util.Collections;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.fail;
@@ -38,7 +44,20 @@ public void serializeDeserializeTest() {
PublicBAOS byteArrayOutputStream = new PublicBAOS();
DataOutputStream outputStream = new DataOutputStream(byteArrayOutputStream);
- DeleteDatabaseProcedure p1 = new DeleteDatabaseProcedure(new TDatabaseSchema("root.sg"), false);
+ TRegionReplicaSet regionReplicaSet =
+ new TRegionReplicaSet(
+ new TConsensusGroupId(TConsensusGroupType.DataRegion, 1),
+ Collections.singletonList(
+ new TDataNodeLocation()
+ .setDataNodeId(1)
+ .setClientRpcEndPoint(new TEndPoint("127.0.0.1", 6667))
+ .setInternalEndPoint(new TEndPoint("127.0.0.1", 10730))
+ .setMPPDataExchangeEndPoint(new TEndPoint("127.0.0.1", 10740))
+ .setDataRegionConsensusEndPoint(new TEndPoint("127.0.0.1", 10760))
+ .setSchemaRegionConsensusEndPoint(new TEndPoint("127.0.0.1", 10750))));
+ DeleteDatabaseProcedure p1 =
+ new DeleteDatabaseProcedure(
+ new TDatabaseSchema("root.sg"), false, Collections.singletonList(regionReplicaSet));
try {
p1.serialize(outputStream);
diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/consensus/SchemaRegionConsensusImpl.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/consensus/SchemaRegionConsensusImpl.java
index f6a1175da1998..e204f568c61a3 100644
--- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/consensus/SchemaRegionConsensusImpl.java
+++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/consensus/SchemaRegionConsensusImpl.java
@@ -24,6 +24,7 @@
import org.apache.iotdb.commons.conf.CommonConfig;
import org.apache.iotdb.commons.conf.CommonDescriptor;
import org.apache.iotdb.commons.consensus.SchemaRegionId;
+import org.apache.iotdb.commons.utils.TestOnly;
import org.apache.iotdb.consensus.ConsensusFactory;
import org.apache.iotdb.consensus.IConsensus;
import org.apache.iotdb.consensus.config.ConsensusConfig;
@@ -52,6 +53,11 @@ public static IConsensus getInstance() {
return SchemaRegionConsensusImplHolder.INSTANCE;
}
+ @TestOnly
+ public static void setInstance(final IConsensus instance) {
+ SchemaRegionConsensusImplHolder.INSTANCE = instance;
+ }
+
public static void reinitializeStatics() {
SchemaRegionConsensusImplHolder.reinitializeStatics();
}
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..a8c0507b2b1be 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
@@ -30,6 +30,7 @@
import org.apache.iotdb.commons.exception.IllegalPathException;
import org.apache.iotdb.commons.exception.MetadataException;
import org.apache.iotdb.commons.log.LoggerPeriodicalLogReducer;
+import org.apache.iotdb.commons.utils.TestOnly;
import org.apache.iotdb.consensus.common.Peer;
import org.apache.iotdb.consensus.exception.ConsensusException;
import org.apache.iotdb.consensus.exception.ConsensusGroupAlreadyExistException;
@@ -47,9 +48,11 @@
import org.slf4j.LoggerFactory;
import java.util.ArrayList;
+import java.util.BitSet;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.locks.ReentrantLock;
import java.util.concurrent.locks.ReentrantReadWriteLock;
/**
@@ -60,13 +63,18 @@ public class DataNodeRegionManager {
private static final Logger LOGGER = LoggerFactory.getLogger(DataNodeRegionManager.class);
- private final SchemaEngine schemaEngine = SchemaEngine.getInstance();
- private final StorageEngine storageEngine = StorageEngine.getInstance();
+ private final SchemaEngine schemaEngine;
+ private final StorageEngine storageEngine;
private final Map schemaRegionLockMap =
new ConcurrentHashMap<>();
private final Map dataRegionLockMap =
new ConcurrentHashMap<>();
+ private final Object deletedRegionGroupFenceLock = new Object();
+ private final BitSet deletedSchemaRegionGroups = new BitSet();
+ private final BitSet deletedDataRegionGroups = new BitSet();
+ private static final int REGION_CREATION_LOCK_COUNT = 256;
+ private final ReentrantLock[] regionCreationLocks = new ReentrantLock[REGION_CREATION_LOCK_COUNT];
private static class DataNodeRegionManagerHolder {
private static final DataNodeRegionManager INSTANCE = new DataNodeRegionManager();
@@ -95,9 +103,24 @@ public void init() {
public void clear() {
schemaRegionLockMap.clear();
dataRegionLockMap.clear();
+ synchronized (deletedRegionGroupFenceLock) {
+ deletedSchemaRegionGroups.clear();
+ deletedDataRegionGroups.clear();
+ }
+ }
+
+ private DataNodeRegionManager() {
+ this(SchemaEngine.getInstance(), StorageEngine.getInstance());
}
- private DataNodeRegionManager() {}
+ @TestOnly
+ DataNodeRegionManager(SchemaEngine schemaEngine, StorageEngine storageEngine) {
+ this.schemaEngine = schemaEngine;
+ this.storageEngine = storageEngine;
+ for (int i = 0; i < REGION_CREATION_LOCK_COUNT; i++) {
+ regionCreationLocks[i] = new ReentrantLock();
+ }
+ }
public ReentrantReadWriteLock getRegionLock(ConsensusGroupId consensusGroupId) {
return consensusGroupId instanceof DataRegionId
@@ -110,9 +133,21 @@ public TSStatus createSchemaRegion(
TSStatus tsStatus;
final SchemaRegionId schemaRegionId =
new SchemaRegionId(regionReplicaSet.getRegionId().getId());
+ final ReentrantLock creationLock = getRegionCreationLock(schemaRegionId);
+ creationLock.lock();
+ boolean localRegionExisted = true;
+ boolean consensusGroupExisted = true;
+ boolean localRegionCreated = false;
try {
- schemaEngine.createSchemaRegion(storageGroup, schemaRegionId);
- schemaRegionLockMap.put(schemaRegionId, new ReentrantReadWriteLock(false));
+ if (isRegionGroupDeleted(schemaRegionId)) {
+ return new TSStatus(TSStatusCode.CREATE_REGION_ERROR.getStatusCode());
+ }
+ localRegionExisted = schemaEngine.getSchemaRegion(schemaRegionId) != null;
+ consensusGroupExisted =
+ SchemaRegionConsensusImpl.getInstance()
+ .getAllConsensusGroupIds()
+ .contains(schemaRegionId);
+ localRegionCreated = schemaEngine.createSchemaRegionIfAbsent(storageGroup, schemaRegionId);
final List peers = new ArrayList<>();
for (final TDataNodeLocation dataNodeLocation : regionReplicaSet.getDataNodeLocations()) {
final TEndPoint endpoint =
@@ -122,6 +157,7 @@ public TSStatus createSchemaRegion(
peers.add(new Peer(schemaRegionId, dataNodeLocation.getDataNodeId(), endpoint));
}
SchemaRegionConsensusImpl.getInstance().createLocalPeer(schemaRegionId, peers);
+ schemaRegionLockMap.putIfAbsent(schemaRegionId, new ReentrantReadWriteLock(false));
tsStatus = new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
} catch (final IllegalPathException e1) {
LOGGER.error(DataNodeMiscMessages.CREATE_SCHEMA_REGION_FAILED_ILLEGAL_PATH, storageGroup);
@@ -138,13 +174,29 @@ public TSStatus createSchemaRegion(
tsStatus.setMessage(
String.format(DataNodeMiscMessages.CREATE_SCHEMA_REGION_FAILED_FMT, e2.getMessage()));
} catch (final ConsensusGroupAlreadyExistException e) {
+ schemaRegionLockMap.putIfAbsent(schemaRegionId, new ReentrantReadWriteLock(false));
tsStatus = new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
tsStatus.setMessage(
String.format(
DataNodeMiscMessages.SCHEMA_REGION_ALREADY_EXISTS_FMT, schemaRegionId.getId()));
} catch (final ConsensusException e) {
+ rollbackSchemaRegionCreation(
+ schemaRegionId, localRegionCreated, consensusGroupExisted, storageGroup);
tsStatus = new TSStatus(TSStatusCode.CREATE_REGION_ERROR.getStatusCode());
tsStatus.setMessage(e.getMessage());
+ } catch (final RuntimeException | OutOfMemoryError e) {
+ rollbackSchemaRegionCreation(
+ schemaRegionId,
+ localRegionCreated
+ || (!localRegionExisted && schemaEngine.getSchemaRegion(schemaRegionId) != null),
+ consensusGroupExisted,
+ storageGroup);
+ LOGGER.error(DataNodeMiscMessages.CREATE_SCHEMA_REGION_FAILED, storageGroup, e.getMessage());
+ tsStatus = new TSStatus(TSStatusCode.CREATE_REGION_ERROR.getStatusCode());
+ tsStatus.setMessage(
+ String.format(DataNodeMiscMessages.CREATE_SCHEMA_REGION_FAILED_FMT, e.getMessage()));
+ } finally {
+ creationLock.unlock();
}
return tsStatus;
}
@@ -152,9 +204,19 @@ public TSStatus createSchemaRegion(
public TSStatus createDataRegion(TRegionReplicaSet regionReplicaSet, String storageGroup) {
TSStatus tsStatus;
DataRegionId dataRegionId = new DataRegionId(regionReplicaSet.getRegionId().getId());
+ final ReentrantLock creationLock = getRegionCreationLock(dataRegionId);
+ creationLock.lock();
+ boolean localRegionExisted = true;
+ boolean consensusGroupExisted = true;
+ boolean localRegionCreated = false;
try {
- storageEngine.createDataRegion(dataRegionId, storageGroup);
- dataRegionLockMap.put(dataRegionId, new ReentrantReadWriteLock(false));
+ if (isRegionGroupDeleted(dataRegionId)) {
+ return new TSStatus(TSStatusCode.CREATE_REGION_ERROR.getStatusCode());
+ }
+ localRegionExisted = storageEngine.getDataRegion(dataRegionId) != null;
+ consensusGroupExisted =
+ DataRegionConsensusImpl.getInstance().getAllConsensusGroupIds().contains(dataRegionId);
+ localRegionCreated = storageEngine.createDataRegionIfAbsent(dataRegionId, storageGroup);
List peers = new ArrayList<>();
for (TDataNodeLocation dataNodeLocation : regionReplicaSet.getDataNodeLocations()) {
TEndPoint endpoint =
@@ -164,6 +226,7 @@ public TSStatus createDataRegion(TRegionReplicaSet regionReplicaSet, String stor
peers.add(new Peer(dataRegionId, dataNodeLocation.getDataNodeId(), endpoint));
}
DataRegionConsensusImpl.getInstance().createLocalPeer(dataRegionId, peers);
+ dataRegionLockMap.putIfAbsent(dataRegionId, new ReentrantReadWriteLock(false));
tsStatus = new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
} catch (DataRegionException e) {
LOGGER.error(DataNodeMiscMessages.CREATE_DATA_REGION_FAILED, storageGroup, e.getMessage());
@@ -171,16 +234,114 @@ public TSStatus createDataRegion(TRegionReplicaSet regionReplicaSet, String stor
tsStatus.setMessage(
String.format(DataNodeMiscMessages.CREATE_DATA_REGION_FAILED_FMT, e.getMessage()));
} catch (ConsensusGroupAlreadyExistException e) {
+ dataRegionLockMap.putIfAbsent(dataRegionId, new ReentrantReadWriteLock(false));
tsStatus = new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
tsStatus.setMessage(
String.format(DataNodeMiscMessages.DATA_REGION_ALREADY_EXISTS_FMT, dataRegionId.getId()));
} catch (ConsensusException e) {
+ rollbackDataRegionCreation(
+ dataRegionId, localRegionCreated, consensusGroupExisted, storageGroup);
tsStatus = new TSStatus(TSStatusCode.CREATE_REGION_ERROR.getStatusCode());
tsStatus.setMessage(e.getMessage());
+ } catch (RuntimeException | OutOfMemoryError e) {
+ rollbackDataRegionCreation(
+ dataRegionId,
+ localRegionCreated
+ || (!localRegionExisted && storageEngine.getDataRegion(dataRegionId) != null),
+ consensusGroupExisted,
+ storageGroup);
+ LOGGER.error(DataNodeMiscMessages.CREATE_DATA_REGION_FAILED, storageGroup, e.getMessage());
+ tsStatus = new TSStatus(TSStatusCode.CREATE_REGION_ERROR.getStatusCode());
+ tsStatus.setMessage(
+ String.format(DataNodeMiscMessages.CREATE_DATA_REGION_FAILED_FMT, e.getMessage()));
+ } finally {
+ creationLock.unlock();
}
return tsStatus;
}
+ private ReentrantLock getRegionCreationLock(ConsensusGroupId regionId) {
+ return regionCreationLocks[
+ (regionId.hashCode() & Integer.MAX_VALUE) % REGION_CREATION_LOCK_COUNT];
+ }
+
+ /**
+ * Fences delayed create RPCs after a whole RegionGroup starts deletion.
+ *
+ * RegionGroup ids are never reused. The caller acquires the same striped lock as creation, so
+ * an already-running creation finishes before deletion starts, while every later creation is
+ * rejected even if it came from an old ConfigNode leader.
+ */
+ public void markRegionGroupDeleted(ConsensusGroupId regionId) {
+ final ReentrantLock creationLock = getRegionCreationLock(regionId);
+ creationLock.lock();
+ try {
+ synchronized (deletedRegionGroupFenceLock) {
+ getDeletedRegionGroupSet(regionId).set(regionId.getId());
+ }
+ } finally {
+ creationLock.unlock();
+ }
+ }
+
+ private boolean isRegionGroupDeleted(ConsensusGroupId regionId) {
+ synchronized (deletedRegionGroupFenceLock) {
+ return getDeletedRegionGroupSet(regionId).get(regionId.getId());
+ }
+ }
+
+ private BitSet getDeletedRegionGroupSet(ConsensusGroupId regionId) {
+ return regionId instanceof DataRegionId ? deletedDataRegionGroups : deletedSchemaRegionGroups;
+ }
+
+ private void rollbackDataRegionCreation(
+ DataRegionId regionId,
+ boolean localRegionCreated,
+ boolean consensusGroupExisted,
+ String storageGroup) {
+ rollbackConsensusPeer(
+ DataRegionConsensusImpl.getInstance(), regionId, consensusGroupExisted, storageGroup);
+ if (localRegionCreated && !consensusGroupExisted) {
+ storageEngine.deleteDataRegion(regionId);
+ dataRegionLockMap.remove(regionId);
+ }
+ }
+
+ private void rollbackSchemaRegionCreation(
+ SchemaRegionId regionId,
+ boolean localRegionCreated,
+ boolean consensusGroupExisted,
+ String storageGroup) {
+ rollbackConsensusPeer(
+ SchemaRegionConsensusImpl.getInstance(), regionId, consensusGroupExisted, storageGroup);
+ if (localRegionCreated && !consensusGroupExisted) {
+ try {
+ schemaEngine.deleteSchemaRegion(regionId);
+ schemaRegionLockMap.remove(regionId);
+ } catch (MetadataException e) {
+ LOGGER.error(
+ DataNodeMiscMessages.CREATE_SCHEMA_REGION_FAILED, storageGroup, e.getMessage());
+ }
+ }
+ }
+
+ private void rollbackConsensusPeer(
+ org.apache.iotdb.consensus.IConsensus consensus,
+ ConsensusGroupId regionId,
+ boolean consensusGroupExisted,
+ String storageGroup) {
+ if (consensusGroupExisted) {
+ return;
+ }
+ try {
+ if (consensus.getAllConsensusGroupIds().contains(regionId)) {
+ consensus.deleteLocalPeer(regionId);
+ }
+ } catch (ConsensusException | RuntimeException | OutOfMemoryError e) {
+ LOGGER.error(DataNodeMiscMessages.CREATE_DATA_REGION_FAILED, storageGroup, e.getMessage());
+ }
+ }
+
public TSStatus createNewRegion(final ConsensusGroupId regionId, final String storageGroup) {
final TSStatus status = new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
LOGGER.info(DataNodeMiscMessages.START_CREATE_NEW_REGION, regionId);
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 a52330073d8de..3222c54b0878a 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
@@ -279,13 +279,23 @@ public void updateSubtreeMeasurementCountForTemplate(final int templateId, final
public synchronized void createSchemaRegion(
final String storageGroup, final SchemaRegionId schemaRegionId) throws MetadataException {
+ createSchemaRegionIfAbsent(storageGroup, schemaRegionId);
+ }
+
+ /**
+ * Atomically creates and registers a SchemaRegion if it is absent.
+ *
+ * @return true only when this invocation created the Region
+ */
+ public synchronized boolean createSchemaRegionIfAbsent(
+ final String storageGroup, final SchemaRegionId schemaRegionId) throws MetadataException {
if (this.schemaRegionMap == null) {
throw new MetadataException(DataNodeSchemaMessages.PEER_IS_SHUTTING_DOWN);
}
final ISchemaRegion schemaRegion = this.schemaRegionMap.get(schemaRegionId);
if (schemaRegion != null) {
if (schemaRegion.getDatabaseFullPath().equals(storageGroup)) {
- return;
+ return false;
} else {
throw new MetadataException(
String.format(
@@ -297,6 +307,7 @@ public synchronized void createSchemaRegion(
}
this.schemaRegionMap.put(
schemaRegionId, createSchemaRegionWithoutExistenceCheck(storageGroup, schemaRegionId));
+ return true;
}
private Callable recoverSchemaRegionTask(
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..74c42fdd2fa11 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
@@ -508,6 +508,15 @@ public DeleteOldRegionPeerTask(
@Override
public void run() {
+ // Negative task ids are reserved for RemoveRegionGroupProcedure. Fence delayed create RPCs
+ // before deleting the peer so an old ConfigNode leader cannot recreate this RegionGroup
+ // after DROP has completed.
+ if (taskId < 0) {
+ DataNodeRegionManager.getInstance()
+ .markRegionGroupDeleted(
+ ConsensusGroupId.Factory.createFromTConsensusGroupId(tRegionId));
+ }
+
// deletePeer: remove the peer from the consensus group
TSStatus runResult = deletePeer();
if (isFailed(runResult)) {
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..19a76c249ccc7 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
@@ -801,13 +801,26 @@ private void makeSureNoOldRegion(DataRegionId regionId) {
// the local storage before adding the corresponding consensusGroup to the consensus layer
public void createDataRegion(DataRegionId regionId, String databaseName)
throws DataRegionException {
+ createDataRegionIfAbsent(regionId, databaseName);
+ }
+
+ /**
+ * Atomically creates and registers a DataRegion if it is absent.
+ *
+ * @return true only when this invocation created the Region
+ */
+ public boolean createDataRegionIfAbsent(DataRegionId regionId, String databaseName)
+ throws DataRegionException {
makeSureNoOldRegion(regionId);
AtomicReference exceptionAtomicReference = new AtomicReference<>(null);
+ AtomicBoolean created = new AtomicBoolean(false);
dataRegionMap.computeIfAbsent(
regionId,
region -> {
try {
- return buildNewDataRegion(databaseName, region);
+ final DataRegion dataRegion = buildNewDataRegion(databaseName, region);
+ created.set(true);
+ return dataRegion;
} catch (DataRegionException e) {
exceptionAtomicReference.set(e);
}
@@ -817,6 +830,7 @@ public void createDataRegion(DataRegionId regionId, String databaseName)
if (exceptionAtomicReference.get() != null) {
throw exceptionAtomicReference.get();
}
+ return created.get();
}
public void deleteDataRegion(DataRegionId regionId) {
diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java
index d2619a9de65bf..153832ddcabee 100644
--- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java
+++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java
@@ -418,68 +418,87 @@ public DataRegion(
config.getDelayAnalyzerConfidenceLevel())
: null;
acquireDirectBufferMemory();
+ ExecutorService createdUpgradeModFileThreadPool = null;
+ DataRegionMetrics createdMetrics = null;
+ boolean tableDiskUsageRegistered = false;
+ try {
+ dataRegionSysDir = SystemFileFactory.INSTANCE.getFile(systemDir, dataRegionIdString);
+ this.tsFileManager =
+ new TsFileManager(databaseName, dataRegionIdString, dataRegionSysDir.getPath());
+ if (dataRegionSysDir.mkdirs()) {
+ logger.info(
+ StorageEngineMessages
+ .STORAGE_LOG_DATABASE_SYSTEM_DIRECTORY_DOESN_T_EXIST_CREATE_IT_9C0E7C68,
+ dataRegionSysDir.getPath());
+ } else if (!dataRegionSysDir.exists()) {
+ logger.error(StorageEngineMessages.CREATE_DB_SYSTEM_DIR_FAILED, dataRegionSysDir.getPath());
+ }
- dataRegionSysDir = SystemFileFactory.INSTANCE.getFile(systemDir, dataRegionIdString);
- this.tsFileManager =
- new TsFileManager(databaseName, dataRegionIdString, dataRegionSysDir.getPath());
- if (dataRegionSysDir.mkdirs()) {
- logger.info(
- StorageEngineMessages
- .STORAGE_LOG_DATABASE_SYSTEM_DIRECTORY_DOESN_T_EXIST_CREATE_IT_9C0E7C68,
- dataRegionSysDir.getPath());
- } else if (!dataRegionSysDir.exists()) {
- logger.error(StorageEngineMessages.CREATE_DB_SYSTEM_DIR_FAILED, dataRegionSysDir.getPath());
- }
-
- lastFlushTimeMap = new HashLastFlushTimeMap();
- upgradeModFileThreadPool =
- IoTDBThreadPoolFactory.newSingleThreadExecutor(
- databaseName + "-" + dataRegionIdString + "-UpgradeMod");
+ lastFlushTimeMap = new HashLastFlushTimeMap();
+ createdUpgradeModFileThreadPool =
+ IoTDBThreadPoolFactory.newSingleThreadExecutor(
+ databaseName + "-" + dataRegionIdString + "-UpgradeMod");
+ upgradeModFileThreadPool = createdUpgradeModFileThreadPool;
- TableDiskUsageIndex.getInstance().registerRegion(this);
+ TableDiskUsageIndex.getInstance().registerRegion(this);
+ tableDiskUsageRegistered = isTableModel;
- // recover tsfiles unless consensus protocol is ratis and storage engine is not ready
- if (config.getDataRegionConsensusProtocolClass().equals(ConsensusFactory.RATIS_CONSENSUS)
- && !StorageEngine.getInstance().isReadyForReadAndWrite()) {
- logger.debug(
- StorageEngineMessages
- .STORAGE_LOG_SKIP_RECOVERING_DATA_REGION_WHEN_CONSENSUS_PROTOCOL_IS_RATIS_43A6A699,
- databaseName,
- dataRegionIdString);
- for (String fileFolder : TierManager.getInstance().getAllFilesFolders()) {
- File dataRegionFolder =
- fsFactory.getFile(fileFolder, databaseName + File.separator + dataRegionIdString);
- try {
- fsFactory.deleteDirectory(dataRegionFolder.getPath());
- } catch (IOException e) {
- logger.error(
- StorageEngineMessages
- .STORAGE_LOG_EXCEPTION_OCCURS_WHEN_DELETING_DATA_REGION_FOLDER_FOR_8ABCF5D1,
- databaseName,
- dataRegionIdString,
- e);
- }
- if (FSUtils.getFSType(dataRegionFolder) == FSType.LOCAL) {
- if (dataRegionFolder.mkdirs()) {
- logger.info(
- StorageEngineMessages
- .STORAGE_LOG_DATA_REGION_DIRECTORY_DOESN_T_EXIST_CREATE_IT_EFB0AE77,
- dataRegionFolder.getPath());
- } else if (!dataRegionFolder.exists()) {
+ // recover tsfiles unless consensus protocol is ratis and storage engine is not ready
+ if (config.getDataRegionConsensusProtocolClass().equals(ConsensusFactory.RATIS_CONSENSUS)
+ && !StorageEngine.getInstance().isReadyForReadAndWrite()) {
+ logger.debug(
+ StorageEngineMessages
+ .STORAGE_LOG_SKIP_RECOVERING_DATA_REGION_WHEN_CONSENSUS_PROTOCOL_IS_RATIS_43A6A699,
+ databaseName,
+ dataRegionIdString);
+ for (String fileFolder : TierManager.getInstance().getAllFilesFolders()) {
+ File dataRegionFolder =
+ fsFactory.getFile(fileFolder, databaseName + File.separator + dataRegionIdString);
+ try {
+ fsFactory.deleteDirectory(dataRegionFolder.getPath());
+ } catch (IOException e) {
logger.error(
- StorageEngineMessages.CREATE_DATA_REGION_DIR_FAILED, dataRegionFolder.getPath());
+ StorageEngineMessages
+ .STORAGE_LOG_EXCEPTION_OCCURS_WHEN_DELETING_DATA_REGION_FOLDER_FOR_8ABCF5D1,
+ databaseName,
+ dataRegionIdString,
+ e);
+ }
+ if (FSUtils.getFSType(dataRegionFolder) == FSType.LOCAL) {
+ if (dataRegionFolder.mkdirs()) {
+ logger.info(
+ StorageEngineMessages
+ .STORAGE_LOG_DATA_REGION_DIRECTORY_DOESN_T_EXIST_CREATE_IT_EFB0AE77,
+ dataRegionFolder.getPath());
+ } else if (!dataRegionFolder.exists()) {
+ logger.error(
+ StorageEngineMessages.CREATE_DATA_REGION_DIR_FAILED, dataRegionFolder.getPath());
+ }
}
}
+ } else {
+ asyncTsFileResourceRecoverTaskList = new ArrayList<>();
+ recover();
}
- } else {
- asyncTsFileResourceRecoverTaskList = new ArrayList<>();
- recover();
- }
- initDiskSelector();
+ initDiskSelector();
- this.metrics = new DataRegionMetrics(this);
- MetricService.getInstance().addMetricSet(metrics);
+ createdMetrics = new DataRegionMetrics(this);
+ this.metrics = createdMetrics;
+ MetricService.getInstance().addMetricSet(metrics);
+ } catch (DataRegionException | RuntimeException | Error e) {
+ if (createdMetrics != null) {
+ MetricService.getInstance().removeMetricSet(createdMetrics);
+ }
+ if (tableDiskUsageRegistered) {
+ TableDiskUsageIndex.getInstance().remove(databaseName, dataRegionId.getId());
+ }
+ if (createdUpgradeModFileThreadPool != null) {
+ createdUpgradeModFileThreadPool.shutdownNow();
+ }
+ releaseDirectBufferMemory();
+ throw e;
+ }
}
@TestOnly
@@ -5248,6 +5267,9 @@ public void markDeleted() {
writeLock("markDeleted");
try {
deleted = true;
+ if (upgradeModFileThreadPool != null) {
+ upgradeModFileThreadPool.shutdownNow();
+ }
releaseDirectBufferMemory();
MetricService.getInstance().removeMetricSet(metrics);
deletedCondition.signalAll();
diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeRegionManagerTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeRegionManagerTest.java
new file mode 100644
index 0000000000000..7227dd69f5925
--- /dev/null
+++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeRegionManagerTest.java
@@ -0,0 +1,191 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.db.protocol.thrift.impl;
+
+import org.apache.iotdb.common.rpc.thrift.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.common.rpc.thrift.TSStatus;
+import org.apache.iotdb.commons.consensus.ConsensusGroupId;
+import org.apache.iotdb.commons.consensus.DataRegionId;
+import org.apache.iotdb.commons.consensus.SchemaRegionId;
+import org.apache.iotdb.consensus.IConsensus;
+import org.apache.iotdb.consensus.exception.ConsensusException;
+import org.apache.iotdb.db.consensus.DataRegionConsensusImpl;
+import org.apache.iotdb.db.consensus.SchemaRegionConsensusImpl;
+import org.apache.iotdb.db.schemaengine.SchemaEngine;
+import org.apache.iotdb.db.storageengine.StorageEngine;
+import org.apache.iotdb.db.storageengine.dataregion.DataRegion;
+import org.apache.iotdb.rpc.TSStatusCode;
+
+import org.junit.After;
+import org.junit.Assert;
+import org.junit.Before;
+import org.junit.Test;
+import org.mockito.Mockito;
+
+import java.util.Collections;
+import java.util.List;
+
+public class DataNodeRegionManagerTest {
+
+ private IConsensus previousDataConsensus;
+ private IConsensus previousSchemaConsensus;
+ private IConsensus dataConsensus;
+ private IConsensus schemaConsensus;
+ private StorageEngine storageEngine;
+ private SchemaEngine schemaEngine;
+ private DataNodeRegionManager regionManager;
+
+ @Before
+ public void setUp() {
+ previousDataConsensus = DataRegionConsensusImpl.getInstance();
+ previousSchemaConsensus = SchemaRegionConsensusImpl.getInstance();
+ dataConsensus = Mockito.mock(IConsensus.class);
+ schemaConsensus = Mockito.mock(IConsensus.class);
+ DataRegionConsensusImpl.setInstance(dataConsensus);
+ SchemaRegionConsensusImpl.setInstance(schemaConsensus);
+ storageEngine = Mockito.mock(StorageEngine.class);
+ schemaEngine = Mockito.mock(SchemaEngine.class);
+ regionManager = new DataNodeRegionManager(schemaEngine, storageEngine);
+ }
+
+ @After
+ public void tearDown() {
+ DataRegionConsensusImpl.setInstance(previousDataConsensus);
+ SchemaRegionConsensusImpl.setInstance(previousSchemaConsensus);
+ }
+
+ @Test
+ public void testDataRegionConsensusOomRollsBackNewLocalState() throws Exception {
+ final DataRegionId regionId = new DataRegionId(1);
+ final List noConsensusGroups = Collections.emptyList();
+ final List partiallyCreatedGroup = Collections.singletonList(regionId);
+ Mockito.when(dataConsensus.getAllConsensusGroupIds())
+ .thenReturn(noConsensusGroups, partiallyCreatedGroup);
+ Mockito.when(storageEngine.getDataRegion(regionId)).thenReturn(null);
+ Mockito.when(storageEngine.createDataRegionIfAbsent(regionId, "root.sg")).thenReturn(true);
+ Mockito.doThrow(new OutOfMemoryError("WAL direct memory exhausted"))
+ .when(dataConsensus)
+ .createLocalPeer(Mockito.eq(regionId), Mockito.anyList());
+
+ final TSStatus status =
+ regionManager.createDataRegion(
+ createReplicaSet(TConsensusGroupType.DataRegion, 1), "root.sg");
+
+ Assert.assertEquals(TSStatusCode.CREATE_REGION_ERROR.getStatusCode(), status.getCode());
+ Mockito.verify(dataConsensus).deleteLocalPeer(regionId);
+ Mockito.verify(storageEngine).deleteDataRegion(regionId);
+ Assert.assertNull(regionManager.getRegionLock(regionId));
+ }
+
+ @Test
+ public void testSchemaRegionConsensusFailureRollsBackNewLocalState() throws Exception {
+ final SchemaRegionId regionId = new SchemaRegionId(2);
+ final List noConsensusGroups = Collections.emptyList();
+ final List partiallyCreatedGroup = Collections.singletonList(regionId);
+ Mockito.when(schemaConsensus.getAllConsensusGroupIds())
+ .thenReturn(noConsensusGroups, partiallyCreatedGroup);
+ Mockito.when(schemaEngine.getSchemaRegion(regionId)).thenReturn(null);
+ Mockito.when(schemaEngine.createSchemaRegionIfAbsent("root.sg", regionId)).thenReturn(true);
+ Mockito.doThrow(new ConsensusException("Ratis create failed"))
+ .when(schemaConsensus)
+ .createLocalPeer(Mockito.eq(regionId), Mockito.anyList());
+
+ final TSStatus status =
+ regionManager.createSchemaRegion(
+ createReplicaSet(TConsensusGroupType.SchemaRegion, 2), "root.sg");
+
+ Assert.assertEquals(TSStatusCode.CREATE_REGION_ERROR.getStatusCode(), status.getCode());
+ Mockito.verify(schemaConsensus).deleteLocalPeer(regionId);
+ Mockito.verify(schemaEngine).deleteSchemaRegion(regionId);
+ Assert.assertNull(regionManager.getRegionLock(regionId));
+ }
+
+ @Test
+ public void testFailedIdempotentRetryDoesNotDeleteExistingDataRegion() throws Exception {
+ final DataRegionId regionId = new DataRegionId(3);
+ Mockito.when(dataConsensus.getAllConsensusGroupIds()).thenReturn(Collections.emptyList());
+ Mockito.when(storageEngine.getDataRegion(regionId)).thenReturn(Mockito.mock(DataRegion.class));
+ Mockito.when(storageEngine.createDataRegionIfAbsent(regionId, "root.sg")).thenReturn(false);
+ Mockito.doThrow(new ConsensusException("consensus unavailable"))
+ .when(dataConsensus)
+ .createLocalPeer(Mockito.eq(regionId), Mockito.anyList());
+
+ final TSStatus status =
+ regionManager.createDataRegion(
+ createReplicaSet(TConsensusGroupType.DataRegion, 3), "root.sg");
+
+ Assert.assertEquals(TSStatusCode.CREATE_REGION_ERROR.getStatusCode(), status.getCode());
+ Mockito.verify(storageEngine, Mockito.never()).deleteDataRegion(regionId);
+ }
+
+ @Test
+ public void testLocalRollbackContinuesWhenConsensusRollbackFails() throws Exception {
+ final DataRegionId regionId = new DataRegionId(4);
+ Mockito.when(dataConsensus.getAllConsensusGroupIds())
+ .thenReturn(Collections.emptyList(), Collections.singletonList(regionId));
+ Mockito.when(storageEngine.getDataRegion(regionId)).thenReturn(null);
+ Mockito.when(storageEngine.createDataRegionIfAbsent(regionId, "root.sg")).thenReturn(true);
+ Mockito.doThrow(new ConsensusException("Ratis create failed"))
+ .when(dataConsensus)
+ .createLocalPeer(Mockito.eq(regionId), Mockito.anyList());
+ Mockito.doThrow(new ConsensusException("Ratis rollback failed"))
+ .when(dataConsensus)
+ .deleteLocalPeer(regionId);
+
+ final TSStatus status =
+ regionManager.createDataRegion(
+ createReplicaSet(TConsensusGroupType.DataRegion, 4), "root.sg");
+
+ Assert.assertEquals(TSStatusCode.CREATE_REGION_ERROR.getStatusCode(), status.getCode());
+ Mockito.verify(storageEngine).deleteDataRegion(regionId);
+ Assert.assertNull(regionManager.getRegionLock(regionId));
+ }
+
+ @Test
+ public void testDeletedRegionGroupRejectsLateCreation() throws Exception {
+ final DataRegionId regionId = new DataRegionId(5);
+ regionManager.markRegionGroupDeleted(regionId);
+
+ final TSStatus status =
+ regionManager.createDataRegion(
+ createReplicaSet(TConsensusGroupType.DataRegion, 5), "root.sg");
+
+ Assert.assertEquals(TSStatusCode.CREATE_REGION_ERROR.getStatusCode(), status.getCode());
+ Mockito.verify(storageEngine, Mockito.never())
+ .createDataRegionIfAbsent(Mockito.any(), Mockito.anyString());
+ Mockito.verify(dataConsensus, Mockito.never())
+ .createLocalPeer(Mockito.any(), Mockito.anyList());
+ }
+
+ private TRegionReplicaSet createReplicaSet(TConsensusGroupType type, int regionId) {
+ final TDataNodeLocation location =
+ new TDataNodeLocation()
+ .setDataNodeId(0)
+ .setInternalEndPoint(new TEndPoint("127.0.0.1", 10730))
+ .setDataRegionConsensusEndPoint(new TEndPoint("127.0.0.1", 10760))
+ .setSchemaRegionConsensusEndPoint(new TEndPoint("127.0.0.1", 10750));
+ return new TRegionReplicaSet(
+ new TConsensusGroupId(type, regionId), Collections.singletonList(location));
+ }
+}