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 33f8699682 [#12576] improvement(core): optimize schema write fencing
(#12577)
33f8699682 is described below
commit 33f8699682ed398462edac6a96329ad1b87612c0
Author: Qi Yu <[email protected]>
AuthorDate: Sun Sep 20 20:22:55 2026 +0800
[#12576] improvement(core): optimize schema write fencing (#12577)
### What changes were proposed in this pull request?
- Add one schema-scoped transaction entry point that acquires the parent
schema lock and runs descendant writes in the same transaction.
- Migrate table, view, fileset, function, model, model-version, topic,
and table-statistic upsert paths to that entry point.
- Replace six complete direct-child metadata loads in the non-cascade
schema emptiness check with one lightweight existence query that stops
after the first active child.
- Preserve missing-schema errors when the fence rejects view or fileset
updates.
- Reconcile the latest fileset OCC and overwrite flow with the schema
fence so its CAS and version writes remain atomic under the parent lock.
- Cover all six direct child types, plus the independently writable
model-version and table-statistic descendants, through real service
paths.
This is a follow-up to #12456, which has now merged.
### Why are the changes needed?
The previous implementation held the schema delete fence correctly, but
schema-scoped services and independently writable descendants had to
assemble or remember the parent-lock transaction pattern themselves.
That made future omissions easy and allowed a table-statistic upsert to
race a schema cascade. Non-cascade deletion also materialized complete
child objects and version information when it only needed a yes/no
answer, extending the time spent under the schema delete lock.
Fix: #12576
### Does this PR introduce _any_ user-facing change?
No.
### How was this patch tested?
- `./gradlew :core:spotlessApply`
- `./gradlew :core:check -PskipITs -PskipDockerTests=true` on embedded
H2.
- Targeted table-statistic schema-fence concurrency coverage on H2,
MySQL, and PostgreSQL with `-PskipDockerTests=false`.
- `git diff --check`
---
.../relational/mapper/SchemaMetaMapper.java | 10 +
.../mapper/SchemaMetaSQLProviderFactory.java | 5 +
.../provider/base/SchemaMetaBaseSQLProvider.java | 31 +++
.../relational/service/FilesetMetaService.java | 181 ++++++++---------
.../relational/service/FunctionMetaService.java | 63 +++---
.../relational/service/ModelMetaService.java | 96 ++++-----
.../service/ModelVersionMetaService.java | 35 ++--
.../relational/service/SchemaMetaService.java | 79 +++----
.../relational/service/StatisticMetaService.java | 44 +++-
.../relational/service/TableMetaService.java | 226 +++++++++------------
.../relational/service/TopicMetaService.java | 68 ++++---
.../relational/service/ViewMetaService.java | 60 +++---
.../base/TestSchemaMetaBaseSQLProvider.java | 61 ++++++
.../service/TestModelVersionMetaService.java | 54 +++++
.../relational/service/TestSchemaMetaService.java | 198 +++++++++++++-----
.../service/TestStatisticMetaService.java | 175 ++++++++++++++++
16 files changed, 904 insertions(+), 482 deletions(-)
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/SchemaMetaMapper.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/SchemaMetaMapper.java
index bb870e44dc..1f6c844c19 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/SchemaMetaMapper.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/SchemaMetaMapper.java
@@ -86,6 +86,16 @@ public interface SchemaMetaMapper {
@SelectProvider(type = SchemaMetaSQLProviderFactory.class, method =
"selectSchemaMetaById")
SchemaPO selectSchemaMetaById(@Param("schemaId") Long schemaId);
+ /**
+ * Returns one when an active table, view, fileset, function, model, or
topic exists in the
+ * schema, and {@code null} otherwise.
+ *
+ * <p>Only a literal is selected because callers need an existence answer,
not complete child
+ * metadata. The final limit also lets the database stop as soon as it finds
the first child.
+ */
+ @SelectProvider(type = SchemaMetaSQLProviderFactory.class, method =
"selectActiveChildBySchemaId")
+ Integer selectActiveChildBySchemaId(@Param("schemaId") Long schemaId);
+
/** Selects and locks an active schema by ID for the current transaction. */
@SelectProvider(
type = SchemaMetaSQLProviderFactory.class,
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/SchemaMetaSQLProviderFactory.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/SchemaMetaSQLProviderFactory.java
index 62c532db54..bcebb42eb9 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/SchemaMetaSQLProviderFactory.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/SchemaMetaSQLProviderFactory.java
@@ -106,6 +106,11 @@ public class SchemaMetaSQLProviderFactory {
return getProvider().selectSchemaMetaById(schemaId);
}
+ /** Returns SQL that checks whether an active child exists in the schema. */
+ public static String selectActiveChildBySchemaId(@Param("schemaId") Long
schemaId) {
+ return getProvider().selectActiveChildBySchemaId(schemaId);
+ }
+
/** Returns SQL that selects and locks an active schema by ID. */
public static String selectSchemaMetaByIdForUpdate(@Param("schemaId") Long
schemaId) {
return getProvider().selectSchemaMetaByIdForUpdate(schemaId);
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/SchemaMetaBaseSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/SchemaMetaBaseSQLProvider.java
index e8d9738764..00a72eb5b0 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/SchemaMetaBaseSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/SchemaMetaBaseSQLProvider.java
@@ -22,7 +22,13 @@ import static
org.apache.gravitino.storage.relational.mapper.SchemaMetaMapper.TA
import java.util.List;
import org.apache.gravitino.storage.relational.mapper.CatalogMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.FilesetMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.FunctionMetaMapper;
import org.apache.gravitino.storage.relational.mapper.MetalakeMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.ModelMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.TableMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.TopicMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.ViewMetaMapper;
import org.apache.gravitino.storage.relational.mapper.provider.DatabaseTimeSQL;
import org.apache.gravitino.storage.relational.po.SchemaPO;
import org.apache.ibatis.annotations.Param;
@@ -183,6 +189,31 @@ public class SchemaMetaBaseSQLProvider {
+ " WHERE schema_id = #{schemaId} AND deleted_at = 0";
}
+ /** Returns SQL that checks whether an active child exists in the schema. */
+ public String selectActiveChildBySchemaId(@Param("schemaId") Long schemaId) {
+ // Each branch returns only the same literal, so UNION ALL avoids
unnecessary duplicate
+ // elimination. LIMIT 1 lets the database stop as soon as any kind of
child is found.
+ return "SELECT 1 FROM "
+ + TableMetaMapper.TABLE_NAME
+ + " WHERE schema_id = #{schemaId} AND deleted_at = 0"
+ + " UNION ALL SELECT 1 FROM "
+ + ViewMetaMapper.TABLE_NAME
+ + " WHERE schema_id = #{schemaId} AND deleted_at = 0"
+ + " UNION ALL SELECT 1 FROM "
+ + FilesetMetaMapper.META_TABLE_NAME
+ + " WHERE schema_id = #{schemaId} AND deleted_at = 0"
+ + " UNION ALL SELECT 1 FROM "
+ + FunctionMetaMapper.TABLE_NAME
+ + " WHERE schema_id = #{schemaId} AND deleted_at = 0"
+ + " UNION ALL SELECT 1 FROM "
+ + ModelMetaMapper.TABLE_NAME
+ + " WHERE schema_id = #{schemaId} AND deleted_at = 0"
+ + " UNION ALL SELECT 1 FROM "
+ + TopicMetaMapper.TABLE_NAME
+ + " WHERE schema_id = #{schemaId} AND deleted_at = 0"
+ + " LIMIT 1";
+ }
+
/** Returns SQL that selects and locks an active schema by ID. */
public String selectSchemaMetaByIdForUpdate(@Param("schemaId") Long
schemaId) {
return selectSchemaMetaById(schemaId) + " FOR UPDATE";
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/FilesetMetaService.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/FilesetMetaService.java
index 22e37f1ce5..67069b4244 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/FilesetMetaService.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/FilesetMetaService.java
@@ -26,7 +26,6 @@ import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
-import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicReference;
import java.util.function.Function;
import java.util.stream.Collectors;
@@ -171,56 +170,51 @@ public class FilesetMetaService {
// The schema lock, metadata row, and every storage-location version row
share one
// transaction. A failure in any later step restores all earlier writes.
- SessionUtils.doMultipleWithCommit(
- // Hold the parent schema row until this transaction ends, so the
fileset cannot be
- // written below a schema that is being dropped.
- () ->
- SchemaMetaService.getInstance()
- .lockSchemaForEntityWrite(
- filesetEntity.nameIdentifier(),
- po.getSchemaId(),
- po.getCatalogId(),
- po.getMetalakeId()),
- () ->
- SessionUtils.doWithoutCommit(
- FilesetMetaMapper.class,
- mapper -> {
- FilesetPO storedPO =
- overwrite
- ?
mapper.selectFilesetMetaBySchemaIdAndNameForUpdate(
- po.getSchemaId(), po.getFilesetName())
- : null;
- if (storedPO == null) {
- mapper.insertFilesetMeta(po);
- return;
- }
-
- // Resolve the natural-key overwrite before building its
snapshot. This keeps
- // the stored ID in both the metadata row and identifier
property without a
- // post-insert JSON/PO rewrite.
- FilesetEntity replacement =
- filesetWithPersistedId(filesetEntity,
storedPO.getFilesetId());
- Long maxStoredVersion =
- SessionUtils.getWithoutCommit(
- FilesetVersionMapper.class,
- versionMapper ->
-
versionMapper.selectMaxFilesetVersion(storedPO.getFilesetId()));
- FilesetPO replacementPO =
- POConverters.updateFilesetPOWithVersion(
- storedPO, replacement, maxStoredVersion);
- Integer updated = mapper.updateFilesetMeta(replacementPO,
storedPO);
- Preconditions.checkState(
- updated != null && updated == 1,
- "The overwritten fileset %s in schema %s changed while
its row was held",
- po.getFilesetName(),
- po.getSchemaId());
- persistedPO.set(replacementPO);
- }),
- () ->
- SessionUtils.doWithoutCommit(
- FilesetVersionMapper.class,
- mapper ->
-
mapper.insertFilesetVersions(persistedPO.get().getFilesetVersionPOs())));
+ SchemaMetaService.getInstance()
+ .doWithSchemaWriteLock(
+ filesetEntity.nameIdentifier(),
+ po.getSchemaId(),
+ po.getCatalogId(),
+ po.getMetalakeId(),
+ () ->
+ SessionUtils.doWithoutCommit(
+ FilesetMetaMapper.class,
+ mapper -> {
+ FilesetPO storedPO =
+ overwrite
+ ?
mapper.selectFilesetMetaBySchemaIdAndNameForUpdate(
+ po.getSchemaId(), po.getFilesetName())
+ : null;
+ if (storedPO == null) {
+ mapper.insertFilesetMeta(po);
+ return;
+ }
+
+ // Resolve the natural-key overwrite before building
its snapshot so the
+ // metadata row and identifier property both keep the
stored ID.
+ FilesetEntity replacement =
+ filesetWithPersistedId(filesetEntity,
storedPO.getFilesetId());
+ Long maxStoredVersion =
+ SessionUtils.getWithoutCommit(
+ FilesetVersionMapper.class,
+ versionMapper ->
+
versionMapper.selectMaxFilesetVersion(storedPO.getFilesetId()));
+ FilesetPO replacementPO =
+ POConverters.updateFilesetPOWithVersion(
+ storedPO, replacement, maxStoredVersion);
+ Integer updated =
mapper.updateFilesetMeta(replacementPO, storedPO);
+ Preconditions.checkState(
+ updated != null && updated == 1,
+ "The overwritten fileset %s in schema %s changed
while its row was held",
+ po.getFilesetName(),
+ po.getSchemaId());
+ persistedPO.set(replacementPO);
+ }),
+ () ->
+ SessionUtils.doWithoutCommit(
+ FilesetVersionMapper.class,
+ mapper ->
+
mapper.insertFilesetVersions(persistedPO.get().getFilesetVersionPOs())));
} catch (RuntimeException re) {
ExceptionUtils.checkSQLException(
re, Entity.EntityType.FILESET,
filesetEntity.nameIdentifier().toString());
@@ -244,27 +238,14 @@ public class FilesetMetaService {
oldFilesetEntity.id());
try {
- FilesetPO newFilesetPO =
- POConverters.updateFilesetPOWithVersion(oldFilesetPO, newEntity,
null);
- if (tryUpdateFileset(newFilesetPO, oldFilesetPO)) {
- return newEntity;
- }
-
- // The metadata CAS also rejects a version that already has an active
stored snapshot. Only
- // that uncommon legacy case needs the MAX(version) round trip; normal
alters finish above.
- Long maxStoredVersion =
- SessionUtils.getWithoutCommit(
- FilesetVersionMapper.class,
- mapper ->
mapper.selectMaxFilesetVersion(oldFilesetPO.getFilesetId()));
- if (maxStoredVersion != null
- && maxStoredVersion >= newFilesetPO.getCurrentVersion()
- && tryUpdateFileset(
- POConverters.updateFilesetPOWithVersion(oldFilesetPO, newEntity,
maxStoredVersion),
- oldFilesetPO)) {
- return newEntity;
- }
-
- throw filesetWriteFailure(identifier, oldFilesetPO);
+ SchemaMetaService.getInstance()
+ .doWithSchemaWriteLock(
+ identifier,
+ oldFilesetPO.getSchemaId(),
+ oldFilesetPO.getCatalogId(),
+ oldFilesetPO.getMetalakeId(),
+ () -> updateFilesetInTransaction(identifier, newEntity,
oldFilesetPO));
+ return newEntity;
} catch (RuntimeException re) {
ExceptionUtils.checkSQLException(
re, Entity.EntityType.FILESET,
newEntity.nameIdentifier().toString());
@@ -476,26 +457,44 @@ public class FilesetMetaService {
() -> filesetWriteFailure(identifier, observedFilesetPO));
}
+ private void updateFilesetInTransaction(
+ NameIdentifier identifier, FilesetEntity newEntity, FilesetPO
oldFilesetPO) {
+ FilesetPO newFilesetPO =
POConverters.updateFilesetPOWithVersion(oldFilesetPO, newEntity, null);
+ if (tryUpdateFileset(newFilesetPO, oldFilesetPO)) {
+ return;
+ }
+
+ // The metadata CAS also rejects a version that already has an active
stored snapshot. Only
+ // that uncommon legacy case needs the MAX(version) round trip; normal
alters finish above.
+ Long maxStoredVersion =
+ SessionUtils.getWithoutCommit(
+ FilesetVersionMapper.class,
+ mapper ->
mapper.selectMaxFilesetVersion(oldFilesetPO.getFilesetId()));
+ if (maxStoredVersion != null
+ && maxStoredVersion >= newFilesetPO.getCurrentVersion()
+ && tryUpdateFileset(
+ POConverters.updateFilesetPOWithVersion(oldFilesetPO, newEntity,
maxStoredVersion),
+ oldFilesetPO)) {
+ return;
+ }
+
+ throw filesetWriteFailure(identifier, oldFilesetPO);
+ }
+
private boolean tryUpdateFileset(FilesetPO newFilesetPO, FilesetPO
oldFilesetPO) {
- AtomicBoolean updated = new AtomicBoolean(false);
- SessionUtils.doMultipleWithCommit(
- () -> {
- Integer updateCount =
- SessionUtils.getWithoutCommit(
- FilesetMetaMapper.class,
- mapper -> mapper.updateFilesetMeta(newFilesetPO,
oldFilesetPO));
- updated.set(updateCount != null && updateCount > 0);
- },
- () -> {
- if (updated.get()) {
- // The metadata row now points to this complete snapshot. It stays
in the same
- // transaction so a failed version insert also restores the
metadata version.
- SessionUtils.doWithoutCommit(
- FilesetVersionMapper.class,
- mapper ->
mapper.insertFilesetVersions(newFilesetPO.getFilesetVersionPOs()));
- }
- });
- return updated.get();
+ Integer updateCount =
+ SessionUtils.getWithoutCommit(
+ FilesetMetaMapper.class,
+ mapper -> mapper.updateFilesetMeta(newFilesetPO, oldFilesetPO));
+ boolean updated = updateCount != null && updateCount > 0;
+ if (updated) {
+ // The metadata row now points to this complete snapshot. The caller's
schema transaction
+ // ensures a failed version insert also restores the metadata version.
+ SessionUtils.doWithoutCommit(
+ FilesetVersionMapper.class,
+ mapper ->
mapper.insertFilesetVersions(newFilesetPO.getFilesetVersionPOs()));
+ }
+ return updated;
}
private FilesetEntity filesetWithPersistedId(FilesetEntity filesetEntity,
Long persistedId) {
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/FunctionMetaService.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/FunctionMetaService.java
index 31fe233741..9ffcaafe66 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/FunctionMetaService.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/FunctionMetaService.java
@@ -116,17 +116,13 @@ public class FunctionMetaService {
fillFunctionPOBuilderParentEntityId(builder, functionEntity.namespace());
FunctionPO po = initializeFunctionPO(functionEntity, builder);
- SessionUtils.doMultipleWithCommit(
- // Hold the parent schema row until this transaction ends, so the
function cannot be
- // written below a schema that is being dropped.
- () ->
- SchemaMetaService.getInstance()
- .lockSchemaForEntityWrite(
- functionEntity.nameIdentifier(),
- po.schemaId(),
- po.catalogId(),
- po.metalakeId()),
- () -> insertFunctionWithoutCommit(functionEntity, po, overwrite));
+ SchemaMetaService.getInstance()
+ .doWithSchemaWriteLock(
+ functionEntity.nameIdentifier(),
+ po.schemaId(),
+ po.catalogId(),
+ po.metalakeId(),
+ () -> insertFunctionWithoutCommit(functionEntity, po,
overwrite));
} catch (RuntimeException re) {
try {
ExceptionUtils.checkSQLException(
@@ -272,29 +268,28 @@ public class FunctionMetaService {
try {
FunctionPO newFunctionPO =
updateFunctionPO(oldFunctionPO, newEntity, newSchemaId,
newCatalogId, newMetalakeId);
- SessionUtils.doMultipleWithCommit(
- () -> {
- if (isSchemaChanged) {
- SchemaMetaService.getInstance()
- .lockSchemaForEntityWrite(
- newEntity.nameIdentifier(), newSchemaId, newCatalogId,
newMetalakeId);
- }
- },
- () -> {
- // function_current_version is the sole OCC token. The root CAS is
the transaction's
- // decision point and must run before the unguarded version-row
insert below.
- int updated =
- SessionUtils.getWithoutCommit(
- FunctionMetaMapper.class,
- mapper -> ops.updatePO(mapper, newFunctionPO,
oldFunctionPO));
- if (updated == 0) {
- throw functionWriteFailure(identifier, oldFunctionPO);
- }
- },
- () ->
- SessionUtils.doWithoutCommit(
- FunctionVersionMetaMapper.class,
- mapper ->
mapper.insertFunctionVersionMeta(newFunctionPO.functionVersionPO())));
+ SchemaMetaService.getInstance()
+ .doWithSchemaWriteLock(
+ newEntity.nameIdentifier(),
+ newSchemaId,
+ newCatalogId,
+ newMetalakeId,
+ () -> {
+ // function_current_version is the sole OCC token. The root
CAS is the transaction's
+ // decision point and must run before the unguarded
version-row insert below.
+ int updated =
+ SessionUtils.getWithoutCommit(
+ FunctionMetaMapper.class,
+ mapper -> ops.updatePO(mapper, newFunctionPO,
oldFunctionPO));
+ if (updated == 0) {
+ throw functionWriteFailure(identifier, oldFunctionPO);
+ }
+ },
+ () ->
+ SessionUtils.doWithoutCommit(
+ FunctionVersionMetaMapper.class,
+ mapper ->
+
mapper.insertFunctionVersionMeta(newFunctionPO.functionVersionPO())));
return newEntity;
} catch (RuntimeException re) {
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/ModelMetaService.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/ModelMetaService.java
index a482cfb3ce..eb8b3d8f70 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/ModelMetaService.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/ModelMetaService.java
@@ -97,26 +97,22 @@ public class ModelMetaService {
fillModelPOBuilderParentEntityId(builder, modelEntity.namespace());
ModelPO po = POConverters.initializeModelPO(modelEntity, builder);
- SessionUtils.doMultipleWithCommit(
- // Hold the parent schema row until this transaction ends, so the
model cannot be
- // written below a schema that is being dropped.
- () ->
- SchemaMetaService.getInstance()
- .lockSchemaForEntityWrite(
- modelEntity.nameIdentifier(),
- po.getSchemaId(),
- po.getCatalogId(),
- po.getMetalakeId()),
- () ->
- SessionUtils.doWithoutCommit(
- ModelMetaMapper.class,
- mapper -> {
- if (overwrite) {
- mapper.insertModelMetaOnDuplicateKeyUpdate(po);
- } else {
- mapper.insertModelMeta(po);
- }
- }));
+ SchemaMetaService.getInstance()
+ .doWithSchemaWriteLock(
+ modelEntity.nameIdentifier(),
+ po.getSchemaId(),
+ po.getCatalogId(),
+ po.getMetalakeId(),
+ () ->
+ SessionUtils.doWithoutCommit(
+ ModelMetaMapper.class,
+ mapper -> {
+ if (overwrite) {
+ mapper.insertModelMetaOnDuplicateKeyUpdate(po);
+ } else {
+ mapper.insertModelMeta(po);
+ }
+ }));
} catch (RuntimeException re) {
ExceptionUtils.checkSQLException(
re, Entity.EntityType.MODEL,
modelEntity.nameIdentifier().toString());
@@ -347,33 +343,39 @@ public class ModelMetaService {
boolean isRenamed = !Objects.equals(oldModelEntity.name(),
newEntity.name());
try {
- SessionUtils.doMultipleWithCommit(
- () -> {
- // This is the first write in the transaction. It succeeds only if
the model still has
- // the concurrency version read above, so an older request cannot
overwrite a newer
- // model or add an incorrect change-log entry.
- int updated =
- SessionUtils.getWithoutCommit(
- ModelMetaMapper.class,
- mapper ->
- mapper.updateModelMeta(
- POConverters.updateModelPO(oldModelPO, newEntity),
oldModelPO));
- if (updated == 0) {
- throw modelWriteFailure(identifier, oldModelPO);
- }
- },
- () -> {
- if (isRenamed) {
- SessionUtils.doWithoutCommit(
- EntityChangeLogMapper.class,
- mapper ->
- mapper.insertEntityChange(
- metalakeName,
- Entity.EntityType.MODEL.name(),
- oldFullName,
- OperateType.ALTER));
- }
- });
+ SchemaMetaService.getInstance()
+ .doWithSchemaWriteLock(
+ identifier,
+ oldModelPO.getSchemaId(),
+ oldModelPO.getCatalogId(),
+ oldModelPO.getMetalakeId(),
+ () -> {
+ // The model CAS is the first child write. It succeeds only if
the concurrency
+ // version
+ // still matches, so an older request cannot overwrite a newer
model or add an
+ // incorrect change-log entry.
+ int updated =
+ SessionUtils.getWithoutCommit(
+ ModelMetaMapper.class,
+ mapper ->
+ mapper.updateModelMeta(
+ POConverters.updateModelPO(oldModelPO,
newEntity), oldModelPO));
+ if (updated == 0) {
+ throw modelWriteFailure(identifier, oldModelPO);
+ }
+ },
+ () -> {
+ if (isRenamed) {
+ SessionUtils.doWithoutCommit(
+ EntityChangeLogMapper.class,
+ mapper ->
+ mapper.insertEntityChange(
+ metalakeName,
+ Entity.EntityType.MODEL.name(),
+ oldFullName,
+ OperateType.ALTER));
+ }
+ });
} catch (RuntimeException re) {
ExceptionUtils.checkSQLException(
re, Entity.EntityType.MODEL, newEntity.nameIdentifier().toString());
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/ModelVersionMetaService.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/ModelVersionMetaService.java
index 4a448b09af..fcc58faf3b 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/ModelVersionMetaService.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/ModelVersionMetaService.java
@@ -175,7 +175,9 @@ public class ModelVersionMetaService {
POConverters.initializeModelVersionAliasRelPO(modelVersionEntity,
modelId);
try {
- SessionUtils.doMultipleWithCommit(
+ doWithSchemaWriteLock(
+ modelIdent,
+ modelPO,
() -> reserveModelVersionRegistration(modelIdent, modelPO),
() ->
SessionUtils.doWithoutCommit(
@@ -229,7 +231,9 @@ public class ModelVersionMetaService {
int modelVersion = observedVersionPOs.get(0).getModelVersion();
try {
- SessionUtils.doMultipleWithCommit(
+ doWithSchemaWriteLock(
+ modelIdent,
+ modelPO,
() -> reserveModelVersionWrite(ident, modelPO, modelVersion),
() -> {
// An alias was resolved to its numeric version above. Delete
every URI row belonging to
@@ -359,7 +363,9 @@ public class ModelVersionMetaService {
isModelVersionUriUpdated(oldModelVersionEntity, newModelVersionEntity);
try {
- SessionUtils.doMultipleWithCommit(
+ doWithSchemaWriteLock(
+ modelIdent,
+ modelPO,
() -> reserveModelVersionUpdate(ident, modelPO, oldModelVersionPOs,
oldAliasRelPOs),
() -> {
int updated;
@@ -439,21 +445,21 @@ public class ModelVersionMetaService {
return !oldUris.equals(newUris);
}
- private void lockSchemaForModelVersionWrite(
- NameIdentifier modelIdentifier, ModelPO observedModelPO) {
+ private void doWithSchemaWriteLock(
+ NameIdentifier modelIdentifier, ModelPO modelPO, Runnable...
modelVersionWriteOperations) {
SchemaMetaService.getInstance()
- .lockSchemaForEntityWrite(
+ .doWithSchemaWriteLock(
modelIdentifier,
- observedModelPO.getSchemaId(),
- observedModelPO.getCatalogId(),
- observedModelPO.getMetalakeId());
+ modelPO.getSchemaId(),
+ modelPO.getCatalogId(),
+ modelPO.getMetalakeId(),
+ modelVersionWriteOperations);
}
private void reserveModelVersionRegistration(
NameIdentifier modelIdentifier, ModelPO observedModelPO) {
- // The schema fence prevents a registration below a schema being dropped.
The model update then
- // allocates the version number and advances the aggregate OCC token in
one row write.
- lockSchemaForModelVersionWrite(modelIdentifier, observedModelPO);
+ // The caller holds the schema fence. Allocate the version number and
advance the aggregate
+ // OCC token in one model-row write.
ModelMetaService.getInstance()
.bumpModelVersionAndLatestVersion(modelIdentifier, observedModelPO);
}
@@ -461,9 +467,8 @@ public class ModelVersionMetaService {
private void reserveModelVersionWrite(
NameIdentifier modelVersionIdentifier, ModelPO observedModelPO, int
observedModelVersion) {
NameIdentifier modelIdentifier =
NameIdentifier.of(modelVersionIdentifier.namespace().levels());
- // Every model-version mutation follows this order: fence the schema, lock
and advance the
- // parent model row, then validate an alias-based identifier while that
row lock is held.
- lockSchemaForModelVersionWrite(modelIdentifier, observedModelPO);
+ // The caller holds the schema fence. Lock and advance the parent model
row, then validate an
+ // alias-based identifier while that row lock is held.
ModelMetaService.getInstance().bumpModelVersion(modelIdentifier,
observedModelPO);
validateModelVersionIdentifier(
modelVersionIdentifier, observedModelPO.getModelId(),
observedModelVersion);
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/SchemaMetaService.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/SchemaMetaService.java
index 796c0b4368..151da5beac 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/SchemaMetaService.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/SchemaMetaService.java
@@ -507,12 +507,30 @@ public class SchemaMetaService {
}
/**
- * Holds the parent schema row while a table, view, fileset, function,
model, model version, or
- * topic is written, so a child cannot be added below a schema that is going
away. The lock is
- * shared, so children of the same schema can still be written in parallel;
dropping the schema
- * takes the row exclusively and therefore waits for them.
+ * Runs schema-scoped writes while holding a shared lock on their parent
schema.
+ *
+ * <p>This method owns the transaction boundary on purpose. If callers
locked the schema in one
+ * transaction and wrote the child in another, the lock would be released
too early and a schema
+ * deletion could slip between those two steps. Keeping the lock and every
supplied operation in
+ * the same transaction makes that mistake impossible for callers of this
entry point.
*/
- void lockSchemaForEntityWrite(
+ void doWithSchemaWriteLock(
+ NameIdentifier entityIdentifier,
+ Long observedSchemaId,
+ Long observedCatalogId,
+ Long observedMetalakeId,
+ Runnable... entityWriteOperations) {
+ Runnable[] transactionOperations = new
Runnable[entityWriteOperations.length + 1];
+ transactionOperations[0] =
+ () ->
+ lockSchemaForEntityWrite(
+ entityIdentifier, observedSchemaId, observedCatalogId,
observedMetalakeId);
+ System.arraycopy(
+ entityWriteOperations, 0, transactionOperations, 1,
entityWriteOperations.length);
+ SessionUtils.doMultipleWithCommit(transactionOperations);
+ }
+
+ private void lockSchemaForEntityWrite(
NameIdentifier entityIdentifier,
Long observedSchemaId,
Long observedCatalogId,
@@ -570,49 +588,18 @@ public class SchemaMetaService {
mapper -> mapper.softDeleteSchemaMetasWithVersion(children)));
}
- /**
- * Checks that nothing is left under the schema. Views and functions are
included: they used to be
- * missing here, which let a non-cascade drop leave their rows behind with
no parent.
- */
+ /** Checks that no active schema or metadata object is left below the
schema. */
private void checkSchemaIsEmpty(NameIdentifier identifier, SchemaPO
schemaPO) {
boolean hasDescendantSchemas =
!listDescendantSchemaPOs(schemaPO).isEmpty();
- boolean hasTables =
- !SessionUtils.getWithoutCommit(
- TableMetaMapper.class,
- mapper ->
mapper.listTablePOsBySchemaId(schemaPO.getSchemaId()))
- .isEmpty();
- boolean hasFilesets =
- !SessionUtils.getWithoutCommit(
- FilesetMetaMapper.class,
- mapper ->
mapper.listFilesetPOsBySchemaId(schemaPO.getSchemaId()))
- .isEmpty();
- boolean hasModels =
- !SessionUtils.getWithoutCommit(
- ModelMetaMapper.class,
- mapper ->
mapper.listModelPOsBySchemaId(schemaPO.getSchemaId()))
- .isEmpty();
- boolean hasTopics =
- !SessionUtils.getWithoutCommit(
- TopicMetaMapper.class,
- mapper ->
mapper.listTopicPOsBySchemaId(schemaPO.getSchemaId()))
- .isEmpty();
- boolean hasViews =
- !SessionUtils.getWithoutCommit(
- ViewMetaMapper.class,
- mapper -> mapper.listViewPOsBySchemaId(schemaPO.getSchemaId()))
- .isEmpty();
- boolean hasFunctions =
- !SessionUtils.getWithoutCommit(
- FunctionMetaMapper.class,
- mapper ->
mapper.listFunctionPOsBySchemaId(schemaPO.getSchemaId()))
- .isEmpty();
- if (hasDescendantSchemas
- || hasTables
- || hasFilesets
- || hasModels
- || hasTopics
- || hasViews
- || hasFunctions) {
+ // A non-cascade delete only needs to know whether any direct child
exists. Asking the database
+ // for one literal avoids building every child PO and loading its version
details while the
+ // schema delete lock is held.
+ boolean hasDirectChild =
+ SessionUtils.getWithoutCommit(
+ SchemaMetaMapper.class,
+ mapper ->
mapper.selectActiveChildBySchemaId(schemaPO.getSchemaId()))
+ != null;
+ if (hasDescendantSchemas || hasDirectChild) {
throw new NonEmptyEntityException(
"Entity %s has sub-entities, you should remove sub-entities first",
identifier);
}
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/StatisticMetaService.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/StatisticMetaService.java
index 31fd741e39..2a429efb8a 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/StatisticMetaService.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/StatisticMetaService.java
@@ -79,9 +79,17 @@ public class StatisticMetaService {
namespacedEntityId.namespaceIds()[0],
namespacedEntityId.entityId(),
NameIdentifierUtil.toMetadataObject(entity, type).type());
- SessionUtils.doWithCommit(
- StatisticMetaMapper.class,
- mapper -> mapper.batchInsertStatisticPOsOnDuplicateKeyUpdate(pos));
+ // Statistics have their own write API, so they do not inherit the schema
fence from the
+ // metadata object update path. Fence schema-scoped targets explicitly to
keep a schema cascade
+ // from deleting the target and then missing this independently committed
statistic upsert.
+ doWithSchemaWriteLockIfNeeded(
+ entity,
+ type,
+ namespacedEntityId,
+ () ->
+ SessionUtils.doWithoutCommit(
+ StatisticMetaMapper.class,
+ mapper ->
mapper.batchInsertStatisticPOsOnDuplicateKeyUpdate(pos)));
}
@Monitored(
@@ -107,4 +115,34 @@ public class StatisticMetaService {
StatisticMetaMapper.class,
mapper -> mapper.deleteStatisticsByLegacyTimeline(legacyTimeline,
limit));
}
+
+ private void doWithSchemaWriteLockIfNeeded(
+ NameIdentifier identifier,
+ Entity.EntityType type,
+ NamespacedEntityId namespacedEntityId,
+ Runnable writeOperation) {
+ long[] namespaceIds = namespacedEntityId.namespaceIds();
+ Long schemaId;
+ switch (type) {
+ case SCHEMA:
+ schemaId = namespacedEntityId.entityId();
+ break;
+ case TABLE:
+ case VIEW:
+ case COLUMN:
+ case FILESET:
+ case TOPIC:
+ case MODEL:
+ case FUNCTION:
+ schemaId = namespaceIds[2];
+ break;
+ default:
+ SessionUtils.doMultipleWithCommit(writeOperation);
+ return;
+ }
+
+ SchemaMetaService.getInstance()
+ .doWithSchemaWriteLock(
+ identifier, schemaId, namespaceIds[1], namespaceIds[0],
writeOperation);
+ }
}
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 7ace373f1d..3342715558 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
@@ -118,80 +118,72 @@ public class TableMetaService {
AtomicReference<TablePO> persistedPO = new AtomicReference<>(po);
// The schema lock, table row, version row, and columns share one
transaction. If any later
// step fails, the earlier inserts are rolled back as well.
- SessionUtils.doMultipleWithCommit(
- // Hold the parent schema row until this transaction ends, so the
table cannot be
- // written below a schema that is being dropped.
- () ->
- SchemaMetaService.getInstance()
- .lockSchemaForEntityWrite(
- tableEntity.nameIdentifier(),
- po.getSchemaId(),
- po.getCatalogId(),
- po.getMetalakeId()),
- () ->
- SessionUtils.doWithoutCommit(
- TableMetaMapper.class,
- mapper -> {
- ops.insertPO(mapper, po, overwrite);
- if (overwrite) {
- // MySQL may resolve the upsert through the active
(schema_id, table_name,
- // deleted_at) key rather than table_id. In that case it
preserves the
- // winner's ID. The upsert already holds that row until
commit, so read the
- // database-derived identity and version back through
the same natural key.
- TablePO storedPO =
- mapper.selectTableMetaBySchemaIdAndName(
- po.getSchemaId(), po.getTableName());
- Preconditions.checkState(
- storedPO != null,
- "The overwritten table %s in schema %s does not
exist",
- po.getTableName(),
- po.getSchemaId());
-
persistedPO.set(tablePOWithPersistedIdentityAndVersions(po, storedPO));
- }
- }),
- () ->
- SessionUtils.doWithoutCommit(
- TableVersionMapper.class,
- mapper -> {
- if (overwrite) {
- TablePO storedPO = persistedPO.get();
- // Retire the version row this overwrite replaces. There
is one only when the
- // upsert updated an existing table: the database then
moved the version from
- // N to N + 1, so the row to retire is N. When the
upsert inserted a brand new
- // table the version is still the initial one and no
earlier row exists.
- if (storedPO.getCurrentVersion() >
POConverters.INIT_VERSION) {
- mapper.softDeleteTableVersionByTableIdAndVersion(
- storedPO.getTableId(),
storedPO.getCurrentVersion() - 1);
- }
- mapper.insertTableVersionOnDuplicateKeyUpdate(storedPO);
- } else {
- mapper.insertTableVersion(po);
- }
- }),
- () -> {
- List<ColumnEntity> columns = tableEntity.columns();
- // We need to delete the columns first if we want to overwrite the
table.
- if (overwrite) {
- TablePO storedPO = persistedPO.get();
- // Overwriting the same table, e.g. re-importing it after an
out-of-band rename, keeps
- // the stored ids of the columns that still exist, so their
tags, owners and
- // privileges stay attached. When the upsert resolved to another
table's row, the
- // ids differ and that table's columns are not inherited.
- if (columns != null
- && po.getTableId().equals(storedPO.getTableId())
- && storedPO.getCurrentVersion() > POConverters.INIT_VERSION)
{
- columns =
- TableColumnMetaService.getInstance()
- .reuseStoredColumnIds(
- storedPO.getTableId(),
storedPO.getCurrentVersion(), columns);
- }
-
TableColumnMetaService.getInstance().deleteColumnsByTableId(storedPO.getTableId());
- }
-
- if (columns != null && !columns.isEmpty()) {
-
TableColumnMetaService.getInstance().insertColumnPOs(persistedPO.get(),
columns);
- }
- });
+ SchemaMetaService.getInstance()
+ .doWithSchemaWriteLock(
+ tableEntity.nameIdentifier(),
+ po.getSchemaId(),
+ po.getCatalogId(),
+ po.getMetalakeId(),
+ () ->
+ SessionUtils.doWithoutCommit(
+ TableMetaMapper.class,
+ mapper -> {
+ ops.insertPO(mapper, po, overwrite);
+ if (overwrite) {
+ // MySQL may preserve the existing table ID during
an upsert. Read the
+ // stored identity and database-generated version
while the row is locked.
+ TablePO storedPO =
+ mapper.selectTableMetaBySchemaIdAndName(
+ po.getSchemaId(), po.getTableName());
+ Preconditions.checkState(
+ storedPO != null,
+ "The overwritten table %s in schema %s does not
exist",
+ po.getTableName(),
+ po.getSchemaId());
+
persistedPO.set(tablePOWithPersistedIdentityAndVersions(po, storedPO));
+ }
+ }),
+ () ->
+ SessionUtils.doWithoutCommit(
+ TableVersionMapper.class,
+ mapper -> {
+ if (overwrite) {
+ TablePO storedPO = persistedPO.get();
+ // An existing table advances from N to N + 1 during
the upsert, so retire
+ // N before recording the new current version. A new
table has no N row.
+ if (storedPO.getCurrentVersion() >
POConverters.INIT_VERSION) {
+ mapper.softDeleteTableVersionByTableIdAndVersion(
+ storedPO.getTableId(),
storedPO.getCurrentVersion() - 1);
+ }
+
mapper.insertTableVersionOnDuplicateKeyUpdate(storedPO);
+ } else {
+ mapper.insertTableVersion(po);
+ }
+ }),
+ () -> {
+ List<ColumnEntity> columns = tableEntity.columns();
+ // We need to delete the columns first if we want to overwrite
the table.
+ if (overwrite) {
+ TablePO storedPO = persistedPO.get();
+ // Reuse stored column IDs when overwriting the same table,
for example after an
+ // out-of-band rename. This preserves their tags, owners,
and privileges. If the
+ // upsert resolves to another table's row, its column IDs
are not inherited.
+ if (columns != null
+ && po.getTableId().equals(storedPO.getTableId())
+ && storedPO.getCurrentVersion() >
POConverters.INIT_VERSION) {
+ columns =
+ TableColumnMetaService.getInstance()
+ .reuseStoredColumnIds(
+ storedPO.getTableId(),
storedPO.getCurrentVersion(), columns);
+ }
+ TableColumnMetaService.getInstance()
+ .deleteColumnsByTableId(storedPO.getTableId());
+ }
+
+ if (columns != null && !columns.isEmpty()) {
+
TableColumnMetaService.getInstance().insertColumnPOs(persistedPO.get(),
columns);
+ }
+ });
} catch (RuntimeException re) {
ExceptionUtils.checkSQLException(
@@ -228,57 +220,41 @@ public class TableMetaService {
POConverters.updateTablePOWithVersionAndSchemaId(oldTablePO,
newTableEntity, newSchemaId);
try {
- SessionUtils.doMultipleWithCommit(
- () -> {
- // Only an update that moves the table to another schema needs a
lock here. The new
- // parent must stay alive until the move commits; locking the old
parent would not
- // protect the table's new location.
- if (isSchemaChanged) {
- SchemaMetaService.getInstance()
- .lockSchemaForEntityWrite(
- newTableEntity.nameIdentifier(),
- newSchemaId,
- oldTablePO.getCatalogId(),
- oldTablePO.getMetalakeId());
- }
- },
- () -> {
- // This update is the decision point for the whole transaction.
current_version is the
- // table's OCC token: if another writer changed the table after we
read it, that writer
- // has already increased the token and this UPDATE changes zero
rows. Throwing here
- // rolls back the transaction before it can touch the version
history or columns.
- int updated =
- SessionUtils.getWithoutCommit(
- TableMetaMapper.class, mapper -> ops.updatePO(mapper,
newTablePO, oldTablePO));
- if (updated == 0) {
- throw tableWriteFailure(identifier, oldTablePO);
- }
- },
- () -> {
- // The table details live in table_version_info, keyed by
(table_id, version), while
- // table_meta only points at the current version. The two rows
have to move together,
- // and the upsert below has no version guard of its own: it
overwrites whatever sits
- // under that key.
- //
- // Say two writers both read version 5 and both want to write 6.
Their version rows
- // carry the same key, (table_id, 6), so whichever runs this
statement second would
- // silently replace the other's details. Ordering this step after
the table_meta CAS is
- // what prevents that: the loser matches no row up there, throws,
and the transaction
- // rolls back before reaching this statement. Only the winner ever
writes version 6.
- SessionUtils.doWithoutCommit(
- TableVersionMapper.class,
- mapper -> {
- mapper.softDeleteTableVersionByTableIdAndVersion(
- oldTablePO.getTableId(), oldTablePO.getCurrentVersion());
- mapper.insertTableVersionOnDuplicateKeyUpdate(newTablePO);
- });
- },
- () -> {
- // Column changes use the same new table version. Keeping this in
the same transaction
- // means a column failure also rolls back table_meta and
table_version_info.
- TableColumnMetaService.getInstance()
- .updateColumnPOsFromTableDiff(oldTableEntity, newTableEntity,
newTablePO);
- });
+ // For a cross-schema rename, the new schema is the parent that must
remain alive. For a
+ // regular update, newSchemaId is the existing parent, so the same entry
point covers both.
+ SchemaMetaService.getInstance()
+ .doWithSchemaWriteLock(
+ newTableEntity.nameIdentifier(),
+ newSchemaId,
+ oldTablePO.getCatalogId(),
+ oldTablePO.getMetalakeId(),
+ () -> {
+ // current_version is the table's OCC token. A zero-row update
means another
+ // writer won, so stop before touching version history or
columns.
+ int updated =
+ SessionUtils.getWithoutCommit(
+ TableMetaMapper.class,
+ mapper -> ops.updatePO(mapper, newTablePO,
oldTablePO));
+ if (updated == 0) {
+ throw tableWriteFailure(identifier, oldTablePO);
+ }
+ },
+ () ->
+ SessionUtils.doWithoutCommit(
+ TableVersionMapper.class,
+ mapper -> {
+ // Only the CAS winner can reach this step, so it is
safe to replace the
+ // details stored under the next table version.
+ mapper.softDeleteTableVersionByTableIdAndVersion(
+ oldTablePO.getTableId(),
oldTablePO.getCurrentVersion());
+
mapper.insertTableVersionOnDuplicateKeyUpdate(newTablePO);
+ }),
+ () -> {
+ // A column failure rolls back the table row and version row
in the same
+ // transaction.
+ TableColumnMetaService.getInstance()
+ .updateColumnPOsFromTableDiff(oldTableEntity,
newTableEntity, newTablePO);
+ });
} catch (RuntimeException re) {
ExceptionUtils.checkSQLException(
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/TopicMetaService.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/TopicMetaService.java
index fbe9c87565..28c94cd763 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/TopicMetaService.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/TopicMetaService.java
@@ -71,26 +71,22 @@ public class TopicMetaService {
fillTopicPOBuilderParentEntityId(builder, topicEntity.namespace());
TopicPO po = POConverters.initializeTopicPOWithVersion(topicEntity,
builder);
- SessionUtils.doMultipleWithCommit(
- // Hold the parent schema row until this transaction ends, so the
topic cannot be
- // written below a schema that is being dropped.
- () ->
- SchemaMetaService.getInstance()
- .lockSchemaForEntityWrite(
- topicEntity.nameIdentifier(),
- po.getSchemaId(),
- po.getCatalogId(),
- po.getMetalakeId()),
- () ->
- SessionUtils.doWithoutCommit(
- TopicMetaMapper.class,
- mapper -> {
- if (overwrite) {
- mapper.insertTopicMetaOnDuplicateKeyUpdate(po);
- } else {
- mapper.insertTopicMeta(po);
- }
- }));
+ SchemaMetaService.getInstance()
+ .doWithSchemaWriteLock(
+ topicEntity.nameIdentifier(),
+ po.getSchemaId(),
+ po.getCatalogId(),
+ po.getMetalakeId(),
+ () ->
+ SessionUtils.doWithoutCommit(
+ TopicMetaMapper.class,
+ mapper -> {
+ if (overwrite) {
+ mapper.insertTopicMetaOnDuplicateKeyUpdate(po);
+ } else {
+ mapper.insertTopicMeta(po);
+ }
+ }));
// TODO: insert topic dataLayout version after supporting it
} catch (RuntimeException re) {
ExceptionUtils.checkSQLException(
@@ -123,19 +119,25 @@ public class TopicMetaService {
try {
TopicPO newTopicPO = POConverters.updateTopicPOWithVersion(oldTopicPO,
newEntity);
- SessionUtils.doMultipleWithCommit(
- () -> {
- // current_version is the decision point for the whole write. Even
if another writer
- // changes the payload and later restores it, that writer still
advances the version,
- // so this stale update changes zero rows.
- int updated =
- SessionUtils.getWithoutCommit(
- TopicMetaMapper.class,
- mapper -> mapper.updateTopicMeta(newTopicPO, oldTopicPO));
- if (updated == 0) {
- throw topicWriteFailure(ident, oldTopicPO);
- }
- });
+ SchemaMetaService.getInstance()
+ .doWithSchemaWriteLock(
+ ident,
+ oldTopicPO.getSchemaId(),
+ oldTopicPO.getCatalogId(),
+ oldTopicPO.getMetalakeId(),
+ () -> {
+ // current_version is the decision point for the whole write.
Even if another writer
+ // changes the payload and later restores it, that writer
still advances the
+ // version,
+ // so this stale update changes zero rows.
+ int updated =
+ SessionUtils.getWithoutCommit(
+ TopicMetaMapper.class,
+ mapper -> mapper.updateTopicMeta(newTopicPO,
oldTopicPO));
+ if (updated == 0) {
+ throw topicWriteFailure(ident, oldTopicPO);
+ }
+ });
} catch (RuntimeException re) {
ExceptionUtils.checkSQLException(
re, Entity.EntityType.TOPIC, newEntity.nameIdentifier().toString());
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/ViewMetaService.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/ViewMetaService.java
index c96f3731d1..614f2c0f26 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/ViewMetaService.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/ViewMetaService.java
@@ -108,17 +108,13 @@ public class ViewMetaService {
try {
ViewPO po = initializeViewPO(viewEntity, builder);
- SessionUtils.doMultipleWithCommit(
- // Hold the parent schema row until this transaction ends, so the
view cannot be
- // written below a schema that is being dropped.
- () ->
- SchemaMetaService.getInstance()
- .lockSchemaForEntityWrite(
- viewEntity.nameIdentifier(),
- po.getSchemaId(),
- po.getCatalogId(),
- po.getMetalakeId()),
- () -> insertViewWithoutCommit(viewEntity, po, overwrite));
+ SchemaMetaService.getInstance()
+ .doWithSchemaWriteLock(
+ viewEntity.nameIdentifier(),
+ po.getSchemaId(),
+ po.getCatalogId(),
+ po.getMetalakeId(),
+ () -> insertViewWithoutCommit(viewEntity, po, overwrite));
} catch (RuntimeException re) {
try {
ExceptionUtils.checkSQLException(
@@ -164,28 +160,26 @@ public class ViewMetaService {
try {
ViewPO newViewPO = updateViewPO(oldViewPO, newEntity);
- SessionUtils.doMultipleWithCommit(
- () -> {
- if (isSchemaChanged) {
- SchemaMetaService.getInstance()
- .lockSchemaForEntityWrite(
- newEntity.nameIdentifier(), newSchemaId, newCatalogId,
newMetalakeId);
- }
- },
- () -> {
- // current_version is the sole OCC token. The root CAS is the
transaction's decision
- // point and must run before the unguarded version-row insert
below.
- int updated =
- SessionUtils.getWithoutCommit(
- ViewMetaMapper.class, mapper -> ops.updatePO(mapper,
newViewPO, oldViewPO));
- if (updated == 0) {
- throw viewWriteFailure(ident, oldViewPO);
- }
- },
- () ->
- SessionUtils.doWithoutCommit(
- ViewVersionInfoMapper.class,
- mapper ->
mapper.insertViewVersionInfo(newViewPO.getViewVersionInfoPO())));
+ SchemaMetaService.getInstance()
+ .doWithSchemaWriteLock(
+ newEntity.nameIdentifier(),
+ newSchemaId,
+ newCatalogId,
+ newMetalakeId,
+ () -> {
+ // current_version is the sole OCC token. The root CAS is the
transaction's decision
+ // point and must run before the unguarded version-row insert
below.
+ int updated =
+ SessionUtils.getWithoutCommit(
+ ViewMetaMapper.class, mapper -> ops.updatePO(mapper,
newViewPO, oldViewPO));
+ if (updated == 0) {
+ throw viewWriteFailure(ident, oldViewPO);
+ }
+ },
+ () ->
+ SessionUtils.doWithoutCommit(
+ ViewVersionInfoMapper.class,
+ mapper ->
mapper.insertViewVersionInfo(newViewPO.getViewVersionInfoPO())));
return newEntity;
} catch (RuntimeException re) {
ExceptionUtils.checkSQLException(
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/mapper/provider/base/TestSchemaMetaBaseSQLProvider.java
b/core/src/test/java/org/apache/gravitino/storage/relational/mapper/provider/base/TestSchemaMetaBaseSQLProvider.java
new file mode 100644
index 0000000000..6cedd927d8
--- /dev/null
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/mapper/provider/base/TestSchemaMetaBaseSQLProvider.java
@@ -0,0 +1,61 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.gravitino.storage.relational.mapper.provider.base;
+
+import java.util.Arrays;
+import java.util.List;
+import org.apache.gravitino.storage.relational.mapper.FilesetMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.FunctionMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.ModelMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.TableMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.TopicMetaMapper;
+import org.apache.gravitino.storage.relational.mapper.ViewMetaMapper;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+class TestSchemaMetaBaseSQLProvider {
+
+ private static final SchemaMetaBaseSQLProvider PROVIDER = new
SchemaMetaBaseSQLProvider();
+
+ @Test
+ void testSelectActiveChildChecksEverySupportedChildType() {
+ String sql = PROVIDER.selectActiveChildBySchemaId(null);
+ List<String> childTables =
+ Arrays.asList(
+ TableMetaMapper.TABLE_NAME,
+ ViewMetaMapper.TABLE_NAME,
+ FilesetMetaMapper.META_TABLE_NAME,
+ FunctionMetaMapper.TABLE_NAME,
+ ModelMetaMapper.TABLE_NAME,
+ TopicMetaMapper.TABLE_NAME);
+
+ childTables.forEach(
+ tableName ->
+ Assertions.assertTrue(
+ sql.contains(
+ "FROM " + tableName + " WHERE schema_id = #{schemaId} AND
deleted_at = 0"),
+ () -> "Missing active-child check for " + tableName + " in: "
+ sql));
+ Assertions.assertEquals(childTables.size() - 1, countOccurrences(sql,
"UNION ALL"));
+ Assertions.assertTrue(sql.endsWith("LIMIT 1"));
+ }
+
+ private static int countOccurrences(String value, String target) {
+ return (value.length() - value.replace(target, "").length()) /
target.length();
+ }
+}
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestModelVersionMetaService.java
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestModelVersionMetaService.java
index 1aa158666f..e67f665dc2 100644
---
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestModelVersionMetaService.java
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestModelVersionMetaService.java
@@ -1711,6 +1711,60 @@ public class TestModelVersionMetaService extends
TestJDBCBackend {
Assertions.assertEquals(ImmutableList.of("second_alias"),
unchanged.aliases());
}
+ @TestTemplate
+ void testAliasConflictRollsBackModelVersionRegistration() throws IOException
{
+ createParentEntities(METALAKE_NAME, CATALOG_NAME, SCHEMA_NAME, AUDIT_INFO);
+ ModelEntity model =
+ createModelEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ MODEL_NS,
+ randomModelName(),
+ "model comment",
+ 0,
+ properties,
+ AUDIT_INFO);
+ ModelMetaService.getInstance().insertModel(model, false);
+ ModelVersionEntity first =
+ createModelVersionEntity(
+ model.nameIdentifier(),
+ 0,
+ ImmutableMap.of(ModelVersion.URI_NAME_UNKNOWN, "first_path"),
+ ImmutableList.of("taken_alias"),
+ "first",
+ properties,
+ AUDIT_INFO);
+ ModelVersionMetaService.getInstance().insertModelVersion(first);
+ ModelPO beforeFailure =
ModelMetaService.getInstance().getModelPOById(model.id());
+ ModelVersionEntity conflicting =
+ createModelVersionEntity(
+ model.nameIdentifier(),
+ 1,
+ ImmutableMap.of(ModelVersion.URI_NAME_UNKNOWN, "second_path"),
+ ImmutableList.of("taken_alias"),
+ "must roll back",
+ properties,
+ AUDIT_INFO);
+
+ Assertions.assertThrows(
+ RuntimeException.class,
+ () ->
ModelVersionMetaService.getInstance().insertModelVersion(conflicting));
+
+ ModelPO afterFailure =
ModelMetaService.getInstance().getModelPOById(model.id());
+ Assertions.assertEquals(beforeFailure.getCurrentVersion(),
afterFailure.getCurrentVersion());
+ Assertions.assertEquals(beforeFailure.getLastVersion(),
afterFailure.getLastVersion());
+ Assertions.assertEquals(
+ beforeFailure.getModelLatestVersion(),
afterFailure.getModelLatestVersion());
+ Assertions.assertTrue(
+ SessionUtils.getWithoutCommit(
+ ModelVersionMetaMapper.class,
+ mapper -> mapper.selectModelVersionMeta(model.id(),
conflicting.version()))
+ .isEmpty());
+ ModelVersionEntity unchanged =
+
ModelVersionMetaService.getInstance().getModelVersionByIdentifier(first.nameIdentifier());
+ Assertions.assertEquals(ImmutableList.of("taken_alias"),
unchanged.aliases());
+ Assertions.assertEquals("first", unchanged.comment());
+ }
+
@TestTemplate
void testDeleteModelVersionsByLegacyTimeline() throws IOException {
createParentEntities(METALAKE_NAME, CATALOG_NAME, SCHEMA_NAME, AUDIT_INFO);
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestSchemaMetaService.java
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestSchemaMetaService.java
index 72d721a4c2..4ea5eea35c 100644
---
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestSchemaMetaService.java
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestSchemaMetaService.java
@@ -31,6 +31,7 @@ import java.time.Instant;
import java.util.Arrays;
import java.util.Collections;
import java.util.List;
+import java.util.Locale;
import java.util.Objects;
import java.util.Set;
import java.util.concurrent.CountDownLatch;
@@ -140,71 +141,81 @@ public class TestSchemaMetaService extends
TestJDBCBackend {
createAndInsertMakeLake(metalakeName);
createAndInsertCatalog(metalakeName, catalogName);
- List<SchemaChildWrite> childWrites =
- Arrays.asList(
- namespace ->
- backend.insert(
- createTableEntity(
- RandomIdGenerator.INSTANCE.nextId(), namespace,
"child_table", AUDIT_INFO),
- false),
- namespace ->
- backend.insert(
- createViewEntity(RandomIdGenerator.INSTANCE.nextId(),
namespace, "child_view"),
- false),
- namespace ->
- backend.insert(
- createFilesetEntity(
- RandomIdGenerator.INSTANCE.nextId(),
- namespace,
- "child_fileset",
- AUDIT_INFO),
- false),
- namespace ->
- backend.insert(
- createFunctionEntity(
- RandomIdGenerator.INSTANCE.nextId(),
- namespace,
- "child_function",
- AUDIT_INFO),
- false),
- namespace ->
- backend.insert(
- createModelEntity(
- RandomIdGenerator.INSTANCE.nextId(),
- namespace,
- "child_model",
- "model comment",
- 0,
- Collections.emptyMap(),
- AUDIT_INFO),
- false),
- namespace ->
- backend.insert(
- createTopicEntity(
- RandomIdGenerator.INSTANCE.nextId(), namespace,
"child_topic", AUDIT_INFO),
- false));
+ for (SchemaChildCase childCase : schemaChildCases()) {
+ SchemaEntity schema =
+ createSchemaEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ NamespaceUtil.ofSchema(metalakeName, catalogName),
+ "schema_for_entity_lock_" +
childCase.type.name().toLowerCase(Locale.ROOT),
+ AUDIT_INFO);
+ backend.insert(schema, false);
+ Namespace childNamespace = Namespace.of(metalakeName, catalogName,
schema.name());
+ assertSchemaChildActionWaitsForConcurrentDelete(
+ schema, () -> childCase.write.run(childNamespace));
+ }
+ }
- for (int index = 0; index < childWrites.size(); index++) {
+ @TestTemplate
+ public void testSchemaChildUpdatesWaitForConcurrentSchemaDelete() throws
Exception {
+ createAndInsertMakeLake(metalakeName);
+ createAndInsertCatalog(metalakeName, catalogName);
+
+ for (SchemaChildCase childCase : schemaChildCases()) {
SchemaEntity schema =
createSchemaEntity(
RandomIdGenerator.INSTANCE.nextId(),
NamespaceUtil.ofSchema(metalakeName, catalogName),
- "schema_for_entity_lock_" + index,
+ "schema_for_entity_update_lock_" +
childCase.type.name().toLowerCase(Locale.ROOT),
AUDIT_INFO);
backend.insert(schema, false);
- assertChildWriteWaitsForConcurrentSchemaDelete(schema,
childWrites.get(index));
+ Namespace childNamespace = Namespace.of(metalakeName, catalogName,
schema.name());
+ childCase.write.run(childNamespace);
+
+ NameIdentifier childIdentifier =
+ NameIdentifier.of(
+ childNamespace, "child_" +
childCase.type.name().toLowerCase(Locale.ROOT));
+ assertSchemaChildActionWaitsForConcurrentDelete(
+ schema, () -> backend.update(childIdentifier, childCase.type, entity
-> entity));
}
}
- private void assertChildWriteWaitsForConcurrentSchemaDelete(
- SchemaEntity schema, SchemaChildWrite childWrite) throws Exception {
+ @TestTemplate
+ public void testSchemaActiveChildExistenceQueryCoversEveryChildType() throws
Exception {
+ createAndInsertMakeLake(metalakeName);
+ createAndInsertCatalog(metalakeName, catalogName);
+
+ for (SchemaChildCase childCase : schemaChildCases()) {
+ SchemaEntity schema =
+ createSchemaEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ NamespaceUtil.ofSchema(metalakeName, catalogName),
+ "schema_for_child_exists_" +
childCase.type.name().toLowerCase(Locale.ROOT),
+ AUDIT_INFO);
+ backend.insert(schema, false);
+
+ Assertions.assertNull(selectActiveSchemaChild(schema.id()));
+ childCase.write.run(Namespace.of(metalakeName, catalogName,
schema.name()));
+ Assertions.assertEquals(1, selectActiveSchemaChild(schema.id()));
+
+ // Each UNION branch must protect the public non-cascade delete path,
not merely return a
+ // value when the mapper is called directly.
+ assertThrows(
+ NonEmptyEntityException.class,
+ () ->
SchemaMetaService.getInstance().deleteSchema(schema.nameIdentifier(), false));
+ SchemaMetaService.getInstance().deleteSchema(schema.nameIdentifier(),
true);
+ Assertions.assertNull(selectActiveSchemaChild(schema.id()));
+ }
+ }
+
+ private void assertSchemaChildActionWaitsForConcurrentDelete(
+ SchemaEntity schema, SchemaChildAction childAction) throws Exception {
SchemaPO observedSchemaPO =
SessionUtils.getWithoutCommit(
SchemaMetaMapper.class, mapper ->
mapper.selectSchemaMetaById(schema.id()));
CountDownLatch schemaDeleteLocked = new CountDownLatch(1);
CountDownLatch allowDeleteCommit = new CountDownLatch(1);
- CountDownLatch entityCreateStarted = new CountDownLatch(1);
+ CountDownLatch childActionStarted = new CountDownLatch(1);
ExecutorService executor = Executors.newFixedThreadPool(2);
Future<Throwable> deleteResult =
executor.submit(
@@ -235,26 +246,26 @@ public class TestSchemaMetaService extends
TestJDBCBackend {
});
try {
assertTrue(schemaDeleteLocked.await(30, TimeUnit.SECONDS));
- Future<Throwable> createResult =
+ Future<Throwable> childActionResult =
executor.submit(
() -> {
- entityCreateStarted.countDown();
+ childActionStarted.countDown();
try {
// Exercise the real JDBCBackend-to-service path. This test
must fail if any
// schema-scoped service forgets to take the parent lock in
its own transaction.
- childWrite.run(Namespace.of(metalakeName, catalogName,
schema.name()));
+ childAction.run();
return null;
} catch (Throwable throwable) {
return throwable;
}
});
- assertTrue(entityCreateStarted.await(30, TimeUnit.SECONDS));
- assertThrows(TimeoutException.class, () -> createResult.get(500,
TimeUnit.MILLISECONDS));
+ assertTrue(childActionStarted.await(30, TimeUnit.SECONDS));
+ assertThrows(TimeoutException.class, () -> childActionResult.get(500,
TimeUnit.MILLISECONDS));
allowDeleteCommit.countDown();
Assertions.assertNull(deleteResult.get(30, TimeUnit.SECONDS));
Assertions.assertInstanceOf(
- NoSuchEntityException.class, createResult.get(30, TimeUnit.SECONDS));
+ NoSuchEntityException.class, childActionResult.get(30,
TimeUnit.SECONDS));
} finally {
allowDeleteCommit.countDown();
executor.shutdownNow();
@@ -1210,8 +1221,85 @@ public class TestSchemaMetaService extends
TestJDBCBackend {
}
}
+ private Integer selectActiveSchemaChild(Long schemaId) {
+ return SessionUtils.getWithoutCommit(
+ SchemaMetaMapper.class, mapper ->
mapper.selectActiveChildBySchemaId(schemaId));
+ }
+
+ private List<SchemaChildCase> schemaChildCases() {
+ return Arrays.asList(
+ new SchemaChildCase(
+ Entity.EntityType.TABLE,
+ namespace ->
+ backend.insert(
+ createTableEntity(
+ RandomIdGenerator.INSTANCE.nextId(), namespace,
"child_table", AUDIT_INFO),
+ false)),
+ new SchemaChildCase(
+ Entity.EntityType.VIEW,
+ namespace ->
+ backend.insert(
+ createViewEntity(RandomIdGenerator.INSTANCE.nextId(),
namespace, "child_view"),
+ false)),
+ new SchemaChildCase(
+ Entity.EntityType.FILESET,
+ namespace ->
+ backend.insert(
+ createFilesetEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ namespace,
+ "child_fileset",
+ AUDIT_INFO),
+ false)),
+ new SchemaChildCase(
+ Entity.EntityType.FUNCTION,
+ namespace ->
+ backend.insert(
+ createFunctionEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ namespace,
+ "child_function",
+ AUDIT_INFO),
+ false)),
+ new SchemaChildCase(
+ Entity.EntityType.MODEL,
+ namespace ->
+ backend.insert(
+ createModelEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ namespace,
+ "child_model",
+ "model comment",
+ 0,
+ Collections.emptyMap(),
+ AUDIT_INFO),
+ false)),
+ new SchemaChildCase(
+ Entity.EntityType.TOPIC,
+ namespace ->
+ backend.insert(
+ createTopicEntity(
+ RandomIdGenerator.INSTANCE.nextId(), namespace,
"child_topic", AUDIT_INFO),
+ false)));
+ }
+
+ @FunctionalInterface
+ private interface SchemaChildAction {
+ void run() throws Exception;
+ }
+
@FunctionalInterface
private interface SchemaChildWrite {
void run(Namespace namespace) throws Exception;
}
+
+ private static class SchemaChildCase {
+ private final Entity.EntityType type;
+ private final SchemaChildWrite write;
+
+ private SchemaChildCase(Entity.EntityType type, SchemaChildWrite write) {
+ this.type = type;
+ this.write = write;
+ }
+ }
}
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestStatisticMetaService.java
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestStatisticMetaService.java
index 8675448ee9..a079abf5d8 100644
---
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestStatisticMetaService.java
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestStatisticMetaService.java
@@ -25,8 +25,16 @@ import java.sql.SQLException;
import java.sql.Statement;
import java.time.Instant;
import java.util.List;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
import org.apache.gravitino.Entity;
+import org.apache.gravitino.NameIdentifier;
import org.apache.gravitino.Namespace;
+import org.apache.gravitino.exceptions.NoSuchEntityException;
import org.apache.gravitino.meta.AuditInfo;
import org.apache.gravitino.meta.BaseMetalake;
import org.apache.gravitino.meta.CatalogEntity;
@@ -40,7 +48,10 @@ import org.apache.gravitino.meta.TopicEntity;
import org.apache.gravitino.stats.StatisticValues;
import org.apache.gravitino.storage.RandomIdGenerator;
import org.apache.gravitino.storage.relational.TestJDBCBackend;
+import org.apache.gravitino.storage.relational.mapper.SchemaMetaMapper;
+import org.apache.gravitino.storage.relational.po.SchemaPO;
import org.apache.gravitino.storage.relational.session.SqlSessionFactoryHelper;
+import org.apache.gravitino.storage.relational.utils.SessionUtils;
import org.apache.ibatis.session.SqlSession;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.TestTemplate;
@@ -48,6 +59,170 @@ import org.junit.jupiter.api.TestTemplate;
public class TestStatisticMetaService extends TestJDBCBackend {
private final StatisticMetaService statisticMetaService =
StatisticMetaService.getInstance();
+ @TestTemplate
+ public void testTableStatisticWriteWaitsForConcurrentSchemaDelete() throws
Exception {
+ String metalakeName = "metalake_for_statistic_schema_fence";
+ String catalogName = "catalog_for_statistic_schema_fence";
+ String schemaName = "schema_for_statistic_schema_fence";
+ AuditInfo auditInfo =
+
AuditInfo.builder().withCreator("creator").withCreateTime(Instant.now()).build();
+ createParentEntities(metalakeName, catalogName, schemaName, auditInfo);
+ SchemaPO observedSchemaPO = selectSchemaPO(metalakeName, catalogName,
schemaName);
+
+ assertTableStatisticWriteWaitsForConcurrentSchemaDelete(
+ metalakeName,
+ catalogName,
+ schemaName,
+ auditInfo,
+ () -> {
+ int deleted =
+ SessionUtils.getWithoutCommit(
+ SchemaMetaMapper.class,
+ mapper ->
+ mapper.softDeleteSchemaMetaBySchemaIdAndVersion(
+ observedSchemaPO.getSchemaId(),
observedSchemaPO.getCurrentVersion()));
+ Assertions.assertEquals(1, deleted);
+ });
+ }
+
+ @TestTemplate
+ public void
testNestedSchemaTableStatisticWriteWaitsForAncestorCascadeDelete() throws
Exception {
+ String metalakeName = "metalake_for_nested_statistic_fence";
+ String catalogName = "catalog_for_nested_statistic_fence";
+ String ancestorName = "fence_anc_a";
+ String nestedName = ancestorName + ":fence_anc_b";
+ AuditInfo auditInfo =
+
AuditInfo.builder().withCreator("creator").withCreateTime(Instant.now()).build();
+ createAndInsertMakeLake(metalakeName);
+ createAndInsertCatalog(metalakeName, catalogName);
+ // Inserting the nested leaf also creates the ancestor row.
+ SchemaMetaService.getInstance()
+ .insertSchema(
+ createSchemaEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ Namespace.of(metalakeName, catalogName),
+ nestedName,
+ auditInfo),
+ false);
+ SchemaPO observedAncestorPO = selectSchemaPO(metalakeName, catalogName,
ancestorName);
+ SchemaPO observedNestedPO = selectSchemaPO(metalakeName, catalogName,
nestedName);
+
+ // The statistic only locks its direct parent, the nested schema. A
cascade delete of the
+ // ancestor soft-deletes every descendant schema row in the same
transaction, so the write
+ // must still wait for it and then fail once the nested schema is gone.
+ assertTableStatisticWriteWaitsForConcurrentSchemaDelete(
+ metalakeName,
+ catalogName,
+ nestedName,
+ auditInfo,
+ () -> {
+ int deletedAncestor =
+ SessionUtils.getWithoutCommit(
+ SchemaMetaMapper.class,
+ mapper ->
+ mapper.softDeleteSchemaMetaBySchemaIdAndVersion(
+ observedAncestorPO.getSchemaId(),
+ observedAncestorPO.getCurrentVersion()));
+ Assertions.assertEquals(1, deletedAncestor);
+ int deletedDescendants =
+ SessionUtils.getWithoutCommit(
+ SchemaMetaMapper.class,
+ mapper ->
mapper.softDeleteSchemaMetasWithVersion(List.of(observedNestedPO)));
+ Assertions.assertEquals(1, deletedDescendants);
+ });
+ }
+
+ private SchemaPO selectSchemaPO(String metalakeName, String catalogName,
String schemaName) {
+ Long schemaId =
+ EntityIdService.getEntityId(
+ NameIdentifier.of(metalakeName, catalogName, schemaName),
Entity.EntityType.SCHEMA);
+ return SessionUtils.getWithoutCommit(
+ SchemaMetaMapper.class, mapper ->
mapper.selectSchemaMetaById(schemaId));
+ }
+
+ /**
+ * Runs {@code schemaDeleteStep} in an uncommitted transaction, then
verifies that a statistic
+ * upsert on a table below {@code schemaName} blocks until that transaction
commits and fails with
+ * {@link NoSuchEntityException} afterwards.
+ */
+ private void assertTableStatisticWriteWaitsForConcurrentSchemaDelete(
+ String metalakeName,
+ String catalogName,
+ String schemaName,
+ AuditInfo auditInfo,
+ Runnable schemaDeleteStep)
+ throws Exception {
+ Long metalakeId =
+ EntityIdService.getEntityId(NameIdentifier.of(metalakeName),
Entity.EntityType.METALAKE);
+ TableEntity table =
+ createTableEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ Namespace.of(metalakeName, catalogName, schemaName),
+ "table",
+ auditInfo);
+ backend.insert(table, false);
+ StatisticEntity statistic =
+ TableStatisticEntity.builder()
+ .withId(RandomIdGenerator.INSTANCE.nextId())
+ .withName("test")
+ .withNamespace(Namespace.of(metalakeName, catalogName, schemaName,
table.name()))
+ .withValue(StatisticValues.longValue(100L))
+ .withAuditInfo(auditInfo)
+ .build();
+
+ CountDownLatch schemaDeleteLocked = new CountDownLatch(1);
+ CountDownLatch allowDeleteCommit = new CountDownLatch(1);
+ CountDownLatch statisticWriteStarted = new CountDownLatch(1);
+ ExecutorService executor = Executors.newFixedThreadPool(2);
+ Future<Throwable> deleteResult =
+ executor.submit(
+ () -> {
+ try {
+ SessionUtils.doMultipleWithCommit(
+ () -> {
+ schemaDeleteStep.run();
+ schemaDeleteLocked.countDown();
+ try {
+ Assertions.assertTrue(allowDeleteCommit.await(30,
TimeUnit.SECONDS));
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ throw new RuntimeException(e);
+ }
+ });
+ return null;
+ } catch (Throwable throwable) {
+ return throwable;
+ }
+ });
+ try {
+ Assertions.assertTrue(schemaDeleteLocked.await(30, TimeUnit.SECONDS));
+ Future<Throwable> statisticWriteResult =
+ executor.submit(
+ () -> {
+ statisticWriteStarted.countDown();
+ try {
+ backend.batchPut(List.of(statistic), true);
+ return null;
+ } catch (Throwable throwable) {
+ return throwable;
+ }
+ });
+ Assertions.assertTrue(statisticWriteStarted.await(30, TimeUnit.SECONDS));
+ Assertions.assertThrows(
+ TimeoutException.class, () -> statisticWriteResult.get(500,
TimeUnit.MILLISECONDS));
+
+ allowDeleteCommit.countDown();
+ Assertions.assertNull(deleteResult.get(30, TimeUnit.SECONDS));
+ Assertions.assertInstanceOf(
+ NoSuchEntityException.class, statisticWriteResult.get(30,
TimeUnit.SECONDS));
+ } finally {
+ allowDeleteCommit.countDown();
+ executor.shutdownNow();
+ }
+
+ Assertions.assertEquals(0, countActiveStats(metalakeId));
+ }
+
@TestTemplate
public void testStatisticsLifeCycle() throws Exception {
String metalakeName = "metalake";