This is an automated email from the ASF dual-hosted git repository.
yuqi1129 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gravitino.git
The following commit(s) were added to refs/heads/main by this push:
new 9f7be9d25e [#12206] improvement(core): decouple the fileset and policy
OCC version from their history version (#13484)
9f7be9d25e is described below
commit 9f7be9d25eab340b02e86e49adde67471e79577c
Author: Qi Yu <[email protected]>
AuthorDate: Tue Sep 29 22:35:40 2026 +0800
[#12206] improvement(core): decouple the fileset and policy OCC version
from their history version (#13484)
### What changes were proposed in this pull request?
`fileset_meta` and `policy_meta` get an additive `occ_version` column.
The CAS predicate and the unconditional increment move to it, and
`current_version` goes back to meaning "the history version this row
points at", advancing only when the fields `*_version_info` stores
actually change.
- **Schema:** the column in `schema-2.0.0-{mysql,postgresql,h2}.sql` and
in the still-open `upgrade-1.3.0-to-2.0.0-*.sql`. `DEFAULT 1` is the
whole backfill, because `occ_version` is only ever compared against
itself on the same row: an upgraded row at `occ_version = 1` with
`current_version = 7` is correct.
- **Write path:** `updateFilesetPOWithVersion` and
`buildNextPolicyPOVersion` always advance `occ_version`; when nothing
stored changed they keep `current_version` / `last_version`, return no
snapshot, and the services skip the version insert.
- **SQL:** the CAS moves to `occ_version` in `updateFilesetMeta`,
`updatePolicyMeta`, `softDeleteFilesetMetasByFilesetId` and
`softDeletePolicyByIdAndVersion`. The unused fileset upsert
`insertFilesetMetaOnDuplicateKeyUpdate` (no production caller; fileset
overwrite uses `updateFilesetMeta`) is removed.
- **Read path: no query changes.** `current_version` moves only in the
statement that inserts the snapshot it moves to, so a live row still
points at exactly one active snapshot.
One subtlety worth pointing at during review: `updateFilesetMeta`'s `NOT
EXISTS` snapshot check is appended **only** when the alter allocates a
version. An alter that allocates none keeps `current_version` where it
is, and the snapshot it points at is supposed to exist, so the check
would otherwise reject every such alter.
### Why are the changes needed?
`current_version` was doing two jobs at once: it is the join key into
`*_version_info`, and it was the value the CAS compared. Once OCC made
the token advance on every alter (#12656, #12782), an audit-only or
rename-only alter had to write a full snapshot just to keep the join
resolvable — for a fileset that is one row per storage location — and
the retention job removed them again. The history table recorded
revisions that revised nothing.
Fix: #12206
**Scope note:** #12206's body is largely stale against `main`. The
full-row comparisons it describes are gone (#12656 for fileset, #12782
for policy), and the `model_meta` insert it calls out already lists
`current_version` / `last_version` (`ModelMetaBaseSQLProvider.java:36`,
`:49`). What remained is the coupling named in the title. `table` /
`view` / `function` share it and are left to a follow-up: their
`current_version` has many more callers.
**Upgrade:** the 1.3.0 → 2.0.0 upgrade is offline
(`docs/how-to-upgrade.md`, Step 1: shut down Gravitino before running
the scripts), so 1.3.0 and 2.0.0 servers never write to the same
database at the same time.
### Does this PR introduce _any_ user-facing change?
No. Neither fileset nor policy version history is reachable through any
public API — `FilesetVersionMapper` and `PolicyVersionMapper` are used
only by the `current_version` join and the retention job, and `grep
getCurrentVersion()` outside `storage/relational` has no hits, so
nothing depends on it increasing monotonically. No API, configuration or
read-path change.
One benign behavioural consequence: because metadata-only alters no
longer consume version numbers, retention evicts genuine content history
more slowly than before.
### How was this patch tested?
`TestFilesetMetaService` 66, `TestPolicyMetaService` 66,
`TestPOConverters` 42, `TestFilesetMetaBaseSQLProvider` 4,
`TestFilesetMetaPostgreSQLProvider` 1, `TestSQLScripts` 12 — 191 tests,
0 failures or skips. The service test templates run against H2, MySQL,
and PostgreSQL. Formatting and compilation checks pass.
New tests:
- `testMetadataOnlyAlterRejectsAStaleUpdate` in both service suites —
commits an audit-only winner after the outer updater reads and before
its CAS, then verifies that the stale audit-only update fails, the
winner remains, and no snapshot is added. Mutation-checked on all three
backends: changing either update CAS back to `current_version` makes its
test fail because no `OptimisticLockException` is thrown.
- `testUpgradeToTwoZeroBackfillsOccVersions` — upgrades live and deleted
fileset/policy metadata with non-initial history versions, verifies
`occ_version = 1`, and verifies that `current_version`, `last_version`,
and deletion state are preserved on all three backends.
- `testAlterWritesASnapshotOnlyWhenStoredContentChanges` — a fileset
with two storage locations: a rename leaves `fileset_version_info` at 2
rows and 1 distinct version while `occ_version` advances; a comment
change takes it to 4 rows and 2 versions; the fileset still reads back
correctly. Added a `countFilesetVersionRows` helper, because the row
count rather than the version count is what a snapshot costs.
- `testRenameWithReorderedPropertiesWritesNoSnapshot` and
`testUpdateFilesetPOVersionComparesPropertiesByValue` — fileset
properties are compared by value, so a rename whose properties went
through a `HashMap` (as `FilesetCatalogOperations` does) writes no
snapshot. Both fail without the fix. Shared comment/properties are
compared once per snapshot. The converter test uses two locations and
also verifies that a property value change or a change to a non-first
location still writes a complete snapshot.
- `testDeleteRejectsAStaleVersionAfterAMetadataOnlyAlter` — pins the
regression this decoupling creates. An audit-only alter no longer moves
`current_version`, so a drop still guarded by it would stop detecting
one and would delete a fileset the caller never observed in its current
state. Mutation-checked: with the delete CAS on `current_version` it
fails with `Expected OptimisticLockException to be thrown, but nothing
was thrown`.
- `testUpdateDropsTheSnapshotCheckWhenNoVersionIsAllocated` — pins the
conditional `NOT EXISTS`.
- `testStalePolicyDeleteAfterMetadataOnlyAlter` — uses the Policy read
back from storage, changes only audit info, and verifies no snapshot is
allocated and a stale delete fails. This exposed unordered
`supportedObjectTypes` serializing in a different order after readback;
content is now compared by value when its JSON strings differ.
- The fileset snapshot test now performs two metadata-only renames
before changing content, and the policy overwrite test first performs a
metadata-only update.
The `filesetSnapshotUnchanged` branch is mutation-checked too: disabling
it makes `testUpdateFilesetPOVersion` fail (`expected: <1> but was:
<2>`).
Two existing tests were rewritten because they asserted the behaviour
this PR reverses: `testMetadataOnlyPolicyAlterCreatesCompleteSnapshot` →
`...AdvancesOnlyTheOccVersion`, and the tail of
`testAlterReportsOptimisticLockConflictAndKeepsWinnerVersion`.
**Coverage limits:** The legacy `maxStoredVersion` path is exercised by
`testAlterSkipsVersionsAlreadyStored` and
`testOverwriteSkipsVersionsAlreadyStored` against the service backends.
### Alternative considered: reuse `last_version`
A [tested experiment
branch](https://github.com/yuqi1129/gravitino/tree/experiment/13484-last-version)
uses `last_version` as the row revision and keeps `current_version` as
the content snapshot pointer. It avoids a new database column: every
write increments `last_version`, while only content changes advance
`current_version`. Consecutive metadata-only writes, the next content
snapshot, stale deletes, and the policy readback case are covered by
tests on that branch.
We kept `occ_version` in this PR because `last_version` previously meant
the latest allocated snapshot version (and was used as a high-water
mark). Reusing it changes an existing column's meaning and needs a
separate compatibility discussion. The experiment branch is available if
we decide to take that route later.
---
.../relational/mapper/FilesetMetaMapper.java | 18 +-
.../mapper/FilesetMetaSQLProviderFactory.java | 11 +-
.../relational/mapper/PolicyMetaMapper.java | 6 +-
.../mapper/PolicyMetaSQLProviderFactory.java | 4 +-
.../provider/base/FilesetMetaBaseSQLProvider.java | 134 +++++-----
.../provider/base/PolicyMetaBaseSQLProvider.java | 25 +-
.../postgresql/FilesetMetaPostgreSQLProvider.java | 44 +---
.../postgresql/PolicyMetaPostgreSQLProvider.java | 4 +-
.../gravitino/storage/relational/po/FilesetPO.java | 33 ++-
.../gravitino/storage/relational/po/PolicyPO.java | 17 ++
.../relational/service/FilesetMetaService.java | 16 +-
.../relational/service/PolicyMetaService.java | 26 +-
.../storage/relational/utils/POConverters.java | 164 ++++++++++--
.../apache/gravitino/storage/TestSQLScripts.java | 54 ++++
.../storage/relational/TestJDBCBackend.java | 26 ++
.../base/TestFilesetMetaBaseSQLProvider.java | 48 ++--
.../TestFilesetMetaPostgreSQLProvider.java | 20 +-
.../relational/service/TestFilesetMetaService.java | 278 ++++++++++++++++++++-
.../relational/service/TestPolicyMetaService.java | 119 ++++++++-
.../storage/relational/utils/TestPOConverters.java | 103 +++++++-
scripts/h2/schema-2.0.0-h2.sql | 2 +
scripts/h2/upgrade-1.3.0-to-2.0.0-h2.sql | 12 +
scripts/mysql/schema-2.0.0-mysql.sql | 2 +
scripts/mysql/upgrade-1.3.0-to-2.0.0-mysql.sql | 12 +
scripts/postgresql/schema-2.0.0-postgresql.sql | 4 +
.../upgrade-1.3.0-to-2.0.0-postgresql.sql | 12 +
26 files changed, 955 insertions(+), 239 deletions(-)
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/FilesetMetaMapper.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/FilesetMetaMapper.java
index fbd2ca9c26..ef18e9fb99 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/FilesetMetaMapper.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/FilesetMetaMapper.java
@@ -73,6 +73,7 @@ public interface FilesetMetaMapper {
@Result(property = "auditInfo", column = "audit_info"),
@Result(property = "currentVersion", column = "current_version"),
@Result(property = "lastVersion", column = "last_version"),
+ @Result(property = "occVersion", column = "occ_version"),
@Result(property = "deletedAt", column = "deleted_at"),
@Result(
property = "filesetVersionPOs",
@@ -95,6 +96,7 @@ public interface FilesetMetaMapper {
@Result(property = "auditInfo", column = "audit_info"),
@Result(property = "currentVersion", column = "current_version"),
@Result(property = "lastVersion", column = "last_version"),
+ @Result(property = "occVersion", column = "occ_version"),
@Result(property = "deletedAt", column = "deleted_at"),
@Result(
property = "filesetVersionPOs",
@@ -122,6 +124,7 @@ public interface FilesetMetaMapper {
@Result(property = "auditInfo", column = "audit_info"),
@Result(property = "currentVersion", column = "current_version"),
@Result(property = "lastVersion", column = "last_version"),
+ @Result(property = "occVersion", column = "occ_version"),
@Result(property = "deletedAt", column = "deleted_at"),
@Result(
property = "filesetVersionPOs",
@@ -150,6 +153,7 @@ public interface FilesetMetaMapper {
@Result(property = "auditInfo", column = "audit_info"),
@Result(property = "currentVersion", column = "current_version"),
@Result(property = "lastVersion", column = "last_version"),
+ @Result(property = "occVersion", column = "occ_version"),
@Result(property = "deletedAt", column = "deleted_at"),
@Result(
property = "filesetVersionPOs",
@@ -175,6 +179,7 @@ public interface FilesetMetaMapper {
@Result(property = "auditInfo", column = "audit_info"),
@Result(property = "currentVersion", column = "current_version"),
@Result(property = "lastVersion", column = "last_version"),
+ @Result(property = "occVersion", column = "occ_version"),
@Result(property = "deletedAt", column = "deleted_at"),
@Result(
property = "filesetVersionPOs",
@@ -210,6 +215,7 @@ public interface FilesetMetaMapper {
@Result(property = "auditInfo", column = "audit_info"),
@Result(property = "currentVersion", column = "current_version"),
@Result(property = "lastVersion", column = "last_version"),
+ @Result(property = "occVersion", column = "occ_version"),
@Result(property = "deletedAt", column = "deleted_at"),
@Result(
property = "filesetVersionPOs",
@@ -231,11 +237,6 @@ public interface FilesetMetaMapper {
@InsertProvider(type = FilesetMetaSQLProviderFactory.class, method =
"insertFilesetMeta")
void insertFilesetMeta(@Param("filesetMeta") FilesetPO filesetPO);
- @InsertProvider(
- type = FilesetMetaSQLProviderFactory.class,
- method = "insertFilesetMetaOnDuplicateKeyUpdate")
- void insertFilesetMetaOnDuplicateKeyUpdate(@Param("filesetMeta") FilesetPO
filesetPO);
-
@UpdateProvider(type = FilesetMetaSQLProviderFactory.class, method =
"updateFilesetMeta")
Integer updateFilesetMeta(
@Param("newFilesetMeta") FilesetPO newFilesetPO,
@@ -257,17 +258,17 @@ public interface FilesetMetaMapper {
Integer softDeleteFilesetMetasBySchemaIds(@Param("schemaIds") List<Long>
schemaIds);
/**
- * Soft-deletes a fileset only if its version has not changed since the
caller read it.
+ * Soft-deletes a fileset only if its OCC version has not changed since the
caller read it.
*
* @param filesetId the fileset ID
- * @param currentVersion the version observed by the caller
+ * @param occVersion the OCC version observed by the caller
* @return the number of deleted rows; zero means the fileset changed or
disappeared
*/
@UpdateProvider(
type = FilesetMetaSQLProviderFactory.class,
method = "softDeleteFilesetMetasByFilesetId")
Integer softDeleteFilesetMetasByFilesetId(
- @Param("filesetId") Long filesetId, @Param("currentVersion") Long
currentVersion);
+ @Param("filesetId") Long filesetId, @Param("occVersion") Long
occVersion);
@DeleteProvider(
type = FilesetMetaSQLProviderFactory.class,
@@ -285,6 +286,7 @@ public interface FilesetMetaMapper {
@Result(property = "auditInfo", column = "audit_info"),
@Result(property = "currentVersion", column = "current_version"),
@Result(property = "lastVersion", column = "last_version"),
+ @Result(property = "occVersion", column = "occ_version"),
@Result(property = "deletedAt", column = "deleted_at"),
@Result(
property = "filesetVersionPOs",
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/FilesetMetaSQLProviderFactory.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/FilesetMetaSQLProviderFactory.java
index b09dd069e8..0b87dc4fb7 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/FilesetMetaSQLProviderFactory.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/FilesetMetaSQLProviderFactory.java
@@ -106,11 +106,6 @@ public class FilesetMetaSQLProviderFactory {
return getProvider().insertFilesetMeta(filesetPO);
}
- public static String insertFilesetMetaOnDuplicateKeyUpdate(
- @Param("filesetMeta") FilesetPO filesetPO) {
- return getProvider().insertFilesetMetaOnDuplicateKeyUpdate(filesetPO);
- }
-
public static String updateFilesetMeta(
@Param("newFilesetMeta") FilesetPO newFilesetPO,
@Param("oldFilesetMeta") FilesetPO oldFilesetPO) {
@@ -133,12 +128,12 @@ public class FilesetMetaSQLProviderFactory {
* Returns SQL that soft-deletes a fileset by ID and expected version.
*
* @param filesetId the fileset ID
- * @param currentVersion the version observed by the caller
+ * @param occVersion the OCC version observed by the caller
* @return the version-checked delete SQL
*/
public static String softDeleteFilesetMetasByFilesetId(
- @Param("filesetId") Long filesetId, @Param("currentVersion") Long
currentVersion) {
- return getProvider().softDeleteFilesetMetasByFilesetId(filesetId,
currentVersion);
+ @Param("filesetId") Long filesetId, @Param("occVersion") Long
occVersion) {
+ return getProvider().softDeleteFilesetMetasByFilesetId(filesetId,
occVersion);
}
public String deleteFilesetMetasByLegacyTimeline(
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/PolicyMetaMapper.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/PolicyMetaMapper.java
index d7c8f5a810..383791342b 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/PolicyMetaMapper.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/PolicyMetaMapper.java
@@ -55,6 +55,7 @@ public interface PolicyMetaMapper {
@Result(property = "auditInfo", column = "audit_info"),
@Result(property = "currentVersion", column = "current_version"),
@Result(property = "lastVersion", column = "last_version"),
+ @Result(property = "occVersion", column = "occ_version"),
@Result(property = "deletedAt", column = "deleted_at"),
@Result(property = "policyVersionPO.id", column = "id"),
@Result(property = "policyVersionPO.metalakeId", column =
"version_metalake_id"),
@@ -94,14 +95,14 @@ public interface PolicyMetaMapper {
* Soft-deletes an active policy when its OCC version still matches.
*
* @param policyId The policy ID.
- * @param currentVersion The version observed by the caller.
+ * @param occVersion The OCC version observed by the caller.
* @return The number of affected rows.
*/
@UpdateProvider(
type = PolicyMetaSQLProviderFactory.class,
method = "softDeletePolicyByIdAndVersion")
Integer softDeletePolicyByIdAndVersion(
- @Param("policyId") Long policyId, @Param("currentVersion") Long
currentVersion);
+ @Param("policyId") Long policyId, @Param("occVersion") Long occVersion);
@UpdateProvider(
type = PolicyMetaSQLProviderFactory.class,
@@ -124,6 +125,7 @@ public interface PolicyMetaMapper {
@Result(property = "auditInfo", column = "audit_info"),
@Result(property = "currentVersion", column = "current_version"),
@Result(property = "lastVersion", column = "last_version"),
+ @Result(property = "occVersion", column = "occ_version"),
@Result(property = "deletedAt", column = "deleted_at")
})
@SelectProvider(
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/PolicyMetaSQLProviderFactory.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/PolicyMetaSQLProviderFactory.java
index 807badd1fb..602d789c8e 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/PolicyMetaSQLProviderFactory.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/PolicyMetaSQLProviderFactory.java
@@ -74,8 +74,8 @@ public class PolicyMetaSQLProviderFactory {
/** Delegates a version-checked policy soft delete. */
public static String softDeletePolicyByIdAndVersion(
- @Param("policyId") Long policyId, @Param("currentVersion") Long
currentVersion) {
- return getProvider().softDeletePolicyByIdAndVersion(policyId,
currentVersion);
+ @Param("policyId") Long policyId, @Param("occVersion") Long occVersion) {
+ return getProvider().softDeletePolicyByIdAndVersion(policyId, occVersion);
}
public static String deletePolicyMetasByLegacyTimeline(
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 285a8200cb..cedfbb654b 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
@@ -23,6 +23,7 @@ import static
org.apache.gravitino.storage.relational.mapper.FilesetMetaMapper.M
import static
org.apache.gravitino.storage.relational.mapper.FilesetMetaMapper.VERSION_TABLE_NAME;
import java.util.List;
+import java.util.Objects;
import org.apache.gravitino.storage.relational.mapper.CatalogMetaMapper;
import org.apache.gravitino.storage.relational.mapper.MetalakeMetaMapper;
import org.apache.gravitino.storage.relational.mapper.SchemaMetaMapper;
@@ -34,7 +35,8 @@ public class FilesetMetaBaseSQLProvider {
public String listFilesetPOsBySchemaId(@Param("schemaId") Long schemaId) {
return "SELECT fm.fileset_id, fm.fileset_name, fm.metalake_id,
fm.catalog_id, fm.schema_id,"
- + " fm.type, fm.audit_info, fm.current_version, fm.last_version,
fm.deleted_at,"
+ + " fm.type, fm.audit_info, fm.current_version, fm.last_version,
fm.occ_version,"
+ + " fm.deleted_at,"
+ " vi.id, vi.metalake_id as version_metalake_id, vi.catalog_id as
version_catalog_id,"
+ " vi.schema_id as version_schema_id, vi.fileset_id as
version_fileset_id,"
+ " vi.version, vi.fileset_comment, vi.properties,
vi.storage_location_name, vi.storage_location,"
@@ -62,6 +64,7 @@ public class FilesetMetaBaseSQLProvider {
fm.audit_info,
fm.current_version,
fm.last_version,
+ fm.occ_version,
fm.deleted_at,
vi.id,
vi.metalake_id as version_metalake_id,
@@ -114,7 +117,8 @@ public class FilesetMetaBaseSQLProvider {
public String listFilesetPOsByFilesetIds(@Param("filesetIds") List<Long>
filesetIds) {
return "<script>"
+ "SELECT fm.fileset_id, fm.fileset_name, fm.metalake_id,
fm.catalog_id, fm.schema_id,"
- + " fm.type, fm.audit_info, fm.current_version, fm.last_version,
fm.deleted_at,"
+ + " fm.type, fm.audit_info, fm.current_version, fm.last_version,
fm.occ_version,"
+ + " fm.deleted_at,"
+ " vi.id, vi.metalake_id as version_metalake_id, vi.catalog_id as
version_catalog_id,"
+ " vi.schema_id as version_schema_id, vi.fileset_id as
version_fileset_id,"
+ " vi.version, vi.fileset_comment, vi.properties,
vi.storage_location_name, vi.storage_location,"
@@ -136,7 +140,8 @@ public class FilesetMetaBaseSQLProvider {
public String selectFilesetMetaBySchemaIdAndName(
@Param("schemaId") Long schemaId, @Param("filesetName") String name) {
return "SELECT fm.fileset_id, fm.fileset_name, fm.metalake_id,
fm.catalog_id, fm.schema_id,"
- + " fm.type, fm.audit_info, fm.current_version, fm.last_version,
fm.deleted_at,"
+ + " fm.type, fm.audit_info, fm.current_version, fm.last_version,
fm.occ_version,"
+ + " fm.deleted_at,"
+ " vi.id, vi.metalake_id as version_metalake_id, vi.catalog_id as
version_catalog_id,"
+ " vi.schema_id as version_schema_id, vi.fileset_id as
version_fileset_id,"
+ " vi.version, vi.fileset_comment, vi.properties,
vi.storage_location_name, vi.storage_location,"
@@ -166,6 +171,7 @@ public class FilesetMetaBaseSQLProvider {
fm.audit_info,
fm.current_version,
fm.last_version,
+ fm.occ_version,
fm.deleted_at,
vi.id,
vi.metalake_id as version_metalake_id,
@@ -210,7 +216,8 @@ public class FilesetMetaBaseSQLProvider {
public String selectFilesetMetaById(@Param("filesetId") Long filesetId) {
return "SELECT fm.fileset_id, fm.fileset_name, fm.metalake_id,
fm.catalog_id, fm.schema_id,"
- + " fm.type, fm.audit_info, fm.current_version, fm.last_version,
fm.deleted_at,"
+ + " fm.type, fm.audit_info, fm.current_version, fm.last_version,
fm.occ_version,"
+ + " fm.deleted_at,"
+ " vi.id, vi.metalake_id as version_metalake_id, vi.catalog_id as
version_catalog_id,"
+ " vi.schema_id as version_schema_id, vi.fileset_id as
version_fileset_id,"
+ " vi.version, vi.fileset_comment, vi.properties,
vi.storage_location_name, vi.storage_location,"
@@ -227,8 +234,8 @@ public class FilesetMetaBaseSQLProvider {
/**
* Returns the active fileset metadata row selected by its natural key.
*
- * <p>An overwrite may match the natural key instead of the incoming ID.
Reading the stored row
- * after the upsert tells dependent version rows which ID and
database-generated version to use.
+ * <p>An overwrite may match the natural key instead of the incoming ID. The
overwrite locks the
+ * stored row with this select, then updates it in place, keeping the stored
ID.
*
* @param schemaId the schema ID
* @param filesetName the fileset name
@@ -240,7 +247,7 @@ public class FilesetMetaBaseSQLProvider {
+ " metalake_id as metalakeId, catalog_id as catalogId, schema_id as
schemaId,"
+ " type as type, audit_info as auditInfo,"
+ " current_version as currentVersion, last_version as lastVersion,"
- + " deleted_at as deletedAt"
+ + " occ_version as occVersion, deleted_at as deletedAt"
+ " FROM "
+ META_TABLE_NAME
+ " WHERE schema_id = #{schemaId} AND fileset_name = #{filesetName}"
@@ -252,7 +259,7 @@ public class FilesetMetaBaseSQLProvider {
+ META_TABLE_NAME
+ " (fileset_id, fileset_name, metalake_id,"
+ " catalog_id, schema_id, type, audit_info,"
- + " current_version, last_version, deleted_at)"
+ + " current_version, last_version, occ_version, deleted_at)"
+ " VALUES ("
+ " #{filesetMeta.filesetId},"
+ " #{filesetMeta.filesetName},"
@@ -263,55 +270,24 @@ public class FilesetMetaBaseSQLProvider {
+ " #{filesetMeta.auditInfo},"
+ " #{filesetMeta.currentVersion},"
+ " #{filesetMeta.lastVersion},"
+ + " #{filesetMeta.occVersion},"
+ " #{filesetMeta.deletedAt}"
+ " )";
}
- public String insertFilesetMetaOnDuplicateKeyUpdate(@Param("filesetMeta")
FilesetPO filesetPO) {
- return "INSERT INTO "
- + META_TABLE_NAME
- + " (fileset_id, fileset_name, metalake_id,"
- + " catalog_id, schema_id, type, audit_info,"
- + " current_version, last_version, deleted_at)"
- + " VALUES ("
- + " #{filesetMeta.filesetId},"
- + " #{filesetMeta.filesetName},"
- + " #{filesetMeta.metalakeId},"
- + " #{filesetMeta.catalogId},"
- + " #{filesetMeta.schemaId},"
- + " #{filesetMeta.type},"
- + " #{filesetMeta.auditInfo},"
- + " #{filesetMeta.currentVersion},"
- + " #{filesetMeta.lastVersion},"
- + " #{filesetMeta.deletedAt}"
- + " )"
- + " ON DUPLICATE KEY UPDATE"
- + " fileset_name = #{filesetMeta.filesetName},"
- + " metalake_id = #{filesetMeta.metalakeId},"
- + " catalog_id = #{filesetMeta.catalogId},"
- + " schema_id = #{filesetMeta.schemaId},"
- + " type = #{filesetMeta.type},"
- + " audit_info = #{filesetMeta.auditInfo},"
- // An overwrite is also a write observed by OCC. Advance from the
stored value instead of
- // resetting the row to the initial version carried by the incoming
create request.
- //
- // Keep current_version last: MySQL evaluates these assignments left
to right against the
- // columns already assigned, while H2 and PostgreSQL evaluate every
right-hand side against
- // the row as it was before the update. Both agree only while
current_version is read
- // before it is assigned.
- + " last_version = current_version + 1,"
- + " current_version = current_version + 1,"
- + " deleted_at = #{filesetMeta.deletedAt}";
- }
-
/**
- * Returns SQL that updates a fileset only while its OCC version is
unchanged and its next
- * snapshot version is free.
+ * Returns SQL that updates a fileset only while its OCC version is
unchanged, and, when the alter
+ * allocates a new snapshot, only while that snapshot version is free.
+ *
+ * <p>{@code occ_version} is the concurrency token, so payload, name, and
audit columns are
+ * deliberately excluded from the predicate. This also detects
change-then-change-back races that
+ * a full-row comparison would miss.
*
- * <p>The version is the concurrency token, so payload, name, and audit
columns are deliberately
- * excluded from the predicate. This also detects change-then-change-back
races that a full-row
- * comparison would miss. The snapshot check detects rows affected by the
legacy overwrite bug
- * without requiring a separate {@code MAX(version)} query on every normal
alter.
+ * <p>The snapshot check detects rows affected by the legacy overwrite bug
without requiring a
+ * separate {@code MAX(version)} query on every normal alter. It is appended
only when the alter
+ * moves {@code current_version}: an alter that changes nothing the version
table stores keeps
+ * that column where it is, and the snapshot it points at is supposed to
exist, so the check would
+ * reject every such alter.
*
* @param newFilesetPO the new fileset values
* @param oldFilesetPO the fileset values and version observed by the caller
@@ -320,25 +296,31 @@ public class FilesetMetaBaseSQLProvider {
public String updateFilesetMeta(
@Param("newFilesetMeta") FilesetPO newFilesetPO,
@Param("oldFilesetMeta") FilesetPO oldFilesetPO) {
- return "UPDATE "
- + META_TABLE_NAME
- + " SET fileset_name = #{newFilesetMeta.filesetName},"
- + " metalake_id = #{newFilesetMeta.metalakeId},"
- + " catalog_id = #{newFilesetMeta.catalogId},"
- + " schema_id = #{newFilesetMeta.schemaId},"
- + " type = #{newFilesetMeta.type},"
- + " audit_info = #{newFilesetMeta.auditInfo},"
- + " current_version = #{newFilesetMeta.currentVersion},"
- + " last_version = #{newFilesetMeta.lastVersion},"
- + " deleted_at = #{newFilesetMeta.deletedAt}"
- + " WHERE fileset_id = #{oldFilesetMeta.filesetId}"
- + " AND current_version = #{oldFilesetMeta.currentVersion}"
- + " AND deleted_at = 0"
- + " AND NOT EXISTS (SELECT 1 FROM "
- + VERSION_TABLE_NAME
- + " fv WHERE fv.fileset_id = #{oldFilesetMeta.filesetId}"
- + " AND fv.version >= #{newFilesetMeta.currentVersion}"
- + " AND fv.deleted_at = 0)";
+ String sql =
+ "UPDATE "
+ + META_TABLE_NAME
+ + " SET fileset_name = #{newFilesetMeta.filesetName},"
+ + " metalake_id = #{newFilesetMeta.metalakeId},"
+ + " catalog_id = #{newFilesetMeta.catalogId},"
+ + " schema_id = #{newFilesetMeta.schemaId},"
+ + " type = #{newFilesetMeta.type},"
+ + " audit_info = #{newFilesetMeta.auditInfo},"
+ + " current_version = #{newFilesetMeta.currentVersion},"
+ + " last_version = #{newFilesetMeta.lastVersion},"
+ + " occ_version = #{newFilesetMeta.occVersion},"
+ + " deleted_at = #{newFilesetMeta.deletedAt}"
+ + " WHERE fileset_id = #{oldFilesetMeta.filesetId}"
+ + " AND occ_version = #{oldFilesetMeta.occVersion}"
+ + " AND deleted_at = 0";
+ if (!Objects.equals(newFilesetPO.getCurrentVersion(),
oldFilesetPO.getCurrentVersion())) {
+ sql +=
+ " AND NOT EXISTS (SELECT 1 FROM "
+ + VERSION_TABLE_NAME
+ + " fv WHERE fv.fileset_id = #{oldFilesetMeta.filesetId}"
+ + " AND fv.version >= #{newFilesetMeta.currentVersion}"
+ + " AND fv.deleted_at = 0)";
+ }
+ return sql;
}
public String softDeleteFilesetMetasByMetalakeId(@Param("metalakeId") Long
metalakeId) {
@@ -372,20 +354,21 @@ public class FilesetMetaBaseSQLProvider {
}
/**
- * Returns SQL that deletes only the fileset version observed by the caller.
+ * Returns SQL that soft-deletes a fileset only while it still carries the
OCC version the caller
+ * observed.
*
* @param filesetId the fileset ID
- * @param currentVersion the version observed by the caller
+ * @param occVersion the OCC version observed by the caller
* @return the version-checked delete SQL
*/
public String softDeleteFilesetMetasByFilesetId(
- @Param("filesetId") Long filesetId, @Param("currentVersion") Long
currentVersion) {
+ @Param("filesetId") Long filesetId, @Param("occVersion") Long
occVersion) {
return "UPDATE "
+ META_TABLE_NAME
+ " SET deleted_at = "
+ DatabaseTimeSQL.MYSQL
+ " WHERE fileset_id = #{filesetId}"
- + " AND current_version = #{currentVersion} AND deleted_at = 0";
+ + " AND occ_version = #{occVersion} AND deleted_at = 0";
}
public String deleteFilesetMetasByLegacyTimeline(
@@ -402,7 +385,8 @@ public class FilesetMetaBaseSQLProvider {
@Param("filesetNames") List<String> filesetNames) {
return "<script>"
+ "SELECT fm.fileset_id, fm.fileset_name, fm.metalake_id,
fm.catalog_id, fm.schema_id,"
- + " fm.type, fm.audit_info, fm.current_version, fm.last_version,
fm.deleted_at,"
+ + " fm.type, fm.audit_info, fm.current_version, fm.last_version,
fm.occ_version,"
+ + " fm.deleted_at,"
+ " vi.id, vi.metalake_id as version_metalake_id, vi.catalog_id as
version_catalog_id,"
+ " vi.schema_id as version_schema_id, vi.fileset_id as
version_fileset_id,"
+ " vi.version, vi.fileset_comment, vi.properties,
vi.storage_location_name, vi.storage_location,"
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 730ce9bed0..94993c6f17 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
@@ -31,7 +31,7 @@ public class PolicyMetaBaseSQLProvider {
public String listPolicyPOsByMetalake(@Param("metalakeName") String
metalakeName) {
return "SELECT pm.policy_id, pm.policy_name, pm.policy_type,
pm.metalake_id,"
- + " pm.audit_info, pm.current_version, pm.last_version,"
+ + " pm.audit_info, pm.current_version, pm.last_version,
pm.occ_version,"
+ " pm.deleted_at, pvi.id, pvi.metalake_id as version_metalake_id,
pvi.policy_id as version_policy_id,"
+ " pvi.version, pvi.policy_comment, pvi.enabled, pvi.content,
pvi.deleted_at as version_deleted_at"
+ " FROM "
@@ -52,7 +52,7 @@ public class PolicyMetaBaseSQLProvider {
@Param("metalakeName") String metalakeName, @Param("policyNames")
List<String> policyNames) {
return "<script>"
+ "SELECT pm.policy_id, pm.policy_name, pm.policy_type,
pm.metalake_id,"
- + " pm.audit_info, pm.current_version, pm.last_version,"
+ + " pm.audit_info, pm.current_version, pm.last_version,
pm.occ_version,"
+ " pm.deleted_at, pvi.id, pvi.metalake_id as version_metalake_id,
pvi.policy_id as version_policy_id,"
+ " pvi.version, pvi.policy_comment, pvi.enabled, pvi.content,
pvi.deleted_at as version_deleted_at"
+ " FROM "
@@ -78,16 +78,16 @@ public class PolicyMetaBaseSQLProvider {
return "INSERT INTO "
+ POLICY_META_TABLE_NAME
+ " (policy_id, policy_name, policy_type, metalake_id,"
- + " audit_info, current_version, last_version, deleted_at)"
+ + " audit_info, current_version, last_version, occ_version,
deleted_at)"
+ " VALUES (#{policyMeta.policyId}, #{policyMeta.policyName},
#{policyMeta.policyType},"
+ " #{policyMeta.metalakeId}, #{policyMeta.auditInfo},
#{policyMeta.currentVersion},"
- + " #{policyMeta.lastVersion}, #{policyMeta.deletedAt})";
+ + " #{policyMeta.lastVersion}, #{policyMeta.occVersion},
#{policyMeta.deletedAt})";
}
public String selectPolicyMetaByMetalakeAndName(
@Param("metalakeName") String metalakeName, @Param("policyName") String
policyName) {
return "SELECT pm.policy_id, pm.policy_name, pm.policy_type,
pm.metalake_id,"
- + " pm.audit_info, pm.current_version, pm.last_version,"
+ + " pm.audit_info, pm.current_version, pm.last_version,
pm.occ_version,"
+ " pm.deleted_at, pvi.id, pvi.metalake_id as version_metalake_id,
pvi.policy_id as version_policy_id,"
+ " pvi.version, pvi.policy_comment, pvi.enabled, pvi.content,
pvi.deleted_at as version_deleted_at"
+ " FROM "
@@ -116,20 +116,21 @@ public class PolicyMetaBaseSQLProvider {
+ " audit_info = #{newPolicyMeta.auditInfo},"
+ " current_version = #{newPolicyMeta.currentVersion},"
+ " last_version = #{newPolicyMeta.lastVersion},"
+ + " occ_version = #{newPolicyMeta.occVersion},"
+ " deleted_at = #{newPolicyMeta.deletedAt}"
+ " WHERE policy_id = #{oldPolicyMeta.policyId}"
- + " AND current_version = #{oldPolicyMeta.currentVersion}"
+ + " AND occ_version = #{oldPolicyMeta.occVersion}"
+ " AND deleted_at = 0";
}
/** Returns SQL that soft-deletes a policy using its stable ID and observed
OCC version. */
public String softDeletePolicyByIdAndVersion(
- @Param("policyId") Long policyId, @Param("currentVersion") Long
currentVersion) {
+ @Param("policyId") Long policyId, @Param("occVersion") Long occVersion) {
return "UPDATE "
+ POLICY_META_TABLE_NAME
+ " SET deleted_at = "
+ DatabaseTimeSQL.MYSQL
- + " WHERE policy_id = #{policyId} AND current_version =
#{currentVersion}"
+ + " WHERE policy_id = #{policyId} AND occ_version = #{occVersion}"
+ " AND deleted_at = 0";
}
@@ -150,7 +151,7 @@ public class PolicyMetaBaseSQLProvider {
public String selectPolicyByPolicyId(@Param("policyId") Long policyId) {
return "SELECT pm.policy_id, pm.policy_name, pm.policy_type,
pm.metalake_id,"
- + " pm.audit_info, pm.current_version, pm.last_version,"
+ + " pm.audit_info, pm.current_version, pm.last_version,
pm.occ_version,"
+ " pm.deleted_at"
+ " FROM "
+ POLICY_META_TABLE_NAME
@@ -182,7 +183,7 @@ public class PolicyMetaBaseSQLProvider {
public String selectPolicyMetaByMetalakeIdAndName(
@Param("metalakeId") Long metalakeId, @Param("policyName") String
policyName) {
return "SELECT pm.policy_id, pm.policy_name, pm.policy_type,
pm.metalake_id,"
- + " pm.audit_info, pm.current_version, pm.last_version,"
+ + " pm.audit_info, pm.current_version, pm.last_version,
pm.occ_version,"
+ " pm.deleted_at"
+ " FROM "
+ POLICY_META_TABLE_NAME
@@ -203,7 +204,7 @@ public class PolicyMetaBaseSQLProvider {
@Param("metalakeName") String metalakeName, @Param("policyNames")
List<String> policyNames) {
return "<script>"
+ "SELECT pm.policy_id, pm.policy_name, pm.policy_type,
pm.metalake_id,"
- + " pm.audit_info, pm.current_version, pm.last_version, pm.deleted_at,"
+ + " pm.audit_info, pm.current_version, pm.last_version,
pm.occ_version, pm.deleted_at,"
+ " pv.id, pv.metalake_id as version_metalake_id, pv.policy_id as
version_policy_id,"
+ " pv.version, pv.policy_comment, pv.enabled, pv.content,
pv.deleted_at as version_deleted_at"
+ " FROM "
@@ -231,7 +232,7 @@ public class PolicyMetaBaseSQLProvider {
*/
private String selectPolicyPOsByPolicyIdsBody() {
return "SELECT pm.policy_id, pm.policy_name, pm.policy_type,
pm.metalake_id,"
- + " pm.audit_info, pm.current_version, pm.last_version,"
+ + " pm.audit_info, pm.current_version, pm.last_version,
pm.occ_version,"
+ " pm.deleted_at"
+ " FROM "
+ POLICY_META_TABLE_NAME
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/FilesetMetaPostgreSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/FilesetMetaPostgreSQLProvider.java
index 2c516e2d96..f65f21fd49 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/FilesetMetaPostgreSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/FilesetMetaPostgreSQLProvider.java
@@ -23,7 +23,6 @@ import static
org.apache.gravitino.storage.relational.mapper.FilesetMetaMapper.M
import java.util.List;
import org.apache.gravitino.storage.relational.mapper.provider.DatabaseTimeSQL;
import
org.apache.gravitino.storage.relational.mapper.provider.base.FilesetMetaBaseSQLProvider;
-import org.apache.gravitino.storage.relational.po.FilesetPO;
import org.apache.ibatis.annotations.Param;
public class FilesetMetaPostgreSQLProvider extends FilesetMetaBaseSQLProvider {
@@ -61,13 +60,13 @@ public class FilesetMetaPostgreSQLProvider extends
FilesetMetaBaseSQLProvider {
}
@Override
- public String softDeleteFilesetMetasByFilesetId(Long filesetId, Long
currentVersion) {
+ public String softDeleteFilesetMetasByFilesetId(Long filesetId, Long
occVersion) {
return "UPDATE "
+ META_TABLE_NAME
+ " SET deleted_at = "
+ DatabaseTimeSQL.POSTGRESQL
+ " WHERE fileset_id = #{filesetId}"
- + " AND current_version = #{currentVersion} AND deleted_at = 0";
+ + " AND occ_version = #{occVersion} AND deleted_at = 0";
}
@Override
@@ -79,43 +78,4 @@ public class FilesetMetaPostgreSQLProvider extends
FilesetMetaBaseSQLProvider {
+ META_TABLE_NAME
+ " WHERE deleted_at > 0 AND deleted_at < #{legacyTimeline} LIMIT
#{limit})";
}
-
- @Override
- public String insertFilesetMetaOnDuplicateKeyUpdate(FilesetPO filesetPO) {
- return "INSERT INTO "
- + META_TABLE_NAME
- + " (fileset_id, fileset_name, metalake_id,"
- + " catalog_id, schema_id, type, audit_info,"
- + " current_version, last_version, deleted_at)"
- + " VALUES ("
- + " #{filesetMeta.filesetId},"
- + " #{filesetMeta.filesetName},"
- + " #{filesetMeta.metalakeId},"
- + " #{filesetMeta.catalogId},"
- + " #{filesetMeta.schemaId},"
- + " #{filesetMeta.type},"
- + " #{filesetMeta.auditInfo},"
- + " #{filesetMeta.currentVersion},"
- + " #{filesetMeta.lastVersion},"
- + " #{filesetMeta.deletedAt}"
- + " )"
- // Overwrite is selected by name, and a create request normally
carries a newly generated
- // ID. Target the natural key so PostgreSQL preserves the ID of the
row being replaced, the
- // same behavior that MySQL and H2 provide for their duplicate-key
upsert.
- + " ON CONFLICT(schema_id, fileset_name, deleted_at) DO UPDATE SET"
- + " fileset_name = #{filesetMeta.filesetName},"
- + " metalake_id = #{filesetMeta.metalakeId},"
- + " catalog_id = #{filesetMeta.catalogId},"
- + " schema_id = #{filesetMeta.schemaId},"
- + " type = #{filesetMeta.type},"
- + " audit_info = #{filesetMeta.auditInfo},"
- // PostgreSQL requires the stored row to be qualified on the update
side of ON CONFLICT.
- + " current_version = "
- + META_TABLE_NAME
- + ".current_version + 1,"
- + " last_version = "
- + META_TABLE_NAME
- + ".current_version + 1,"
- + " deleted_at = #{filesetMeta.deletedAt}";
- }
}
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/PolicyMetaPostgreSQLProvider.java
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/PolicyMetaPostgreSQLProvider.java
index e2c5770da2..2bd850d1f3 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/PolicyMetaPostgreSQLProvider.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/PolicyMetaPostgreSQLProvider.java
@@ -26,12 +26,12 @@ import
org.apache.gravitino.storage.relational.mapper.provider.base.PolicyMetaBa
public class PolicyMetaPostgreSQLProvider extends PolicyMetaBaseSQLProvider {
@Override
- public String softDeletePolicyByIdAndVersion(Long policyId, Long
currentVersion) {
+ public String softDeletePolicyByIdAndVersion(Long policyId, Long occVersion)
{
return "UPDATE "
+ POLICY_META_TABLE_NAME
+ " SET deleted_at = "
+ DatabaseTimeSQL.POSTGRESQL
- + " WHERE policy_id = #{policyId} AND current_version =
#{currentVersion}"
+ + " WHERE policy_id = #{policyId} AND occ_version = #{occVersion}"
+ " AND deleted_at = 0";
}
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/po/FilesetPO.java
b/core/src/main/java/org/apache/gravitino/storage/relational/po/FilesetPO.java
index 5a06e29bd4..07a4bdf541 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/po/FilesetPO.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/po/FilesetPO.java
@@ -32,6 +32,7 @@ public class FilesetPO {
private String auditInfo;
private Long currentVersion;
private Long lastVersion;
+ private Long occVersion;
private Long deletedAt;
private List<FilesetVersionPO> filesetVersionPOs;
@@ -71,6 +72,15 @@ public class FilesetPO {
return lastVersion;
}
+ /**
+ * Returns the concurrency token advanced by every metadata update.
+ *
+ * @return the optimistic concurrency version
+ */
+ public Long getOccVersion() {
+ return occVersion;
+ }
+
public Long getDeletedAt() {
return deletedAt;
}
@@ -97,6 +107,7 @@ public class FilesetPO {
&& Objects.equal(getAuditInfo(), filesetPO.getAuditInfo())
&& Objects.equal(getCurrentVersion(), filesetPO.getCurrentVersion())
&& Objects.equal(getLastVersion(), filesetPO.getLastVersion())
+ && Objects.equal(getOccVersion(), filesetPO.getOccVersion())
&& Objects.equal(getDeletedAt(), filesetPO.getDeletedAt())
&& Objects.equal(getFilesetVersionPOs(),
filesetPO.getFilesetVersionPOs());
}
@@ -113,6 +124,7 @@ public class FilesetPO {
getAuditInfo(),
getCurrentVersion(),
getLastVersion(),
+ getOccVersion(),
getDeletedAt(),
getFilesetVersionPOs());
}
@@ -169,6 +181,17 @@ public class FilesetPO {
return this;
}
+ /**
+ * Sets the concurrency token, independently of the content snapshot
version.
+ *
+ * @param occVersion the optimistic concurrency version
+ * @return this builder
+ */
+ public FilesetPO.Builder withOccVersion(Long occVersion) {
+ filesetPO.occVersion = occVersion;
+ return this;
+ }
+
public FilesetPO.Builder withDeletedAt(Long deletedAt) {
filesetPO.deletedAt = deletedAt;
return this;
@@ -203,10 +226,16 @@ public class FilesetPO {
Preconditions.checkArgument(filesetPO.auditInfo != null, "Audit info is
required");
Preconditions.checkArgument(filesetPO.currentVersion != null, "Current
version is required");
Preconditions.checkArgument(filesetPO.lastVersion != null, "Last version
is required");
+ Preconditions.checkArgument(filesetPO.occVersion != null, "OCC version
is required");
Preconditions.checkArgument(filesetPO.deletedAt != null, "Deleted at is
required");
+ // An alter that leaves every stored field untouched allocates no
snapshot and keeps
+ // current_version pointing at the one already stored, so the list is
empty rather than null.
+ //
+ // This used to reject an empty list, which also stopped a create from
storing a row with no
+ // snapshot for current_version to resolve. That case cannot arise: a
create builds one
+ // snapshot per storage location, and FilesetEntity.validate rejects an
entity that has none.
Preconditions.checkArgument(
- filesetPO.filesetVersionPOs != null &&
!filesetPO.filesetVersionPOs.isEmpty(),
- "Fileset version is required");
+ filesetPO.filesetVersionPOs != null, "Fileset version is required");
}
public FilesetPO build() {
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/po/PolicyPO.java
b/core/src/main/java/org/apache/gravitino/storage/relational/po/PolicyPO.java
index 6e2cf0dc81..76e0fb4f9a 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/po/PolicyPO.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/po/PolicyPO.java
@@ -31,6 +31,7 @@ public class PolicyPO {
private String auditInfo;
private Long currentVersion;
private Long lastVersion;
+ private Long occVersion;
private Long deletedAt;
private PolicyVersionPO policyVersionPO;
@@ -54,6 +55,7 @@ public class PolicyPO {
&& Objects.equal(auditInfo, policyPO.auditInfo)
&& Objects.equal(currentVersion, policyPO.currentVersion)
&& Objects.equal(lastVersion, policyPO.lastVersion)
+ && Objects.equal(occVersion, policyPO.occVersion)
&& Objects.equal(policyVersionPO, policyPO.policyVersionPO)
&& Objects.equal(deletedAt, policyPO.deletedAt);
}
@@ -68,6 +70,7 @@ public class PolicyPO {
auditInfo,
currentVersion,
lastVersion,
+ occVersion,
policyVersionPO,
deletedAt);
}
@@ -80,6 +83,7 @@ public class PolicyPO {
private String auditInfo;
private Long currentVersion;
private Long lastVersion;
+ private Long occVersion;
private Long deletedAt;
private PolicyVersionPO policyVersionPO;
@@ -118,6 +122,17 @@ public class PolicyPO {
return this;
}
+ /**
+ * Sets the concurrency token, independently of the content snapshot
version.
+ *
+ * @param occVersion the optimistic concurrency version
+ * @return this builder
+ */
+ public Builder withOccVersion(Long occVersion) {
+ this.occVersion = occVersion;
+ return this;
+ }
+
public Builder withDeletedAt(Long deletedAt) {
this.deletedAt = deletedAt;
return this;
@@ -143,6 +158,7 @@ public class PolicyPO {
policyPO.auditInfo = auditInfo;
policyPO.currentVersion = currentVersion;
policyPO.lastVersion = lastVersion;
+ policyPO.occVersion = occVersion;
policyPO.deletedAt = deletedAt;
policyPO.policyVersionPO = policyVersionPO;
return policyPO;
@@ -155,6 +171,7 @@ public class PolicyPO {
Preconditions.checkArgument(policyType != null, "Policy type is
required");
Preconditions.checkArgument(currentVersion != null, "Current version is
required");
Preconditions.checkArgument(lastVersion != null, "Last version is
required");
+ Preconditions.checkArgument(occVersion != null, "OCC version is
required");
Preconditions.checkArgument(deletedAt != null, "Deleted at is required");
Preconditions.checkArgument(auditInfo != null, "Audit info is required");
Preconditions.checkArgument(policyVersionPO != null, "Policy version is
required");
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/FilesetMetaService.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/FilesetMetaService.java
index 29642f9106..72a0a5be2e 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/FilesetMetaService.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/FilesetMetaService.java
@@ -198,6 +198,8 @@ public class FilesetMetaService {
FilesetVersionMapper.class,
versionMapper ->
versionMapper.selectMaxFilesetVersion(storedPO.getFilesetId()));
+ // storedPO carries no snapshot rows, so an overwrite
always allocates a
+ // new version and writes its snapshot.
FilesetPO replacementPO =
POConverters.updateFilesetPOWithVersion(
storedPO, replacement, maxStoredVersion);
@@ -452,7 +454,7 @@ public class FilesetMetaService {
FilesetMetaMapper.class,
mapper ->
mapper.softDeleteFilesetMetasByFilesetId(
- observedFilesetPO.getFilesetId(),
observedFilesetPO.getCurrentVersion())),
+ observedFilesetPO.getFilesetId(),
observedFilesetPO.getOccVersion())),
() -> filesetWriteFailure(identifier, observedFilesetPO));
}
@@ -463,6 +465,13 @@ public class FilesetMetaService {
return;
}
+ // The snapshot check is in the statement only when the alter allocates a
version, so an alter
+ // that allocates none can have failed for one reason: it lost the OCC
race. Its observed
+ // version is fixed, so retrying would repeat the same comparison and fail
again.
+ if
(newFilesetPO.getCurrentVersion().equals(oldFilesetPO.getCurrentVersion())) {
+ throw filesetWriteFailure(identifier, oldFilesetPO);
+ }
+
// The metadata CAS also rejects a version that already has an active
stored snapshot. Only
// that uncommon legacy case needs the MAX(version) round trip; normal
alters finish above.
Long maxStoredVersion =
@@ -486,9 +495,12 @@ public class FilesetMetaService {
FilesetMetaMapper.class,
mapper -> mapper.updateFilesetMeta(newFilesetPO, oldFilesetPO));
boolean updated = updateCount != null && updateCount > 0;
- if (updated) {
+ if (updated && !newFilesetPO.getFilesetVersionPOs().isEmpty()) {
// The metadata row now points to this complete snapshot. The caller's
schema transaction
// ensures a failed version insert also restores the metadata version.
+ //
+ // An alter that changed nothing the version table stores allocates no
snapshot and leaves
+ // current_version alone, so the row keeps pointing at the snapshot it
already had.
SessionUtils.doWithoutCommit(
FilesetVersionMapper.class,
mapper ->
mapper.insertFilesetVersions(newFilesetPO.getFilesetVersionPOs()));
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/service/PolicyMetaService.java
b/core/src/main/java/org/apache/gravitino/storage/relational/service/PolicyMetaService.java
index 2dc8e1be43..2f30649b0b 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/service/PolicyMetaService.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/service/PolicyMetaService.java
@@ -149,10 +149,7 @@ public class PolicyMetaService {
POConverters.updatePolicyPOWithVersion(oldPolicyPO,
updatedPolicyEntity);
SessionUtils.doMultipleWithCommit(
() -> updatePolicyRootWithVersion(ident, oldPolicyPO, newPolicyPO),
- () ->
- SessionUtils.doWithoutCommit(
- PolicyVersionMapper.class,
- mapper ->
mapper.insertPolicyVersion(newPolicyPO.getPolicyVersionPO())));
+ () -> insertPolicyVersionIfAllocated(oldPolicyPO, newPolicyPO));
} catch (RuntimeException re) {
ExceptionUtils.checkSQLException(
re, Entity.EntityType.POLICY,
updatedPolicyEntity.nameIdentifier().toString());
@@ -380,9 +377,26 @@ public class PolicyMetaService {
NameIdentifier observedIdentifier =
NameIdentifier.of(policyEntity.namespace(),
existingPolicyPO.getPolicyName());
updatePolicyRootWithVersion(observedIdentifier, existingPolicyPO,
replacementPolicyPO);
+ insertPolicyVersionIfAllocated(existingPolicyPO, replacementPolicyPO);
+ }
+
+ /**
+ * Inserts the new content snapshot, unless the write allocated none.
+ *
+ * <p>A write that leaves comment, enabled and content untouched keeps
{@code current_version}
+ * where it is, and the row still points at the snapshot it already had.
Inserting that snapshot
+ * again would collide with the unique key over (policy_id, version,
deleted_at).
+ *
+ * @param oldPolicyPO the row being replaced
+ * @param newPolicyPO the replacement, carrying the snapshot it points at
+ */
+ private void insertPolicyVersionIfAllocated(PolicyPO oldPolicyPO, PolicyPO
newPolicyPO) {
+ if
(newPolicyPO.getCurrentVersion().equals(oldPolicyPO.getCurrentVersion())) {
+ return;
+ }
SessionUtils.doWithoutCommit(
PolicyVersionMapper.class,
- mapper ->
mapper.insertPolicyVersion(replacementPolicyPO.getPolicyVersionPO()));
+ mapper ->
mapper.insertPolicyVersion(newPolicyPO.getPolicyVersionPO()));
}
private void insertNewPolicyWithoutCommit(PolicyPO policyPO) {
@@ -457,7 +471,7 @@ public class PolicyMetaService {
PolicyMetaMapper.class,
mapper ->
mapper.softDeletePolicyByIdAndVersion(
- observedPolicyPO.getPolicyId(),
observedPolicyPO.getCurrentVersion())),
+ observedPolicyPO.getPolicyId(),
observedPolicyPO.getOccVersion())),
() -> policyWriteFailure(identifier, observedPolicyPO));
}
diff --git
a/core/src/main/java/org/apache/gravitino/storage/relational/utils/POConverters.java
b/core/src/main/java/org/apache/gravitino/storage/relational/utils/POConverters.java
index 71025147dc..9dfe048e2c 100644
---
a/core/src/main/java/org/apache/gravitino/storage/relational/utils/POConverters.java
+++
b/core/src/main/java/org/apache/gravitino/storage/relational/utils/POConverters.java
@@ -64,6 +64,7 @@ import org.apache.gravitino.meta.TagEntity;
import org.apache.gravitino.meta.TopicEntity;
import org.apache.gravitino.meta.UserEntity;
import org.apache.gravitino.policy.Policy;
+import org.apache.gravitino.policy.PolicyContent;
import org.apache.gravitino.rel.Column;
import org.apache.gravitino.rel.Table;
import org.apache.gravitino.rel.expressions.Expression;
@@ -687,6 +688,7 @@ public class POConverters {
.withAuditInfo(JsonUtils.anyFieldMapper().writeValueAsString(filesetEntity.auditInfo()))
.withCurrentVersion(INIT_VERSION)
.withLastVersion(INIT_VERSION)
+ .withOccVersion(INIT_VERSION)
.withDeletedAt(DEFAULT_DELETED_AT)
.withFilesetVersionPOs(filesetVersionPOs)
.build();
@@ -696,23 +698,35 @@ public class POConverters {
}
/**
- * Update FilesetPO version
+ * Updates fileset metadata, advancing OCC and reusing a known unchanged
content snapshot.
*
* @param oldFilesetPO the existing {@link FilesetPO} containing the current
and last version data
* @param newFileset the {@link FilesetEntity} with updated metadata and
storage locations
* @param maxStoredVersion the highest version the fileset still has a
stored snapshot for, or
- * {@code null} when it has none
- * @return {@code FilesetPO} object with updated version
+ * {@code null} when it has not been queried or no active snapshot exists
+ * @return the updated fileset row, carrying only newly allocated snapshot
rows
* @throws RuntimeException if JSON serialization of properties fails
*/
public static FilesetPO updateFilesetPOWithVersion(
FilesetPO oldFilesetPO, FilesetEntity newFileset, @Nullable Long
maxStoredVersion) {
try {
- // Every successful fileset alter advances the OCC token. The current
version is also the
- // value used by reads to find the fileset details, so even a rename or
audit-only change
- // needs a complete snapshot at the new version. Alters that change
nothing therefore still
- // write one row per storage location; the version retention job is what
removes them again.
- //
+ // Every successful fileset alter advances the OCC token, which is the
value the CAS compares.
+ Long occVersion = oldFilesetPO.getOccVersion() + 1;
+ String props =
JsonUtils.anyFieldMapper().writeValueAsString(newFileset.properties());
+
+ // The current version is a different thing: it is the join key reads
use to find the fileset
+ // details, so it may only ever point at a version that has a stored
snapshot. An alter that
+ // leaves every stored field untouched, such as a rename or an
audit-only change, therefore
+ // keeps it where it is and writes no snapshot at all.
+ if (filesetSnapshotUnchanged(oldFilesetPO, newFileset, props)) {
+ return newFilesetPOBuilder(oldFilesetPO, newFileset)
+ .withCurrentVersion(oldFilesetPO.getCurrentVersion())
+ .withLastVersion(oldFilesetPO.getLastVersion())
+ .withOccVersion(occVersion)
+ .withFilesetVersionPOs(Collections.emptyList())
+ .build();
+ }
+
// The stored snapshots are taken into account as well, because a
fileset written before the
// version reset was fixed can carry snapshots newer than the version
its metadata row
// records. Starting from the metadata row alone would rebuild a version
that already exists
@@ -723,7 +737,6 @@ public class POConverters {
previousVersion = Math.max(previousVersion, maxStoredVersion);
}
Long currentVersion = previousVersion + 1;
- String props =
JsonUtils.anyFieldMapper().writeValueAsString(newFileset.properties());
List<FilesetVersionPO> newFilesetVersionPOs =
newFileset.storageLocations().entrySet().stream()
.map(
@@ -741,17 +754,10 @@ public class POConverters {
.withDeletedAt(DEFAULT_DELETED_AT)
.build())
.collect(Collectors.toList());
- return FilesetPO.builder()
- .withFilesetId(newFileset.id())
- .withFilesetName(newFileset.name())
- .withMetalakeId(oldFilesetPO.getMetalakeId())
- .withCatalogId(oldFilesetPO.getCatalogId())
- .withSchemaId(oldFilesetPO.getSchemaId())
- .withType(newFileset.filesetType().name())
-
.withAuditInfo(JsonUtils.anyFieldMapper().writeValueAsString(newFileset.auditInfo()))
+ return newFilesetPOBuilder(oldFilesetPO, newFileset)
.withCurrentVersion(currentVersion)
.withLastVersion(currentVersion)
- .withDeletedAt(DEFAULT_DELETED_AT)
+ .withOccVersion(occVersion)
.withFilesetVersionPOs(newFilesetVersionPOs)
.build();
} catch (JsonProcessingException e) {
@@ -759,8 +765,74 @@ public class POConverters {
}
}
+ private static FilesetPO.Builder newFilesetPOBuilder(
+ FilesetPO oldFilesetPO, FilesetEntity newFileset) throws
JsonProcessingException {
+ return FilesetPO.builder()
+ .withFilesetId(newFileset.id())
+ .withFilesetName(newFileset.name())
+ .withMetalakeId(oldFilesetPO.getMetalakeId())
+ .withCatalogId(oldFilesetPO.getCatalogId())
+ .withSchemaId(oldFilesetPO.getSchemaId())
+ .withType(newFileset.filesetType().name())
+
.withAuditInfo(JsonUtils.anyFieldMapper().writeValueAsString(newFileset.auditInfo()))
+ .withDeletedAt(DEFAULT_DELETED_AT);
+ }
+
/**
- * Builds the next complete policy metadata and content snapshot.
+ * Tells whether an alter leaves every field {@code fileset_version_info}
stores untouched.
+ *
+ * <p>Compares exactly the persisted columns: comment, properties, and the
storage locations.
+ * Properties are compared by value, because the same map can serialize in a
different key order
+ * after a read/write round trip. This decides only whether to write a
snapshot, never whether a
+ * concurrent write happened, which is what the OCC version is for; a wrong
{@code false} costs
+ * one redundant snapshot, the behaviour every alter used to have.
+ *
+ * @param oldFilesetPO the row being replaced, carrying the snapshot its
current version points at
+ * @param newFileset the updated fileset
+ * @param newProperties the updated properties, already serialized
+ * @return true when no stored field changed and no new snapshot is needed
+ */
+ private static boolean filesetSnapshotUnchanged(
+ FilesetPO oldFilesetPO, FilesetEntity newFileset, String newProperties) {
+ List<FilesetVersionPO> storedVersions =
oldFilesetPO.getFilesetVersionPOs();
+ if (storedVersions == null || storedVersions.isEmpty()) {
+ // Nothing to point at, so the alter has to write a snapshot whatever it
changed.
+ return false;
+ }
+ Map<String, String> storedLocations =
+ storedVersions.stream()
+ .collect(
+ Collectors.toMap(
+ FilesetVersionPO::getLocationName,
FilesetVersionPO::getStorageLocation));
+ if (!storedLocations.equals(newFileset.storageLocations())) {
+ return false;
+ }
+ // Rows for the same snapshot share the comment and properties; only
locations differ.
+ FilesetVersionPO snapshot = storedVersions.get(0);
+ return Objects.equals(snapshot.getFilesetComment(), newFileset.comment())
+ && filesetPropertiesUnchanged(
+ snapshot.getProperties(), newProperties, newFileset.properties());
+ }
+
+ private static boolean filesetPropertiesUnchanged(
+ String storedProperties, String newProperties, Map<String, String>
newPropertyMap) {
+ if (Objects.equals(storedProperties, newProperties)) {
+ return true;
+ }
+
+ // The alter path copies properties into a HashMap, so the same map can
serialize in a
+ // different key order than the stored snapshot. Compare the maps so a
rename does not
+ // allocate a redundant snapshot.
+ try {
+ return Objects.equals(
+ JsonUtils.anyFieldMapper().readValue(storedProperties, Map.class),
newPropertyMap);
+ } catch (JsonProcessingException e) {
+ throw new RuntimeException("Failed to deserialize fileset properties:",
e);
+ }
+ }
+
+ /**
+ * Updates policy metadata, advancing OCC and reusing a known unchanged
content snapshot.
*
* <p>The row keeps the ID it already has: {@code oldPolicyPO} is the row
being replaced, and its
* ID is what the version snapshots and every relation row point at. An
alter cannot change the
@@ -770,7 +842,7 @@ public class POConverters {
*
* @param oldPolicyPO The policy row observed by the caller.
* @param newPolicy The policy values to persist.
- * @return The policy row and version snapshot at the next monotonic version.
+ * @return The updated policy row and the content snapshot it points at.
*/
public static PolicyPO updatePolicyPOWithVersion(PolicyPO oldPolicyPO,
PolicyEntity newPolicy) {
try {
@@ -788,7 +860,7 @@ public class POConverters {
}
/**
- * Builds the next policy version from values that were serialized before
acquiring a row lock.
+ * Updates a policy from values that were serialized before acquiring a row
lock.
*
* <p>This overload is used by overwrite: the initialized replacement
already contains the
* serialized audit and content values, so advancing the locked row does not
repeat CPU-bound JSON
@@ -796,7 +868,7 @@ public class POConverters {
*
* @param oldPolicyPO The locked policy row being replaced.
* @param replacementPolicyPO The initialized replacement values.
- * @return The policy row and version snapshot at the next monotonic version.
+ * @return The updated policy row and the content snapshot it points at.
*/
public static PolicyPO updatePolicyPOWithVersion(
PolicyPO oldPolicyPO, PolicyPO replacementPolicyPO) {
@@ -1529,6 +1601,7 @@ public class POConverters {
.withAuditInfo(JsonUtils.anyFieldMapper().writeValueAsString(policyEntity.auditInfo()))
.withCurrentVersion(INIT_VERSION)
.withLastVersion(INIT_VERSION)
+ .withOccVersion(INIT_VERSION)
.withDeletedAt(DEFAULT_DELETED_AT)
.withPolicyVersionPO(policyVersionPO)
.build();
@@ -1794,6 +1867,33 @@ public class POConverters {
String policyComment,
boolean enabled,
String content) {
+ // Every successful policy alter advances the OCC token, which is the
value the CAS compares.
+ Long occVersion = oldPolicyPO.getOccVersion() + 1;
+
+ // The current version is the join key reads use to find the policy
content, so it may only
+ // point at a version that has a stored snapshot. An alter that leaves
comment, enabled and
+ // content untouched keeps it where it is and writes no snapshot, so the
row keeps pointing at
+ // the one it already has.
+ PolicyVersionPO storedVersionPO = oldPolicyPO.getPolicyVersionPO();
+ if (storedVersionPO != null
+ && Objects.equals(storedVersionPO.getPolicyComment(), policyComment)
+ && storedVersionPO.isEnabled() == enabled
+ && Objects.equals(oldPolicyPO.getPolicyType(), policyType)
+ && policyContentUnchanged(storedVersionPO.getContent(), content,
policyType)) {
+ return PolicyPO.builder()
+ .withPolicyId(oldPolicyPO.getPolicyId())
+ .withPolicyName(policyName)
+ .withPolicyType(policyType)
+ .withMetalakeId(oldPolicyPO.getMetalakeId())
+ .withAuditInfo(auditInfo)
+ .withCurrentVersion(oldPolicyPO.getCurrentVersion())
+ .withLastVersion(oldPolicyPO.getLastVersion())
+ .withOccVersion(occVersion)
+ .withDeletedAt(DEFAULT_DELETED_AT)
+ .withPolicyVersionPO(storedVersionPO)
+ .build();
+ }
+
Long nextVersion = Math.max(oldPolicyPO.getCurrentVersion(),
oldPolicyPO.getLastVersion()) + 1;
PolicyVersionPO newPolicyVersionPO =
PolicyVersionPO.builder()
@@ -1813,11 +1913,31 @@ public class POConverters {
.withAuditInfo(auditInfo)
.withCurrentVersion(nextVersion)
.withLastVersion(nextVersion)
+ .withOccVersion(occVersion)
.withDeletedAt(DEFAULT_DELETED_AT)
.withPolicyVersionPO(newPolicyVersionPO)
.build();
}
+ private static boolean policyContentUnchanged(
+ String storedContent, String newContent, String policyType) {
+ if (Objects.equals(storedContent, newContent)) {
+ return true;
+ }
+
+ // Sets in policy content can serialize in a different order after a
read/write round trip.
+ // Compare the content objects so an audit-only update does not allocate a
redundant snapshot.
+ try {
+ Class<? extends PolicyContent> contentClass =
+ Policy.BuiltInType.fromPolicyType(policyType).contentClass();
+ return Objects.equals(
+ JsonUtils.anyFieldMapper().readValue(storedContent, contentClass),
+ JsonUtils.anyFieldMapper().readValue(newContent, contentClass));
+ } catch (JsonProcessingException e) {
+ throw new RuntimeException("Failed to deserialize policy content:", e);
+ }
+ }
+
private static ModelVersionAliasRelPO createAliasRelPO(Long modelId, int
version, String alias) {
return ModelVersionAliasRelPO.builder()
.withModelVersion(version)
diff --git
a/core/src/test/java/org/apache/gravitino/storage/TestSQLScripts.java
b/core/src/test/java/org/apache/gravitino/storage/TestSQLScripts.java
index c758614778..0997134152 100644
--- a/core/src/test/java/org/apache/gravitino/storage/TestSQLScripts.java
+++ b/core/src/test/java/org/apache/gravitino/storage/TestSQLScripts.java
@@ -201,6 +201,60 @@ public class TestSQLScripts extends TestJDBCBackend {
}
}
+ /** Verifies the OCC backfill preserves existing history versions on live
and deleted rows. */
+ @TestTemplate
+ public void testUpgradeToTwoZeroBackfillsOccVersions() throws SQLException,
IOException {
+ String gravitinoHome = System.getenv("GRAVITINO_HOME");
+ Assertions.assertNotNull(gravitinoHome, "GRAVITINO_HOME environment
variable is not set");
+ Path scriptDir = Path.of(gravitinoHome, "scripts",
backendType.toLowerCase());
+ String suffix = "-" + backendType.toLowerCase() + ".sql";
+ dropAllTables();
+ executeScript(scriptDir.resolve("schema-1.3.0" + suffix).toFile());
+
+ try (SqlSession sqlSession =
+
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+ Connection connection = sqlSession.getConnection();
+ Statement statement = connection.createStatement()) {
+ statement.execute(
+ "INSERT INTO fileset_meta (fileset_id, fileset_name, metalake_id,
catalog_id, schema_id,"
+ + " type, audit_info, current_version, last_version, deleted_at)
VALUES"
+ + " (1, 'live', 1, 1, 1, 'MANAGED', '{}', 7, 9, 0),"
+ + " (2, 'deleted', 1, 1, 1, 'MANAGED', '{}', 5, 5, 100)");
+ statement.execute(
+ "INSERT INTO policy_meta (policy_id, policy_name, policy_type,
metalake_id, audit_info,"
+ + " current_version, last_version, deleted_at) VALUES"
+ + " (1, 'live', 'custom', 1, '{}', 7, 9, 0),"
+ + " (2, 'deleted', 'custom', 1, '{}', 5, 5, 100)");
+ }
+
+ executeScript(scriptDir.resolve("upgrade-1.3.0-to-2.0.0" +
suffix).toFile());
+
+ try (SqlSession sqlSession =
+
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+ Connection connection = sqlSession.getConnection();
+ Statement statement = connection.createStatement()) {
+ for (String table : List.of("fileset_meta", "policy_meta")) {
+ try (ResultSet rows =
+ statement.executeQuery(
+ "SELECT current_version, last_version, occ_version, deleted_at
FROM "
+ + table
+ + " ORDER BY deleted_at")) {
+ Assertions.assertTrue(rows.next());
+ Assertions.assertEquals(7, rows.getLong("current_version"));
+ Assertions.assertEquals(9, rows.getLong("last_version"));
+ Assertions.assertEquals(1, rows.getLong("occ_version"));
+ Assertions.assertEquals(0, rows.getLong("deleted_at"));
+ Assertions.assertTrue(rows.next());
+ Assertions.assertEquals(5, rows.getLong("current_version"));
+ Assertions.assertEquals(5, rows.getLong("last_version"));
+ Assertions.assertEquals(1, rows.getLong("occ_version"));
+ Assertions.assertEquals(100, rows.getLong("deleted_at"));
+ Assertions.assertFalse(rows.next());
+ }
+ }
+ }
+ }
+
private void executeScript(File scriptFile) throws IOException, SQLException
{
List<String> ddls = extractStatements(scriptFile.toPath());
try (SqlSession sqlSession =
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/TestJDBCBackend.java
b/core/src/test/java/org/apache/gravitino/storage/relational/TestJDBCBackend.java
index fa1c2fc554..5c878525c2 100644
---
a/core/src/test/java/org/apache/gravitino/storage/relational/TestJDBCBackend.java
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/TestJDBCBackend.java
@@ -295,6 +295,32 @@ public abstract class TestJDBCBackend {
return versionDeletedTime;
}
+ /**
+ * Counts the rows a fileset holds in {@code fileset_version_info}.
+ *
+ * <p>A snapshot is one row per storage location, so the row count, unlike
the number of distinct
+ * versions, shows what an alter that allocates a version actually costs.
+ *
+ * @param filesetId the fileset ID
+ * @return the number of active and deleted version rows
+ */
+ protected int countFilesetVersionRows(Long filesetId) {
+ try (SqlSession sqlSession =
+
SqlSessionFactoryHelper.getInstance().getSqlSessionFactory().openSession(true);
+ Connection connection = sqlSession.getConnection();
+ Statement statement = connection.createStatement();
+ ResultSet rs =
+ statement.executeQuery(
+ String.format(
+ "SELECT COUNT(*) AS row_count FROM fileset_version_info
WHERE fileset_id = %d",
+ filesetId))) {
+ rs.next();
+ return rs.getInt("row_count");
+ } catch (SQLException e) {
+ throw new RuntimeException("SQL execution failed", e);
+ }
+ }
+
protected Map<Integer, Long> listPolicyVersions(Long policyId) {
Map<Integer, Long> versionDeletedTime = new HashMap<>();
try (SqlSession sqlSession =
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/mapper/provider/base/TestFilesetMetaBaseSQLProvider.java
b/core/src/test/java/org/apache/gravitino/storage/relational/mapper/provider/base/TestFilesetMetaBaseSQLProvider.java
index cdbfcffd37..23e45b41a0 100644
---
a/core/src/test/java/org/apache/gravitino/storage/relational/mapper/provider/base/TestFilesetMetaBaseSQLProvider.java
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/mapper/provider/base/TestFilesetMetaBaseSQLProvider.java
@@ -18,35 +18,29 @@
*/
package org.apache.gravitino.storage.relational.mapper.provider.base;
+import org.apache.gravitino.storage.relational.po.FilesetPO;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
+import org.mockito.Mockito;
class TestFilesetMetaBaseSQLProvider {
private static final FilesetMetaBaseSQLProvider PROVIDER = new
FilesetMetaBaseSQLProvider();
- @Test
- void testOverwriteAdvancesStoredVersion() {
- String sql = PROVIDER.insertFilesetMetaOnDuplicateKeyUpdate(null);
- String updateClause = sql.substring(sql.indexOf(" ON DUPLICATE KEY
UPDATE"));
-
- Assertions.assertTrue(updateClause.contains("last_version =
current_version + 1"));
- Assertions.assertTrue(updateClause.contains("current_version =
current_version + 1"));
- Assertions.assertTrue(
- updateClause.indexOf("last_version =") <
updateClause.indexOf("current_version ="));
- Assertions.assertFalse(
- updateClause.contains("current_version =
#{filesetMeta.currentVersion}"));
- Assertions.assertFalse(updateClause.contains("last_version =
#{filesetMeta.lastVersion}"));
- }
-
@Test
void testUpdateUsesVersionCasAndRejectsAnOccupiedSnapshotVersion() {
- String sql = PROVIDER.updateFilesetMeta(null, null);
+ // An alter that allocates a new version keeps the snapshot check.
+ FilesetPO allocating = Mockito.mock(FilesetPO.class);
+ Mockito.when(allocating.getCurrentVersion()).thenReturn(4L);
+ FilesetPO stored = Mockito.mock(FilesetPO.class);
+ Mockito.when(stored.getCurrentVersion()).thenReturn(3L);
+
+ String sql = PROVIDER.updateFilesetMeta(allocating, stored);
String whereClause = sql.substring(sql.indexOf(" WHERE"));
Assertions.assertEquals(
" WHERE fileset_id = #{oldFilesetMeta.filesetId}"
- + " AND current_version = #{oldFilesetMeta.currentVersion}"
+ + " AND occ_version = #{oldFilesetMeta.occVersion}"
+ " AND deleted_at = 0"
+ " AND NOT EXISTS (SELECT 1 FROM fileset_version_info fv"
+ " WHERE fv.fileset_id = #{oldFilesetMeta.filesetId}"
@@ -55,11 +49,31 @@ class TestFilesetMetaBaseSQLProvider {
whereClause);
}
+ @Test
+ void testUpdateDropsTheSnapshotCheckWhenNoVersionIsAllocated() {
+ // An alter that changes nothing the version table stores keeps
current_version where it is,
+ // and the snapshot it points at is supposed to exist, so the check would
reject every such
+ // alter.
+ FilesetPO unchanged = Mockito.mock(FilesetPO.class);
+ Mockito.when(unchanged.getCurrentVersion()).thenReturn(3L);
+ FilesetPO stored = Mockito.mock(FilesetPO.class);
+ Mockito.when(stored.getCurrentVersion()).thenReturn(3L);
+
+ String sql = PROVIDER.updateFilesetMeta(unchanged, stored);
+
+ Assertions.assertFalse(sql.contains("NOT EXISTS"));
+ Assertions.assertTrue(
+ sql.endsWith(
+ " WHERE fileset_id = #{oldFilesetMeta.filesetId}"
+ + " AND occ_version = #{oldFilesetMeta.occVersion}"
+ + " AND deleted_at = 0"));
+ }
+
@Test
void testDirectDeleteUsesVersionCas() {
String sql = PROVIDER.softDeleteFilesetMetasByFilesetId(null, null);
- Assertions.assertTrue(sql.contains("AND current_version =
#{currentVersion}"));
+ Assertions.assertTrue(sql.contains("AND occ_version = #{occVersion}"));
Assertions.assertTrue(sql.endsWith("AND deleted_at = 0"));
}
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TestFilesetMetaPostgreSQLProvider.java
b/core/src/test/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TestFilesetMetaPostgreSQLProvider.java
index 49753bfdea..35a879b26c 100644
---
a/core/src/test/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TestFilesetMetaPostgreSQLProvider.java
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/mapper/provider/postgresql/TestFilesetMetaPostgreSQLProvider.java
@@ -18,34 +18,16 @@
*/
package org.apache.gravitino.storage.relational.mapper.provider.postgresql;
-import org.apache.gravitino.storage.relational.mapper.FilesetMetaMapper;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
class TestFilesetMetaPostgreSQLProvider {
- @Test
- void testOverwriteAdvancesStoredVersion() {
- String sql = new
FilesetMetaPostgreSQLProvider().insertFilesetMetaOnDuplicateKeyUpdate(null);
- String conflictClause = sql.substring(sql.indexOf(" ON CONFLICT"));
-
- Assertions.assertTrue(
- conflictClause.startsWith(" ON CONFLICT(schema_id, fileset_name,
deleted_at)"));
- Assertions.assertTrue(
- conflictClause.contains(
- "current_version = " + FilesetMetaMapper.META_TABLE_NAME +
".current_version + 1"));
- Assertions.assertTrue(
- conflictClause.contains(
- "last_version = " + FilesetMetaMapper.META_TABLE_NAME +
".current_version + 1"));
-
Assertions.assertFalse(conflictClause.contains("#{filesetMeta.currentVersion}"));
-
Assertions.assertFalse(conflictClause.contains("#{filesetMeta.lastVersion}"));
- }
-
@Test
void testDirectDeleteUsesVersionCas() {
String sql = new
FilesetMetaPostgreSQLProvider().softDeleteFilesetMetasByFilesetId(null, null);
- Assertions.assertTrue(sql.contains("AND current_version =
#{currentVersion}"));
+ Assertions.assertTrue(sql.contains("AND occ_version = #{occVersion}"));
Assertions.assertTrue(sql.endsWith("AND deleted_at = 0"));
}
}
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestFilesetMetaService.java
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestFilesetMetaService.java
index 737e998520..bb95266903 100644
---
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestFilesetMetaService.java
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestFilesetMetaService.java
@@ -21,6 +21,7 @@ package org.apache.gravitino.storage.relational.service;
import static org.apache.gravitino.file.Fileset.LOCATION_NAME_UNKNOWN;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNotEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -33,6 +34,7 @@ import java.sql.ResultSet;
import java.sql.SQLException;
import java.sql.Statement;
import java.time.Instant;
+import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
@@ -53,6 +55,7 @@ import org.apache.gravitino.exceptions.NoSuchEntityException;
import org.apache.gravitino.exceptions.OptimisticLockException;
import org.apache.gravitino.file.Fileset;
import org.apache.gravitino.integration.test.util.GravitinoITUtils;
+import org.apache.gravitino.json.JsonUtils;
import org.apache.gravitino.meta.AuditInfo;
import org.apache.gravitino.meta.FilesetEntity;
import org.apache.gravitino.meta.SchemaEntity;
@@ -401,6 +404,177 @@ public class TestFilesetMetaService extends
TestJDBCBackend {
.build();
}
+ @TestTemplate
+ public void testAlterWritesASnapshotOnlyWhenStoredContentChanges() throws
IOException {
+ String filesetName =
GravitinoITUtils.genRandomName("tst_fs_snapshot_cost");
+ NameIdentifier filesetIdent =
+ NameIdentifier.of(metalakeName, catalogName, schemaName, filesetName);
+ // Two storage locations, because a snapshot costs one row per location:
that multiplier is
+ // what made a no-op alter expensive.
+ FilesetEntity fileset =
+ FilesetEntity.builder()
+ .withId(RandomIdGenerator.INSTANCE.nextId())
+ .withName(filesetName)
+ .withNamespace(NamespaceUtil.ofFileset(metalakeName, catalogName,
schemaName))
+ .withFilesetType(Fileset.Type.MANAGED)
+ .withStorageLocations(ImmutableMap.of("first", "/tmp-a", "second",
"/tmp-b"))
+ .withComment("comment-v1")
+ .withProperties(ImmutableMap.of("k", "v1"))
+ .withAuditInfo(AUDIT_INFO)
+ .build();
+ FilesetMetaService.getInstance().insertFileset(fileset, false);
+ FilesetPO initialPO = getFilesetPO(fileset.id());
+ assertEquals(1, listFilesetVersions(fileset.id()).size());
+ assertEquals(2, countFilesetVersionRows(fileset.id()));
+
+ // A rename touches nothing fileset_version_info stores.
+ String renamed = filesetName + "_renamed";
+ FilesetMetaService.getInstance()
+ .updateFileset(
+ filesetIdent,
+ entity -> {
+ FilesetEntity current = (FilesetEntity) entity;
+ return FilesetEntity.builder()
+ .withId(current.id())
+ .withName(renamed)
+ .withNamespace(current.namespace())
+ .withFilesetType(current.filesetType())
+ .withStorageLocations(current.storageLocations())
+ .withComment(current.comment())
+ .withProperties(current.properties())
+ .withAuditInfo(current.auditInfo())
+ .build();
+ });
+
+ FilesetPO afterRename = getFilesetPO(fileset.id());
+ assertEquals(renamed, afterRename.getFilesetName());
+ assertEquals(initialPO.getOccVersion() + 1,
afterRename.getOccVersion().longValue());
+ assertEquals(initialPO.getCurrentVersion(),
afterRename.getCurrentVersion());
+ assertEquals(initialPO.getLastVersion(), afterRename.getLastVersion());
+ assertEquals(1, listFilesetVersions(fileset.id()).size());
+ assertEquals(2, countFilesetVersionRows(fileset.id()));
+
+ // A second rename still writes no snapshot, while every write advances
the OCC token.
+ String renamedAgain = renamed + "_again";
+ FilesetMetaService.getInstance()
+ .updateFileset(
+ NameIdentifier.of(metalakeName, catalogName, schemaName, renamed),
+ entity -> {
+ FilesetEntity current = (FilesetEntity) entity;
+ return FilesetEntity.builder()
+ .withId(current.id())
+ .withName(renamedAgain)
+ .withNamespace(current.namespace())
+ .withFilesetType(current.filesetType())
+ .withStorageLocations(current.storageLocations())
+ .withComment(current.comment())
+ .withProperties(current.properties())
+ .withAuditInfo(current.auditInfo())
+ .build();
+ });
+ FilesetPO afterSecondRename = getFilesetPO(fileset.id());
+ assertEquals(afterRename.getOccVersion() + 1,
afterSecondRename.getOccVersion().longValue());
+ assertEquals(afterRename.getCurrentVersion(),
afterSecondRename.getCurrentVersion());
+ assertEquals(afterRename.getLastVersion(),
afterSecondRename.getLastVersion());
+ assertEquals(1, listFilesetVersions(fileset.id()).size());
+ assertEquals(2, countFilesetVersionRows(fileset.id()));
+
+ // Changing the comment does change stored content, so it allocates a
snapshot per location.
+ NameIdentifier renamedIdent =
+ NameIdentifier.of(metalakeName, catalogName, schemaName, renamedAgain);
+ FilesetMetaService.getInstance()
+ .updateFileset(
+ renamedIdent,
+ entity -> {
+ FilesetEntity current = (FilesetEntity) entity;
+ return FilesetEntity.builder()
+ .withId(current.id())
+ .withName(current.name())
+ .withNamespace(current.namespace())
+ .withFilesetType(current.filesetType())
+ .withStorageLocations(current.storageLocations())
+ .withComment("comment-v2")
+ .withProperties(current.properties())
+ .withAuditInfo(current.auditInfo())
+ .build();
+ });
+
+ FilesetPO afterComment = getFilesetPO(fileset.id());
+ assertEquals(afterSecondRename.getOccVersion() + 1,
afterComment.getOccVersion().longValue());
+ assertEquals(
+ afterSecondRename.getCurrentVersion() + 1,
afterComment.getCurrentVersion().longValue());
+ assertEquals(afterComment.getCurrentVersion(),
afterComment.getLastVersion());
+ assertEquals(2, listFilesetVersions(fileset.id()).size());
+ assertEquals(4, countFilesetVersionRows(fileset.id()));
+
+ // Whatever the mix of alters, the row still resolves to exactly one
active snapshot.
+ FilesetEntity readBack =
FilesetMetaService.getInstance().getFilesetByIdentifier(renamedIdent);
+ assertEquals("comment-v2", readBack.comment());
+ assertEquals("/tmp-a", readBack.storageLocations().get("first"));
+ assertEquals("/tmp-b", readBack.storageLocations().get("second"));
+ }
+
+ @TestTemplate
+ public void testRenameWithReorderedPropertiesWritesNoSnapshot() throws
IOException {
+ String filesetName =
GravitinoITUtils.genRandomName("tst_fs_reordered_props");
+ NameIdentifier filesetIdent =
+ NameIdentifier.of(metalakeName, catalogName, schemaName, filesetName);
+ // FilesetCatalogOperations copies the stored properties into a HashMap
before an alter. For
+ // these keys the HashMap iterates in a different order than the stored
one.
+ Map<String, String> properties =
+ ImmutableMap.of("team", "data", "retention", "7d",
"gravitino.identifier", "id");
+ Map<String, String> reordered = new HashMap<>(properties);
+ assertNotEquals(
+ JsonUtils.anyFieldMapper().writeValueAsString(properties),
+ JsonUtils.anyFieldMapper().writeValueAsString(reordered),
+ "The test needs keys whose HashMap order differs from the stored
order");
+
+ FilesetEntity fileset =
+ FilesetEntity.builder()
+ .withId(RandomIdGenerator.INSTANCE.nextId())
+ .withName(filesetName)
+ .withNamespace(NamespaceUtil.ofFileset(metalakeName, catalogName,
schemaName))
+ .withFilesetType(Fileset.Type.MANAGED)
+ .withStorageLocations(ImmutableMap.of("first", "/tmp-a", "second",
"/tmp-b"))
+ .withComment("comment")
+ .withProperties(properties)
+ .withAuditInfo(AUDIT_INFO)
+ .build();
+ FilesetMetaService.getInstance().insertFileset(fileset, false);
+ FilesetPO initialPO = getFilesetPO(fileset.id());
+
+ String renamed = filesetName + "_renamed";
+ FilesetMetaService.getInstance()
+ .updateFileset(
+ filesetIdent,
+ entity -> {
+ FilesetEntity current = (FilesetEntity) entity;
+ return FilesetEntity.builder()
+ .withId(current.id())
+ .withName(renamed)
+ .withNamespace(current.namespace())
+ .withFilesetType(current.filesetType())
+ .withStorageLocations(current.storageLocations())
+ .withComment(current.comment())
+ .withProperties(new HashMap<>(current.properties()))
+ .withAuditInfo(current.auditInfo())
+ .build();
+ });
+
+ FilesetPO afterRename = getFilesetPO(fileset.id());
+ assertEquals(renamed, afterRename.getFilesetName());
+ assertEquals(initialPO.getOccVersion() + 1,
afterRename.getOccVersion().longValue());
+ assertEquals(initialPO.getCurrentVersion(),
afterRename.getCurrentVersion());
+ assertEquals(1, listFilesetVersions(fileset.id()).size());
+ assertEquals(2, countFilesetVersionRows(fileset.id()));
+ assertEquals(
+ properties,
+ FilesetMetaService.getInstance()
+ .getFilesetByIdentifier(
+ NameIdentifier.of(metalakeName, catalogName, schemaName,
renamed))
+ .properties());
+ }
+
@TestTemplate
public void testAlterReportsOptimisticLockConflictAndKeepsWinnerVersion()
throws IOException {
String filesetName = GravitinoITUtils.genRandomName("tst_fs_conflict");
@@ -445,7 +619,7 @@ public class TestFilesetMetaService extends TestJDBCBackend
{
filesetIdent,
e -> {
// Commit another alter after the outer call has read
its snapshot. The
- // outer write must then lose the current_version
comparison.
+ // outer write must then lose the occ_version comparison.
updateFilesetUnchecked(
filesetIdent,
entity ->
@@ -466,10 +640,12 @@ public class TestFilesetMetaService extends
TestJDBCBackend {
Assertions.assertEquals("/tmp",
persistedEntity.storageLocations().get(LOCATION_NAME_UNKNOWN));
Assertions.assertNotEquals(updatedFilesetEntity, persistedEntity);
FilesetPO currentPO = getFilesetPO(filesetEntity.id());
- Assertions.assertEquals(
- initialPO.getCurrentVersion() + 1,
currentPO.getCurrentVersion().longValue());
+ // The alter that won changed only the audit info, which
fileset_version_info does not store,
+ // so it advanced the OCC token alone and the row still points at its
original snapshot.
+ Assertions.assertEquals(initialPO.getOccVersion() + 1,
currentPO.getOccVersion().longValue());
+ Assertions.assertEquals(initialPO.getCurrentVersion(),
currentPO.getCurrentVersion());
Assertions.assertEquals(currentPO.getCurrentVersion(),
currentPO.getLastVersion());
- Assertions.assertEquals(2, listFilesetVersions(filesetEntity.id()).size());
+ Assertions.assertEquals(1, listFilesetVersions(filesetEntity.id()).size());
}
@TestTemplate
@@ -653,6 +829,100 @@ public class TestFilesetMetaService extends
TestJDBCBackend {
assertVersionActive(versions, 2);
}
+ /** Verifies that an audit-only winner fences a stale audit-only update
without a new snapshot. */
+ @TestTemplate
+ public void testMetadataOnlyAlterRejectsAStaleUpdate() throws IOException {
+ FilesetEntity fileset =
+ createFilesetEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ NamespaceUtil.ofFileset(metalakeName, catalogName, schemaName),
+ GravitinoITUtils.genRandomName("tst_fs_metadata_conflict"),
+ AUDIT_INFO,
+ "/tmp");
+ FilesetMetaService service = FilesetMetaService.getInstance();
+ service.insertFileset(fileset, false);
+ FilesetPO initialPO = getFilesetPO(fileset.id());
+ AuditInfo winningAudit =
+
AuditInfo.builder().withCreator("winning-updater").withCreateTime(Instant.now()).build();
+ AuditInfo staleAudit =
+
AuditInfo.builder().withCreator("stale-updater").withCreateTime(Instant.now()).build();
+
+ assertThrows(
+ OptimisticLockException.class,
+ () ->
+ updateFilesetUnchecked(
+ fileset.nameIdentifier(),
+ current -> {
+ // Commit the winner after the outer update reads, before
its CAS executes.
+ updateFilesetUnchecked(
+ current.nameIdentifier(),
+ winner ->
+ copyFileset(
+ winner,
+ winner.id(),
+ winner.name(),
+ winner.comment(),
+ "/tmp",
+ winningAudit));
+ return copyFileset(
+ current, current.id(), current.name(),
current.comment(), "/tmp", staleAudit);
+ }));
+
+ FilesetEntity stored =
service.getFilesetByIdentifier(fileset.nameIdentifier());
+ assertEquals(winningAudit, stored.auditInfo());
+ assertEquals(fileset.comment(), stored.comment());
+ FilesetPO afterConflict = getFilesetPO(fileset.id());
+ assertEquals(initialPO.getOccVersion() + 1,
afterConflict.getOccVersion().longValue());
+ assertEquals(initialPO.getCurrentVersion(),
afterConflict.getCurrentVersion());
+ assertEquals(initialPO.getLastVersion(), afterConflict.getLastVersion());
+ assertEquals(1, countFilesetVersionRows(fileset.id()));
+ }
+
+ @TestTemplate
+ public void testDeleteRejectsAStaleVersionAfterAMetadataOnlyAlter() throws
IOException {
+ String filesetName = GravitinoITUtils.genRandomName("tst_fs_stale_delete");
+ FilesetEntity fileset =
+ createFilesetEntity(
+ RandomIdGenerator.INSTANCE.nextId(),
+ NamespaceUtil.ofFileset(metalakeName, catalogName, schemaName),
+ filesetName,
+ AUDIT_INFO,
+ "/tmp");
+ FilesetMetaService.getInstance().insertFileset(fileset, false);
+ FilesetPO stalePO = getFilesetPO(fileset.id());
+
+ // An audit-only alter advances occ_version and deliberately leaves
current_version alone. A
+ // drop still guarded by current_version would not notice it and would
delete a fileset the
+ // caller never observed in its current state.
+ AuditInfo laterAudit =
+
AuditInfo.builder().withCreator("later-updater").withCreateTime(Instant.now()).build();
+ updateFilesetUnchecked(
+ fileset.nameIdentifier(),
+ entity ->
+ createFilesetEntity(
+ entity.id(), entity.namespace(), entity.name(), laterAudit,
"/tmp"));
+
+ FilesetPO afterAlter = getFilesetPO(fileset.id());
+ Assertions.assertEquals(stalePO.getCurrentVersion(),
afterAlter.getCurrentVersion());
+ Assertions.assertEquals(stalePO.getOccVersion() + 1,
afterAlter.getOccVersion().longValue());
+
+ Assertions.assertThrows(
+ OptimisticLockException.class,
+ () ->
+ SessionUtils.doMultipleWithCommit(
+ () ->
+ FilesetMetaService.getInstance()
+ .deleteFilesetWithVersion(fileset.nameIdentifier(),
stalePO)));
+
+ // The stale drop stopped before removing anything.
+ Assertions.assertEquals(
+ laterAudit.creator(),
+ FilesetMetaService.getInstance()
+ .getFilesetByIdentifier(fileset.nameIdentifier())
+ .auditInfo()
+ .creator());
+ }
+
@TestTemplate
public void testDeleteReportsNoSuchWhenDeletedConcurrently() throws
IOException {
String filesetName =
GravitinoITUtils.genRandomName("tst_fs_double_delete");
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestPolicyMetaService.java
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestPolicyMetaService.java
index fb523b7f3a..06b44444f6 100644
---
a/core/src/test/java/org/apache/gravitino/storage/relational/service/TestPolicyMetaService.java
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/service/TestPolicyMetaService.java
@@ -18,8 +18,10 @@
*/
package org.apache.gravitino.storage.relational.service;
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNotEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -433,7 +435,7 @@ public class TestPolicyMetaService extends TestJDBCBackend {
}
@TestTemplate
- public void testMetadataOnlyPolicyAlterCreatesCompleteSnapshot() throws
IOException {
+ public void testMetadataOnlyPolicyAlterAdvancesOnlyTheOccVersion() throws
IOException {
createAndInsertMakeLake(METALAKE_NAME);
PolicyMetaService policyMetaService = PolicyMetaService.getInstance();
PolicyEntity policy =
@@ -452,14 +454,20 @@ public class TestPolicyMetaService extends
TestJDBCBackend {
policyMetaService.updatePolicy(policy.nameIdentifier(), ignored ->
metadataOnlyUpdate);
PolicyPO updatedPO = getPolicyPO(policy.nameIdentifier());
- assertEquals(initialPO.getCurrentVersion() + 1,
updatedPO.getCurrentVersion().longValue());
- assertEquals(updatedPO.getCurrentVersion(), updatedPO.getLastVersion());
+ // The audit info is the only thing that changed, and policy_version_info
does not store it, so
+ // the alter advances the OCC token alone and writes no snapshot.
+ assertEquals(initialPO.getOccVersion() + 1,
updatedPO.getOccVersion().longValue());
+ assertEquals(initialPO.getCurrentVersion(), updatedPO.getCurrentVersion());
+ assertEquals(initialPO.getLastVersion(), updatedPO.getLastVersion());
+ assertNotEquals(initialPO.getAuditInfo(), updatedPO.getAuditInfo());
+
+ // The row still points at the snapshot it already had, and reads still
resolve it.
assertEquals(updatedPO.getCurrentVersion(),
updatedPO.getPolicyVersionPO().getVersion());
assertEquals(policy.comment(),
updatedPO.getPolicyVersionPO().getPolicyComment());
assertEquals(policy.enabled(), updatedPO.getPolicyVersionPO().isEnabled());
assertEquals(
initialPO.getPolicyVersionPO().getContent(),
updatedPO.getPolicyVersionPO().getContent());
- assertEquals(2, listPolicyVersions(policy.id()).size());
+ assertEquals(1, listPolicyVersions(policy.id()).size());
}
@TestTemplate
@@ -475,12 +483,29 @@ public class TestPolicyMetaService extends
TestJDBCBackend {
policyMetaService.insertPolicy(policy, false);
PolicyPO initialPO = getPolicyPO(policy.nameIdentifier());
+ PolicyEntity metadataUpdate =
+ copyPolicy(
+ policy,
+ policy.name(),
+ policy.comment(),
+ AuditInfo.builder()
+ .withCreator("updated-creator")
+ .withCreateTime(Instant.now())
+ .build());
+ policyMetaService.updatePolicy(policy.nameIdentifier(), ignored ->
metadataUpdate);
+ PolicyPO afterMetadataUpdate = getPolicyPO(policy.nameIdentifier());
+ assertEquals(initialPO.getCurrentVersion(),
afterMetadataUpdate.getCurrentVersion());
+ assertEquals(initialPO.getLastVersion(),
afterMetadataUpdate.getLastVersion());
+ assertEquals(initialPO.getOccVersion() + 1,
afterMetadataUpdate.getOccVersion().longValue());
+
PolicyEntity replacement = copyPolicy(policy,
"policy_overwrite_occ_renamed", "replacement");
policyMetaService.insertPolicy(replacement, true);
PolicyPO overwrittenPO = getPolicyPO(replacement.nameIdentifier());
assertEquals(initialPO.getCurrentVersion() + 1,
overwrittenPO.getCurrentVersion().longValue());
assertEquals(overwrittenPO.getCurrentVersion(),
overwrittenPO.getLastVersion());
+ assertEquals(
+ afterMetadataUpdate.getOccVersion() + 1,
overwrittenPO.getOccVersion().longValue());
assertEquals(2, listPolicyVersions(policy.id()).size());
assertEquals(
replacement,
policyMetaService.getPolicyByIdentifier(replacement.nameIdentifier()));
@@ -560,6 +585,92 @@ public class TestPolicyMetaService extends TestJDBCBackend
{
listPolicyVersions(policy.id()).values().stream().filter(v ->
v.longValue() == 0L).count());
}
+ /** Verifies that an audit-only winner fences a stale audit-only update
without a new snapshot. */
+ @TestTemplate
+ public void testMetadataOnlyAlterRejectsAStaleUpdate() throws IOException {
+ createAndInsertMakeLake(METALAKE_NAME);
+ PolicyMetaService service = PolicyMetaService.getInstance();
+ PolicyEntity policy =
+ createPolicy(
+ RandomIdGenerator.INSTANCE.nextId(),
+ NamespaceUtil.ofPolicy(METALAKE_NAME),
+ "policy_metadata_conflict",
+ AUDIT_INFO);
+ service.insertPolicy(policy, false);
+ PolicyPO initialPO = getPolicyPO(policy.nameIdentifier());
+ AuditInfo winningAudit =
+
AuditInfo.builder().withCreator("winning-updater").withCreateTime(Instant.now()).build();
+ AuditInfo staleAudit =
+
AuditInfo.builder().withCreator("stale-updater").withCreateTime(Instant.now()).build();
+
+ assertThrows(
+ OptimisticLockException.class,
+ () ->
+ service.updatePolicy(
+ policy.nameIdentifier(),
+ entity -> {
+ PolicyEntity current = (PolicyEntity) entity;
+ // Commit the winner after the outer update reads, before
its CAS executes.
+ assertDoesNotThrow(
+ () ->
+ service.updatePolicy(
+ current.nameIdentifier(),
+ winner ->
+ copyPolicy(
+ (PolicyEntity) winner,
+ current.name(),
+ current.comment(),
+ winningAudit)));
+ return copyPolicy(current, current.name(),
current.comment(), staleAudit);
+ }));
+
+ PolicyEntity stored =
service.getPolicyByIdentifier(policy.nameIdentifier());
+ assertEquals(winningAudit, stored.auditInfo());
+ assertEquals(policy.comment(), stored.comment());
+ PolicyPO afterConflict = getPolicyPO(policy.nameIdentifier());
+ assertEquals(initialPO.getOccVersion() + 1,
afterConflict.getOccVersion().longValue());
+ assertEquals(initialPO.getCurrentVersion(),
afterConflict.getCurrentVersion());
+ assertEquals(initialPO.getLastVersion(), afterConflict.getLastVersion());
+ assertEquals(1, listPolicyVersions(policy.id()).size());
+ }
+
+ @TestTemplate
+ public void testStalePolicyDeleteAfterMetadataOnlyAlter() throws IOException
{
+ createAndInsertMakeLake(METALAKE_NAME);
+ PolicyMetaService policyMetaService = PolicyMetaService.getInstance();
+ PolicyEntity policy =
+ createPolicy(
+ RandomIdGenerator.INSTANCE.nextId(),
+ NamespaceUtil.ofPolicy(METALAKE_NAME),
+ "policy_metadata_delete_occ",
+ AUDIT_INFO);
+ policyMetaService.insertPolicy(policy, false);
+ PolicyPO stalePO = getPolicyPO(policy.nameIdentifier());
+ assertEquals(
+ policy.content(),
+
policyMetaService.getPolicyByIdentifier(policy.nameIdentifier()).content());
+
+ AuditInfo updatedAudit =
+
AuditInfo.builder().withCreator("updated-creator").withCreateTime(Instant.now()).build();
+ policyMetaService.updatePolicy(
+ policy.nameIdentifier(),
+ entity ->
+ copyPolicy(
+ (PolicyEntity) entity,
+ ((PolicyEntity) entity).name(),
+ ((PolicyEntity) entity).comment(),
+ updatedAudit));
+
+ PolicyPO afterAlter = getPolicyPO(policy.nameIdentifier());
+ assertEquals(stalePO.getCurrentVersion(), afterAlter.getCurrentVersion());
+ assertEquals(stalePO.getOccVersion() + 1,
afterAlter.getOccVersion().longValue());
+ assertThrows(
+ OptimisticLockException.class,
+ () -> policyMetaService.deletePolicy(policy.nameIdentifier(),
stalePO));
+ assertTrue(backend.exists(policy.nameIdentifier(),
Entity.EntityType.POLICY));
+ assertEquals(1, listPolicyVersions(policy.id()).size());
+ }
+
@TestTemplate
public void testPolicyCreateIsFencedByParentMetalake() {
PolicyMetaService policyMetaService = PolicyMetaService.getInstance();
diff --git
a/core/src/test/java/org/apache/gravitino/storage/relational/utils/TestPOConverters.java
b/core/src/test/java/org/apache/gravitino/storage/relational/utils/TestPOConverters.java
index 4a23e64bc6..dae145b560 100644
---
a/core/src/test/java/org/apache/gravitino/storage/relational/utils/TestPOConverters.java
+++
b/core/src/test/java/org/apache/gravitino/storage/relational/utils/TestPOConverters.java
@@ -21,7 +21,9 @@ package org.apache.gravitino.storage.relational.utils;
import static org.apache.gravitino.file.Fileset.LOCATION_NAME_UNKNOWN;
import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotEquals;
import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
import com.fasterxml.jackson.core.JsonProcessingException;
import com.google.common.collect.ImmutableList;
@@ -35,6 +37,7 @@ import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
import java.util.HashMap;
+import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
import java.util.stream.Collectors;
@@ -844,35 +847,109 @@ public class TestPOConverters {
assertEquals(updatedFileset.storageLocations(), storageLocations);
assertEquals(2, updatePO1.getCurrentVersion());
assertEquals(2, updatePO1.getLastVersion());
+ assertEquals(2, updatePO1.getOccVersion());
assertEquals(2, updatePO1.getFilesetVersionPOs().get(0).getVersion());
Map<String, String> updatedProperties =
JsonUtils.anyFieldMapper()
.readValue(updatePO1.getFilesetVersionPOs().get(0).getProperties(), Map.class);
assertEquals("value1", updatedProperties.get("key"));
- // Metadata-only changes must also advance the OCC token. Reads join the
version table through
- // current_version, so the converter writes the unchanged content as a new
complete snapshot.
+ // A rename changes nothing that fileset_version_info stores, so it
advances the OCC token
+ // alone. The history version stays where it is and no snapshot is
written, which leaves the
+ // metadata row still pointing at the snapshot reads join through
current_version.
FilesetPO updatePO2 = POConverters.updateFilesetPOWithVersion(initPO,
updatedFileset1, null);
- Map<String, String> storageLocations2 =
- updatePO2.getFilesetVersionPOs().stream()
- .collect(
- Collectors.toMap(
- FilesetVersionPO::getLocationName,
FilesetVersionPO::getStorageLocation));
- assertEquals(filesetEntity.storageLocation(),
storageLocations2.get(LOCATION_NAME_UNKNOWN));
- assertEquals(filesetEntity.storageLocations(), storageLocations2);
- assertEquals(2, updatePO2.getCurrentVersion());
- assertEquals(2, updatePO2.getLastVersion());
- assertEquals(2, updatePO2.getFilesetVersionPOs().get(0).getVersion());
assertEquals("test1", updatePO2.getFilesetName());
+ assertEquals(2, updatePO2.getOccVersion());
+ assertEquals(1, updatePO2.getCurrentVersion());
+ assertEquals(1, updatePO2.getLastVersion());
+ assertTrue(updatePO2.getFilesetVersionPOs().isEmpty());
// A snapshot stored above the version the metadata row records must not
be rebuilt: the next
// version starts above every snapshot the fileset still owns.
FilesetPO updatePO3 = POConverters.updateFilesetPOWithVersion(initPO,
updatedFileset, 7L);
assertEquals(8, updatePO3.getCurrentVersion());
assertEquals(8, updatePO3.getLastVersion());
+ assertEquals(2, updatePO3.getOccVersion());
assertEquals(8, updatePO3.getFilesetVersionPOs().get(0).getVersion());
}
+ @Test
+ public void testUpdateFilesetPOVersionComparesPropertiesByValue() throws
JsonProcessingException {
+ // The stored snapshot keeps this key order. For these keys a HashMap
iterates differently.
+ Map<String, String> properties = new LinkedHashMap<>();
+ properties.put("team", "data");
+ properties.put("retention", "7d");
+ properties.put("gravitino.identifier", "id");
+ Map<String, String> reordered = new HashMap<>(properties);
+ assertNotEquals(
+ JsonUtils.anyFieldMapper().writeValueAsString(properties),
+ JsonUtils.anyFieldMapper().writeValueAsString(reordered));
+
+ Namespace namespace = NamespaceUtil.ofFileset("test_metalake",
"test_catalog", "test_schema");
+ Map<String, String> locations =
+ ImmutableMap.of("first", "hdfs://localhost/first", "second",
"hdfs://localhost/second");
+ FilesetEntity original =
+ createFileset(1L, "test", namespace, "this is test",
"hdfs://localhost/first", properties);
+ FilesetEntity.Builder filesetBuilder =
+ FilesetEntity.builder()
+ .withId(original.id())
+ .withName(original.name())
+ .withNamespace(namespace)
+ .withFilesetType(original.filesetType())
+ .withStorageLocations(locations)
+ .withComment(original.comment())
+ .withProperties(properties)
+ .withAuditInfo(original.auditInfo());
+ FilesetEntity filesetEntity = filesetBuilder.build();
+ FilesetEntity renamedFileset =
+ filesetBuilder.withName("test1").withProperties(reordered).build();
+
+ FilesetPO.Builder builder =
+
FilesetPO.builder().withMetalakeId(1L).withCatalogId(1L).withSchemaId(1L);
+ FilesetPO initPO =
POConverters.initializeFilesetPOWithVersion(filesetEntity, builder);
+
+ // Same properties in a different order: the rename writes no snapshot.
+ FilesetPO renamedPO = POConverters.updateFilesetPOWithVersion(initPO,
renamedFileset, null);
+ assertEquals(2, renamedPO.getOccVersion());
+ assertEquals(1, renamedPO.getCurrentVersion());
+ assertTrue(renamedPO.getFilesetVersionPOs().isEmpty());
+ assertEquals(2, initPO.getFilesetVersionPOs().size());
+
+ // A real property change still writes a complete snapshot for both
locations.
+ Map<String, String> changedProperties = new HashMap<>(reordered);
+ changedProperties.put("retention", "14d");
+ FilesetPO changedPO =
+ POConverters.updateFilesetPOWithVersion(
+ initPO, filesetBuilder.withProperties(changedProperties).build(),
null);
+ assertEquals(2, changedPO.getCurrentVersion());
+ assertEquals(2, changedPO.getOccVersion());
+ assertEquals(2, changedPO.getFilesetVersionPOs().size());
+ for (FilesetVersionPO version : changedPO.getFilesetVersionPOs()) {
+ assertEquals(
+ changedProperties,
+ JsonUtils.anyFieldMapper().readValue(version.getProperties(),
Map.class));
+ }
+
+ // Comparing shared fields once must not skip changes to any location,
including a row
+ // other than the first one used for the shared fields.
+ String changedLocationName =
initPO.getFilesetVersionPOs().get(1).getLocationName();
+ Map<String, String> changedLocations = new HashMap<>(locations);
+ changedLocations.put(changedLocationName, "hdfs://localhost/changed");
+ FilesetPO changedLocationPO =
+ POConverters.updateFilesetPOWithVersion(
+ initPO,
+
filesetBuilder.withProperties(reordered).withStorageLocations(changedLocations).build(),
+ null);
+ assertEquals(2, changedLocationPO.getCurrentVersion());
+ assertEquals(2, changedLocationPO.getOccVersion());
+ assertEquals(
+ changedLocations,
+ changedLocationPO.getFilesetVersionPOs().stream()
+ .collect(
+ Collectors.toMap(
+ FilesetVersionPO::getLocationName,
FilesetVersionPO::getStorageLocation)));
+ }
+
@Test
public void testUpdateGroupPOVersionUsesCurrentVersion() {
AuditInfo auditInfo =
@@ -1771,6 +1848,7 @@ public class TestPOConverters {
.withAuditInfo(JsonUtils.anyFieldMapper().writeValueAsString(auditInfo))
.withCurrentVersion(1L)
.withLastVersion(1L)
+ .withOccVersion(1L)
.withDeletedAt(0L)
.withFilesetVersionPOs(ImmutableList.of(filesetVersionPO))
.build();
@@ -1836,6 +1914,7 @@ public class TestPOConverters {
.withAuditInfo(JsonUtils.anyFieldMapper().writeValueAsString(auditInfo))
.withCurrentVersion(1L)
.withLastVersion(1L)
+ .withOccVersion(1L)
.withDeletedAt(0L)
.withPolicyVersionPO(policyVersionPO)
.build();
diff --git a/scripts/h2/schema-2.0.0-h2.sql b/scripts/h2/schema-2.0.0-h2.sql
index e01f35e812..373b8558bf 100644
--- a/scripts/h2/schema-2.0.0-h2.sql
+++ b/scripts/h2/schema-2.0.0-h2.sql
@@ -120,6 +120,7 @@ CREATE TABLE IF NOT EXISTS `fileset_meta` (
`audit_info` CLOB NOT NULL COMMENT 'fileset audit info',
`current_version` INT UNSIGNED NOT NULL DEFAULT 1 COMMENT 'fileset current
version',
`last_version` INT UNSIGNED NOT NULL DEFAULT 1 COMMENT 'fileset last
version',
+ `occ_version` INT UNSIGNED NOT NULL DEFAULT 1 COMMENT 'fileset optimistic
concurrency version',
`deleted_at` BIGINT(20) UNSIGNED NOT NULL DEFAULT 0 COMMENT 'fileset
deleted at',
PRIMARY KEY (fileset_id),
CONSTRAINT uk_sid_fn_del UNIQUE (schema_id, fileset_name, deleted_at),
@@ -398,6 +399,7 @@ CREATE TABLE IF NOT EXISTS `policy_meta` (
`audit_info` CLOB NOT NULL COMMENT 'policy audit info',
`current_version` INT UNSIGNED NOT NULL DEFAULT 1 COMMENT 'policy current
version',
`last_version` INT UNSIGNED NOT NULL DEFAULT 1 COMMENT 'policy last
version',
+ `occ_version` INT UNSIGNED NOT NULL DEFAULT 1 COMMENT 'policy optimistic
concurrency version',
`deleted_at` BIGINT(20) UNSIGNED NOT NULL DEFAULT 0 COMMENT 'policy
deleted at',
PRIMARY KEY (`policy_id`),
UNIQUE KEY `uk_mi_pn_del` (`metalake_id`, `policy_name`, `deleted_at`)
diff --git a/scripts/h2/upgrade-1.3.0-to-2.0.0-h2.sql
b/scripts/h2/upgrade-1.3.0-to-2.0.0-h2.sql
index 2b1e63bf06..f34df6f3e8 100644
--- a/scripts/h2/upgrade-1.3.0-to-2.0.0-h2.sql
+++ b/scripts/h2/upgrade-1.3.0-to-2.0.0-h2.sql
@@ -114,3 +114,15 @@ UPDATE `owner_meta`
AND d.`metadata_object_id` = `owner_meta`.`metadata_object_id`
AND d.`metadata_object_type` = `owner_meta`.`metadata_object_type`
);
+
+-- Separate the optimistic-concurrency token from the history version for
fileset and policy.
+-- Until now `current_version` served as both: it is the join key into
`*_version_info` and the
+-- value the CAS compares, so every alter had to advance it and write a
snapshot even when nothing
+-- in that snapshot changed. `occ_version` takes over the CAS;
`current_version` again advances
+-- only when the stored snapshot changes. The default is the whole backfill,
because `occ_version`
+-- is only ever compared against itself on the same row.
+ALTER TABLE `fileset_meta`
+ ADD COLUMN `occ_version` INT UNSIGNED NOT NULL DEFAULT 1 COMMENT 'fileset
optimistic concurrency version' AFTER `last_version`;
+
+ALTER TABLE `policy_meta`
+ ADD COLUMN `occ_version` INT UNSIGNED NOT NULL DEFAULT 1 COMMENT 'policy
optimistic concurrency version' AFTER `last_version`;
diff --git a/scripts/mysql/schema-2.0.0-mysql.sql
b/scripts/mysql/schema-2.0.0-mysql.sql
index 68e9f9cf23..0edbb0f1f2 100644
--- a/scripts/mysql/schema-2.0.0-mysql.sql
+++ b/scripts/mysql/schema-2.0.0-mysql.sql
@@ -114,6 +114,7 @@ CREATE TABLE IF NOT EXISTS `fileset_meta` (
`audit_info` MEDIUMTEXT NOT NULL COMMENT 'fileset audit info',
`current_version` INT UNSIGNED NOT NULL DEFAULT 1 COMMENT 'fileset current
version',
`last_version` INT UNSIGNED NOT NULL DEFAULT 1 COMMENT 'fileset last
version',
+ `occ_version` INT UNSIGNED NOT NULL DEFAULT 1 COMMENT 'fileset optimistic
concurrency version',
`deleted_at` BIGINT(20) UNSIGNED NOT NULL DEFAULT 0 COMMENT 'fileset
deleted at',
PRIMARY KEY (`fileset_id`),
UNIQUE KEY `fileset_meta_uk_sid_fn_del` (`schema_id`, `fileset_name`,
`deleted_at`),
@@ -389,6 +390,7 @@ CREATE TABLE IF NOT EXISTS `policy_meta` (
`audit_info` MEDIUMTEXT NOT NULL COMMENT 'policy audit info',
`current_version` INT UNSIGNED NOT NULL DEFAULT 1 COMMENT 'policy current
version',
`last_version` INT UNSIGNED NOT NULL DEFAULT 1 COMMENT 'policy last
version',
+ `occ_version` INT UNSIGNED NOT NULL DEFAULT 1 COMMENT 'policy optimistic
concurrency version',
`deleted_at` BIGINT(20) UNSIGNED NOT NULL DEFAULT 0 COMMENT 'policy
deleted at',
PRIMARY KEY (`policy_id`),
UNIQUE KEY `uk_mi_pn_del` (`metalake_id`, `policy_name`, `deleted_at`)
diff --git a/scripts/mysql/upgrade-1.3.0-to-2.0.0-mysql.sql
b/scripts/mysql/upgrade-1.3.0-to-2.0.0-mysql.sql
index 4bad49e12d..121271190b 100644
--- a/scripts/mysql/upgrade-1.3.0-to-2.0.0-mysql.sql
+++ b/scripts/mysql/upgrade-1.3.0-to-2.0.0-mysql.sql
@@ -197,3 +197,15 @@ UPDATE `owner_meta` o
SET o.`deleted_at` = ((UNIX_TIMESTAMP() * 1000.0) + EXTRACT(MICROSECOND
FROM CURRENT_TIMESTAMP(3)) / 1000),
o.`updated_at` = ((UNIX_TIMESTAMP() * 1000.0) + EXTRACT(MICROSECOND
FROM CURRENT_TIMESTAMP(3)) / 1000)
WHERE o.`deleted_at` = 0 AND o.`id` <> d.keep_id;
+
+-- Separate the optimistic-concurrency token from the history version for
fileset and policy.
+-- Until now `current_version` served as both: it is the join key into
`*_version_info` and the
+-- value the CAS compares, so every alter had to advance it and write a
snapshot even when nothing
+-- in that snapshot changed. `occ_version` takes over the CAS;
`current_version` again advances
+-- only when the stored snapshot changes. The default is the whole backfill,
because `occ_version`
+-- is only ever compared against itself on the same row.
+ALTER TABLE `fileset_meta`
+ ADD COLUMN `occ_version` INT UNSIGNED NOT NULL DEFAULT 1 COMMENT 'fileset
optimistic concurrency version' AFTER `last_version`;
+
+ALTER TABLE `policy_meta`
+ ADD COLUMN `occ_version` INT UNSIGNED NOT NULL DEFAULT 1 COMMENT 'policy
optimistic concurrency version' AFTER `last_version`;
diff --git a/scripts/postgresql/schema-2.0.0-postgresql.sql
b/scripts/postgresql/schema-2.0.0-postgresql.sql
index 0d4bab39f7..3fe920d0eb 100644
--- a/scripts/postgresql/schema-2.0.0-postgresql.sql
+++ b/scripts/postgresql/schema-2.0.0-postgresql.sql
@@ -194,6 +194,7 @@ CREATE TABLE IF NOT EXISTS fileset_meta (
audit_info TEXT NOT NULL,
current_version INT NOT NULL DEFAULT 1,
last_version INT NOT NULL DEFAULT 1,
+ occ_version INT NOT NULL DEFAULT 1,
deleted_at BIGINT NOT NULL DEFAULT 0,
PRIMARY KEY (fileset_id),
UNIQUE (schema_id, fileset_name, deleted_at)
@@ -212,6 +213,7 @@ COMMENT ON COLUMN fileset_meta.type IS 'fileset type';
COMMENT ON COLUMN fileset_meta.audit_info IS 'fileset audit info';
COMMENT ON COLUMN fileset_meta.current_version IS 'fileset current version';
COMMENT ON COLUMN fileset_meta.last_version IS 'fileset last version';
+COMMENT ON COLUMN fileset_meta.occ_version IS 'fileset optimistic concurrency
version';
COMMENT ON COLUMN fileset_meta.deleted_at IS 'fileset deleted at';
@@ -688,6 +690,7 @@ CREATE TABLE IF NOT EXISTS policy_meta (
audit_info TEXT NOT NULL,
current_version INT NOT NULL DEFAULT 1,
last_version INT NOT NULL DEFAULT 1,
+ occ_version INT NOT NULL DEFAULT 1,
deleted_at BIGINT NOT NULL DEFAULT 0,
PRIMARY KEY (policy_id),
UNIQUE (metalake_id, policy_name, deleted_at)
@@ -701,6 +704,7 @@ COMMENT ON COLUMN policy_meta.metalake_id IS 'metalake id';
COMMENT ON COLUMN policy_meta.audit_info IS 'policy audit info';
COMMENT ON COLUMN policy_meta.current_version IS 'policy current version';
COMMENT ON COLUMN policy_meta.last_version IS 'policy last version';
+COMMENT ON COLUMN policy_meta.occ_version IS 'policy optimistic concurrency
version';
COMMENT ON COLUMN policy_meta.deleted_at IS 'policy deleted at';
diff --git a/scripts/postgresql/upgrade-1.3.0-to-2.0.0-postgresql.sql
b/scripts/postgresql/upgrade-1.3.0-to-2.0.0-postgresql.sql
index 774cec6b8e..06d0412acf 100644
--- a/scripts/postgresql/upgrade-1.3.0-to-2.0.0-postgresql.sql
+++ b/scripts/postgresql/upgrade-1.3.0-to-2.0.0-postgresql.sql
@@ -163,3 +163,15 @@ UPDATE owner_meta
AND d.metadata_object_id = owner_meta.metadata_object_id
AND d.metadata_object_type = owner_meta.metadata_object_type
);
+
+-- Separate the optimistic-concurrency token from the history version for
fileset and policy.
+-- Until now current_version served as both: it is the join key into
*_version_info and the value
+-- the CAS compares, so every alter had to advance it and write a snapshot
even when nothing in
+-- that snapshot changed. occ_version takes over the CAS; current_version
again advances only when
+-- the stored snapshot changes. The default is the whole backfill, because
occ_version is only ever
+-- compared against itself on the same row.
+ALTER TABLE fileset_meta ADD COLUMN occ_version INT NOT NULL DEFAULT 1;
+COMMENT ON COLUMN fileset_meta.occ_version IS 'fileset optimistic concurrency
version';
+
+ALTER TABLE policy_meta ADD COLUMN occ_version INT NOT NULL DEFAULT 1;
+COMMENT ON COLUMN policy_meta.occ_version IS 'policy optimistic concurrency
version';