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 76140de5b453527fbc26284fab82ddf53f4db4c4 Author: yuqi <[email protected]> AuthorDate: Fri Jul 24 15:30:25 2026 +0800 [#12166] improvement(core): version-checked soft-delete (drop CAS) for view, function, model view/function/model already read the full PO on delete, so pass its version to the soft-delete and check it in the WHERE (base + PostgreSQL). function uses function_current_version; view threads the version through its private deleteView helper; model deletes by (schemaId, modelName) + current_version. Part of #12166 (drop CAS). --- .../gravitino/storage/relational/mapper/FunctionMetaMapper.java | 4 +++- .../relational/mapper/FunctionMetaSQLProviderFactory.java | 6 ++++-- .../gravitino/storage/relational/mapper/ModelMetaMapper.java | 4 +++- .../storage/relational/mapper/ModelMetaSQLProviderFactory.java | 7 +++++-- .../gravitino/storage/relational/mapper/ViewMetaMapper.java | 3 ++- .../storage/relational/mapper/ViewMetaSQLProviderFactory.java | 5 +++-- .../mapper/provider/base/FunctionMetaBaseSQLProvider.java | 9 +++++++-- .../mapper/provider/base/ModelMetaBaseSQLProvider.java | 8 ++++++-- .../relational/mapper/provider/base/ViewMetaBaseSQLProvider.java | 7 +++++-- .../provider/postgresql/FunctionMetaPostgreSQLProvider.java | 9 +++++++-- .../mapper/provider/postgresql/ModelMetaPostgreSQLProvider.java | 8 ++++++-- .../mapper/provider/postgresql/ViewMetaPostgreSQLProvider.java | 7 +++++-- .../storage/relational/service/FunctionMetaService.java | 4 +++- .../gravitino/storage/relational/service/ModelMetaService.java | 3 ++- .../gravitino/storage/relational/service/ViewMetaService.java | 8 +++++--- 15 files changed, 66 insertions(+), 26 deletions(-) diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/FunctionMetaMapper.java b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/FunctionMetaMapper.java index 100c2683c3..753aad3a43 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/FunctionMetaMapper.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/FunctionMetaMapper.java @@ -151,7 +151,9 @@ public interface FunctionMetaMapper { @UpdateProvider( type = FunctionMetaSQLProviderFactory.class, method = "softDeleteFunctionMetaByFunctionId") - Integer softDeleteFunctionMetaByFunctionId(@Param("functionId") Long functionId); + Integer softDeleteFunctionMetaByFunctionId( + @Param("functionId") Long functionId, + @Param("functionCurrentVersion") Integer functionCurrentVersion); @UpdateProvider( type = FunctionMetaSQLProviderFactory.class, diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/FunctionMetaSQLProviderFactory.java b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/FunctionMetaSQLProviderFactory.java index 2e2ad59129..c3babbb9be 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/FunctionMetaSQLProviderFactory.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/FunctionMetaSQLProviderFactory.java @@ -95,8 +95,10 @@ public class FunctionMetaSQLProviderFactory { return getProvider().selectFunctionIdBySchemaIdAndFunctionName(schemaId, functionName); } - public static String softDeleteFunctionMetaByFunctionId(@Param("functionId") Long functionId) { - return getProvider().softDeleteFunctionMetaByFunctionId(functionId); + public static String softDeleteFunctionMetaByFunctionId( + @Param("functionId") Long functionId, + @Param("functionCurrentVersion") Integer functionCurrentVersion) { + return getProvider().softDeleteFunctionMetaByFunctionId(functionId, functionCurrentVersion); } public static String softDeleteFunctionMetasByCatalogId(@Param("catalogId") Long catalogId) { diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ModelMetaMapper.java b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ModelMetaMapper.java index 9fc338ad98..579bfc1465 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ModelMetaMapper.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ModelMetaMapper.java @@ -79,7 +79,9 @@ public interface ModelMetaMapper { type = ModelMetaSQLProviderFactory.class, method = "softDeleteModelMetaBySchemaIdAndModelName") Integer softDeleteModelMetaBySchemaIdAndModelName( - @Param("schemaId") Long schemaId, @Param("modelName") String modelName); + @Param("schemaId") Long schemaId, + @Param("modelName") String modelName, + @Param("currentVersion") Long currentVersion); @UpdateProvider( type = ModelMetaSQLProviderFactory.class, diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ModelMetaSQLProviderFactory.java b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ModelMetaSQLProviderFactory.java index d16a5c485b..014b2089ce 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ModelMetaSQLProviderFactory.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ModelMetaSQLProviderFactory.java @@ -98,8 +98,11 @@ public class ModelMetaSQLProviderFactory { } public static String softDeleteModelMetaBySchemaIdAndModelName( - @Param("schemaId") Long schemaId, @Param("modelName") String modelName) { - return getProvider().softDeleteModelMetaBySchemaIdAndModelName(schemaId, modelName); + @Param("schemaId") Long schemaId, + @Param("modelName") String modelName, + @Param("currentVersion") Long currentVersion) { + return getProvider() + .softDeleteModelMetaBySchemaIdAndModelName(schemaId, modelName, currentVersion); } public static String softDeleteModelMetasByCatalogId(@Param("catalogId") Long catalogId) { diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ViewMetaMapper.java b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ViewMetaMapper.java index 813d646c15..ac54b2df0c 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ViewMetaMapper.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ViewMetaMapper.java @@ -128,7 +128,8 @@ public interface ViewMetaMapper { @Param("newViewMeta") ViewPO newViewPO, @Param("oldViewMeta") ViewPO oldViewPO); @UpdateProvider(type = ViewMetaSQLProviderFactory.class, method = "softDeleteViewMetasByViewId") - Integer softDeleteViewMetasByViewId(@Param("viewId") Long viewId); + Integer softDeleteViewMetasByViewId( + @Param("viewId") Long viewId, @Param("currentVersion") Long currentVersion); @UpdateProvider( type = ViewMetaSQLProviderFactory.class, diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ViewMetaSQLProviderFactory.java b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ViewMetaSQLProviderFactory.java index cb25d743d3..44af0443db 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ViewMetaSQLProviderFactory.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/ViewMetaSQLProviderFactory.java @@ -98,8 +98,9 @@ public class ViewMetaSQLProviderFactory { return getProvider().updateViewMeta(newViewPO, oldViewPO); } - public static String softDeleteViewMetasByViewId(@Param("viewId") Long viewId) { - return getProvider().softDeleteViewMetasByViewId(viewId); + public static String softDeleteViewMetasByViewId( + @Param("viewId") Long viewId, @Param("currentVersion") Long currentVersion) { + return getProvider().softDeleteViewMetasByViewId(viewId, currentVersion); } public static String softDeleteViewMetasByMetalakeId(@Param("metalakeId") Long metalakeId) { diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/FunctionMetaBaseSQLProvider.java b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/FunctionMetaBaseSQLProvider.java index ecd2d68b94..162333130b 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/FunctionMetaBaseSQLProvider.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/FunctionMetaBaseSQLProvider.java @@ -242,12 +242,17 @@ public class FunctionMetaBaseSQLProvider { + " WHERE schema_id = #{schemaId} AND function_name = #{functionName} AND deleted_at = 0"; } - public String softDeleteFunctionMetaByFunctionId(@Param("functionId") Long functionId) { + public String softDeleteFunctionMetaByFunctionId( + @Param("functionId") Long functionId, + @Param("functionCurrentVersion") Integer functionCurrentVersion) { return "UPDATE " + TABLE_NAME + " SET deleted_at = (UNIX_TIMESTAMP() * 1000.0)" + " + EXTRACT(MICROSECOND FROM CURRENT_TIMESTAMP(3)) / 1000" - + " WHERE function_id = #{functionId} AND deleted_at = 0"; + // OCC: version-checked delete (0 rows = stale version; the service returns false). + + " WHERE function_id = #{functionId}" + + " AND function_current_version = #{functionCurrentVersion}" + + " AND deleted_at = 0"; } public String softDeleteFunctionMetasByCatalogId(@Param("catalogId") Long catalogId) { diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/ModelMetaBaseSQLProvider.java b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/ModelMetaBaseSQLProvider.java index 561be20e90..b4836c2e6f 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/ModelMetaBaseSQLProvider.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/ModelMetaBaseSQLProvider.java @@ -207,12 +207,16 @@ public class ModelMetaBaseSQLProvider { } public String softDeleteModelMetaBySchemaIdAndModelName( - @Param("schemaId") Long schemaId, @Param("modelName") String modelName) { + @Param("schemaId") Long schemaId, + @Param("modelName") String modelName, + @Param("currentVersion") Long currentVersion) { return "UPDATE " + ModelMetaMapper.TABLE_NAME + " SET deleted_at = (UNIX_TIMESTAMP() * 1000.0)" + " + EXTRACT(MICROSECOND FROM CURRENT_TIMESTAMP(3)) / 1000" - + " WHERE schema_id = #{schemaId} AND model_name = #{modelName} AND deleted_at = 0"; + // OCC: version-checked delete (0 rows = stale version; the service returns false). + + " WHERE schema_id = #{schemaId} AND model_name = #{modelName}" + + " AND current_version = #{currentVersion} AND deleted_at = 0"; } public String softDeleteModelMetasByCatalogId(@Param("catalogId") Long catalogId) { diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/ViewMetaBaseSQLProvider.java b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/ViewMetaBaseSQLProvider.java index 4910a7327e..3d653bde99 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/ViewMetaBaseSQLProvider.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/ViewMetaBaseSQLProvider.java @@ -214,12 +214,15 @@ public class ViewMetaBaseSQLProvider { + "</script>"; } - public String softDeleteViewMetasByViewId(@Param("viewId") Long viewId) { + public String softDeleteViewMetasByViewId( + @Param("viewId") Long viewId, @Param("currentVersion") Long currentVersion) { return "UPDATE " + TABLE_NAME + " SET deleted_at = (UNIX_TIMESTAMP() * 1000.0)" + " + EXTRACT(MICROSECOND FROM CURRENT_TIMESTAMP(3)) / 1000" - + " WHERE view_id = #{viewId} AND deleted_at = 0"; + // OCC: version-checked delete (0 rows = stale version; the service returns false). + + " WHERE view_id = #{viewId} AND current_version = #{currentVersion}" + + " AND deleted_at = 0"; } public String softDeleteViewMetasByMetalakeId(@Param("metalakeId") Long metalakeId) { diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/FunctionMetaPostgreSQLProvider.java b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/FunctionMetaPostgreSQLProvider.java index 205dd47eb9..f9abfb2096 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/FunctionMetaPostgreSQLProvider.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/FunctionMetaPostgreSQLProvider.java @@ -101,11 +101,16 @@ public class FunctionMetaPostgreSQLProvider extends FunctionMetaBaseSQLProvider } @Override - public String softDeleteFunctionMetaByFunctionId(@Param("functionId") Long functionId) { + public String softDeleteFunctionMetaByFunctionId( + @Param("functionId") Long functionId, + @Param("functionCurrentVersion") Integer functionCurrentVersion) { return "UPDATE " + FunctionMetaMapper.TABLE_NAME + " SET deleted_at = CAST(EXTRACT(EPOCH FROM CURRENT_TIMESTAMP) * 1000 AS BIGINT)" - + " WHERE function_id = #{functionId} AND deleted_at = 0"; + // OCC: version-checked delete (see the base provider for rationale). + + " WHERE function_id = #{functionId}" + + " AND function_current_version = #{functionCurrentVersion}" + + " AND deleted_at = 0"; } @Override diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/ModelMetaPostgreSQLProvider.java b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/ModelMetaPostgreSQLProvider.java index b70cdf0324..c873fb3820 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/ModelMetaPostgreSQLProvider.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/ModelMetaPostgreSQLProvider.java @@ -50,11 +50,15 @@ public class ModelMetaPostgreSQLProvider extends ModelMetaBaseSQLProvider { @Override public String softDeleteModelMetaBySchemaIdAndModelName( - @Param("schemaId") Long schemaId, @Param("modelName") String modelName) { + @Param("schemaId") Long schemaId, + @Param("modelName") String modelName, + @Param("currentVersion") Long currentVersion) { return "UPDATE " + ModelMetaMapper.TABLE_NAME + " SET deleted_at = CAST(EXTRACT(EPOCH FROM CURRENT_TIMESTAMP) * 1000 AS BIGINT)" - + " WHERE schema_id = #{schemaId} AND model_name = #{modelName} AND deleted_at = 0"; + // OCC: version-checked delete (see the base provider for rationale). + + " WHERE schema_id = #{schemaId} AND model_name = #{modelName}" + + " AND current_version = #{currentVersion} AND deleted_at = 0"; } @Override diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/ViewMetaPostgreSQLProvider.java b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/ViewMetaPostgreSQLProvider.java index 76e80226ec..f14d3524bc 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/ViewMetaPostgreSQLProvider.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/ViewMetaPostgreSQLProvider.java @@ -94,11 +94,14 @@ public class ViewMetaPostgreSQLProvider extends ViewMetaBaseSQLProvider { } @Override - public String softDeleteViewMetasByViewId(@Param("viewId") Long viewId) { + public String softDeleteViewMetasByViewId( + @Param("viewId") Long viewId, @Param("currentVersion") Long currentVersion) { return "UPDATE " + TABLE_NAME + " SET deleted_at = CAST(EXTRACT(EPOCH FROM CURRENT_TIMESTAMP) * 1000 AS BIGINT)" - + " WHERE view_id = #{viewId} AND deleted_at = 0"; + // OCC: version-checked delete (see the base provider for rationale). + + " WHERE view_id = #{viewId} AND current_version = #{currentVersion}" + + " AND deleted_at = 0"; } @Override 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 963bb64b3c..80f13b5d59 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 @@ -149,7 +149,9 @@ public class FunctionMetaService { functionDeletedCount.set( SessionUtils.getWithoutCommit( FunctionMetaMapper.class, - mapper -> mapper.softDeleteFunctionMetaByFunctionId(functionId))), + mapper -> + mapper.softDeleteFunctionMetaByFunctionId( + functionId, functionPO.functionCurrentVersion()))), // delete function versions, owner rels, and securable object rels after meta deletion () -> { 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 d492371024..3d17ab4e3a 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 @@ -154,7 +154,8 @@ public class ModelMetaService { SessionUtils.getWithoutCommit( ModelMetaMapper.class, mapper -> - mapper.softDeleteModelMetaBySchemaIdAndModelName(schemaId, ident.name()))), + mapper.softDeleteModelMetaBySchemaIdAndModelName( + schemaId, ident.name(), modelPO.getCurrentVersion()))), () -> SessionUtils.doWithoutCommit( OwnerMetaMapper.class, 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 a305e05645..794cfb9cd4 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 @@ -200,7 +200,7 @@ public class ViewMetaService { String viewFullName = NameIdentifierUtil.ofView(metalakeName, catalogName, schemaName, viewPO.getViewName()) .toString(); - return deleteView(viewPO.getViewId(), metalakeName, viewFullName); + return deleteView(viewPO.getViewId(), viewPO.getCurrentVersion(), metalakeName, viewFullName); } @Monitored( @@ -224,13 +224,15 @@ public class ViewMetaService { return ops; } - private boolean deleteView(Long viewId, String metalakeName, String viewFullName) { + private boolean deleteView( + Long viewId, Long currentVersion, String metalakeName, String viewFullName) { AtomicInteger deleteResult = new AtomicInteger(0); SessionUtils.doMultipleWithCommit( () -> deleteResult.set( SessionUtils.getWithoutCommit( - ViewMetaMapper.class, mapper -> mapper.softDeleteViewMetasByViewId(viewId))), + ViewMetaMapper.class, + mapper -> mapper.softDeleteViewMetasByViewId(viewId, currentVersion))), () -> { if (deleteResult.get() > 0) { SessionUtils.doWithoutCommit(
