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

yuqi1129 pushed a commit to branch feat/12166-occ-drop-cas
in repository https://gitbox.apache.org/repos/asf/gravitino.git

commit e6245150e79e595da984553d840c24503e4ab58e
Author: yuqi <[email protected]>
AuthorDate: Fri Jul 24 16:15:37 2026 +0800

    [#12166] improvement(core): version-checked soft-delete (drop CAS) for 
metalake and catalog
    
    These delete services resolved only the entity id, so fetch the PO for the 
version
    (getMetalakePOByName selectMetalakeMetaByName / getCatalogPOByName; both 
throw
    NoSuch, preserving behavior), pass current_version to the versioned 
self-delete, and
    return based on the captured row count (Option A) instead of always-true. 
Cascade
    soft-deletes by parent id are unchanged. base + PostgreSQL providers 
updated.
    
    Part of #12166 (drop CAS).
---
 .../relational/mapper/CatalogMetaMapper.java       |  3 ++-
 .../mapper/CatalogMetaSQLProviderFactory.java      |  5 ++--
 .../relational/mapper/MetalakeMetaMapper.java      |  3 ++-
 .../mapper/MetalakeMetaSQLProviderFactory.java     |  5 ++--
 .../provider/base/CatalogMetaBaseSQLProvider.java  |  7 ++++--
 .../provider/base/MetalakeMetaBaseSQLProvider.java |  7 ++++--
 .../postgresql/CatalogMetaPostgreSQLProvider.java  |  6 +++--
 .../postgresql/MetalakeMetaPostgreSQLProvider.java |  6 +++--
 .../relational/service/CatalogMetaService.java     | 24 ++++++++++++-------
 .../relational/service/MetalakeMetaService.java    | 28 +++++++++++++++-------
 10 files changed, 63 insertions(+), 31 deletions(-)

diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/CatalogMetaMapper.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/CatalogMetaMapper.java
index 9f19d1a6a8..f1f1f3231a 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/CatalogMetaMapper.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/CatalogMetaMapper.java
@@ -89,7 +89,8 @@ public interface CatalogMetaMapper {
   @UpdateProvider(
       type = CatalogMetaSQLProviderFactory.class,
       method = "softDeleteCatalogMetasByCatalogId")
-  Integer softDeleteCatalogMetasByCatalogId(@Param("catalogId") Long 
catalogId);
+  Integer softDeleteCatalogMetasByCatalogId(
+      @Param("catalogId") Long catalogId, @Param("currentVersion") Long 
currentVersion);
 
   @UpdateProvider(
       type = CatalogMetaSQLProviderFactory.class,
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/CatalogMetaSQLProviderFactory.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/CatalogMetaSQLProviderFactory.java
index c3a7954a25..506c7874b9 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/CatalogMetaSQLProviderFactory.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/CatalogMetaSQLProviderFactory.java
@@ -109,8 +109,9 @@ public class CatalogMetaSQLProviderFactory {
     return getProvider().updateCatalogMeta(newCatalogPO, oldCatalogPO);
   }
 
-  public static String softDeleteCatalogMetasByCatalogId(@Param("catalogId") 
Long catalogId) {
-    return getProvider().softDeleteCatalogMetasByCatalogId(catalogId);
+  public static String softDeleteCatalogMetasByCatalogId(
+      @Param("catalogId") Long catalogId, @Param("currentVersion") Long 
currentVersion) {
+    return getProvider().softDeleteCatalogMetasByCatalogId(catalogId, 
currentVersion);
   }
 
   public static String softDeleteCatalogMetasByMetalakeId(@Param("metalakeId") 
Long metalakeId) {
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/MetalakeMetaMapper.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/MetalakeMetaMapper.java
index f705c283ce..33787a71a5 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/MetalakeMetaMapper.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/MetalakeMetaMapper.java
@@ -73,7 +73,8 @@ public interface MetalakeMetaMapper {
   @UpdateProvider(
       type = MetalakeMetaSQLProviderFactory.class,
       method = "softDeleteMetalakeMetaByMetalakeId")
-  Integer softDeleteMetalakeMetaByMetalakeId(@Param("metalakeId") Long 
metalakeId);
+  Integer softDeleteMetalakeMetaByMetalakeId(
+      @Param("metalakeId") Long metalakeId, @Param("currentVersion") Long 
currentVersion);
 
   @DeleteProvider(
       type = MetalakeMetaSQLProviderFactory.class,
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/MetalakeMetaSQLProviderFactory.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/MetalakeMetaSQLProviderFactory.java
index eba26f9e02..6d95b3a1e4 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/MetalakeMetaSQLProviderFactory.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/MetalakeMetaSQLProviderFactory.java
@@ -88,8 +88,9 @@ public class MetalakeMetaSQLProviderFactory {
     return getProvider().updateMetalakeMeta(newMetalakePO, oldMetalakePO);
   }
 
-  public static String softDeleteMetalakeMetaByMetalakeId(@Param("metalakeId") 
Long metalakeId) {
-    return getProvider().softDeleteMetalakeMetaByMetalakeId(metalakeId);
+  public static String softDeleteMetalakeMetaByMetalakeId(
+      @Param("metalakeId") Long metalakeId, @Param("currentVersion") Long 
currentVersion) {
+    return getProvider().softDeleteMetalakeMetaByMetalakeId(metalakeId, 
currentVersion);
   }
 
   public static String deleteMetalakeMetasByLegacyTimeline(
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/CatalogMetaBaseSQLProvider.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/CatalogMetaBaseSQLProvider.java
index 6c486db95b..c38db5ed1e 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/CatalogMetaBaseSQLProvider.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/CatalogMetaBaseSQLProvider.java
@@ -211,12 +211,15 @@ public class CatalogMetaBaseSQLProvider {
         + " AND deleted_at = 0";
   }
 
-  public String softDeleteCatalogMetasByCatalogId(@Param("catalogId") Long 
catalogId) {
+  public String softDeleteCatalogMetasByCatalogId(
+      @Param("catalogId") Long catalogId, @Param("currentVersion") Long 
currentVersion) {
     return "UPDATE "
         + TABLE_NAME
         + " SET deleted_at = (UNIX_TIMESTAMP() * 1000.0)"
         + " + EXTRACT(MICROSECOND FROM CURRENT_TIMESTAMP(3)) / 1000"
-        + " WHERE catalog_id = #{catalogId} AND deleted_at = 0";
+        // OCC: version-checked delete (0 rows = stale version; the service 
returns false).
+        + " WHERE catalog_id = #{catalogId} AND current_version = 
#{currentVersion}"
+        + " AND deleted_at = 0";
   }
 
   public String softDeleteCatalogMetasByMetalakeId(@Param("metalakeId") Long 
metalakeId) {
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/MetalakeMetaBaseSQLProvider.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/MetalakeMetaBaseSQLProvider.java
index 4892e8816b..683a1f6527 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/MetalakeMetaBaseSQLProvider.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/MetalakeMetaBaseSQLProvider.java
@@ -148,12 +148,15 @@ public class MetalakeMetaBaseSQLProvider {
         + " AND deleted_at = 0";
   }
 
-  public String softDeleteMetalakeMetaByMetalakeId(@Param("metalakeId") Long 
metalakeId) {
+  public String softDeleteMetalakeMetaByMetalakeId(
+      @Param("metalakeId") Long metalakeId, @Param("currentVersion") Long 
currentVersion) {
     return "UPDATE "
         + TABLE_NAME
         + " SET deleted_at = (UNIX_TIMESTAMP() * 1000.0)"
         + " + EXTRACT(MICROSECOND FROM CURRENT_TIMESTAMP(3)) / 1000"
-        + " WHERE metalake_id = #{metalakeId} AND deleted_at = 0";
+        // OCC: version-checked delete (0 rows = stale version; the service 
returns false).
+        + " WHERE metalake_id = #{metalakeId} AND current_version = 
#{currentVersion}"
+        + " AND deleted_at = 0";
   }
 
   public String deleteMetalakeMetasByLegacyTimeline(
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/CatalogMetaPostgreSQLProvider.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/CatalogMetaPostgreSQLProvider.java
index 2c79fa2060..c5cde0bbcf 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/CatalogMetaPostgreSQLProvider.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/CatalogMetaPostgreSQLProvider.java
@@ -26,11 +26,13 @@ import org.apache.ibatis.annotations.Param;
 
 public class CatalogMetaPostgreSQLProvider extends CatalogMetaBaseSQLProvider {
   @Override
-  public String softDeleteCatalogMetasByCatalogId(Long catalogId) {
+  public String softDeleteCatalogMetasByCatalogId(Long catalogId, Long 
currentVersion) {
     return "UPDATE "
         + TABLE_NAME
         + " SET deleted_at = CAST(EXTRACT(EPOCH FROM CURRENT_TIMESTAMP) * 1000 
AS BIGINT)"
-        + " WHERE catalog_id = #{catalogId} AND deleted_at = 0";
+        // OCC: version-checked delete (see the base provider for rationale).
+        + " WHERE catalog_id = #{catalogId} AND current_version = 
#{currentVersion}"
+        + " AND deleted_at = 0";
   }
 
   @Override
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/MetalakeMetaPostgreSQLProvider.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/MetalakeMetaPostgreSQLProvider.java
index 2aeb5f9b67..501a7b9efe 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/MetalakeMetaPostgreSQLProvider.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/MetalakeMetaPostgreSQLProvider.java
@@ -26,11 +26,13 @@ import org.apache.ibatis.annotations.Param;
 
 public class MetalakeMetaPostgreSQLProvider extends 
MetalakeMetaBaseSQLProvider {
   @Override
-  public String softDeleteMetalakeMetaByMetalakeId(Long metalakeId) {
+  public String softDeleteMetalakeMetaByMetalakeId(Long metalakeId, Long 
currentVersion) {
     return "UPDATE "
         + TABLE_NAME
         + " SET deleted_at = CAST(EXTRACT(EPOCH FROM CURRENT_TIMESTAMP) * 1000 
AS BIGINT)"
-        + " WHERE metalake_id = #{metalakeId} AND deleted_at = 0";
+        // OCC: version-checked delete (see the base provider for rationale).
+        + " WHERE metalake_id = #{metalakeId} AND current_version = 
#{currentVersion}"
+        + " AND deleted_at = 0";
   }
 
   @Override
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/service/CatalogMetaService.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/service/CatalogMetaService.java
index a85e1b85ba..c98b6eb9d3 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/service/CatalogMetaService.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/service/CatalogMetaService.java
@@ -270,15 +270,20 @@ public class CatalogMetaService {
     NameIdentifierUtil.checkCatalog(identifier);
 
     String catalogName = identifier.name();
-    long catalogId = EntityIdService.getEntityId(identifier, 
Entity.EntityType.CATALOG);
     String metalakeName = identifier.namespace().level(0);
+    CatalogPO catalogPO = getCatalogPOByName(metalakeName, catalogName);
+    Long catalogId = catalogPO.getCatalogId();
+    Long currentVersion = catalogPO.getCurrentVersion();
+    AtomicInteger deleteCount = new AtomicInteger(0);
 
     if (cascade) {
       SessionUtils.doMultipleWithCommit(
           () ->
-              SessionUtils.doWithoutCommit(
-                  CatalogMetaMapper.class,
-                  mapper -> 
mapper.softDeleteCatalogMetasByCatalogId(catalogId)),
+              deleteCount.set(
+                  SessionUtils.getWithoutCommit(
+                      CatalogMetaMapper.class,
+                      mapper ->
+                          mapper.softDeleteCatalogMetasByCatalogId(catalogId, 
currentVersion))),
           () ->
               SessionUtils.doWithoutCommit(
                   SchemaMetaMapper.class,
@@ -369,9 +374,11 @@ public class CatalogMetaService {
       }
       SessionUtils.doMultipleWithCommit(
           () ->
-              SessionUtils.doWithoutCommit(
-                  CatalogMetaMapper.class,
-                  mapper -> 
mapper.softDeleteCatalogMetasByCatalogId(catalogId)),
+              deleteCount.set(
+                  SessionUtils.getWithoutCommit(
+                      CatalogMetaMapper.class,
+                      mapper ->
+                          mapper.softDeleteCatalogMetasByCatalogId(catalogId, 
currentVersion))),
           () ->
               SessionUtils.doWithoutCommit(
                   OwnerMetaMapper.class,
@@ -412,7 +419,8 @@ public class CatalogMetaService {
           });
     }
 
-    return true;
+    // OCC: false when the catalog's version changed between read and delete.
+    return deleteCount.get() > 0;
   }
 
   @Monitored(
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/service/MetalakeMetaService.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/service/MetalakeMetaService.java
index 607dab00c1..e6e71dc399 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/service/MetalakeMetaService.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/service/MetalakeMetaService.java
@@ -218,14 +218,21 @@ public class MetalakeMetaService {
       baseMetricName = "deleteMetalake")
   public boolean deleteMetalake(NameIdentifier ident, boolean cascade) {
     NameIdentifierUtil.checkMetalake(ident);
-    Long metalakeId = getMetalakeIdByName(ident.name());
-    if (metalakeId != null) {
+    MetalakePO metalakePO =
+        SessionUtils.getWithoutCommit(
+            MetalakeMetaMapper.class, mapper -> 
mapper.selectMetalakeMetaByName(ident.name()));
+    AtomicInteger deleteCount = new AtomicInteger(0);
+    if (metalakePO != null) {
+      Long metalakeId = metalakePO.getMetalakeId();
+      Long currentVersion = metalakePO.getCurrentVersion();
       if (cascade) {
         SessionUtils.doMultipleWithCommit(
             () ->
-                SessionUtils.doWithoutCommit(
-                    MetalakeMetaMapper.class,
-                    mapper -> 
mapper.softDeleteMetalakeMetaByMetalakeId(metalakeId)),
+                deleteCount.set(
+                    SessionUtils.getWithoutCommit(
+                        MetalakeMetaMapper.class,
+                        mapper ->
+                            
mapper.softDeleteMetalakeMetaByMetalakeId(metalakeId, currentVersion))),
             () ->
                 SessionUtils.doWithoutCommit(
                     CatalogMetaMapper.class,
@@ -354,9 +361,11 @@ public class MetalakeMetaService {
         }
         SessionUtils.doMultipleWithCommit(
             () ->
-                SessionUtils.doWithoutCommit(
-                    MetalakeMetaMapper.class,
-                    mapper -> 
mapper.softDeleteMetalakeMetaByMetalakeId(metalakeId)),
+                deleteCount.set(
+                    SessionUtils.getWithoutCommit(
+                        MetalakeMetaMapper.class,
+                        mapper ->
+                            
mapper.softDeleteMetalakeMetaByMetalakeId(metalakeId, currentVersion))),
             () ->
                 SessionUtils.doWithoutCommit(
                     UserRoleRelMapper.class,
@@ -417,7 +426,8 @@ public class MetalakeMetaService {
             });
       }
     }
-    return true;
+    // OCC: false when the metalake was absent, or its version changed between 
read and delete.
+    return deleteCount.get() > 0;
   }
 
   @Monitored(

Reply via email to