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 b92ead1eea7897255ae3c9aca6d4c684266f005a Author: yuqi <[email protected]> AuthorDate: Thu Jul 23 22:40:41 2026 +0800 [#12166] improvement(core): slim topic UPDATE WHERE to version CAS Reduce updateTopicMeta WHERE from a full-row compare to id + current_version + deleted_at (base + PostgreSQL providers). With current_version now monotonic, the version check is a real compare-and-set; the old field-by-field / JSON-byte match was redundant and fragile. Add a DB-level CAS test: current version -> 1 row and version advances; stale version -> 0 rows. Reference slice for the per-entity WHERE-slim pass. Part of #12166. --- .../provider/base/TopicMetaBaseSQLProvider.java | 12 ++---- .../postgresql/TopicMetaPostgreSQLProvider.java | 11 +----- .../relational/service/TestTopicMetaService.java | 44 +++++++++++++++++++++- 3 files changed, 47 insertions(+), 20 deletions(-) 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 8d508b351b..375f28427b 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 @@ -233,17 +233,11 @@ public class TopicMetaBaseSQLProvider { + " current_version = #{newTopicMeta.currentVersion}," + " last_version = #{newTopicMeta.lastVersion}," + " deleted_at = #{newTopicMeta.deletedAt}" + // OCC: compare-and-set on the version alone. Because a successful update always raises + // current_version, matching id + current_version uniquely identifies the row the caller + // read; the old full-row comparison was redundant and fragile (JSON byte equality). + " WHERE topic_id = #{oldTopicMeta.topicId}" - + " AND topic_name = #{oldTopicMeta.topicName}" - + " AND metalake_id = #{oldTopicMeta.metalakeId}" - + " AND catalog_id = #{oldTopicMeta.catalogId}" - + " AND schema_id = #{oldTopicMeta.schemaId}" - + " AND (comment = #{oldTopicMeta.comment}" - + " OR (comment IS NULL and #{oldTopicMeta.comment} IS NULL))" - + " AND properties = #{oldTopicMeta.properties}" - + " AND audit_info = #{oldTopicMeta.auditInfo}" + " AND current_version = #{oldTopicMeta.currentVersion}" - + " AND last_version = #{oldTopicMeta.lastVersion}" + " AND deleted_at = 0"; } 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 711c951934..d9ca98bbfc 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 @@ -42,18 +42,9 @@ public class TopicMetaPostgreSQLProvider extends TopicMetaBaseSQLProvider { + " current_version = #{newTopicMeta.currentVersion}," + " last_version = #{newTopicMeta.lastVersion}," + " deleted_at = #{newTopicMeta.deletedAt}" + // OCC: compare-and-set on the version alone (see the base provider for rationale). + " WHERE topic_id = #{oldTopicMeta.topicId}" - + " AND topic_name = #{oldTopicMeta.topicName}" - + " AND metalake_id = #{oldTopicMeta.metalakeId}" - + " AND catalog_id = #{oldTopicMeta.catalogId}" - + " AND schema_id = #{oldTopicMeta.schemaId}" - + " AND (comment = #{oldTopicMeta.comment}" - + " OR (CAST(comment AS VARCHAR) IS NULL" - + " AND CAST(#{oldTopicMeta.comment} AS VARCHAR) IS NULL))" - + " AND properties = #{oldTopicMeta.properties}" - + " AND audit_info = #{oldTopicMeta.auditInfo}" + " AND current_version = #{oldTopicMeta.currentVersion}" - + " AND last_version = #{oldTopicMeta.lastVersion}" + " AND deleted_at = 0"; } 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 a8a3c6fc1f..d12b18b841 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 @@ -33,11 +33,17 @@ import org.apache.gravitino.EntityAlreadyExistsException; import org.apache.gravitino.NameIdentifier; import org.apache.gravitino.Namespace; import org.apache.gravitino.exceptions.NoSuchEntityException; +import org.apache.gravitino.meta.SchemaEntity; import org.apache.gravitino.meta.TopicEntity; import org.apache.gravitino.storage.RandomIdGenerator; import org.apache.gravitino.storage.relational.TestJDBCBackend; +import org.apache.gravitino.storage.relational.mapper.TopicMetaMapper; +import org.apache.gravitino.storage.relational.po.TopicPO; +import org.apache.gravitino.storage.relational.session.SqlSessionFactoryHelper; +import org.apache.gravitino.storage.relational.utils.POConverters; import org.apache.gravitino.utils.NameIdentifierUtil; import org.apache.gravitino.utils.NamespaceUtil; +import org.apache.ibatis.session.SqlSession; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.TestTemplate; @@ -46,12 +52,48 @@ public class TestTopicMetaService extends TestJDBCBackend { private final String metalakeName = "metalake_for_topic_test"; private final String catalogName = "catalog_for_topic_test"; private final String schemaName = "schema_for_topic_test"; + private long schemaId; @BeforeEach public void prepare() throws IOException { createAndInsertMakeLake(metalakeName); createAndInsertCatalog(metalakeName, catalogName); - createAndInsertSchema(metalakeName, catalogName, schemaName); + SchemaEntity schema = createAndInsertSchema(metalakeName, catalogName, schemaName); + schemaId = schema.id(); + } + + @TestTemplate + public void testUpdateTopicMetaIsVersionCas() throws IOException { + TopicEntity topic = + createTopicEntity( + RandomIdGenerator.INSTANCE.nextId(), + NamespaceUtil.ofTopic(metalakeName, catalogName, schemaName), + "cas_topic", + AUDIT_INFO); + backend.insert(topic, false); + + TopicEntity updated = + createTopicEntity( + topic.id(), + NamespaceUtil.ofTopic(metalakeName, catalogName, schemaName), + "cas_topic", + AUDIT_INFO); + + try (SqlSession session = + SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true)) { + TopicMetaMapper mapper = session.getMapper(TopicMetaMapper.class); + TopicPO oldPO = mapper.selectTopicMetaBySchemaIdAndName(schemaId, "cas_topic"); + Assertions.assertEquals(1L, oldPO.getCurrentVersion()); + + // A write against the current version succeeds and advances the version (1 -> 2). + TopicPO newPO = POConverters.updateTopicPOWithVersion(oldPO, updated); + Assertions.assertEquals(2L, newPO.getCurrentVersion()); + Assertions.assertEquals(1, mapper.updateTopicMeta(newPO, oldPO)); + + // A second write reusing the now-stale snapshot (version 1) matches 0 rows: version CAS lost. + TopicPO staleNewPO = POConverters.updateTopicPOWithVersion(oldPO, updated); + Assertions.assertEquals(0, mapper.updateTopicMeta(staleNewPO, oldPO)); + } } @TestTemplate
