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 52a5727d49b127a417e9d73a54bf4dbb5af3c97b
Author: yuqi <[email protected]>
AuthorDate: Fri Jul 24 14:24:00 2026 +0800

    [#12166] improvement(core): version-checked soft-delete (drop CAS) for 
table and fileset
    
    Same pattern as topic: add a current_version param to 
softDeleteTableMetasByTableId
    and softDeleteFilesetMetasByFilesetId (base + PostgreSQL), check it in the 
WHERE, and
    have deleteTable/deleteFileset pass the version they read. Cascade paths 
unchanged.
    
    Part of #12166 (drop CAS).
---
 .../gravitino/storage/relational/mapper/FilesetMetaMapper.java     | 3 ++-
 .../storage/relational/mapper/FilesetMetaSQLProviderFactory.java   | 5 +++--
 .../gravitino/storage/relational/mapper/TableMetaMapper.java       | 3 ++-
 .../storage/relational/mapper/TableMetaSQLProviderFactory.java     | 5 +++--
 .../mapper/provider/base/FilesetMetaBaseSQLProvider.java           | 7 +++++--
 .../relational/mapper/provider/base/TableMetaBaseSQLProvider.java  | 7 +++++--
 .../mapper/provider/postgresql/FilesetMetaPostgreSQLProvider.java  | 6 ++++--
 .../mapper/provider/postgresql/TableMetaPostgreSQLProvider.java    | 6 ++++--
 .../gravitino/storage/relational/service/FilesetMetaService.java   | 4 +++-
 .../gravitino/storage/relational/service/TableMetaService.java     | 4 +++-
 10 files changed, 34 insertions(+), 16 deletions(-)

diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/FilesetMetaMapper.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/FilesetMetaMapper.java
index fcbbc66c0b..fa6c6bc795 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/FilesetMetaMapper.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/FilesetMetaMapper.java
@@ -246,7 +246,8 @@ public interface FilesetMetaMapper {
   @UpdateProvider(
       type = FilesetMetaSQLProviderFactory.class,
       method = "softDeleteFilesetMetasByFilesetId")
-  Integer softDeleteFilesetMetasByFilesetId(@Param("filesetId") Long 
filesetId);
+  Integer softDeleteFilesetMetasByFilesetId(
+      @Param("filesetId") Long filesetId, @Param("currentVersion") Long 
currentVersion);
 
   @DeleteProvider(
       type = FilesetMetaSQLProviderFactory.class,
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/FilesetMetaSQLProviderFactory.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/FilesetMetaSQLProviderFactory.java
index 07aaebb1cb..74d4fe9d69 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/FilesetMetaSQLProviderFactory.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/FilesetMetaSQLProviderFactory.java
@@ -117,8 +117,9 @@ public class FilesetMetaSQLProviderFactory {
     return getProvider().softDeleteFilesetMetasBySchemaIds(schemaIds);
   }
 
-  public String softDeleteFilesetMetasByFilesetId(@Param("filesetId") Long 
filesetId) {
-    return getProvider().softDeleteFilesetMetasByFilesetId(filesetId);
+  public String softDeleteFilesetMetasByFilesetId(
+      @Param("filesetId") Long filesetId, @Param("currentVersion") Long 
currentVersion) {
+    return getProvider().softDeleteFilesetMetasByFilesetId(filesetId, 
currentVersion);
   }
 
   public String deleteFilesetMetasByLegacyTimeline(
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TableMetaMapper.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TableMetaMapper.java
index acb682916d..52207794ee 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TableMetaMapper.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TableMetaMapper.java
@@ -90,7 +90,8 @@ public interface TableMetaMapper {
   @UpdateProvider(
       type = TableMetaSQLProviderFactory.class,
       method = "softDeleteTableMetasByTableId")
-  Integer softDeleteTableMetasByTableId(@Param("tableId") Long tableId);
+  Integer softDeleteTableMetasByTableId(
+      @Param("tableId") Long tableId, @Param("currentVersion") Long 
currentVersion);
 
   @UpdateProvider(
       type = TableMetaSQLProviderFactory.class,
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TableMetaSQLProviderFactory.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TableMetaSQLProviderFactory.java
index c69fd9a191..bdaa485496 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TableMetaSQLProviderFactory.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TableMetaSQLProviderFactory.java
@@ -104,8 +104,9 @@ public class TableMetaSQLProviderFactory {
     return getProvider().updateTableMeta(newTablePO, oldTablePO, newSchemaId);
   }
 
-  public static String softDeleteTableMetasByTableId(@Param("tableId") Long 
tableId) {
-    return getProvider().softDeleteTableMetasByTableId(tableId);
+  public static String softDeleteTableMetasByTableId(
+      @Param("tableId") Long tableId, @Param("currentVersion") Long 
currentVersion) {
+    return getProvider().softDeleteTableMetasByTableId(tableId, 
currentVersion);
   }
 
   public static String softDeleteTableMetasByMetalakeId(@Param("metalakeId") 
Long metalakeId) {
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/FilesetMetaBaseSQLProvider.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/FilesetMetaBaseSQLProvider.java
index 64f8de39be..dc4b491f74 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/FilesetMetaBaseSQLProvider.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/FilesetMetaBaseSQLProvider.java
@@ -332,12 +332,15 @@ public class FilesetMetaBaseSQLProvider {
         + "</script>";
   }
 
-  public String softDeleteFilesetMetasByFilesetId(@Param("filesetId") Long 
filesetId) {
+  public String softDeleteFilesetMetasByFilesetId(
+      @Param("filesetId") Long filesetId, @Param("currentVersion") Long 
currentVersion) {
     return "UPDATE "
         + META_TABLE_NAME
         + " SET deleted_at = (UNIX_TIMESTAMP() * 1000.0)"
         + " + EXTRACT(MICROSECOND FROM CURRENT_TIMESTAMP(3)) / 1000"
-        + " WHERE fileset_id = #{filesetId} AND deleted_at = 0";
+        // OCC: version-checked delete (0 rows = stale version; the service 
returns false).
+        + " WHERE fileset_id = #{filesetId} AND current_version = 
#{currentVersion}"
+        + " AND deleted_at = 0";
   }
 
   public String deleteFilesetMetasByLegacyTimeline(
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TableMetaBaseSQLProvider.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TableMetaBaseSQLProvider.java
index 7b55f0967f..09d91185d0 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TableMetaBaseSQLProvider.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TableMetaBaseSQLProvider.java
@@ -245,12 +245,15 @@ public class TableMetaBaseSQLProvider {
         + " AND deleted_at = 0";
   }
 
-  public String softDeleteTableMetasByTableId(@Param("tableId") Long tableId) {
+  public String softDeleteTableMetasByTableId(
+      @Param("tableId") Long tableId, @Param("currentVersion") Long 
currentVersion) {
     return "UPDATE "
         + TABLE_NAME
         + " SET deleted_at = (UNIX_TIMESTAMP() * 1000.0)"
         + " + EXTRACT(MICROSECOND FROM CURRENT_TIMESTAMP(3)) / 1000"
-        + " WHERE table_id = #{tableId} AND deleted_at = 0";
+        // OCC: version-checked delete (0 rows = stale version; the service 
returns false).
+        + " WHERE table_id = #{tableId} AND current_version = 
#{currentVersion}"
+        + " AND deleted_at = 0";
   }
 
   public String softDeleteTableMetasByMetalakeId(@Param("metalakeId") Long 
metalakeId) {
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/FilesetMetaPostgreSQLProvider.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/FilesetMetaPostgreSQLProvider.java
index ebb4d1731f..0487d8e785 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/FilesetMetaPostgreSQLProvider.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/FilesetMetaPostgreSQLProvider.java
@@ -57,11 +57,13 @@ public class FilesetMetaPostgreSQLProvider extends 
FilesetMetaBaseSQLProvider {
   }
 
   @Override
-  public String softDeleteFilesetMetasByFilesetId(Long filesetId) {
+  public String softDeleteFilesetMetasByFilesetId(Long filesetId, Long 
currentVersion) {
     return "UPDATE "
         + META_TABLE_NAME
         + " SET deleted_at = CAST(EXTRACT(EPOCH FROM CURRENT_TIMESTAMP) * 1000 
AS BIGINT)"
-        + " WHERE fileset_id = #{filesetId} AND deleted_at = 0";
+        // OCC: version-checked delete (see the base provider for rationale).
+        + " WHERE fileset_id = #{filesetId} AND current_version = 
#{currentVersion}"
+        + " AND deleted_at = 0";
   }
 
   @Override
diff --git 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TableMetaPostgreSQLProvider.java
 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TableMetaPostgreSQLProvider.java
index 7add18edac..bebd522cdc 100644
--- 
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TableMetaPostgreSQLProvider.java
+++ 
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TableMetaPostgreSQLProvider.java
@@ -56,11 +56,13 @@ public class TableMetaPostgreSQLProvider extends 
TableMetaBaseSQLProvider {
   }
 
   @Override
-  public String softDeleteTableMetasByTableId(Long tableId) {
+  public String softDeleteTableMetasByTableId(Long tableId, Long 
currentVersion) {
     return "UPDATE "
         + TABLE_NAME
         + " SET deleted_at = CAST(EXTRACT(EPOCH FROM CURRENT_TIMESTAMP) * 1000 
AS BIGINT)"
-        + " WHERE table_id = #{tableId} AND deleted_at = 0";
+        // OCC: version-checked delete (see the base provider for rationale).
+        + " WHERE table_id = #{tableId} AND current_version = 
#{currentVersion}"
+        + " AND deleted_at = 0";
   }
 
   @Override
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 9c1f677b4b..3a01def13b 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
@@ -326,7 +326,9 @@ public class FilesetMetaService {
             deleteResult.set(
                 SessionUtils.getWithoutCommit(
                     FilesetMetaMapper.class,
-                    mapper -> 
mapper.softDeleteFilesetMetasByFilesetId(filesetId))),
+                    mapper ->
+                        mapper.softDeleteFilesetMetasByFilesetId(
+                            filesetId, filesetPO.getCurrentVersion()))),
         () -> {
           if (deleteResult.get() > 0) {
             SessionUtils.doWithoutCommit(
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 b23f6e9f49..0d10b3ce80 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
@@ -269,7 +269,9 @@ public class TableMetaService {
             deleteResult.set(
                 SessionUtils.getWithoutCommit(
                     TableMetaMapper.class,
-                    mapper -> 
mapper.softDeleteTableMetasByTableId(tablePO.getTableId()))),
+                    mapper ->
+                        mapper.softDeleteTableMetasByTableId(
+                            tablePO.getTableId(), 
tablePO.getCurrentVersion()))),
         () -> {
           if (deleteResult.get() > 0) {
             SessionUtils.doWithoutCommit(

Reply via email to