This is an automated email from the ASF dual-hosted git repository.
CRZbulabula pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/master by this push:
new d61ef998215 Add several verifications for filtering out databases of
the table model (#18592)
d61ef998215 is described below
commit d61ef9982157007d2fb40c3bcfde904db0b54374
Author: libo <[email protected]>
AuthorDate: Tue Sep 8 12:04:27 2026 +0800
Add several verifications for filtering out databases of the table model
(#18592)
---
...ataPartitionTableIntegrityCheckProcedureIT.java | 127 ++++++++++++++++++++-
.../DataPartitionTableIntegrityCheckProcedure.java | 16 ++-
2 files changed, 137 insertions(+), 6 deletions(-)
diff --git
a/integration-test/src/test/java/org/apache/iotdb/confignode/it/partition/DataPartitionTableIntegrityCheckProcedureIT.java
b/integration-test/src/test/java/org/apache/iotdb/confignode/it/partition/DataPartitionTableIntegrityCheckProcedureIT.java
index da54c39479e..f2390a6c608 100644
---
a/integration-test/src/test/java/org/apache/iotdb/confignode/it/partition/DataPartitionTableIntegrityCheckProcedureIT.java
+++
b/integration-test/src/test/java/org/apache/iotdb/confignode/it/partition/DataPartitionTableIntegrityCheckProcedureIT.java
@@ -19,11 +19,19 @@
package org.apache.iotdb.confignode.it.partition;
+import org.apache.iotdb.common.rpc.thrift.TConsensusGroupId;
+import org.apache.iotdb.common.rpc.thrift.TSeriesPartitionSlot;
+import org.apache.iotdb.common.rpc.thrift.TTimePartitionSlot;
+import org.apache.iotdb.commons.client.sync.SyncConfigNodeIServiceClient;
import org.apache.iotdb.commons.enums.RepairDataPartitionTableProgressState;
+import org.apache.iotdb.confignode.rpc.thrift.TDataPartitionReq;
+import org.apache.iotdb.confignode.rpc.thrift.TDataPartitionTableResp;
import org.apache.iotdb.it.env.EnvFactory;
import org.apache.iotdb.it.framework.IoTDBTestRunner;
import org.apache.iotdb.itbase.category.ClusterIT;
import org.apache.iotdb.itbase.category.LocalStandaloneIT;
+import org.apache.iotdb.itbase.env.BaseEnv;
+import org.apache.iotdb.rpc.TSStatusCode;
import org.junit.After;
import org.junit.Assert;
@@ -41,6 +49,8 @@ import java.sql.Statement;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
+import java.util.Map;
+import java.util.TreeMap;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
@@ -55,6 +65,15 @@ public class DataPartitionTableIntegrityCheckProcedureIT {
private static final Logger LOGGER =
LoggerFactory.getLogger(DataPartitionTableIntegrityCheckProcedureIT.class);
+ private static final String TABLE_DATABASE = "repair_table_db";
+ private static final String TABLE_NAME = "table1";
+ private static final long TIME_PARTITION_INTERVAL = 604_800_000L;
+ private static final long TABLE_TTL = 7 * TIME_PARTITION_INTERVAL;
+ private static final long CURRENT_TIME_PARTITION_START =
+ System.currentTimeMillis() / TIME_PARTITION_INTERVAL *
TIME_PARTITION_INTERVAL;
+ private static final long EXPIRED_TIME_PARTITION_START =
+ CURRENT_TIME_PARTITION_START - TABLE_TTL * 2;
+
@Before
public void setUp() {
EnvFactory.getEnv()
@@ -63,7 +82,16 @@ public class DataPartitionTableIntegrityCheckProcedureIT {
.setConfigNodeConsensusProtocolClass(RATIS_CONSENSUS)
.setSchemaRegionConsensusProtocolClass(RATIS_CONSENSUS)
.setDataRegionConsensusProtocolClass(RATIS_CONSENSUS)
- .setDataReplicationFactor(1);
+ .setDataReplicationFactor(1)
+ .setTimePartitionInterval(TIME_PARTITION_INTERVAL);
+ EnvFactory.getEnv()
+ .getConfig()
+ .getConfigNodeCommonConfig()
+ .setTTLCheckInterval(TimeUnit.MILLISECONDS.toMillis(500));
+ EnvFactory.getEnv()
+ .getConfig()
+ .getDataNodeCommonConfig()
+ .setTTLCheckInterval(TimeUnit.MINUTES.toMillis(10));
EnvFactory.getEnv().initClusterEnvironment(1, 1);
}
@@ -142,6 +170,103 @@ public class DataPartitionTableIntegrityCheckProcedureIT {
}
}
+ @Test
+ public void testRepairDataPartitionTableIgnoresTableModelDatabase() throws
Exception {
+ final TDataPartitionReq dataPartitionReq = new TDataPartitionReq();
+ dataPartitionReq.putToPartitionSlotsMap(TABLE_DATABASE, new TreeMap<>());
+
+ try (final Connection tableConnection =
+ EnvFactory.getEnv().getConnection(BaseEnv.TABLE_SQL_DIALECT);
+ final Statement tableStatement = tableConnection.createStatement()) {
+ tableStatement.execute(String.format("CREATE DATABASE %s",
TABLE_DATABASE));
+ tableStatement.execute(String.format("USE %s", TABLE_DATABASE));
+ tableStatement.execute(
+ String.format("CREATE TABLE %s (device_id STRING TAG, value INT64
FIELD)", TABLE_NAME));
+ tableStatement.execute(
+ String.format(
+ "INSERT INTO %s(time, device_id, value) VALUES (%d, 'd1', 1)",
+ TABLE_NAME, EXPIRED_TIME_PARTITION_START));
+ tableStatement.execute(
+ String.format(
+ "INSERT INTO %s(time, device_id, value) VALUES (%d, 'd1', 2)",
+ TABLE_NAME, CURRENT_TIME_PARTITION_START));
+ tableStatement.execute("FLUSH");
+
+ Assert.assertTrue(
+ "The expired time partition must exist before TTL cleanup",
+ containsTimePartition(dataPartitionReq,
EXPIRED_TIME_PARTITION_START));
+
+ tableStatement.execute(
+ String.format("ALTER TABLE %s SET PROPERTIES TTL=%d", TABLE_NAME,
TABLE_TTL));
+ waitUntilTimePartitionRemoved(dataPartitionReq,
EXPIRED_TIME_PARTITION_START);
+ Assert.assertTrue(
+ "The current table-model time partition must survive TTL cleanup",
+ containsTimePartition(dataPartitionReq,
CURRENT_TIME_PARTITION_START));
+
+ // Stop periodic cleanup from hiding an incorrectly restored partition
after repair.
+ tableStatement.execute(String.format("ALTER TABLE %s SET PROPERTIES
TTL='INF'", TABLE_NAME));
+ }
+
+ try (final Connection treeConnection =
+ EnvFactory.getEnv().getConnection(BaseEnv.TREE_SQL_DIALECT);
+ final Statement treeStatement = treeConnection.createStatement()) {
+ treeStatement.execute("REPAIR DATA PARTITION TABLE");
+ waitForRepairCompletion(treeStatement);
+ }
+
+ Assert.assertFalse(
+ "Repair must not restore a table-model time partition removed by TTL
cleanup",
+ containsTimePartition(dataPartitionReq, EXPIRED_TIME_PARTITION_START));
+ Assert.assertTrue(
+ "Repair must preserve the current table-model time partition",
+ containsTimePartition(dataPartitionReq, CURRENT_TIME_PARTITION_START));
+ }
+
+ private static boolean containsTimePartition(
+ final TDataPartitionReq request, final long timePartitionStart) throws
Exception {
+ try (final SyncConfigNodeIServiceClient client =
+ (SyncConfigNodeIServiceClient)
EnvFactory.getEnv().getLeaderConfigNodeConnection()) {
+ final TDataPartitionTableResp response =
client.getDataPartitionTable(request);
+ Assert.assertEquals(
+ TSStatusCode.SUCCESS_STATUS.getStatusCode(),
response.getStatus().getCode());
+ final Map<TSeriesPartitionSlot, Map<TTimePartitionSlot,
List<TConsensusGroupId>>>
+ seriesPartitionTable =
response.getDataPartitionTable().get(TABLE_DATABASE);
+ return seriesPartitionTable != null
+ && seriesPartitionTable.values().stream()
+ .anyMatch(
+ timePartitionTable ->
+ timePartitionTable.containsKey(new
TTimePartitionSlot(timePartitionStart)));
+ }
+ }
+
+ private static void waitUntilTimePartitionRemoved(
+ final TDataPartitionReq request, final long timePartitionStart) throws
Exception {
+ for (int retry = 0; retry < 120; retry++) {
+ if (!containsTimePartition(request, timePartitionStart)) {
+ return;
+ }
+ TimeUnit.SECONDS.sleep(1);
+ }
+ Assert.fail("The expired table-model time partition was not cleaned within
the timeout");
+ }
+
+ private static void waitForRepairCompletion(final Statement statement)
throws Exception {
+ TimeUnit.SECONDS.sleep(1);
+ for (int retry = 0; retry < 120; retry++) {
+ try (final ResultSet resultSet =
+ statement.executeQuery("SHOW REPAIR DATA PARTITION TABLE PROGRESS"))
{
+ Assert.assertTrue(resultSet.next());
+ if (RepairDataPartitionTableProgressState.IDLE
+ .name()
+ .equals(resultSet.getString("Status"))) {
+ return;
+ }
+ }
+ TimeUnit.SECONDS.sleep(1);
+ }
+ Assert.fail("The data partition table repair did not complete within the
timeout");
+ }
+
private static void assertRepairProgress(
final Statement statement,
final String expectedStatus,
diff --git
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/partition/DataPartitionTableIntegrityCheckProcedure.java
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/partition/DataPartitionTableIntegrityCheckProcedure.java
index f4cfb24e64c..0fc0f5d5b9d 100644
---
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/partition/DataPartitionTableIntegrityCheckProcedure.java
+++
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/partition/DataPartitionTableIntegrityCheckProcedure.java
@@ -30,6 +30,7 @@ import
org.apache.iotdb.commons.enums.RepairDataPartitionTableProgressState;
import org.apache.iotdb.commons.partition.DataPartitionTable;
import org.apache.iotdb.commons.partition.DatabaseScopedDataPartitionTable;
import org.apache.iotdb.commons.partition.SeriesPartitionTable;
+import org.apache.iotdb.commons.utils.PathUtils;
import org.apache.iotdb.commons.utils.TimePartitionUtils;
import org.apache.iotdb.confignode.client.sync.CnToDnSyncRequestType;
import org.apache.iotdb.confignode.client.sync.SyncDataNodeClientPool;
@@ -303,7 +304,9 @@ public class DataPartitionTableIntegrityCheckProcedure
// Merge with existing timeslots (take minimum)
for (Map.Entry<String, Long> entry : nodeTimeslots.entrySet()) {
- earliestTimeslots.merge(entry.getKey(), entry.getValue(), Math::min);
+ if (!PathUtils.isTableModelDatabase(entry.getKey())) {
+ earliestTimeslots.merge(entry.getKey(), entry.getValue(),
Math::min);
+ }
}
if (LOG.isDebugEnabled()) {
@@ -362,16 +365,19 @@ public class DataPartitionTableIntegrityCheckProcedure
String database = entry.getKey();
long earliestTimeslot = entry.getValue();
+ if (PathUtils.isTableModelDatabase(database)) {
+ continue;
+ }
+
// Get current DataPartitionTable from ConfigManager
Map<String, Map<TSeriesPartitionSlot, Map<TTimePartitionSlot,
List<TConsensusGroupId>>>>
localDataPartitionTable = getLocalDataPartitionTable(env, database);
// Check if ConfigNode has a data partition that is associated with the
earliestTimeslot
if ((localDataPartitionTable == null
- || localDataPartitionTable.isEmpty()
- || localDataPartitionTable.get(database) == null
- || localDataPartitionTable.get(database).isEmpty())
- && database.startsWith("root.")) {
+ || localDataPartitionTable.isEmpty()
+ || localDataPartitionTable.get(database) == null
+ || localDataPartitionTable.get(database).isEmpty())) {
databasesWithLostDataPartition.add(database);
LOG.warn(
ProcedureMessages