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));