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

Reply via email to