Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -99,6 +99,7 @@ public enum TSStatusCode {
COLUMN_ALREADY_EXISTS(552),
TABLE_IS_LOST(553),
TABLE_INCOMPATIBLE(554),
TABLE_IN_PRE_DELETE(555),
ONLY_LOGICAL_VIEW(560),

// Storage Engine
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
import org.apache.iotdb.common.rpc.thrift.TSStatus;
import org.apache.iotdb.commons.conf.CommonDescriptor;
import org.apache.iotdb.commons.exception.MetadataException;
import org.apache.iotdb.commons.exception.table.TableInDeletionException;
import org.apache.iotdb.commons.path.PartialPath;
import org.apache.iotdb.commons.path.PathPatternTree;
import org.apache.iotdb.commons.schema.SchemaConstant;
Expand Down Expand Up @@ -1540,13 +1541,26 @@ public Optional<Pair<TsTable, TableNodeStatus>> getTableAndStatusIfExists(
return clusterSchemaInfo.getTsTableIfExists(database, tableName);
}

public Optional<TsTable> getTableWithUsingStatusIfExists(
final String database, final String tableName) throws MetadataException {
final Optional<Pair<TsTable, TableNodeStatus>> tableAndStatus =
getTableAndStatusIfExists(database, tableName);
if (tableAndStatus.isEmpty()) {
return Optional.empty();
}
if (TableNodeStatus.PRE_DELETE == tableAndStatus.get().getRight()) {
throw new TableInDeletionException(database, tableName);
}
return Optional.of(tableAndStatus.get().getLeft());
}

public synchronized Pair<TSStatus, TsTable> tableColumnCheckForColumnExtension(
final String database,
final String tableName,
final List<TsTableColumnSchema> columnSchemaList,
final boolean isTableView)
throws MetadataException {
final TsTable originalTable = getTableIfExists(database, tableName).orElse(null);
final TsTable originalTable = getTableWithUsingStatusIfExists(database, tableName).orElse(null);

if (Objects.isNull(originalTable)) {
return new Pair<>(
Expand Down Expand Up @@ -1601,7 +1615,7 @@ public synchronized Pair<TSStatus, TsTable> tableColumnCheckForColumnAltering(
final TSDataType dataType,
final boolean isGeneratedByPipe)
throws MetadataException {
final TsTable originalTable = getTableIfExists(database, tableName).orElse(null);
final TsTable originalTable = getTableWithUsingStatusIfExists(database, tableName).orElse(null);

if (Objects.isNull(originalTable)) {
return new Pair<>(
Expand Down Expand Up @@ -1638,7 +1652,7 @@ public synchronized Pair<TSStatus, TsTable> tableColumnCheckForColumnRenaming(
final String newName,
final boolean isTableView)
throws MetadataException {
final TsTable originalTable = getTableIfExists(database, tableName).orElse(null);
final TsTable originalTable = getTableWithUsingStatusIfExists(database, tableName).orElse(null);

if (Objects.isNull(originalTable)) {
return new Pair<>(
Expand Down Expand Up @@ -1691,7 +1705,7 @@ public synchronized Pair<TSStatus, TsTable> tableCheckForRenaming(
final String newName,
final boolean isTableView)
throws MetadataException {
final TsTable originalTable = getTableIfExists(database, tableName).orElse(null);
final TsTable originalTable = getTableWithUsingStatusIfExists(database, tableName).orElse(null);

if (Objects.isNull(originalTable)) {
return new Pair<>(
Expand Down Expand Up @@ -1776,7 +1790,7 @@ public synchronized Pair<TSStatus, TsTable> updateTableProperties(
final Map<String, String> updatedProperties,
final boolean isTableView)
throws MetadataException {
final TsTable originalTable = getTableIfExists(database, tableName).orElse(null);
final TsTable originalTable = getTableWithUsingStatusIfExists(database, tableName).orElse(null);

if (Objects.isNull(originalTable)) {
return new Pair<>(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1352,9 +1352,11 @@ public ShowTableResp showTables(final ShowTablePlan plan) {
})
.collect(Collectors.toList())
: tableModelMTree
.getAllUsingTablesUnderSpecificDatabase(
.getAllTablesUnderSpecificDatabase(
getQualifiedDatabasePartialPath(plan.getDatabase()))
.stream()
.filter(pair -> pair.getRight() != TableNodeStatus.PRE_CREATE)
.map(Pair::getLeft)
.map(
tsTable ->
new TTableInfo(
Expand Down Expand Up @@ -1447,7 +1449,7 @@ public DescTableResp descTable(final DescTablePlan plan) {
}
return new DescTableResp(
StatusUtils.OK,
tableModelMTree.getUsingTableSchema(databasePath, plan.getTableName()),
tableModelMTree.getTableSchemaForDesc(databasePath, plan.getTableName()),
null,
null);
} catch (final MetadataException e) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
import org.apache.iotdb.commons.exception.SemanticException;
import org.apache.iotdb.commons.exception.table.ColumnNotExistsException;
import org.apache.iotdb.commons.exception.table.TableAlreadyExistsException;
import org.apache.iotdb.commons.exception.table.TableInDeletionException;
import org.apache.iotdb.commons.exception.table.TableNotExistsException;
import org.apache.iotdb.commons.path.PartialPath;
import org.apache.iotdb.commons.path.PathPatternTree;
Expand Down Expand Up @@ -796,7 +797,7 @@ public void setTableComment(
final String comment,
final boolean isView)
throws MetadataException {
final TsTable table = getTable(database, tableName);
final TsTable table = getTableWithUsingStatus(database, tableName).getTable();
final Optional<Pair<TSStatus, TsTable>> check =
ClusterSchemaManager.checkTable4View(database.getTailNode(), table, isView);
if (check.isPresent()) {
Expand All @@ -816,9 +817,8 @@ public void setTableColumnComment(
final @Nonnull String columnName,
final @Nullable String comment)
throws MetadataException {
final TsTable table = getTable(database, tableName);

final TsTableColumnSchema columnSchema = table.getColumnSchema(columnName);
final TsTableColumnSchema columnSchema =
getTableWithUsingStatus(database, tableName).getTable().getColumnSchema(columnName);

if (Objects.isNull(columnSchema)) {
throw new ColumnNotExistsException(
Expand Down Expand Up @@ -988,14 +988,13 @@ public boolean preDeleteColumn(
final String columnName,
final boolean isView)
throws MetadataException, SemanticException {
final ConfigTableNode node = getTableNode(database, tableName);
final ConfigTableNode node = getTableWithUsingStatus(database, tableName);
final Optional<Pair<TSStatus, TsTable>> check =
ClusterSchemaManager.checkTable4View(database.getTailNode(), node.getTable(), isView);
if (check.isPresent()) {
throw new SemanticException(
check.get().getLeft().getMessage(), check.get().getLeft().getCode());
}

final TsTableColumnSchema columnSchema = node.getTable().getColumnSchema(columnName);
if (Objects.isNull(columnSchema)) {
throw new ColumnNotExistsException(
Expand Down Expand Up @@ -1024,7 +1023,7 @@ public void commitDeleteColumn(
public void preAlterColumnDataType(
PartialPath database, String tableName, String columnName, TSDataType dataType)
throws MetadataException {
final ConfigTableNode node = getTableNode(database, tableName);
final ConfigTableNode node = getTableWithUsingStatus(database, tableName);
final TsTableColumnSchema columnSchema = node.getTable().getColumnSchema(columnName);

if (Objects.isNull(columnSchema)) {
Expand Down Expand Up @@ -1089,6 +1088,36 @@ public TsTable getUsingTableSchema(final PartialPath database, final String tabl
return newTable;
}

/**
* The schema to report in DESC. Unlike {@link #getUsingTableSchema}, a column whose deletion is
* still pending is kept, so that the user can see which columns exist and are only waiting for
* the pending procedure. A pending data type change is applied on top.
*/
public TsTable getTableSchemaForDesc(final PartialPath database, final String tableName)
throws MetadataException {
final ConfigTableNode node = getTableNode(database, tableName);
if (node.getPreAlteredColumns().isEmpty()) {
return node.getTable();
}
final TsTable table = new TsTable(node.getTable());
node.getPreAlteredColumns()
.forEach(
(columnName, dataType) -> {
final TsTableColumnSchema columnSchema = table.getColumnSchema(columnName);
if (columnSchema == null) {
return;
}
columnSchema.setDataType(dataType);
if (columnSchema instanceof FieldColumnSchema) {
final FieldColumnSchema fieldColumnSchema = (FieldColumnSchema) columnSchema;
fieldColumnSchema.setEncoding(
SchemaUtils.getDataTypeCompatibleEncoding(
dataType, fieldColumnSchema.getEncoding()));
}
});
return table;
}

public TableSchemaDetails getTableSchemaDetails(
final PartialPath database, final String tableName) throws MetadataException {
final ConfigTableNode node = getTableNode(database, tableName);
Expand Down Expand Up @@ -1130,6 +1159,19 @@ private ConfigTableNode getTableNode(final PartialPath database, final String ta
return ((ConfigTableNode) databaseNode.getChild(tableName));
}

/**
* A table in the pre-delete status is about to be dropped, pre-create status is a temporary
* status, ignore it.
*/
private ConfigTableNode getTableWithUsingStatus(
final PartialPath database, final String tableName) throws MetadataException {
final ConfigTableNode tableNode = getTableNode(database, tableName);
if (tableNode.getStatus() == TableNodeStatus.PRE_DELETE) {
throw new TableInDeletionException(database.getFullPath(), tableName);
}
return tableNode;
}

// endregion

// region Serialization and Deserialization
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -148,7 +148,7 @@ private void checkTableExistence(final ConfigNodeProcedureEnv env) {
try {
if (!env.getConfigManager()
.getClusterSchemaManager()
.getTableIfExists(database, tableName)
.getTableWithUsingStatusIfExists(database, tableName)
.isPresent()) {
setFailure(
new ProcedureException(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,18 +18,27 @@
*/
package org.apache.iotdb.confignode.manager;

import org.apache.iotdb.commons.exception.table.TableInDeletionException;
import org.apache.iotdb.commons.schema.table.TsTable;
import org.apache.iotdb.commons.schema.table.TsTableInternalRPCUtil;
import org.apache.iotdb.confignode.consensus.request.ConfigPhysicalPlanType;
import org.apache.iotdb.confignode.consensus.request.write.database.DatabaseSchemaPlan;
import org.apache.iotdb.confignode.consensus.request.write.table.CommitCreateTablePlan;
import org.apache.iotdb.confignode.consensus.request.write.table.PreCreateTablePlan;
import org.apache.iotdb.confignode.consensus.request.write.table.PreDeleteTablePlan;
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.tsfile.enums.TSDataType;
import org.apache.tsfile.utils.Pair;
import org.junit.Assert;
import org.junit.Test;
import org.mockito.Mockito;

import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
import java.util.List;
Expand Down Expand Up @@ -89,4 +98,84 @@ public void testGetAllTableInfoForDataNodeActivationWithDeletedDatabase() {
Assert.assertEquals(Collections.singleton("test"), tableInfo.right.keySet());
Assert.assertTrue(tableInfo.right.get("test").isEmpty());
}

@Test
public void testGetTableWithUsingStatusIfExists() throws Exception {
final String database = "root.pre_delete_manager_test";
final String table = "table1";
final ClusterSchemaInfo clusterSchemaInfo = new ClusterSchemaInfo();
final ClusterSchemaManager clusterSchemaManager = managerOf(clusterSchemaInfo);
clusterSchemaInfo.createDatabase(
new DatabaseSchemaPlan(
ConfigPhysicalPlanType.CreateDatabase,
new TDatabaseSchema(database).setIsTableModel(true)));

// A missing table yields an empty result instead of an exception.
Assert.assertFalse(
clusterSchemaManager.getTableWithUsingStatusIfExists(database, table).isPresent());

clusterSchemaInfo.preCreateTable(new PreCreateTablePlan(database, new TsTable(table)));
clusterSchemaInfo.commitCreateTable(new CommitCreateTablePlan(database, table));

// A table in the using status is returned.
Assert.assertTrue(
clusterSchemaManager.getTableWithUsingStatusIfExists(database, table).isPresent());

clusterSchemaInfo.preDeleteTable(new PreDeleteTablePlan(database, table));

// A table in the pre-delete status is rejected with the dedicated status code.
final TableInDeletionException exception =
Assert.assertThrows(
TableInDeletionException.class,
() -> clusterSchemaManager.getTableWithUsingStatusIfExists(database, table));
Assert.assertEquals(TSStatusCode.TABLE_IN_PRE_DELETE.getStatusCode(), exception.getErrorCode());
}

@Test
public void testTableChecksRejectTableInPreDelete() throws Exception {
final String database = "root.pre_delete_manager_test";
final String table = "table1";
final ClusterSchemaInfo clusterSchemaInfo = new ClusterSchemaInfo();
final ClusterSchemaManager clusterSchemaManager = managerOf(clusterSchemaInfo);
clusterSchemaInfo.createDatabase(
new DatabaseSchemaPlan(
ConfigPhysicalPlanType.CreateDatabase,
new TDatabaseSchema(database).setIsTableModel(true)));
clusterSchemaInfo.preCreateTable(new PreCreateTablePlan(database, new TsTable(table)));
clusterSchemaInfo.commitCreateTable(new CommitCreateTablePlan(database, table));
clusterSchemaInfo.preDeleteTable(new PreDeleteTablePlan(database, table));

// Every check that guards a table procedure must reject the table, so that no procedure keeps
// modifying a table that is being deleted.
Assert.assertThrows(
TableInDeletionException.class,
() ->
clusterSchemaManager.tableColumnCheckForColumnExtension(
database, table, new ArrayList<>(), false));
Assert.assertThrows(
TableInDeletionException.class,
() ->
clusterSchemaManager.tableColumnCheckForColumnAltering(
database, table, "field", TSDataType.INT32, false));
Assert.assertThrows(
TableInDeletionException.class,
() ->
clusterSchemaManager.tableColumnCheckForColumnRenaming(
database, table, "field", "field2", false));
Assert.assertThrows(
TableInDeletionException.class,
() -> clusterSchemaManager.tableCheckForRenaming(database, table, "table2", false));
Assert.assertThrows(
TableInDeletionException.class,
() ->
clusterSchemaManager.updateTableProperties(
database, table, new HashMap<>(), new HashMap<>(), false));
}

private static ClusterSchemaManager managerOf(final ClusterSchemaInfo clusterSchemaInfo) {
return new ClusterSchemaManager(
Mockito.mock(IManager.class),
clusterSchemaInfo,
Mockito.mock(ClusterSchemaQuotaStatistics.class));
}
}
Loading
Loading