Skip to content

Commit daf9dd8

Browse files
authored
Reject duplicate schema region IDs during recovery (#18754)
1 parent 94694d9 commit daf9dd8

4 files changed

Lines changed: 85 additions & 1 deletion

File tree

‎iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeSchemaMessages.java‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -44,6 +44,9 @@ public final class DataNodeSchemaMessages {
4444
public static final String PEER_IS_SHUTTING_DOWN = "Peer is shutting down now.";
4545
public static final String SCHEMA_REGION_DUPLICATED =
4646
"SchemaRegion [%s] is duplicated between [%s] and [%s], and the former one has been recovered.";
47+
public static final String
48+
EXCEPTION_CANNOT_RECOVER_DUPLICATED_SCHEMAREGION_ARG_FOUND_IN_DATABASES_ARG_AND_ARG_D9F60E05 =
49+
"Cannot recover duplicated SchemaRegion [%s] found in databases [%s] and [%s].";
4750

4851
// ======================== MemSchemaEngineStatistics ========================
4952

‎iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeSchemaMessages.java‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,9 @@ public final class DataNodeSchemaMessages {
4242
public static final String PEER_IS_SHUTTING_DOWN = "节点正在关闭中。";
4343
public static final String SCHEMA_REGION_DUPLICATED =
4444
"SchemaRegion [%s] 在 [%s] 和 [%s] 之间重复,前者已被恢复。";
45+
public static final String
46+
EXCEPTION_CANNOT_RECOVER_DUPLICATED_SCHEMAREGION_ARG_FOUND_IN_DATABASES_ARG_AND_ARG_D9F60E05 =
47+
"无法恢复重复的 SchemaRegion [%s],其同时存在于数据库 [%s] 和 [%s] 中。";
4548

4649
// ======================== MemSchemaEngineStatistics 相关消息 ========================
4750

‎iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/SchemaEngine.java‎

Lines changed: 22 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -184,6 +184,7 @@ public static Map<String, List<SchemaRegionId>> getLocalSchemaRegionInfo() {
184184
private void initSchemaRegion() {
185185
// recover SchemaRegion concurrently
186186
final Map<String, List<SchemaRegionId>> localSchemaRegionInfo = getLocalSchemaRegionInfo();
187+
validateNoDuplicatedSchemaRegionId(localSchemaRegionInfo);
187188
final ExecutorService schemaRegionRecoverPools =
188189
IoTDBThreadPoolFactory.newFixedThreadPool(
189190
Runtime.getRuntime().availableProcessors(),
@@ -207,6 +208,27 @@ private void initSchemaRegion() {
207208
schemaRegionRecoverPools.shutdown();
208209
}
209210

211+
static void validateNoDuplicatedSchemaRegionId(
212+
final Map<String, List<SchemaRegionId>> localSchemaRegionInfo) {
213+
final Map<SchemaRegionId, String> databaseBySchemaRegionId = new HashMap<>();
214+
localSchemaRegionInfo.forEach(
215+
(database, schemaRegionIds) -> {
216+
for (final SchemaRegionId schemaRegionId : schemaRegionIds) {
217+
final String existingDatabase =
218+
databaseBySchemaRegionId.putIfAbsent(schemaRegionId, database);
219+
if (existingDatabase != null) {
220+
throw new IllegalStateException(
221+
String.format(
222+
DataNodeSchemaMessages
223+
.EXCEPTION_CANNOT_RECOVER_DUPLICATED_SCHEMAREGION_ARG_FOUND_IN_DATABASES_ARG_AND_ARG_D9F60E05,
224+
schemaRegionId,
225+
existingDatabase,
226+
database));
227+
}
228+
}
229+
});
230+
}
231+
210232
private void initSchemaEngineStatistics() {
211233
if (CommonDescriptor.getInstance().getConfig().getSchemaEngineMode().equals("Memory")) {
212234
schemaEngineStatistics = new MemSchemaEngineStatistics();
@@ -306,7 +328,6 @@ private Callable<ISchemaRegion> recoverSchemaRegionTask(
306328
return () -> {
307329
long timeRecord = System.currentTimeMillis();
308330
try {
309-
// TODO: handle duplicated regionId across different database
310331
ISchemaRegion schemaRegion =
311332
createSchemaRegionWithoutExistenceCheck(storageGroup, schemaRegionId);
312333
timeRecord = System.currentTimeMillis() - timeRecord;
Lines changed: 57 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,57 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one
3+
* or more contributor license agreements. See the NOTICE file
4+
* distributed with this work for additional information
5+
* regarding copyright ownership. The ASF licenses this file
6+
* to you under the Apache License, Version 2.0 (the
7+
* "License"); you may not use this file except in compliance
8+
* with the License. You may obtain a copy of the License at
9+
*
10+
* http://www.apache.org/licenses/LICENSE-2.0
11+
*
12+
* Unless required by applicable law or agreed to in writing,
13+
* software distributed under the License is distributed on an
14+
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
15+
* KIND, either express or implied. See the License for the
16+
* specific language governing permissions and limitations
17+
* under the License.
18+
*/
19+
package org.apache.iotdb.db.schemaengine;
20+
21+
import org.apache.iotdb.commons.consensus.SchemaRegionId;
22+
23+
import org.junit.Assert;
24+
import org.junit.Test;
25+
26+
import java.util.Arrays;
27+
import java.util.Collections;
28+
import java.util.HashMap;
29+
import java.util.List;
30+
import java.util.Map;
31+
32+
public class SchemaEngineTest {
33+
34+
@Test
35+
public void uniqueSchemaRegionIdsPassRecoveryValidation() {
36+
final Map<String, List<SchemaRegionId>> localSchemaRegionInfo = new HashMap<>();
37+
localSchemaRegionInfo.put(
38+
"root.db1", Arrays.asList(new SchemaRegionId(0), new SchemaRegionId(1)));
39+
localSchemaRegionInfo.put("root.db2", Collections.singletonList(new SchemaRegionId(2)));
40+
41+
SchemaEngine.validateNoDuplicatedSchemaRegionId(localSchemaRegionInfo);
42+
}
43+
44+
@Test
45+
public void duplicatedSchemaRegionIdIsRejectedBeforeRecovery() {
46+
final Map<String, List<SchemaRegionId>> localSchemaRegionInfo = new HashMap<>();
47+
localSchemaRegionInfo.put("root.db1", Collections.singletonList(new SchemaRegionId(0)));
48+
localSchemaRegionInfo.put("root.db2", Collections.singletonList(new SchemaRegionId(0)));
49+
50+
final IllegalStateException exception =
51+
Assert.assertThrows(
52+
IllegalStateException.class,
53+
() -> SchemaEngine.validateNoDuplicatedSchemaRegionId(localSchemaRegionInfo));
54+
Assert.assertTrue(exception.getMessage().contains("root.db1"));
55+
Assert.assertTrue(exception.getMessage().contains("root.db2"));
56+
}
57+
}

0 commit comments

Comments
 (0)