This is an automated email from the ASF dual-hosted git repository.

yuqi1129 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gravitino.git


The following commit(s) were added to refs/heads/main by this push:
     new 4aa6f2170c [#13502] fix(core): Replace a stale table registration 
instead of inheriting it (#13503)
4aa6f2170c is described below

commit 4aa6f2170cb119030460245c0592746423bfac11
Author: Qi Yu <[email protected]>
AuthorDate: Tue Sep 29 20:04:19 2026 +0800

    [#13502] fix(core): Replace a stale table registration instead of 
inheriting it (#13503)
    
    ### What changes were proposed in this pull request?
    
    - `TableMetaService.insertTable(overwrite = true)`: before the upsert,
    in the same schema-locked transaction, delete a row that holds the
    table's name under a different id, together with its dependents
    (columns, version, tag/owner/securable relations, statistics). It reuses
    the drop path's CAS (`deleteTableWithVersion`) and cleanup
    (`deleteTableDependents`). An overwrite with the same id, e.g. a
    re-import after an out-of-band rename, is unchanged.
    - `TableOperationDispatcher.importTable`: overwrite only when the id
    comes from the table's properties. A generated id identifies nothing, so
    a concurrent import of the same table keeps its row: the plain insert
    conflicts and `loadTable` reloads it, as before.
    
    ### Why are the changes needed?
    
    A table dropped outside Gravitino leaves its registration in the store.
    Creating a table with the same name again goes through
    `store.put(entity, true)` with a new id:
    
    - On MySQL/H2, `ON DUPLICATE KEY UPDATE` matches the `(schema_id,
    table_name, deleted_at)` unique key and keeps the stale `table_id`, so
    the new table inherits the old table's tags, owner, privileges and
    statistics.
    - On PostgreSQL, `ON CONFLICT (table_id)` doesn't cover the name key, so
    the insert fails and the store write is lost.
    
    Without the `importTable` change, the stale-row replacement would let
    the last of two concurrent imports of a table without a stored id delete
    the first import's row.
    
    Schema, topic and view creates follow the same `put(overwrite)` pattern
    and are tracked separately in #13303. Privileges already pushed to
    name-based authorization plugins (e.g. Ranger) for the old table are not
    touched, as with any out-of-band drop.
    
    Fix: #13502
    
    ### Does this PR introduce _any_ user-facing change?
    
    Yes. A table recreated under the name of a stale registration now gets a
    new identity and no longer inherits the old table's tags, owner,
    privileges or statistics. No API or configuration changes.
    
    ### How was this patch tested?
    
    - `TestTableMetaService.testOverwriteWithDifferentIdRetiresStaleTable`:
    fails without the fix on all three backends (H2/MySQL keep the stale id,
    PostgreSQL throws `EntityAlreadyExistsException`) and passes with it. It
    replaces `testNaturalKeyOverwriteUsesPersistedTableId`, which asserted
    the old id-preserving behavior.
    -
    
`TestTableOperationDispatcher.testImportWithoutStoredIdDoesNotOverwriteConcurrentImport`:
    an import with a generated id writes without overwrite.
    - Ran `TestTableMetaService` locally on H2, MySQL and PostgreSQL, and
    `TestTableOperationDispatcher` on H2. The full `:core` suite is left to
    CI.
---
 .../catalog/TableOperationDispatcher.java          |  6 +-
 .../relational/service/TableMetaService.java       | 41 ++++++++++++-
 .../catalog/TestTableOperationDispatcher.java      | 49 +++++++++++++++
 .../relational/service/TestTableMetaService.java   | 70 ++++++++++++++--------
 4 files changed, 138 insertions(+), 28 deletions(-)

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 afcef4b385..77d2a5a8b7 100644
--- 
a/core/src/main/java/org/apache/gravitino/catalog/TableOperationDispatcher.java
+++ 
b/core/src/main/java/org/apache/gravitino/catalog/TableOperationDispatcher.java
@@ -579,7 +579,11 @@ public class TableOperationDispatcher extends 
OperationDispatcher implements Tab
             .withAuditInfo(audit)
             .build();
     try {
-      store.put(tableEntity, true);
+      // Overwrite only with the id stored in the catalog: it identifies the 
table, so a row under
+      // the same name with another id is stale and gets replaced. A generated 
id identifies
+      // nothing, so a row that appeared meanwhile, e.g. from a concurrent 
import on another node,
+      // must win. The plain insert then conflicts, and loadTable reloads that 
row.
+      store.put(tableEntity, stringId != null);
     } catch (EntityAlreadyExistsException e) {
       throw e;
     } catch (Exception e) {
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/service/TableMetaService.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/service/TableMetaService.java
index ecbd87aa4a..62a1d0a4c9 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/service/TableMetaService.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/service/TableMetaService.java
@@ -50,10 +50,13 @@ import 
org.apache.gravitino.storage.relational.utils.POConverters;
 import org.apache.gravitino.storage.relational.utils.SessionUtils;
 import org.apache.gravitino.utils.NameIdentifierUtil;
 import org.apache.gravitino.utils.NamespaceUtil;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
 
 /** The service class for table metadata. It provides the basic database 
operations for table. */
 public class TableMetaService {
 
+  private static final Logger LOG = 
LoggerFactory.getLogger(TableMetaService.class);
   private static final TableMetaService INSTANCE = new TableMetaService();
   private BasePOStorageOps<TablePO, TableMetaMapper> ops;
 
@@ -123,13 +126,18 @@ public class TableMetaService {
               po.getSchemaId(),
               po.getCatalogId(),
               po.getMetalakeId(),
+              () -> {
+                if (overwrite) {
+                  deleteStaleTableWithSameName(tableEntity.nameIdentifier(), 
po);
+                }
+              },
               () ->
                   SessionUtils.doWithoutCommit(
                       TableMetaMapper.class,
                       mapper -> {
                         ops.insertPO(mapper, po, overwrite);
                         if (overwrite) {
-                          // MySQL may preserve the existing table ID during 
an upsert. Read the
+                          // The upsert may update the existing row with the 
same table ID. Read the
                           // stored identity and database-generated version 
while the row is locked.
                           TablePO storedPO =
                               mapper.selectTableMetaBySchemaIdAndName(
@@ -410,9 +418,36 @@ public class TableMetaService {
     builder.withSchemaId(namespacedEntityId.entityId());
   }
 
+  /**
+   * Deletes the table stored under the same name as {@code po} when it has a 
different ID.
+   *
+   * <p>Such a row is a stale registration, for example a table dropped 
outside Gravitino and then
+   * created again. Upserting over it would keep the stale ID on MySQL and H2, 
so the new table
+   * would inherit the old table's tags, owner, privileges and statistics, and 
would fail on the
+   * name's unique key on PostgreSQL. The caller must hold the schema write 
lock and run this in the
+   * same transaction as the insert.
+   */
+  private void deleteStaleTableWithSameName(NameIdentifier identifier, TablePO 
po) {
+    TablePO storedPO =
+        SessionUtils.getWithoutCommit(
+            TableMetaMapper.class,
+            mapper -> 
mapper.selectTableMetaBySchemaIdAndName(po.getSchemaId(), po.getTableName()));
+    if (storedPO == null || storedPO.getTableId().equals(po.getTableId())) {
+      return;
+    }
+
+    LOG.warn(
+        "Replacing stale registration of table {} with ID {} by the table with 
ID {}",
+        identifier,
+        storedPO.getTableId(),
+        po.getTableId());
+    deleteTableWithVersion(identifier, storedPO);
+    deleteTableDependents(storedPO);
+  }
+
   private TablePO tablePOWithPersistedIdentityAndVersions(TablePO incomingPO, 
TablePO persistedPO) {
-    // The upsert derives the version inside the database and may preserve an 
existing table ID, so
-    // its dependent rows must carry the identity and versions the database 
ended up with.
+    // The upsert derives the version inside the database, so its dependent 
rows must carry the
+    // versions the database ended up with.
     return TablePO.builder(incomingPO)
         .withTableId(persistedPO.getTableId())
         .withCurrentVersion(persistedPO.getCurrentVersion())
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 6df3796d16..27f7fb1902 100644
--- 
a/core/src/test/java/org/apache/gravitino/catalog/TestTableOperationDispatcher.java
+++ 
b/core/src/test/java/org/apache/gravitino/catalog/TestTableOperationDispatcher.java
@@ -32,7 +32,9 @@ import static org.mockito.Mockito.doAnswer;
 import static org.mockito.Mockito.doReturn;
 import static org.mockito.Mockito.doThrow;
 import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.reset;
+import static org.mockito.Mockito.verify;
 
 import com.google.common.collect.ImmutableMap;
 import java.io.IOException;
@@ -75,8 +77,11 @@ import org.apache.gravitino.meta.TableEntity;
 import org.apache.gravitino.rel.Column;
 import org.apache.gravitino.rel.Table;
 import org.apache.gravitino.rel.TableChange;
+import org.apache.gravitino.rel.expressions.distributions.Distributions;
 import org.apache.gravitino.rel.expressions.literals.Literals;
+import org.apache.gravitino.rel.expressions.sorts.SortOrder;
 import org.apache.gravitino.rel.expressions.transforms.Transform;
+import org.apache.gravitino.rel.indexes.Indexes;
 import org.apache.gravitino.rel.types.Types;
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.BeforeAll;
@@ -375,6 +380,50 @@ public class TestTableOperationDispatcher extends 
TestOperationDispatcher {
     Assertions.assertEquals("comment", loadedTable.comment());
   }
 
+  @Test
+  public void testImportWithoutStoredIdDoesNotOverwriteConcurrentImport() 
throws IOException {
+    Namespace tableNs = Namespace.of(metalake, catalog, "schema_import_no_id");
+    Map<String, String> props = ImmutableMap.of("k1", "v1", "k2", "v2");
+    
schemaOperationDispatcher.createSchema(NameIdentifier.of(tableNs.levels()), 
"comment", props);
+    NameIdentifier tableIdent = NameIdentifier.of(tableNs, 
"table_import_no_id");
+    Column[] columns =
+        new Column[] {
+          TestColumn.builder()
+              .withName("col1")
+              .withPosition(0)
+              .withType(Types.StringType.get())
+              .build()
+        };
+
+    // Create the table outside Gravitino without a stored Gravitino id, so 
loading imports it
+    // under a freshly generated id.
+    TestCatalog testCatalog =
+        (TestCatalog)
+            catalogManager.loadCatalogAndWrap(NameIdentifier.of(metalake, 
catalog)).catalog();
+    ((TestCatalogOperations) testCatalog.ops())
+        .createTable(
+            tableIdent,
+            columns,
+            "comment",
+            props,
+            new Transform[0],
+            Distributions.NONE,
+            new SortOrder[0],
+            Indexes.EMPTY_INDEXES);
+
+    // A generated id carries no identity, so the import must not overwrite a 
registration another
+    // node may have written for the same table meanwhile. A plain insert 
conflicts instead, and
+    // loadTable reloads the winner's entity.
+    reset(entityStore);
+    try {
+      tableOperationDispatcher.loadTable(tableIdent);
+      verify(entityStore).put(any(TableEntity.class), eq(false));
+      verify(entityStore, never()).put(any(TableEntity.class), eq(true));
+    } finally {
+      reset(entityStore);
+    }
+  }
+
   @Test
   public void testConcurrentImportTableFailsOnMismatchedIdentifier() throws 
IOException {
     Namespace tableNs = Namespace.of(metalake, catalog, 
"schemaConcurrentMismatch");
diff --git 
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestTableMetaService.java
 
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestTableMetaService.java
index 2ff71674de..6eded3f650 100644
--- 
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestTableMetaService.java
+++ 
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestTableMetaService.java
@@ -37,6 +37,7 @@ import java.util.function.Function;
 import java.util.stream.Collectors;
 import org.apache.gravitino.Entity;
 import org.apache.gravitino.EntityAlreadyExistsException;
+import org.apache.gravitino.MetadataObject;
 import org.apache.gravitino.NameIdentifier;
 import org.apache.gravitino.Namespace;
 import org.apache.gravitino.exceptions.NoSuchEntityException;
@@ -46,6 +47,7 @@ import org.apache.gravitino.meta.CatalogEntity;
 import org.apache.gravitino.meta.ColumnEntity;
 import org.apache.gravitino.meta.SchemaEntity;
 import org.apache.gravitino.meta.TableEntity;
+import org.apache.gravitino.meta.TagEntity;
 import org.apache.gravitino.rel.Table;
 import org.apache.gravitino.rel.expressions.NamedReference;
 import org.apache.gravitino.rel.expressions.distributions.Distribution;
@@ -66,6 +68,7 @@ import 
org.apache.gravitino.storage.relational.TestJDBCBackend;
 import org.apache.gravitino.storage.relational.mapper.EntityChangeLogMapper;
 import org.apache.gravitino.storage.relational.mapper.SchemaMetaMapper;
 import org.apache.gravitino.storage.relational.mapper.TableMetaMapper;
+import 
org.apache.gravitino.storage.relational.mapper.TagMetadataObjectRelMapper;
 import org.apache.gravitino.storage.relational.po.SchemaPO;
 import org.apache.gravitino.storage.relational.po.TablePO;
 import org.apache.gravitino.storage.relational.po.cache.EntityChangeRecord;
@@ -74,7 +77,6 @@ import 
org.apache.gravitino.storage.relational.utils.SessionUtils;
 import org.apache.gravitino.utils.NameIdentifierUtil;
 import org.apache.gravitino.utils.NamespaceUtil;
 import org.junit.jupiter.api.Assertions;
-import org.junit.jupiter.api.Assumptions;
 import org.junit.jupiter.api.TestTemplate;
 import org.junit.jupiter.api.function.Executable;
 
@@ -268,43 +270,53 @@ public class TestTableMetaService extends TestJDBCBackend 
{
   }
 
   @TestTemplate
-  public void testNaturalKeyOverwriteUsesPersistedTableId() throws IOException 
{
-    // PostgreSQL's upsert targets table_id and rejects a different ID on the 
natural key before
-    // readback. This regression covers MySQL/H2 ON DUPLICATE KEY, which can 
choose either key.
-    Assumptions.assumeFalse("postgresql".equalsIgnoreCase(backendType));
+  public void testOverwriteWithDifferentIdRetiresStaleTable() throws 
IOException {
+    // A row stored under the same name but with another ID is a stale 
registration, e.g. the table
+    // was dropped outside Gravitino and then recreated. The new table must 
not take over the old
+    // row's ID, or it would inherit the old table's tags, policies, owner and 
privileges.
     createParentEntities(metalakeName, catalogName, schemaName, AUDIT_INFO);
     Namespace tableNamespace = NamespaceUtil.ofTable(metalakeName, 
catalogName, schemaName);
-    TableEntity original =
+    TableEntity stale =
         TableEntity.builder()
             .withId(RandomIdGenerator.INSTANCE.nextId())
-            .withName("table_natural_key_overwrite")
+            .withName("table_stale_registration")
             .withNamespace(tableNamespace)
-            .withColumns(List.of(column("original_column", 
Types.IntegerType.get())))
-            .withComment("original")
+            .withColumns(List.of(column("stale_column", 
Types.IntegerType.get())))
+            .withComment("stale")
             .withAuditInfo(AUDIT_INFO)
             .build();
-    TableMetaService.getInstance().insertTable(original, false);
-    TablePO beforeOverwrite = getTablePO(original.id());
-    TableEntity replacement =
+    TableMetaService.getInstance().insertTable(stale, false);
+    TagEntity tag = createAndInsertTagEntity("tag_on_stale_table", "comment", 
metalakeName);
+    TagMetaService.getInstance()
+        .associateTagsWithMetadataObject(
+            stale.nameIdentifier(),
+            Entity.EntityType.TABLE,
+            new NameIdentifier[] {tag.nameIdentifier()},
+            new NameIdentifier[0]);
+    Assertions.assertEquals(1, countActiveTagRels(stale.id()));
+
+    TableEntity recreated =
         TableEntity.builder()
             .withId(RandomIdGenerator.INSTANCE.nextId())
-            .withName(original.name())
+            .withName(stale.name())
             .withNamespace(tableNamespace)
-            .withColumns(List.of(column("replacement_column", 
Types.StringType.get())))
-            .withComment("replacement")
+            .withColumns(List.of(column("recreated_column", 
Types.StringType.get())))
+            .withComment("recreated")
             .withAuditInfo(AUDIT_INFO)
             .build();
-
-    TableMetaService.getInstance().insertTable(replacement, true);
+    TableMetaService.getInstance().insertTable(recreated, true);
 
     TableEntity stored =
-        
TableMetaService.getInstance().getTableByIdentifier(original.nameIdentifier());
-    TablePO afterOverwrite = getTablePO(original.id());
-    Assertions.assertEquals(original.id(), stored.id());
-    Assertions.assertEquals("replacement", stored.comment());
-    Assertions.assertEquals("replacement_column", 
stored.columns().get(0).name());
-    Assertions.assertEquals(
-        beforeOverwrite.getCurrentVersion() + 1, 
afterOverwrite.getCurrentVersion());
+        
TableMetaService.getInstance().getTableByIdentifier(recreated.nameIdentifier());
+    Assertions.assertEquals(recreated.id(), stored.id());
+    Assertions.assertEquals("recreated", stored.comment());
+    Assertions.assertEquals("recreated_column", 
stored.columns().get(0).name());
+    Assertions.assertTrue(
+        SessionUtils.getWithoutCommit(
+                TableMetaMapper.class, mapper -> 
mapper.listTablePOsByTableIds(List.of(stale.id())))
+            .isEmpty());
+    Assertions.assertEquals(0, countActiveTagRels(stale.id()));
+    Assertions.assertEquals(0, countActiveTagRels(recreated.id()));
   }
 
   @TestTemplate
@@ -900,6 +912,16 @@ public class TestTableMetaService extends TestJDBCBackend {
         });
   }
 
+  private int countActiveTagRels(long metadataObjectId) {
+    return SessionUtils.getWithoutCommit(
+        TagMetadataObjectRelMapper.class,
+        mapper ->
+            mapper
+                .listTagPOsByMetadataObjectIdAndType(
+                    metadataObjectId, MetadataObject.Type.TABLE.name())
+                .size());
+  }
+
   private TablePO getTablePO(long tableId) {
     return SessionUtils.getWithoutCommit(
         TableMetaMapper.class, mapper -> 
mapper.listTablePOsByTableIds(List.of(tableId)).get(0));

Reply via email to