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 feea8f4c45a86ca59bd17042c6fbc618fc21d990 Author: yuqi <[email protected]> AuthorDate: Fri Jul 24 10:08:36 2026 +0800 [#12166] improvement(core): version-checked soft-delete (drop CAS) for topic Add a current_version parameter to softDeleteTopicMetasByTopicId and check it in the WHERE (base + PostgreSQL), so a drop carrying a stale version matches 0 rows instead of deleting whatever live row exists. deleteTopic passes the version it read; the mapper returns 0/1. DB-level test: stale version -> 0 rows and topic stays; current version -> 1 row and topic gone. Reference slice for the drop-CAS pass (PR2 of #12166). Cascade paths (by schema/ catalog id) are unchanged since a parent drop is authoritative. --- .../storage/relational/mapper/TopicMetaMapper.java | 3 ++- .../mapper/TopicMetaSQLProviderFactory.java | 5 +++-- .../provider/base/TopicMetaBaseSQLProvider.java | 8 +++++-- .../postgresql/TopicMetaPostgreSQLProvider.java | 6 ++++-- .../relational/service/TopicMetaService.java | 4 +++- .../relational/service/TestTopicMetaService.java | 25 ++++++++++++++++++++++ 6 files changed, 43 insertions(+), 8 deletions(-) diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TopicMetaMapper.java b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TopicMetaMapper.java index fd447015d8..f9ce3773ad 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TopicMetaMapper.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TopicMetaMapper.java @@ -82,7 +82,8 @@ public interface TopicMetaMapper { @UpdateProvider( type = TopicMetaSQLProviderFactory.class, method = "softDeleteTopicMetasByTopicId") - Integer softDeleteTopicMetasByTopicId(@Param("topicId") Long topicId); + Integer softDeleteTopicMetasByTopicId( + @Param("topicId") Long topicId, @Param("currentVersion") Long currentVersion); @UpdateProvider( type = TopicMetaSQLProviderFactory.class, diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TopicMetaSQLProviderFactory.java b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TopicMetaSQLProviderFactory.java index 32d6c1b52e..935ca804eb 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TopicMetaSQLProviderFactory.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/TopicMetaSQLProviderFactory.java @@ -103,8 +103,9 @@ public class TopicMetaSQLProviderFactory { return getProvider().selectTopicIdBySchemaIdAndName(schemaId, name); } - public static String softDeleteTopicMetasByTopicId(@Param("topicId") Long topicId) { - return getProvider().softDeleteTopicMetasByTopicId(topicId); + public static String softDeleteTopicMetasByTopicId( + @Param("topicId") Long topicId, @Param("currentVersion") Long currentVersion) { + return getProvider().softDeleteTopicMetasByTopicId(topicId, currentVersion); } public static String softDeleteTopicMetasByCatalogId(@Param("catalogId") Long catalogId) { diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TopicMetaBaseSQLProvider.java b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TopicMetaBaseSQLProvider.java index 375f28427b..32a4a8b28f 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TopicMetaBaseSQLProvider.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TopicMetaBaseSQLProvider.java @@ -249,12 +249,16 @@ public class TopicMetaBaseSQLProvider { + " AND deleted_at = 0"; } - public String softDeleteTopicMetasByTopicId(@Param("topicId") Long topicId) { + public String softDeleteTopicMetasByTopicId( + @Param("topicId") Long topicId, @Param("currentVersion") Long currentVersion) { return "UPDATE " + TABLE_NAME + " SET deleted_at = (UNIX_TIMESTAMP() * 1000.0)" + " + EXTRACT(MICROSECOND FROM CURRENT_TIMESTAMP(3)) / 1000" - + " WHERE topic_id = #{topicId} AND deleted_at = 0"; + // OCC: version-checked delete. 0 rows means the caller's version is stale (the row was + // updated or already deleted concurrently); the service re-reads to tell the two apart. + + " WHERE topic_id = #{topicId} AND current_version = #{currentVersion}" + + " AND deleted_at = 0"; } public String softDeleteTopicMetasByCatalogId(@Param("catalogId") Long catalogId) { diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TopicMetaPostgreSQLProvider.java b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TopicMetaPostgreSQLProvider.java index d9ca98bbfc..eec34eaefb 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TopicMetaPostgreSQLProvider.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TopicMetaPostgreSQLProvider.java @@ -49,11 +49,13 @@ public class TopicMetaPostgreSQLProvider extends TopicMetaBaseSQLProvider { } @Override - public String softDeleteTopicMetasByTopicId(Long topicId) { + public String softDeleteTopicMetasByTopicId(Long topicId, Long currentVersion) { return "UPDATE " + TABLE_NAME + " SET deleted_at = CAST(EXTRACT(EPOCH FROM CURRENT_TIMESTAMP) * 1000 AS BIGINT)" - + " WHERE topic_id = #{topicId} AND deleted_at = 0"; + // OCC: version-checked delete (see the base provider for rationale). + + " WHERE topic_id = #{topicId} AND current_version = #{currentVersion}" + + " AND deleted_at = 0"; } @Override diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/service/TopicMetaService.java b/core/src/main/java/org/apache/gravitino/storage/relational/service/TopicMetaService.java index 85bc2dbbe0..6ab7e9458e 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/service/TopicMetaService.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/service/TopicMetaService.java @@ -306,7 +306,9 @@ public class TopicMetaService { deleteResult.set( SessionUtils.getWithoutCommit( TopicMetaMapper.class, - mapper -> mapper.softDeleteTopicMetasByTopicId(topicId))), + mapper -> + mapper.softDeleteTopicMetasByTopicId( + topicId, topicPO.getCurrentVersion()))), () -> { if (deleteResult.get() > 0) { SessionUtils.doWithoutCommit( diff --git a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestTopicMetaService.java b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestTopicMetaService.java index d12b18b841..fc156dd410 100644 --- a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestTopicMetaService.java +++ b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestTopicMetaService.java @@ -96,6 +96,31 @@ public class TestTopicMetaService extends TestJDBCBackend { } } + @TestTemplate + public void testSoftDeleteTopicIsVersionCas() throws IOException { + TopicEntity topic = + createTopicEntity( + RandomIdGenerator.INSTANCE.nextId(), + NamespaceUtil.ofTopic(metalakeName, catalogName, schemaName), + "drop_cas_topic", + AUDIT_INFO); + backend.insert(topic, false); + long topicId = topic.id(); + + try (SqlSession session = + SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true)) { + TopicMetaMapper mapper = session.getMapper(TopicMetaMapper.class); + + // A drop carrying a stale version matches 0 rows and leaves the topic live. + Assertions.assertEquals(0, mapper.softDeleteTopicMetasByTopicId(topicId, 999L)); + assertTrue(backend.exists(topic.nameIdentifier(), Entity.EntityType.TOPIC)); + + // A drop carrying the current version (1) soft-deletes exactly one row. + Assertions.assertEquals(1, mapper.softDeleteTopicMetasByTopicId(topicId, 1L)); + assertFalse(backend.exists(topic.nameIdentifier(), Entity.EntityType.TOPIC)); + } + } + @TestTemplate public void testInsertAlreadyExistsException() throws IOException { TopicEntity topic =
