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 e91b0882567182cb35207278350b9827f9ff2567 Author: yuqi <[email protected]> AuthorDate: Thu Jul 23 23:11:11 2026 +0800 [#12166] improvement(core): slim table/function UPDATE WHERE to version CAS; keep fileset/policy full-row table and function always raise their version on update, so reduce their WHERE to id + current_version + deleted_at (function base + PostgreSQL). view was already version-only. fileset and policy use CONDITIONAL versioning: current_version only bumps on a versioned-field change, so an audit-only concurrent update would slip past a version-only CAS (proven by TestFilesetMetaService's no-rows conflict test). Keep their full-row compare as the OCC guard; add a comment explaining why. Part of #12166. Remaining: model schema migration. --- .../mapper/provider/base/FilesetMetaBaseSQLProvider.java | 3 +++ .../mapper/provider/base/FunctionMetaBaseSQLProvider.java | 8 +------- .../mapper/provider/base/PolicyMetaBaseSQLProvider.java | 3 +++ .../relational/mapper/provider/base/TableMetaBaseSQLProvider.java | 7 +------ .../provider/postgresql/FunctionMetaPostgreSQLProvider.java | 8 +------- 5 files changed, 9 insertions(+), 20 deletions(-) diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/FilesetMetaBaseSQLProvider.java b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/FilesetMetaBaseSQLProvider.java index 8a1ad653c8..64f8de39be 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/FilesetMetaBaseSQLProvider.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/FilesetMetaBaseSQLProvider.java @@ -287,6 +287,9 @@ public class FilesetMetaBaseSQLProvider { + " current_version = #{newFilesetMeta.currentVersion}," + " last_version = #{newFilesetMeta.lastVersion}," + " deleted_at = #{newFilesetMeta.deletedAt}" + // Fileset uses CONDITIONAL versioning: current_version only bumps on a versioned-field + // change, so a version-only CAS would miss an audit-only concurrent update. Keep the + // full-row compare, which is the actual OCC guard here. + " WHERE fileset_id = #{oldFilesetMeta.filesetId}" + " AND fileset_name = #{oldFilesetMeta.filesetName}" + " AND metalake_id = #{oldFilesetMeta.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 ebb000edef..ecd2d68b94 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 @@ -302,15 +302,9 @@ public class FunctionMetaBaseSQLProvider { + " function_latest_version = #{newFunctionMeta.functionLatestVersion}," + " audit_info = #{newFunctionMeta.auditInfo}," + " deleted_at = #{newFunctionMeta.deletedAt}" + // OCC: compare-and-set on the version alone (function_current_version is monotonic). + " WHERE function_id = #{oldFunctionMeta.functionId}" - + " AND function_name = #{oldFunctionMeta.functionName}" - + " AND metalake_id = #{oldFunctionMeta.metalakeId}" - + " AND catalog_id = #{oldFunctionMeta.catalogId}" - + " AND schema_id = #{oldFunctionMeta.schemaId}" - + " AND function_type = #{oldFunctionMeta.functionType}" + " AND function_current_version = #{oldFunctionMeta.functionCurrentVersion}" - + " AND function_latest_version = #{oldFunctionMeta.functionLatestVersion}" - + " AND audit_info = #{oldFunctionMeta.auditInfo}" + " AND deleted_at = 0"; } } diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/PolicyMetaBaseSQLProvider.java b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/PolicyMetaBaseSQLProvider.java index c0f6b21f6d..ad5ac707c3 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/PolicyMetaBaseSQLProvider.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/PolicyMetaBaseSQLProvider.java @@ -134,6 +134,9 @@ public class PolicyMetaBaseSQLProvider { + " current_version = #{newPolicyMeta.currentVersion}," + " last_version = #{newPolicyMeta.lastVersion}," + " deleted_at = #{newPolicyMeta.deletedAt}" + // Policy uses CONDITIONAL versioning (like fileset): current_version only bumps on a + // versioned-field change, so a version-only CAS would miss an audit-only concurrent + // update. Keep the full-row compare, which is the actual OCC guard here. + " WHERE policy_id = #{oldPolicyMeta.policyId}" + " AND policy_name = #{oldPolicyMeta.policyName}" + " AND policy_type = #{oldPolicyMeta.policyType}" diff --git a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TableMetaBaseSQLProvider.java b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TableMetaBaseSQLProvider.java index 06684724fa..7b55f0967f 100644 --- a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TableMetaBaseSQLProvider.java +++ b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/TableMetaBaseSQLProvider.java @@ -239,14 +239,9 @@ public class TableMetaBaseSQLProvider { + " current_version = #{newTableMeta.currentVersion}," + " last_version = #{newTableMeta.lastVersion}," + " deleted_at = #{newTableMeta.deletedAt}" + // OCC: compare-and-set on the version alone (current_version is monotonic on update). + " WHERE table_id = #{oldTableMeta.tableId}" - + " AND table_name = #{oldTableMeta.tableName}" - + " AND metalake_id = #{oldTableMeta.metalakeId}" - + " AND catalog_id = #{oldTableMeta.catalogId}" - + " AND schema_id = #{oldTableMeta.schemaId}" - + " AND audit_info = #{oldTableMeta.auditInfo}" + " AND current_version = #{oldTableMeta.currentVersion}" - + " AND last_version = #{oldTableMeta.lastVersion}" + " AND deleted_at = 0"; } 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 6117bdc5c0..205dd47eb9 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 @@ -164,15 +164,9 @@ public class FunctionMetaPostgreSQLProvider extends FunctionMetaBaseSQLProvider + " function_latest_version = #{newFunctionMeta.functionLatestVersion}," + " audit_info = #{newFunctionMeta.auditInfo}," + " deleted_at = #{newFunctionMeta.deletedAt}" + // OCC: compare-and-set on the version alone (function_current_version is monotonic). + " WHERE function_id = #{oldFunctionMeta.functionId}" - + " AND function_name = #{oldFunctionMeta.functionName}" - + " AND metalake_id = #{oldFunctionMeta.metalakeId}" - + " AND catalog_id = #{oldFunctionMeta.catalogId}" - + " AND schema_id = #{oldFunctionMeta.schemaId}" - + " AND function_type = #{oldFunctionMeta.functionType}" + " AND function_current_version = #{oldFunctionMeta.functionCurrentVersion}" - + " AND function_latest_version = #{oldFunctionMeta.functionLatestVersion}" - + " AND audit_info = #{oldFunctionMeta.auditInfo}" + " AND deleted_at = 0"; } }
