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(
