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 897af74faf0ed063c4d1fe761bf4ed4d31def2ff Author: yuqi <[email protected]> AuthorDate: Fri Jul 24 17:31:13 2026 +0800 [#12166] improvement(core): version-checked soft-delete (drop CAS) for schema (non-cascade) Add softDeleteSchemaMetaBySchemaIdAndVersion (base + PostgreSQL) and use it for the non-cascade single-schema delete, which reads schemaPO and now passes its current_version, returning based on the row count (Option A). The hierarchical cascade path deletes multiple descendant schemas via a name-batch, so it stays unversioned (a single-version CAS does not apply there). Part of #12166 (drop CAS). --- .../storage/relational/mapper/SchemaMetaMapper.java | 6 ++++++ .../relational/mapper/SchemaMetaSQLProviderFactory.java | 5 +++++ .../mapper/provider/base/SchemaMetaBaseSQLProvider.java | 12 ++++++++++++ .../provider/postgresql/SchemaMetaPostgreSQLProvider.java | 10 ++++++++++ .../storage/relational/service/SchemaMetaService.java | 13 +++++++++---- 5 files changed, 42 insertions(+), 4 deletions(-) diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/SchemaMetaMapper.java b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/SchemaMetaMapper.java index 1c9b5286b2..e8310e977c 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/SchemaMetaMapper.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/SchemaMetaMapper.java @@ -107,6 +107,12 @@ public interface SchemaMetaMapper { method = "softDeleteSchemaMetasBySchemaIds") Integer softDeleteSchemaMetasBySchemaIds(@Param("schemaIds") List<Long> schemaIds); + @UpdateProvider( + type = SchemaMetaSQLProviderFactory.class, + method = "softDeleteSchemaMetaBySchemaIdAndVersion") + Integer softDeleteSchemaMetaBySchemaIdAndVersion( + @Param("schemaId") Long schemaId, @Param("currentVersion") Long currentVersion); + @UpdateProvider( type = SchemaMetaSQLProviderFactory.class, method = "softDeleteSchemaMetasByMetalakeId") diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/SchemaMetaSQLProviderFactory.java b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/SchemaMetaSQLProviderFactory.java index acc2717026..c65b23ad14 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/SchemaMetaSQLProviderFactory.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/SchemaMetaSQLProviderFactory.java @@ -120,6 +120,11 @@ public class SchemaMetaSQLProviderFactory { return getProvider().softDeleteSchemaMetasBySchemaIds(schemaIds); } + public static String softDeleteSchemaMetaBySchemaIdAndVersion( + @Param("schemaId") Long schemaId, @Param("currentVersion") Long currentVersion) { + return getProvider().softDeleteSchemaMetaBySchemaIdAndVersion(schemaId, currentVersion); + } + public static String softDeleteSchemaMetasByMetalakeId(@Param("metalakeId") Long metalakeId) { return getProvider().softDeleteSchemaMetasByMetalakeId(metalakeId); } diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/SchemaMetaBaseSQLProvider.java b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/SchemaMetaBaseSQLProvider.java index 0a2f98b178..232eb5cf22 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/SchemaMetaBaseSQLProvider.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/SchemaMetaBaseSQLProvider.java @@ -292,6 +292,18 @@ public class SchemaMetaBaseSQLProvider { + "</script>"; } + public String softDeleteSchemaMetaBySchemaIdAndVersion( + @Param("schemaId") Long schemaId, @Param("currentVersion") Long currentVersion) { + return "UPDATE " + + TABLE_NAME + + " SET deleted_at = (UNIX_TIMESTAMP() * 1000.0)" + + " + EXTRACT(MICROSECOND FROM CURRENT_TIMESTAMP(3)) / 1000" + // OCC: version-checked delete for the non-cascade single-schema drop (0 rows = stale + // version; the service returns false). The hierarchical cascade batch stays unversioned. + + " WHERE schema_id = #{schemaId} AND current_version = #{currentVersion}" + + " AND deleted_at = 0"; + } + public String softDeleteSchemaMetasByMetalakeId(@Param("metalakeId") Long metalakeId) { return "UPDATE " + TABLE_NAME diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/SchemaMetaPostgreSQLProvider.java b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/SchemaMetaPostgreSQLProvider.java index c55ca530ca..fb6344a357 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/SchemaMetaPostgreSQLProvider.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/SchemaMetaPostgreSQLProvider.java @@ -117,6 +117,16 @@ public class SchemaMetaPostgreSQLProvider extends SchemaMetaBaseSQLProvider { + "</script>"; } + @Override + public String softDeleteSchemaMetaBySchemaIdAndVersion(Long schemaId, Long currentVersion) { + return "UPDATE " + + TABLE_NAME + + " SET deleted_at = CAST(EXTRACT(EPOCH FROM CURRENT_TIMESTAMP) * 1000 AS BIGINT)" + // OCC: version-checked delete for the non-cascade single-schema drop. + + " WHERE schema_id = #{schemaId} AND current_version = #{currentVersion}" + + " AND deleted_at = 0"; + } + @Override public String softDeleteSchemaMetasByMetalakeId(Long metalakeId) { return "UPDATE " diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/service/SchemaMetaService.java b/core/src/main/java/org/apache/gravitino/storage/relational/service/SchemaMetaService.java index 39d820b42d..82f464447b 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/service/SchemaMetaService.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/service/SchemaMetaService.java @@ -411,12 +411,15 @@ public class SchemaMetaService { "Entity %s has sub-entities, you should remove sub-entities first", identifier); } - List<Long> singleSchemaId = Collections.singletonList(schemaId); + int[] schemaDeletedCount = new int[] {0}; SessionUtils.doMultipleWithCommit( () -> - SessionUtils.doWithoutCommit( - SchemaMetaMapper.class, - mapper -> mapper.softDeleteSchemaMetasBySchemaIds(singleSchemaId)), + schemaDeletedCount[0] = + SessionUtils.getWithoutCommit( + SchemaMetaMapper.class, + mapper -> + mapper.softDeleteSchemaMetaBySchemaIdAndVersion( + schemaId, schemaPO.getCurrentVersion())), () -> SessionUtils.doWithoutCommit( OwnerMetaMapper.class, @@ -455,6 +458,8 @@ public class SchemaMetaService { schemaFullName, OperateType.DROP)); }); + // OCC: false when the schema's version changed between read and delete. + return schemaDeletedCount[0] > 0; } return true; }
