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(
