This is an automated email from the ASF dual-hosted git repository.
jerryshao pushed a commit to branch branch-1.3
in repository https://gitbox.apache.org/repos/asf/gravitino.git
The following commit(s) were added to refs/heads/branch-1.3 by this push:
new ec157e2f71 [Cherry-pick to branch-1.3] [#13304][#13305] fix(core):
Keep column tags consistent across column drop and rename (#13307) (#13308)
ec157e2f71 is described below
commit ec157e2f71504e48fec0d38344c0791192a073e4
Author: Jerry Shao <[email protected]>
AuthorDate: Fri Sep 18 19:17:04 2026 +0800
[Cherry-pick to branch-1.3] [#13304][#13305] fix(core): Keep column tags
consistent across column drop and rename (#13307) (#13308)
Cherry-pick of #13307 to `branch-1.3`.
**Differences from #13307:**
- The orphan relation GC change and its test are not included, because
`branch-1.3` has no `OrphanedMetadataObjectRelationService`. Tag
relations of columns dropped before this fix stay in the database but
are hidden from tag object lists.
- The new batch soft-delete SQL uses `branch-1.3`'s inline timestamp
expressions instead of `DatabaseTimeSQL`.
- `AlterTableCatalogResult` only carries the column name changes
(`branch-1.3` has no `tableEntityBeforeRename`).
`./gradlew :core:test -PskipITs` passes on `branch-1.3` locally (H2).
---
### What changes were proposed in this pull request?
Keep column tags consistent when a column is dropped or renamed.
**Dropped columns (#13304)**
- `MetadataObjectService.getColumnObjectsFullName` returns `null` for a
column whose latest row is `DELETE`. A tag's object list therefore no
longer shows dropped columns. This also covers relations left over from
before this fix.
- `TableColumnMetaService.updateColumnPOsFromTableDiff` now soft-deletes
the tag and owner relations of dropped columns. This runs in the same
transaction as the table update, in batches of 1000 ids, using a new
`softDeleteTagMetadataObjectRelsByMetadataObjects` mapper method.
Policies cannot be attached to columns, so there is nothing to clean up
for them.
**Renamed columns (#13305)**
- `TableOperationDispatcher.alterTable` replays the normalized
`TableChange`s (`resolveColumnNameChanges`) to work out each existing
top-level column's new name. It then uses that when matching stored
columns to the catalog's columns, so a renamed column keeps its id and
its tags. The replay handles:
- chained renames,
- swaps through a temporary name,
- a rename plus adding a new column with the old name,
- a drop plus re-adding the same name in one change (the new column gets
a new id).
Nested fields are ignored.
- Stored columns that were renamed or dropped claim their catalog names
before untouched ones. Without this, a stale stored column (e.g. dropped
outside Gravitino) with the same name as a rename target could take that
name, depending on `HashMap` order, and the renamed column would lose
its id.
- `selectColumnIdByTableIdAndName` matches the name against each
column's latest row only. A renamed column's old name, and a dropped
column's name, no longer resolve to it.
- The load path (`updateColumnsIfNecessaryWhenLoad`) now matches columns
again against the entity passed to `store.update`, which is the latest
stored one. Before, it wrote a column list computed before taking the
lock. A concurrent load could therefore undo a rename that an in-flight
`alterTable` had just stored, which dropped the column's tags. This
narrows the load/alter race but does not close it: a load that read the
catalog before a concurrent alter and the store after it still works
from a stale catalog snapshot, and two concurrent `alterTable`s can
interleave the same way. Closing it needs `alterTable` and `loadTable`
to exclude each other (e.g. a WRITE tree lock in `alterTable`), which is
tracked in #13303.
### Why are the changes needed?
- Dropping a tagged column left it in `GET /tags/{tag}/objects`, while
the column itself reported `NoSuchMetadataObjectException`. The two
views disagreed.
- Renaming a tagged column through `alterTable` on an external catalog
(Hive, Iceberg, …) stored the rename as a drop plus a new column with a
new id. The column lost its tags, and the tag still listed the old name.
Fix: #13304, #13305
Part of #13303
### Does this PR introduce _any_ user-facing change?
No API or configuration change. Behaviour changes:
- A column renamed through Gravitino keeps its tags and owner.
- Dropped columns no longer appear in a tag's object list.
- The old name of a renamed column no longer resolves.
Columns renamed outside Gravitino, e.g. directly in Spark or Hive, still
can't be told apart from a drop plus an add. Their relations are now
soft-deleted instead of being left dangling.
### How was this patch tested?
- **Unit tests** (H2 locally):
- `TestTagMetaService.testTagRelationsFollowColumnDropAndRename`
- `TestTableColumnMetaService` (old-name lookup, dropped-column full
name, batched relation cleanup over more than one batch)
- `TestTableOperationDispatcher`:
`testAlterTableKeepsColumnIdsAcrossRenames`,
`testResolveColumnNameChanges`,
`testLoadDoesNotUndoConcurrentColumnRename`,
`testRenameIntoNameOfStaleStoredColumn`
- **Integration tests** (Hive, Docker):
`TagIT.testDroppedColumnIsRemovedFromTagObjects` and
`TagIT.testRenamedColumnKeepsItsTags` replay the reported REST steps.
- `./gradlew :core:test -PskipITs` passes locally. The MySQL/PostgreSQL
backends and `TagIT` still need to run in CI.
🤖 Generated with [Claude Code](https://claude.com/claude-code)
---------
Co-authored-by: Claude Opus 5 <[email protected]>
---
.../gravitino/client/integration/test/TagIT.java | 93 ++++++++
.../catalog/TableOperationDispatcher.java | 152 +++++++++++-
.../mapper/TagMetadataObjectRelMapper.java | 7 +
.../TagMetadataObjectRelSQLProviderFactory.java | 7 +
.../provider/base/TableColumnBaseSQLProvider.java | 27 ++-
.../base/TagMetadataObjectRelBaseSQLProvider.java | 17 ++
.../TagMetadataObjectRelPostgreSQLProvider.java | 17 ++
.../relational/service/MetadataObjectService.java | 7 +
.../relational/service/TableColumnMetaService.java | 31 +++
.../catalog/TestTableOperationDispatcher.java | 263 +++++++++++++++++++++
.../service/TestTableColumnMetaService.java | 114 +++++++++
.../relational/service/TestTagMetaService.java | 135 +++++++++++
12 files changed, 846 insertions(+), 24 deletions(-)
diff --git
a/clients/client-java/src/test/java/org/apache/gravitino/client/integration/test/TagIT.java
b/clients/client-java/src/test/java/org/apache/gravitino/client/integration/test/TagIT.java
index ffa59da404..412bd6329a 100644
---
a/clients/client-java/src/test/java/org/apache/gravitino/client/integration/test/TagIT.java
+++
b/clients/client-java/src/test/java/org/apache/gravitino/client/integration/test/TagIT.java
@@ -31,6 +31,7 @@ import org.apache.gravitino.Schema;
import org.apache.gravitino.client.GravitinoMetalake;
import org.apache.gravitino.dto.tag.MetadataObjectDTO;
import org.apache.gravitino.exceptions.NoSuchTagException;
+import org.apache.gravitino.exceptions.NotFoundException;
import org.apache.gravitino.exceptions.TagAlreadyAssociatedException;
import org.apache.gravitino.exceptions.TagAlreadyExistsException;
import org.apache.gravitino.function.Function;
@@ -50,6 +51,7 @@ import org.apache.gravitino.rel.Column;
import org.apache.gravitino.rel.Dialects;
import org.apache.gravitino.rel.SQLRepresentation;
import org.apache.gravitino.rel.Table;
+import org.apache.gravitino.rel.TableChange;
import org.apache.gravitino.rel.View;
import org.apache.gravitino.rel.types.Types;
import org.apache.gravitino.tag.Tag;
@@ -791,6 +793,97 @@ public class TagIT extends BaseIT {
Assertions.assertEquals(column.name(),
tag4.associatedObjects().objects()[0].name());
}
+ @Test
+ public void testDroppedColumnIsRemovedFromTagObjects() {
+ NameIdentifier tableIdent =
createColumnTestTable("tag_it_drop_column_table");
+ try {
+ Tag tag =
+ metalake.createTag(
+ GravitinoITUtils.genRandomName("tag_it_drop_column_tag"),
+ "comment",
+ Collections.emptyMap());
+ Column c1 = loadColumn(tableIdent, "c1");
+ c1.supportsTags().associateTags(new String[] {tag.name()}, null);
+ Assertions.assertArrayEquals(new String[] {tag.name()},
c1.supportsTags().listTags());
+
+ relationalCatalog
+ .asTableCatalog()
+ .alterTable(tableIdent, TableChange.deleteColumn(new String[]
{"c1"}, true));
+
+ // The dropped column no longer resolves, and the tag no longer lists
it. The server reports
+ // NoSuchMetadataObjectException; the client's tag error handler
surfaces it as
+ // NotFoundException.
+ Assertions.assertThrows(NotFoundException.class, () ->
c1.supportsTags().listTags());
+ Assertions.assertEquals(0,
metalake.getTag(tag.name()).associatedObjects().count());
+
+ // A new column with the same name does not inherit the dropped column's
tag.
+ relationalCatalog
+ .asTableCatalog()
+ .alterTable(
+ tableIdent, TableChange.addColumn(new String[] {"c1"},
Types.IntegerType.get()));
+ Assertions.assertEquals(0, loadColumn(tableIdent,
"c1").supportsTags().listTags().length);
+ Assertions.assertEquals(0,
metalake.getTag(tag.name()).associatedObjects().count());
+ } finally {
+ relationalCatalog.asTableCatalog().dropTable(tableIdent);
+ }
+ }
+
+ @Test
+ public void testRenamedColumnKeepsItsTags() {
+ NameIdentifier tableIdent =
createColumnTestTable("tag_it_rename_column_table");
+ try {
+ Tag tag =
+ metalake.createTag(
+ GravitinoITUtils.genRandomName("tag_it_rename_column_tag"),
+ "comment",
+ Collections.emptyMap());
+ Column c1 = loadColumn(tableIdent, "c1");
+ c1.supportsTags().associateTags(new String[] {tag.name()}, null);
+
+ relationalCatalog
+ .asTableCatalog()
+ .alterTable(tableIdent, TableChange.renameColumn(new String[]
{"c1"}, "c1_new"));
+
+ // The tag follows the column to its new name.
+ Column renamed = loadColumn(tableIdent, "c1_new");
+ Assertions.assertArrayEquals(new String[] {tag.name()},
renamed.supportsTags().listTags());
+
Assertions.assertFalse(renamed.supportsTags().getTag(tag.name()).inherited().get());
+
+ // The tag lists the column under its new name only, and the old name no
longer resolves.
+ MetadataObject[] objects =
metalake.getTag(tag.name()).associatedObjects().objects();
+ Assertions.assertEquals(1, objects.length);
+ Assertions.assertEquals(MetadataObject.Type.COLUMN, objects[0].type());
+ Assertions.assertEquals(
+ String.join(".", relationalCatalog.name(), schema.name(),
tableIdent.name(), "c1_new"),
+ objects[0].fullName());
+ Assertions.assertThrows(NotFoundException.class, () ->
c1.supportsTags().listTags());
+ } finally {
+ relationalCatalog.asTableCatalog().dropTable(tableIdent);
+ }
+ }
+
+ private NameIdentifier createColumnTestTable(String prefix) {
+ NameIdentifier tableIdent =
+ NameIdentifier.of(schema.name(),
GravitinoITUtils.genRandomName(prefix));
+ relationalCatalog
+ .asTableCatalog()
+ .createTable(
+ tableIdent,
+ new Column[] {
+ Column.of("c1", Types.IntegerType.get()), Column.of("c2",
Types.StringType.get())
+ },
+ "comment",
+ Collections.emptyMap());
+ return tableIdent;
+ }
+
+ private Column loadColumn(NameIdentifier tableIdent, String columnName) {
+ return
Arrays.stream(relationalCatalog.asTableCatalog().loadTable(tableIdent).columns())
+ .filter(c -> c.name().equals(columnName))
+ .findFirst()
+ .get();
+ }
+
@Test
public void testAssociateAndDeleteTags() {
Tag tag1 =
diff --git
a/core/src/main/java/org/apache/gravitino/catalog/TableOperationDispatcher.java
b/core/src/main/java/org/apache/gravitino/catalog/TableOperationDispatcher.java
index 3cdcd77b43..52297b0383 100644
---
a/core/src/main/java/org/apache/gravitino/catalog/TableOperationDispatcher.java
+++
b/core/src/main/java/org/apache/gravitino/catalog/TableOperationDispatcher.java
@@ -25,12 +25,17 @@ import static
org.apache.gravitino.rel.expressions.transforms.Transforms.EMPTY_T
import static
org.apache.gravitino.utils.NameIdentifierUtil.getCatalogIdentifier;
import static
org.apache.gravitino.utils.NameIdentifierUtil.getSchemaIdentifier;
+import com.google.common.annotations.VisibleForTesting;
import com.google.common.base.Objects;
import com.google.common.base.Preconditions;
import com.google.common.collect.Lists;
import java.time.Instant;
+import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
+import java.util.Comparator;
+import java.util.HashMap;
+import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
@@ -38,6 +43,7 @@ import java.util.function.Function;
import java.util.function.Supplier;
import java.util.stream.Collectors;
import java.util.stream.IntStream;
+import javax.annotation.Nullable;
import org.apache.commons.lang3.tuple.Pair;
import org.apache.gravitino.EntityAlreadyExistsException;
import org.apache.gravitino.EntityStore;
@@ -270,19 +276,21 @@ public class TableOperationDispatcher extends
OperationDispatcher implements Tab
nameIdentifierForLock.equals(ident) ? LockType.READ : LockType.WRITE,
() -> {
NameIdentifier catalogIdent = getCatalogIdentifier(ident);
- TableCatalogResult catalogResult =
+ AlterTableCatalogResult catalogResult =
doWithCatalog(
catalogIdent,
catalog -> {
validateAlterProperties(
catalog, HasPropertyMetadata::tablePropertiesMetadata,
changes);
boolean managed = isManagedEntity(catalog,
Capability.Scope.TABLE);
+ TableChange[] normalizedChanges =
+ applyCapabilities(catalog.capabilities(), changes);
Table table =
catalog.doWithTableOps(
- tableOps ->
- tableOps.alterTable(
- ident,
applyCapabilities(catalog.capabilities(), changes)));
- return snapshotTable(catalog, table, managed);
+ tableOps -> tableOps.alterTable(ident,
normalizedChanges));
+ return new AlterTableCatalogResult(
+ snapshotTable(catalog, table, managed),
+ resolveColumnNameChanges(normalizedChanges));
},
NoSuchTableException.class,
IllegalArgumentException.class);
@@ -324,7 +332,8 @@ public class TableOperationDispatcher extends
OperationDispatcher implements Tab
// Update the columns
Pair<Boolean, List<ColumnEntity>>
columnsUpdateResult =
- updateColumnsIfNecessary(alteredTable,
tableEntity);
+ updateColumnsIfNecessary(
+ alteredTable, tableEntity,
catalogResult.columnNameChanges);
return TableEntity.builder()
.withId(tableEntity.id())
@@ -752,8 +761,83 @@ public class TableOperationDispatcher extends
OperationDispatcher implements Tab
return String.join(", ", differences);
}
+ /**
+ * Works out what each top-level column that existed before an alter is
called after it.
+ *
+ * <p>Stored columns are matched to the catalog's columns by name, so
without this a renamed
+ * column would look like a dropped column plus a new one, and lose its id
and everything attached
+ * to it. The changes are replayed in order, so chained renames resolve to
the final name. A
+ * column that is not in the returned map keeps its name, and a column
mapped to {@code null} was
+ * dropped, so a new column that reuses its name is not mistaken for it.
+ *
+ * @param changes the changes applied to the table, after capability
normalization
+ * @return the new name of each renamed column, or {@code null} for each
dropped column, keyed by
+ * the column's name before the alter
+ */
+ @VisibleForTesting
+ static Map<String, String> resolveColumnNameChanges(TableChange... changes) {
+ // Original name -> current name, or null once the original column is
dropped.
+ Map<String, String> originalToCurrent = new HashMap<>();
+ // Current names of columns added by these changes.
+ Set<String> addedColumns = new HashSet<>();
+
+ for (TableChange change : changes) {
+ if (change instanceof TableChange.AddColumn) {
+ String[] fieldName = ((TableChange.AddColumn) change).fieldName();
+ if (fieldName.length == 1) {
+ addedColumns.add(fieldName[0]);
+ }
+
+ } else if (change instanceof TableChange.RenameColumn) {
+ TableChange.RenameColumn rename = (TableChange.RenameColumn) change;
+ if (rename.fieldName().length != 1) {
+ continue;
+ }
+ String from = rename.fieldName()[0];
+ String to = rename.getNewName();
+ String original = originalColumnNamed(originalToCurrent, addedColumns,
from);
+ if (original != null) {
+ originalToCurrent.put(original, to);
+ } else if (addedColumns.remove(from)) {
+ addedColumns.add(to);
+ }
+
+ } else if (change instanceof TableChange.DeleteColumn) {
+ String[] fieldName = ((TableChange.DeleteColumn) change).fieldName();
+ if (fieldName.length != 1) {
+ continue;
+ }
+ String original = originalColumnNamed(originalToCurrent, addedColumns,
fieldName[0]);
+ if (original != null) {
+ originalToCurrent.put(original, null);
+ } else {
+ addedColumns.remove(fieldName[0]);
+ }
+ }
+ }
+
+ originalToCurrent.entrySet().removeIf(e ->
e.getKey().equals(e.getValue()));
+ return originalToCurrent;
+ }
+
+ /** Returns the original name of the pre-existing column currently called
{@code name}. */
+ @Nullable
+ private static String originalColumnNamed(
+ Map<String, String> originalToCurrent, Set<String> addedColumns, String
name) {
+ for (Map.Entry<String, String> entry : originalToCurrent.entrySet()) {
+ if (name.equals(entry.getValue())) {
+ return entry.getKey();
+ }
+ }
+ // A name no change has touched yet still belongs to the pre-existing
column of that name.
+ if (!originalToCurrent.containsKey(name) && !addedColumns.contains(name)) {
+ return name;
+ }
+ return null;
+ }
+
private Pair<Boolean, List<ColumnEntity>> updateColumnsIfNecessary(
- Table tableFromCatalog, TableEntity tableFromGravitino) {
+ Table tableFromCatalog, TableEntity tableFromGravitino, Map<String,
String> nameChanges) {
if (tableFromCatalog == null || tableFromGravitino == null) {
LOG.warn(
"Cannot update columns for table when altering because table or
table entity is "
@@ -775,9 +859,28 @@ public class TableOperationDispatcher extends
OperationDispatcher implements Tab
// Check if columns need to be updated in Gravitino store
List<ColumnEntity> columnsToInsert = Lists.newArrayList();
+ Set<String> matchedCatalogColumns = new HashSet<>();
boolean columnsNeedsUpdate = false;
- for (Map.Entry<String, ColumnEntity> entry :
columnsFromTableEntity.entrySet()) {
- Pair<Integer, Column> columnPair =
columnsFromCatalogTable.get(entry.getKey());
+ // Renamed and dropped columns claim their catalog names first. Otherwise
a stale stored column
+ // (e.g. dropped outside Gravitino) that has the same name as a rename
target could be visited
+ // first and take that name, and the renamed column would lose its id.
+ List<Map.Entry<String, ColumnEntity>> storedColumns =
+ new ArrayList<>(columnsFromTableEntity.entrySet());
+ storedColumns.sort(Comparator.comparing(e ->
!nameChanges.containsKey(e.getKey())));
+ for (Map.Entry<String, ColumnEntity> entry : storedColumns) {
+ // Follow renames so the stored column keeps its id under its new name.
+ String catalogColumnName =
+ nameChanges.containsKey(entry.getKey())
+ ? nameChanges.get(entry.getKey())
+ : entry.getKey();
+ Pair<Integer, Column> columnPair =
+ catalogColumnName == null ||
matchedCatalogColumns.contains(catalogColumnName)
+ ? null
+ : columnsFromCatalogTable.get(catalogColumnName);
+ if (columnPair != null) {
+ matchedCatalogColumns.add(catalogColumnName);
+ }
+
if (columnPair == null) {
LOG.debug(
"Column {} is not found in the table from underlying source, it
will be removed"
@@ -824,7 +927,7 @@ public class TableOperationDispatcher extends
OperationDispatcher implements Tab
// Check if there are new columns in the table from the underlying source
for (Map.Entry<String, Pair<Integer, Column>> entry :
columnsFromCatalogTable.entrySet()) {
- if (!columnsFromTableEntity.containsKey(entry.getKey())) {
+ if (!matchedCatalogColumns.contains(entry.getKey())) {
LOG.debug(
"Column {} of table: {} is found in the table from underlying
source but not in the table "
+ "entity, it will be added to the table entity",
@@ -852,14 +955,23 @@ public class TableOperationDispatcher extends
OperationDispatcher implements Tab
NameIdentifier tableIdent, EntityCombinedTable combinedTable) {
Pair<Boolean, List<ColumnEntity>> columnsUpdateResult =
updateColumnsIfNecessary(
- combinedTable.tableFromCatalog(),
combinedTable.tableFromGravitino());
+ combinedTable.tableFromCatalog(),
+ combinedTable.tableFromGravitino(),
+ Collections.emptyMap());
// No need to update the columns
if (!columnsUpdateResult.getLeft()) {
return combinedTable.tableFromGravitino();
}
- // Update the columns in the Gravitino store
+ // Update the columns in the Gravitino store. The diff above only decides
whether a write is
+ // needed: it was computed before taking the lock, and a concurrent alter
may have changed the
+ // stored columns since, e.g. renamed a column while keeping its id. So
the columns are matched
+ // again against the entity being updated, which is the latest stored one.
This narrows the
+ // load/alter race but does not close it: if this load read the catalog
before a concurrent
+ // alter and the store after it, the catalog snapshot is the stale side.
Closing that needs
+ // alterTable and loadTable to exclude each other, which is tracked
separately.
+ Table tableFromCatalog = combinedTable.tableFromCatalog();
return TreeLockUtils.doWithTreeLock(
tableIdent,
LockType.WRITE,
@@ -878,7 +990,10 @@ public class TableOperationDispatcher extends
OperationDispatcher implements Tab
.withNamespace(entity.namespace())
.withComment(entity.comment())
.withProperties(entity.properties())
- .withColumns(columnsUpdateResult.getRight())
+ .withColumns(
+ updateColumnsIfNecessary(
+ tableFromCatalog, entity,
Collections.emptyMap())
+ .getRight())
.withPartitioning(entity.partitioning())
.withDistribution(entity.distribution())
.withSortOrders(entity.sortOrders())
@@ -909,6 +1024,17 @@ public class TableOperationDispatcher extends
OperationDispatcher implements Tab
}
}
+ private static final class AlterTableCatalogResult extends
TableCatalogResult {
+
+ private final Map<String, String> columnNameChanges;
+
+ private AlterTableCatalogResult(
+ TableCatalogResult tableResult, Map<String, String> columnNameChanges)
{
+ super(tableResult.table, tableResult.managed,
tableResult.hiddenProperties);
+ this.columnNameChanges = columnNameChanges;
+ }
+ }
+
private static final class CreateTableCatalogResult extends
TableCatalogResult {
private final long id;
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TagMetadataObjectRelMapper.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TagMetadataObjectRelMapper.java
index 14d6d87f7b..fda0ac5819 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TagMetadataObjectRelMapper.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TagMetadataObjectRelMapper.java
@@ -82,6 +82,13 @@ public interface TagMetadataObjectRelMapper {
@Param("metadataObjectId") Long metadataObjectId,
@Param("metadataObjectType") String metadataObjectType);
+ @UpdateProvider(
+ type = TagMetadataObjectRelSQLProviderFactory.class,
+ method = "softDeleteTagMetadataObjectRelsByMetadataObjects")
+ void softDeleteTagMetadataObjectRelsByMetadataObjects(
+ @Param("metadataObjectIds") List<Long> metadataObjectIds,
+ @Param("metadataObjectType") String metadataObjectType);
+
@UpdateProvider(
type = TagMetadataObjectRelSQLProviderFactory.class,
method = "softDeleteTagMetadataObjectRelsByCatalogId")
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TagMetadataObjectRelSQLProviderFactory.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TagMetadataObjectRelSQLProviderFactory.java
index fd2f49e645..46cfaa2f2d 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TagMetadataObjectRelSQLProviderFactory.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TagMetadataObjectRelSQLProviderFactory.java
@@ -102,6 +102,13 @@ public class TagMetadataObjectRelSQLProviderFactory {
.softDeleteTagMetadataObjectRelsByMetadataObject(metadataObjectId,
metadataObjectType);
}
+ public static String softDeleteTagMetadataObjectRelsByMetadataObjects(
+ @Param("metadataObjectIds") List<Long> metadataObjectIds,
+ @Param("metadataObjectType") String metadataObjectType) {
+ return getProvider()
+ .softDeleteTagMetadataObjectRelsByMetadataObjects(metadataObjectIds,
metadataObjectType);
+ }
+
public static String softDeleteTagMetadataObjectRelsByCatalogId(
@Param("catalogId") Long catalogId) {
return getProvider().softDeleteTagMetadataObjectRelsByCatalogId(catalogId);
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TableColumnBaseSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TableColumnBaseSQLProvider.java
index 6044fe9fd8..a1243a9cd9 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TableColumnBaseSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TableColumnBaseSQLProvider.java
@@ -121,18 +121,23 @@ public class TableColumnBaseSQLProvider {
public String selectColumnIdByTableIdAndName(
@Param("tableId") Long tableId, @Param("columnName") String name) {
- return "SELECT"
- + " CASE"
- + " WHEN column_op_type = 3 THEN NULL"
- + " ELSE column_id"
- + " END"
- + " FROM "
+ // Match the name against each column's latest row only. A column keeps
its id when it is
+ // renamed, so its older rows still carry the old name; those must not
resolve. A dropped column
+ // has a DELETE row as its latest row. Dropping and re-adding a column
with the same name in one
+ // change gives two column ids with rows at the same version, and only the
live one matches.
+ // The latest rows are found with one pass over the table's rows, so the
cost does not grow with
+ // how often a column was rewritten.
+ return "SELECT c.column_id FROM "
+ + TableColumnMapper.COLUMN_TABLE_NAME
+ + " c JOIN ("
+ + " SELECT column_id, MAX(table_version) AS max_version FROM "
+ TableColumnMapper.COLUMN_TABLE_NAME
- + " WHERE table_id = #{tableId} AND column_name = #{columnName} AND
deleted_at = 0"
- // Update a column will generate two records with the same version,
one with op_type = 2
- // (update) and another with op_type = 3 (delete). We should not
return NULL if both records
- // exist with the same version, otherwise the caller will think the
column does not exist.
- + " ORDER BY table_version DESC, column_op_type ASC, id DESC LIMIT 1";
+ + " WHERE table_id = #{tableId} AND deleted_at = 0 GROUP BY column_id)
latest"
+ + " ON c.column_id = latest.column_id AND c.table_version =
latest.max_version"
+ + " WHERE c.table_id = #{tableId} AND c.column_name = #{columnName}"
+ + " AND c.deleted_at = 0 AND c.column_op_type <> "
+ + ColumnPO.ColumnOpType.DELETE.value()
+ + " ORDER BY c.table_version DESC, c.id DESC LIMIT 1";
}
public String selectColumnPOById(@Param("columnId") Long columnId) {
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TagMetadataObjectRelBaseSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TagMetadataObjectRelBaseSQLProvider.java
index 1f3a727066..b617ef9b68 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TagMetadataObjectRelBaseSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TagMetadataObjectRelBaseSQLProvider.java
@@ -166,6 +166,23 @@ public class TagMetadataObjectRelBaseSQLProvider {
+ " AND metadata_object_type = #{metadataObjectType}";
}
+ public String softDeleteTagMetadataObjectRelsByMetadataObjects(
+ @Param("metadataObjectIds") List<Long> metadataObjectIds,
+ @Param("metadataObjectType") String metadataObjectType) {
+ return "<script>"
+ + "UPDATE "
+ + TagMetadataObjectRelMapper.TAG_METADATA_OBJECT_RELATION_TABLE_NAME
+ + " SET deleted_at = (UNIX_TIMESTAMP() * 1000.0)"
+ + " + EXTRACT(MICROSECOND FROM CURRENT_TIMESTAMP(3)) / 1000"
+ + " WHERE deleted_at = 0 AND metadata_object_type =
#{metadataObjectType}"
+ + " AND metadata_object_id IN ("
+ + "<foreach collection='metadataObjectIds' item='metadataObjectId'
separator=','>"
+ + "#{metadataObjectId}"
+ + "</foreach>"
+ + ")"
+ + "</script>";
+ }
+
public String softDeleteTagMetadataObjectRelsByCatalogId(@Param("catalogId")
Long catalogId) {
return " UPDATE "
+ TagMetadataObjectRelMapper.TAG_METADATA_OBJECT_RELATION_TABLE_NAME
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TagMetadataObjectRelPostgreSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TagMetadataObjectRelPostgreSQLProvider.java
index 992e105ee6..48097d7ad8 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TagMetadataObjectRelPostgreSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TagMetadataObjectRelPostgreSQLProvider.java
@@ -73,6 +73,23 @@ public class TagMetadataObjectRelPostgreSQLProvider extends
TagMetadataObjectRel
+ " AND metadata_object_type = #{metadataObjectType}";
}
+ @Override
+ public String softDeleteTagMetadataObjectRelsByMetadataObjects(
+ @Param("metadataObjectIds") List<Long> metadataObjectIds,
+ @Param("metadataObjectType") String metadataObjectType) {
+ return "<script>"
+ + "UPDATE "
+ + TAG_METADATA_OBJECT_RELATION_TABLE_NAME
+ + " SET deleted_at = CAST(EXTRACT(EPOCH FROM CURRENT_TIMESTAMP) * 1000
AS BIGINT)"
+ + " WHERE deleted_at = 0 AND metadata_object_type =
#{metadataObjectType}"
+ + " AND metadata_object_id IN ("
+ + "<foreach collection='metadataObjectIds' item='metadataObjectId'
separator=','>"
+ + "#{metadataObjectId}"
+ + "</foreach>"
+ + ")"
+ + "</script>";
+ }
+
@Override
public String softDeleteTagMetadataObjectRelsByCatalogId(@Param("catalogId")
Long catalogId) {
return " UPDATE "
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/MetadataObjectService.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/MetadataObjectService.java
index 634f63be01..051fe1e720 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/MetadataObjectService.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/MetadataObjectService.java
@@ -472,6 +472,13 @@ public class MetadataObjectService {
columnPOs.forEach(
columnPO -> {
+ // A dropped column keeps a live row whose op type is DELETE, so it
must be reported as
+ // deleted instead of returning its last name.
+ if (columnPO.getColumnOpType() ==
ColumnPO.ColumnOpType.DELETE.value()) {
+ columnIdAndNameMap.put(columnPO.getColumnId(), null);
+ return;
+ }
+
// since the table can be deleted, we need to check the null value,
// and when the table is deleted, we will set fullName of column to
// null
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/TableColumnMetaService.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/TableColumnMetaService.java
index 925c12e87d..62653e23f7 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/TableColumnMetaService.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/TableColumnMetaService.java
@@ -28,12 +28,16 @@ import java.util.Map;
import java.util.function.Function;
import java.util.stream.Collectors;
import org.apache.gravitino.Entity;
+import org.apache.gravitino.MetadataObject;
import org.apache.gravitino.exceptions.NoSuchEntityException;
import org.apache.gravitino.meta.ColumnEntity;
import org.apache.gravitino.meta.TableEntity;
import org.apache.gravitino.metrics.Monitored;
+import org.apache.gravitino.storage.relational.mapper.OwnerMetaMapper;
import org.apache.gravitino.storage.relational.mapper.TableColumnMapper;
+import
org.apache.gravitino.storage.relational.mapper.TagMetadataObjectRelMapper;
import org.apache.gravitino.storage.relational.po.ColumnPO;
+import org.apache.gravitino.storage.relational.po.OwnerRelForDeletion;
import org.apache.gravitino.storage.relational.po.TablePO;
import org.apache.gravitino.storage.relational.utils.POConverters;
import org.apache.gravitino.storage.relational.utils.SessionUtils;
@@ -174,10 +178,12 @@ public class TableColumnMetaService {
}
// Mark the columns to DELETE if they are not existed in new columns.
+ List<Long> deletedColumnIds = Lists.newArrayList();
for (ColumnEntity oldColumn : oldColumns.values()) {
if (!newColumns.containsKey(oldColumn.id())) {
columnPOsToInsert.add(
POConverters.initializeColumnPO(newTablePO, oldColumn,
ColumnPO.ColumnOpType.DELETE));
+ deletedColumnIds.add(oldColumn.id());
}
}
@@ -194,6 +200,31 @@ public class TableColumnMetaService {
}
insertColumnPOsInBatches(columnPOsToInsert);
+ deleteColumnRelations(deletedColumnIds);
+ }
+
+ private void deleteColumnRelations(List<Long> columnIds) {
+ // A dropped column keeps its rows, so nothing else removes the relations
that reference it.
+ // This runs in the table update transaction, so the relations go away
with the column. A wide
+ // table can drop many columns at once, so the ids are deleted in batches.
Tags are the main
+ // relation on columns; the owner API does not reject columns either, so
owner rows are removed
+ // too. Policies cannot be attached to columns, so there are no policy
relations to remove.
+ String columnType = MetadataObject.Type.COLUMN.name();
+ Lists.partition(columnIds, COLUMN_INSERT_BATCH_SIZE)
+ .forEach(
+ batch -> {
+ SessionUtils.doWithoutCommit(
+ TagMetadataObjectRelMapper.class,
+ mapper ->
+
mapper.softDeleteTagMetadataObjectRelsByMetadataObjects(batch, columnType));
+ SessionUtils.doWithoutCommit(
+ OwnerMetaMapper.class,
+ mapper ->
+ mapper.batchSoftDeleteOwnerRelByMetadataObjects(
+ batch.stream()
+ .map(id -> new OwnerRelForDeletion(id,
columnType))
+ .collect(Collectors.toList())));
+ });
}
private void insertColumnPOsInBatches(List<ColumnPO> columnPOs) {
diff --git
a/core/src/test/java/org/apache/gravitino/catalog/TestTableOperationDispatcher.java
b/core/src/test/java/org/apache/gravitino/catalog/TestTableOperationDispatcher.java
index 9977b132a1..c576682425 100644
---
a/core/src/test/java/org/apache/gravitino/catalog/TestTableOperationDispatcher.java
+++
b/core/src/test/java/org/apache/gravitino/catalog/TestTableOperationDispatcher.java
@@ -28,6 +28,7 @@ import static
org.apache.gravitino.TestBasePropertiesMetadata.COMMENT_KEY;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyBoolean;
import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.doAnswer;
import static org.mockito.Mockito.doReturn;
import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.mock;
@@ -42,6 +43,7 @@ import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Optional;
+import java.util.concurrent.atomic.AtomicBoolean;
import java.util.function.Function;
import java.util.function.Supplier;
import java.util.stream.Collectors;
@@ -820,6 +822,267 @@ public class TestTableOperationDispatcher extends
TestOperationDispatcher {
testColumnAndColumnEntities(alteredTable6.columns(),
tableEntity6.columns());
}
+ @Test
+ public void testAlterTableKeepsColumnIdsAcrossRenames() throws IOException {
+ Namespace tableNs = Namespace.of(metalake, catalog,
"schema_column_rename");
+ Map<String, String> props = ImmutableMap.of("k1", "v1", "k2", "v2");
+
schemaOperationDispatcher.createSchema(NameIdentifier.of(tableNs.levels()),
"comment", props);
+ NameIdentifier tableIdent = NameIdentifier.of(tableNs,
"table_column_rename");
+ Column[] columns =
+ new Column[] {
+ TestColumn.builder()
+ .withName("a")
+ .withPosition(0)
+ .withType(Types.IntegerType.get())
+ .build(),
+ TestColumn.builder()
+ .withName("b")
+ .withPosition(1)
+ .withType(Types.IntegerType.get())
+ .build(),
+ TestColumn.builder()
+ .withName("c")
+ .withPosition(2)
+ .withType(Types.IntegerType.get())
+ .build()
+ };
+ tableOperationDispatcher.createTable(tableIdent, columns, "comment",
props, new Transform[0]);
+ Map<String, Long> ids = columnIds(tableIdent);
+
+ // A rename keeps the column id.
+ tableOperationDispatcher.alterTable(
+ tableIdent, TableChange.renameColumn(new String[] {"a"}, "a1"));
+ Map<String, Long> afterRename = columnIds(tableIdent);
+ Assertions.assertEquals(ids.get("a"), afterRename.get("a1"));
+ Assertions.assertFalse(afterRename.containsKey("a"));
+
+ // Swapping two names through a temporary name swaps the ids with them.
+ tableOperationDispatcher.alterTable(
+ tableIdent,
+ TableChange.renameColumn(new String[] {"b"}, "tmp"),
+ TableChange.renameColumn(new String[] {"c"}, "b"),
+ TableChange.renameColumn(new String[] {"tmp"}, "c"));
+ Map<String, Long> afterSwap = columnIds(tableIdent);
+ Assertions.assertEquals(ids.get("b"), afterSwap.get("c"));
+ Assertions.assertEquals(ids.get("c"), afterSwap.get("b"));
+
+ // A new column that takes a renamed column's old name gets a new id.
+ tableOperationDispatcher.alterTable(
+ tableIdent,
+ TableChange.renameColumn(new String[] {"a1"}, "a2"),
+ TableChange.addColumn(new String[] {"a1"}, Types.IntegerType.get()));
+ Map<String, Long> afterRenameAndAdd = columnIds(tableIdent);
+ Assertions.assertEquals(ids.get("a"), afterRenameAndAdd.get("a2"));
+ Assertions.assertFalse(ids.containsValue(afterRenameAndAdd.get("a1")));
+
+ // Dropping a column and adding one with the same name in one change gives
a new id.
+ tableOperationDispatcher.alterTable(
+ tableIdent,
+ TableChange.deleteColumn(new String[] {"b"}, false),
+ TableChange.addColumn(new String[] {"b"}, Types.IntegerType.get()));
+ Map<String, Long> afterDropAndAdd = columnIds(tableIdent);
+ Assertions.assertNotEquals(afterSwap.get("b"), afterDropAndAdd.get("b"));
+ Assertions.assertEquals(ids.get("b"), afterDropAndAdd.get("c"));
+ }
+
+ @Test
+ public void testLoadDoesNotUndoConcurrentColumnRename() throws IOException {
+ Namespace tableNs = Namespace.of(metalake, catalog, "schema_load_race");
+ Map<String, String> props = ImmutableMap.of("k1", "v1", "k2", "v2");
+
schemaOperationDispatcher.createSchema(NameIdentifier.of(tableNs.levels()),
"comment", props);
+ NameIdentifier tableIdent = NameIdentifier.of(tableNs, "table_load_race");
+ Column[] columns =
+ new Column[] {
+ TestColumn.builder()
+ .withName("c1")
+ .withPosition(0)
+ .withType(Types.IntegerType.get())
+ .build(),
+ TestColumn.builder()
+ .withName("c2")
+ .withPosition(1)
+ .withType(Types.IntegerType.get())
+ .build()
+ };
+ tableOperationDispatcher.createTable(tableIdent, columns, "comment",
props, new Transform[0]);
+ long c1Id = columnIds(tableIdent).get("c1");
+
+ // An alter renames c1 in the catalog. A load then sees c1_new in the
catalog but c1 in the
+ // store, and decides to replace the column before the alter has written
the store.
+ TestCatalog testCatalog =
+ (TestCatalog)
+ catalogManager.loadCatalogAndWrap(NameIdentifier.of(metalake,
catalog)).catalog();
+ ((TestCatalogOperations) testCatalog.ops())
+ .alterTable(tableIdent, TableChange.renameColumn(new String[] {"c1"},
"c1_new"));
+
+ // The alter's store write, which keeps the column id, lands just before
the load's write.
+ AtomicBoolean alterWritten = new AtomicBoolean(false);
+ doAnswer(
+ invocation -> {
+ if (alterWritten.compareAndSet(false, true)) {
+ entityStore.update(
+ tableIdent,
+ TableEntity.class,
+ TABLE,
+ (TableEntity old) -> renameStoredColumn(old, "c1",
"c1_new"));
+ }
+ return invocation.callRealMethod();
+ })
+ .when(entityStore)
+ .update(any(), any(), any(), any());
+ try {
+ tableOperationDispatcher.loadTable(tableIdent);
+ } finally {
+ reset(entityStore);
+ }
+
+ Assertions.assertTrue(alterWritten.get());
+ Map<String, Long> ids = columnIds(tableIdent);
+ Assertions.assertEquals(c1Id, ids.get("c1_new"));
+ Assertions.assertFalse(ids.containsKey("c1"));
+ }
+
+ private static TableEntity renameStoredColumn(TableEntity table, String
from, String to) {
+ List<ColumnEntity> columns =
+ table.columns().stream()
+ .map(
+ c ->
+ c.name().equals(from)
+ ? ColumnEntity.builder()
+ .withId(c.id())
+ .withName(to)
+ .withPosition(c.position())
+ .withDataType(c.dataType())
+ .withComment(c.comment())
+ .withNullable(c.nullable())
+ .withAutoIncrement(c.autoIncrement())
+ .withDefaultValue(c.defaultValue())
+ .withAuditInfo((AuditInfo) c.auditInfo())
+ .build()
+ : c)
+ .collect(Collectors.toList());
+ return TableEntity.builder()
+ .withId(table.id())
+ .withName(table.name())
+ .withNamespace(table.namespace())
+ .withComment(table.comment())
+ .withProperties(table.properties())
+ .withColumns(columns)
+ .withPartitioning(table.partitioning())
+ .withDistribution(table.distribution())
+ .withSortOrders(table.sortOrders())
+ .withIndexes(table.indexes())
+ .withAuditInfo(table.auditInfo())
+ .build();
+ }
+
+ @Test
+ public void testRenameIntoNameOfStaleStoredColumn() throws IOException {
+ Namespace tableNs = Namespace.of(metalake, catalog, "schema_stale_column");
+ Map<String, String> props = ImmutableMap.of("k1", "v1", "k2", "v2");
+
schemaOperationDispatcher.createSchema(NameIdentifier.of(tableNs.levels()),
"comment", props);
+ NameIdentifier tableIdent = NameIdentifier.of(tableNs,
"table_stale_column");
+ // The names are chosen so that the stale column "b" comes before "z" in
the stored columns'
+ // HashMap iteration order.
+ Column[] columns =
+ new Column[] {
+ TestColumn.builder()
+ .withName("z")
+ .withPosition(0)
+ .withType(Types.IntegerType.get())
+ .build(),
+ TestColumn.builder()
+ .withName("b")
+ .withPosition(1)
+ .withType(Types.IntegerType.get())
+ .build(),
+ TestColumn.builder()
+ .withName("c")
+ .withPosition(2)
+ .withType(Types.IntegerType.get())
+ .build()
+ };
+ tableOperationDispatcher.createTable(tableIdent, columns, "comment",
props, new Transform[0]);
+ Map<String, Long> ids = columnIds(tableIdent);
+
+ // Drop b outside Gravitino, so the store still has it, then rename z to b
through Gravitino.
+ TestCatalog testCatalog =
+ (TestCatalog)
+ catalogManager.loadCatalogAndWrap(NameIdentifier.of(metalake,
catalog)).catalog();
+ ((TestCatalogOperations) testCatalog.ops())
+ .alterTable(tableIdent, TableChange.deleteColumn(new String[] {"b"},
false));
+ tableOperationDispatcher.alterTable(
+ tableIdent, TableChange.renameColumn(new String[] {"z"}, "b"));
+
+ Map<String, Long> afterRename = columnIds(tableIdent);
+ Assertions.assertEquals(ids.get("z"), afterRename.get("b"));
+ Assertions.assertEquals(ids.get("c"), afterRename.get("c"));
+ Assertions.assertEquals(2, afterRename.size());
+ }
+
+ @Test
+ public void testResolveColumnNameChanges() {
+ Assertions.assertEquals(
+ ImmutableMap.of("a", "c"),
+ TableOperationDispatcher.resolveColumnNameChanges(
+ TableChange.renameColumn(new String[] {"a"}, "b"),
+ TableChange.renameColumn(new String[] {"b"}, "c")));
+
+ // Renaming back to the original name is not a change.
+ Assertions.assertEquals(
+ ImmutableMap.of(),
+ TableOperationDispatcher.resolveColumnNameChanges(
+ TableChange.renameColumn(new String[] {"a"}, "b"),
+ TableChange.renameColumn(new String[] {"b"}, "a")));
+
+ Map<String, String> dropped =
+ TableOperationDispatcher.resolveColumnNameChanges(
+ TableChange.deleteColumn(new String[] {"a"}, false),
+ TableChange.addColumn(new String[] {"a"}, Types.IntegerType.get()),
+ TableChange.renameColumn(new String[] {"a"}, "b"));
+ Assertions.assertEquals(1, dropped.size());
+ Assertions.assertTrue(dropped.containsKey("a"));
+ Assertions.assertNull(dropped.get("a"));
+
+ // A column renamed and then dropped is dropped under its original name.
+ Map<String, String> renamedThenDropped =
+ TableOperationDispatcher.resolveColumnNameChanges(
+ TableChange.renameColumn(new String[] {"a"}, "b"),
+ TableChange.deleteColumn(new String[] {"b"}, false));
+ Assertions.assertEquals(1, renamedThenDropped.size());
+ Assertions.assertTrue(renamedThenDropped.containsKey("a"));
+ Assertions.assertNull(renamedThenDropped.get("a"));
+
+ // A column added and dropped in the same change never existed before it.
+ Assertions.assertEquals(
+ ImmutableMap.of(),
+ TableOperationDispatcher.resolveColumnNameChanges(
+ TableChange.addColumn(new String[] {"x"}, Types.IntegerType.get()),
+ TableChange.deleteColumn(new String[] {"x"}, false)));
+
+ // Dropping a missing column with ifExists only marks that name as
dropped, and a later rename
+ // of another column is still tracked.
+ Map<String, String> missingThenRenamed =
+ TableOperationDispatcher.resolveColumnNameChanges(
+ TableChange.deleteColumn(new String[] {"missing"}, true),
+ TableChange.renameColumn(new String[] {"a"}, "b"));
+ Assertions.assertEquals(2, missingThenRenamed.size());
+ Assertions.assertNull(missingThenRenamed.get("missing"));
+ Assertions.assertEquals("b", missingThenRenamed.get("a"));
+
+ // Nested fields are not columns of their own.
+ Assertions.assertEquals(
+ ImmutableMap.of(),
+ TableOperationDispatcher.resolveColumnNameChanges(
+ TableChange.renameColumn(new String[] {"s", "x"}, "y"),
+ TableChange.deleteColumn(new String[] {"s", "z"}, false)));
+ }
+
+ private Map<String, Long> columnIds(NameIdentifier tableIdent) throws
IOException {
+ return entityStore.get(tableIdent, TABLE,
TableEntity.class).columns().stream()
+ .collect(Collectors.toMap(ColumnEntity::name, ColumnEntity::id));
+ }
+
@Test
public void testCreateAndAlterTableWithColumn() throws IOException {
Namespace tableNs = Namespace.of(metalake, catalog, "schema101");
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestTableColumnMetaService.java
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestTableColumnMetaService.java
index bf406b5a3a..377cbf8801 100644
---
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestTableColumnMetaService.java
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestTableColumnMetaService.java
@@ -24,6 +24,7 @@ import java.sql.Connection;
import java.sql.PreparedStatement;
import java.sql.ResultSet;
import java.sql.SQLException;
+import java.sql.Statement;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
@@ -40,7 +41,9 @@ import org.apache.gravitino.storage.RandomIdGenerator;
import org.apache.gravitino.storage.relational.TestJDBCBackend;
import org.apache.gravitino.storage.relational.mapper.TableColumnMapper;
import org.apache.gravitino.storage.relational.po.ColumnPO;
+import org.apache.gravitino.storage.relational.session.SqlSessionFactoryHelper;
import org.apache.gravitino.storage.relational.session.SqlSessions;
+import org.apache.ibatis.session.SqlSession;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.TestTemplate;
@@ -186,6 +189,101 @@ public class TestTableColumnMetaService extends
TestJDBCBackend {
compareTwoColumns(createdTable.columns(), retrievedTable.columns());
}
+ @TestTemplate
+ public void testDropColumnsRemovesTheirRelationsInBatches() throws Exception
{
+ String catalogName = "catalog1";
+ String schemaName = "schema1";
+ createParentEntities(METALAKE_NAME, catalogName, schemaName, AUDIT_INFO);
+
+ // More dropped columns than one cleanup batch holds.
+ int columnCount = 1502;
+ List<ColumnEntity> columns = new ArrayList<>(columnCount);
+ for (int i = 0; i < columnCount; i++) {
+ columns.add(
+ ColumnEntity.builder()
+ .withId(RandomIdGenerator.INSTANCE.nextId())
+ .withName("column_" + i)
+ .withPosition(i)
+ .withDataType(Types.IntegerType.get())
+ .withNullable(true)
+ .withAutoIncrement(false)
+ .withAuditInfo(AUDIT_INFO)
+ .build());
+ }
+ TableEntity table =
+ TableEntity.builder()
+ .withId(RandomIdGenerator.INSTANCE.nextId())
+ .withName("table_drop_columns")
+ .withNamespace(Namespace.of(METALAKE_NAME, catalogName,
schemaName))
+ .withColumns(columns)
+ .withAuditInfo(AUDIT_INFO)
+ .build();
+ TableMetaService.getInstance().insertTable(table, false);
+
+ long firstDropped = columns.get(0).id();
+ long lastDropped = columns.get(columnCount - 2).id();
+ long kept = columns.get(columnCount - 1).id();
+ for (long columnId : new long[] {firstDropped, lastDropped, kept}) {
+ insertColumnRelations(columnId);
+ }
+
+ TableEntity updated =
+ TableEntity.builder()
+ .withId(table.id())
+ .withName(table.name())
+ .withNamespace(table.namespace())
+ .withColumns(Lists.newArrayList(columns.get(columnCount - 1)))
+ .withAuditInfo(AUDIT_INFO)
+ .build();
+ TableMetaService.getInstance()
+ .updateTable(table.nameIdentifier(), (TableEntity old) -> updated);
+
+ Assertions.assertEquals(0, countActiveColumnRelations(firstDropped));
+ Assertions.assertEquals(0, countActiveColumnRelations(lastDropped));
+ Assertions.assertEquals(2, countActiveColumnRelations(kept));
+ }
+
+ private void insertColumnRelations(long columnId) throws SQLException {
+ try (SqlSession session =
+
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+ Connection connection = session.getConnection();
+ Statement statement = connection.createStatement()) {
+ statement.executeUpdate(
+ "INSERT INTO owner_meta (metalake_id, metadata_object_id,
metadata_object_type,"
+ + " owner_id, owner_type, audit_info, current_version,
last_version, deleted_at,"
+ + " updated_at) VALUES (1, "
+ + columnId
+ + ", 'COLUMN', 1, 'USER', '{}', 0, 0, 0, 0)");
+ statement.executeUpdate(
+ "INSERT INTO tag_relation_meta (tag_id, metadata_object_id,
metadata_object_type,"
+ + " audit_info, current_version, last_version, deleted_at)
VALUES (1, "
+ + columnId
+ + ", 'COLUMN', '{}', 0, 0, 0)");
+ }
+ }
+
+ private int countActiveColumnRelations(long columnId) throws SQLException {
+ int count = 0;
+ try (SqlSession session =
+
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+ Connection connection = session.getConnection();
+ Statement statement = connection.createStatement()) {
+ for (String table : new String[] {"owner_meta", "tag_relation_meta"}) {
+ try (ResultSet resultSet =
+ statement.executeQuery(
+ "SELECT COUNT(*) FROM "
+ + table
+ + " WHERE metadata_object_type = 'COLUMN' AND deleted_at =
0"
+ + " AND metadata_object_id = "
+ + columnId)) {
+ Assertions.assertTrue(resultSet.next());
+ count += resultSet.getInt(1);
+ }
+ }
+ }
+ return count;
+ }
+
@TestTemplate
public void testUpdateTable() throws IOException {
String catalogName = "catalog1";
@@ -536,6 +634,12 @@ public class TestTableColumnMetaService extends
TestJDBCBackend {
TableColumnMetaService.getInstance()
.getColumnIdByTableIdAndName(retrievedTable.id(),
updatedColumn.name());
Assertions.assertEquals(updatedColumn.id(), updatedColumnId);
+ // The column keeps its id on rename, so its old name must no longer
resolve to it.
+ Assertions.assertThrows(
+ NoSuchEntityException.class,
+ () ->
+ TableColumnMetaService.getInstance()
+ .getColumnIdByTableIdAndName(retrievedTable.id(),
column.name()));
ColumnPO updatedColumnPO =
TableColumnMetaService.getInstance().getColumnPOById(updatedColumn.id());
@@ -564,6 +668,16 @@ public class TestTableColumnMetaService extends
TestJDBCBackend {
Assertions.assertThrows(
NoSuchEntityException.class,
() ->
TableColumnMetaService.getInstance().getColumnPOById(updatedColumn.id()));
+ // Neither name of the dropped column resolves, and it has no full name
any more.
+ Assertions.assertThrows(
+ NoSuchEntityException.class,
+ () ->
+ TableColumnMetaService.getInstance()
+ .getColumnIdByTableIdAndName(retrievedTable.id(),
column.name()));
+ Map<Long, String> fullNames =
+
MetadataObjectService.getColumnObjectsFullName(Lists.newArrayList(updatedColumn.id()));
+ Assertions.assertTrue(fullNames.containsKey(updatedColumn.id()));
+ Assertions.assertNull(fullNames.get(updatedColumn.id()));
}
@TestTemplate
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestTagMetaService.java
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestTagMetaService.java
index 0add50e617..35cd9d7d2b 100644
---
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestTagMetaService.java
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestTagMetaService.java
@@ -1112,6 +1112,141 @@ public class TestTagMetaService extends TestJDBCBackend
{
Assertions.assertEquals(21, countAllTagRel(tagEntity1.id()));
}
+ @TestTemplate
+ public void testTagRelationsFollowColumnDropAndRename() throws IOException {
+ BaseMetalake metalake =
+ createBaseMakeLake(RandomIdGenerator.INSTANCE.nextId(), METALAKE_NAME,
AUDIT_INFO);
+ backend.insert(metalake, false);
+ CatalogEntity catalog =
+ createCatalog(
+ RandomIdGenerator.INSTANCE.nextId(),
+ Namespace.of(METALAKE_NAME),
+ "catalog1",
+ AUDIT_INFO);
+ backend.insert(catalog, false);
+ SchemaEntity schema =
+ createSchemaEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ Namespace.of(METALAKE_NAME, catalog.name()),
+ "schema1",
+ AUDIT_INFO);
+ backend.insert(schema, false);
+
+ ColumnEntity dropped = newColumn("c1", 0);
+ ColumnEntity renamed = newColumn("c2", 1);
+ ColumnEntity kept = newColumn("c3", 2);
+ TableEntity table =
+ TableEntity.builder()
+ .withId(RandomIdGenerator.INSTANCE.nextId())
+ .withName("table1")
+ .withNamespace(Namespace.of(METALAKE_NAME, catalog.name(),
schema.name()))
+ .withColumns(Lists.newArrayList(dropped, renamed, kept))
+ .withAuditInfo(AUDIT_INFO)
+ .build();
+ backend.insert(table, false);
+
+ TagMetaService tagMetaService = TagMetaService.getInstance();
+ TagEntity tag =
+ TagEntity.builder()
+ .withId(RandomIdGenerator.INSTANCE.nextId())
+ .withName("tag1")
+ .withNamespace(NamespaceUtil.ofTag(METALAKE_NAME))
+ .withAuditInfo(AUDIT_INFO)
+ .build();
+ tagMetaService.insertTag(tag, false);
+ for (ColumnEntity column : Lists.newArrayList(dropped, renamed, kept)) {
+ tagMetaService.associateTagsWithMetadataObject(
+ columnIdent(table, column.name()),
+ Entity.EntityType.COLUMN,
+ new NameIdentifier[] {tag.nameIdentifier()},
+ new NameIdentifier[0]);
+ }
+
+ // Drop c1 and rename c2 to c2_new in one update. A rename keeps the
column id.
+ ColumnEntity renamedColumn = copyColumn(renamed, "c2_new", renamed.id());
+ TableMetaService.getInstance()
+ .updateTable(
+ table.nameIdentifier(), (TableEntity old) -> withColumns(old,
renamedColumn, kept));
+
+ List<GenericEntity> objects =
+
tagMetaService.listAssociatedMetadataObjectsForTag(tag.nameIdentifier());
+ Assertions.assertEquals(2, objects.size());
+ Assertions.assertTrue(
+ containsGenericEntity(objects, "catalog1.schema1.table1.c2_new",
Entity.EntityType.COLUMN));
+ Assertions.assertTrue(
+ containsGenericEntity(objects, "catalog1.schema1.table1.c3",
Entity.EntityType.COLUMN));
+ // The dropped column's relation row is removed, not just hidden.
+ Assertions.assertEquals(2, countActiveTagRel(tag.id()));
+
+ Assertions.assertEquals(
+ 1,
+ tagMetaService
+ .listTagsForMetadataObject(columnIdent(table, "c2_new"),
Entity.EntityType.COLUMN)
+ .size());
+ // Neither the old name of the renamed column nor the dropped column
resolves any more.
+ Assertions.assertThrows(
+ NoSuchEntityException.class,
+ () ->
+ tagMetaService.listTagsForMetadataObject(
+ columnIdent(table, "c2"), Entity.EntityType.COLUMN));
+ Assertions.assertThrows(
+ NoSuchEntityException.class,
+ () ->
+ tagMetaService.listTagsForMetadataObject(
+ columnIdent(table, "c1"), Entity.EntityType.COLUMN));
+
+ // A new column that reuses the dropped column's name does not inherit its
tag.
+ ColumnEntity readded = copyColumn(dropped, "c1",
RandomIdGenerator.INSTANCE.nextId());
+ TableMetaService.getInstance()
+ .updateTable(
+ table.nameIdentifier(),
+ (TableEntity old) -> withColumns(old, renamedColumn, kept,
readded));
+ Assertions.assertTrue(
+ tagMetaService
+ .listTagsForMetadataObject(columnIdent(table, "c1"),
Entity.EntityType.COLUMN)
+ .isEmpty());
+ Assertions.assertEquals(
+ 2,
tagMetaService.listAssociatedMetadataObjectsForTag(tag.nameIdentifier()).size());
+ }
+
+ private static ColumnEntity newColumn(String name, int position) {
+ return ColumnEntity.builder()
+ .withId(RandomIdGenerator.INSTANCE.nextId())
+ .withName(name)
+ .withPosition(position)
+ .withAutoIncrement(false)
+ .withNullable(true)
+ .withDataType(Types.IntegerType.get())
+ .withAuditInfo(AUDIT_INFO)
+ .build();
+ }
+
+ private static ColumnEntity copyColumn(ColumnEntity column, String name,
long id) {
+ return ColumnEntity.builder()
+ .withId(id)
+ .withName(name)
+ .withPosition(column.position())
+ .withAutoIncrement(column.autoIncrement())
+ .withNullable(column.nullable())
+ .withDataType(column.dataType())
+ .withAuditInfo(AUDIT_INFO)
+ .build();
+ }
+
+ private static TableEntity withColumns(TableEntity table, ColumnEntity...
columns) {
+ return TableEntity.builder()
+ .withId(table.id())
+ .withName(table.name())
+ .withNamespace(table.namespace())
+ .withColumns(Lists.newArrayList(columns))
+ .withAuditInfo(table.auditInfo())
+ .build();
+ }
+
+ private static NameIdentifier columnIdent(TableEntity table, String
columnName) {
+ return
NameIdentifier.of(Namespace.fromString(table.nameIdentifier().toString()),
columnName);
+ }
+
@TestTemplate
public void testGetTagIdByTagNameWhenTagNotFound() throws IOException {
createAndInsertMakeLake(METALAKE_NAME);