jerryshao commented on code in PR #13484:
URL: https://github.com/apache/gravitino/pull/13484#discussion_r4119862474
##########
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`;
Review Comment:
[Question] Is a rolling upgrade a supported path for 1.3.0 to 2.0.0? If it
is, this change opens a window where OCC stops working across node versions.
After this script runs, a 1.3.0 node keeps comparing and advancing
`current_version` (its `updateFilesetMeta` /
`softDeleteFilesetMetasByFilesetId` CAS), while a 2.0.0 node compares
`occ_version`. Two concurrent alters, one from each version, would both find
their own token unchanged and both commit — neither notices the other. The
backfill itself is fine (`DEFAULT 1` is only ever compared against the same
row, as the comment says); the issue is only the mixed-binary window.
If the upgrade is expected to be offline, or the scripts are documented as
"stop, upgrade, start", then nothing to do and this is just worth a line in the
PR description. If rolling upgrades are supported, this probably needs a note
in the upgrade docs at minimum.
Verified by: reading the converted CAS statements in this run
(`FilesetMetaBaseSQLProvider.java:347`, `:411`,
`PolicyMetaBaseSQLProvider.java:122`, `:133`) against their pre-PR
`current_version` form in `git diff origin/main...HEAD`.
##########
core/src/main/java/org/apache/gravitino/storage/relational/utils/POConverters.java:
##########
@@ -741,24 +754,67 @@ public static FilesetPO updateFilesetPOWithVersion(
.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) {
throw new RuntimeException("Failed to serialize json object:", e);
}
}
+ 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);
+ }
+
+ /**
+ * Tells whether an alter leaves every field {@code fileset_version_info}
stores untouched.
+ *
+ * <p>Compares exactly the persisted columns: comment, the serialized
properties, and the storage
+ * locations. Comparing the serialized properties rather than the map keeps
the answer aligned
+ * with what a snapshot would actually hold, so a map that serializes
identically is correctly
+ * reported as unchanged. 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;
+ }
+ return storedVersions.stream()
+ .allMatch(
+ version ->
+ Objects.equals(version.getFilesetComment(),
newFileset.comment())
+ && Objects.equals(version.getProperties(), newProperties));
Review Comment:
[Important] Comparing the *serialized* properties makes the snapshot skip
order-sensitive, and the production rename path does not preserve that order —
so the alter this PR is about still writes a snapshot per storage location for
many filesets.
The chain: a create stores properties serialized from an `ImmutableMap` with
`gravitino.identifier` appended **last**
(`StringIdentifier.newPropertiesWithId`, `StringIdentifier.java:126-129`). A
read deserializes that JSON into a `LinkedHashMap`, preserving the stored order
(`POConverters.fromFilesetPO`, `POConverters.java:891-893`). The alter path
then copies it into a `HashMap` — `Map<String, String> props =
Maps.newHashMap(filesetEntity.properties())`
(`FilesetCatalogOperations.java:1284-1287`) — and
`FilesetEntity.Builder.withProperties` keeps that map as is
(`FilesetEntity.java:308-311`). Re-serializing iterates in hash order, so
`newProperties` is a different string from `version.getProperties()` even
though nothing changed, `filesetSnapshotUnchanged` returns false, and the
rename allocates a full snapshot again.
Whether it bites depends on the key hashes, which is the unpleasant part:
with only `gravitino.identifier` it always matches, but with two realistic user
properties it is a coin toss. I measured it: copying `{userKey1, userKey2,
gravitino.identifier}` into a `HashMap` reorders 34 of the 120 key pairs drawn
from plausible fileset property names (`owner`, `format`, `team`, `retention`,
`credential-providers`, ...), and 1176 of 2000 random 1-4 key sets. Correctness
is not at risk — a wrong `false` costs one redundant snapshot, as the Javadoc
says — but the optimization silently does not fire on the main production path,
which is the PR's stated goal.
Suggested fix: compare by value rather than by serialized text, the same way
`policyContentUnchanged` (`POConverters.java:1906-1924`) was changed in 3c714bc
for exactly this readback problem — e.g. deserialize `version.getProperties()`
into a `Map<String, String>` and compare it with `newFileset.properties()`, or
canonicalize both sides with sorted keys. Then extend
`testAlterWritesASnapshotOnlyWhenStoredContentChanges` to rename a fileset
whose properties have been round-tripped through a `HashMap`, since the current
test passes `current.properties()` through unchanged and therefore cannot catch
this.
Verified by: reading the four files in the chain in this run
(`POConverters.java:796-816`, `:891-893`,
`FilesetCatalogOperations.java:1282-1326`, `StringIdentifier.java:109-130`,
`FilesetEntity.java:308-311`) and measuring `HashMap` iteration order against
the stored order with a standalone JDK program over realistic and random key
sets.
##########
core/src/main/java/org/apache/gravitino/storage/relational/service/FilesetMetaService.java:
##########
@@ -209,11 +210,16 @@ public void insertFileset(FilesetEntity filesetEntity,
boolean overwrite) throws
po.getSchemaId());
persistedPO.set(replacementPO);
}),
- () ->
- SessionUtils.doWithoutCommit(
- FilesetVersionMapper.class,
- mapper ->
-
mapper.insertFilesetVersions(persistedPO.get().getFilesetVersionPOs())));
+ () -> {
+ // An overwrite that replaces a row with identical stored
content allocates no
+ // snapshot, and an empty batch insert is not valid SQL.
+ List<FilesetVersionPO> versionPOs =
persistedPO.get().getFilesetVersionPOs();
+ if (versionPOs.isEmpty()) {
Review Comment:
[Nit] This guard cannot trigger, and the comment above it describes a case
the overwrite path cannot produce.
`storedPO` comes from `selectFilesetMetaBySchemaIdAndNameForUpdate`, whose
SQL selects metadata columns only and no `fileset_version_info` rows
(`FilesetMetaBaseSQLProvider.java:245-258`). `filesetSnapshotUnchanged` returns
false as soon as that list is null (`POConverters.java:798-801`), so
`updateFilesetPOWithVersion` always takes the allocating branch for an
overwrite and `getFilesetVersionPOs()` is never empty here.
Two honest options: drop the branch and the comment, or make it real by
fetching the stored snapshot in the overwrite path (the
`selectMaxFilesetVersion` round trip is already there) so an overwrite with
identical content also skips its snapshot. The second is probably follow-up
work; either way the comment should not claim a case the code cannot reach.
Verified by: reading the overwrite path (`FilesetMetaService.java:176-222`),
the for-update SQL, and the null-list branch of `filesetSnapshotUnchanged` in
this run.
##########
core/src/main/java/org/apache/gravitino/storage/relational/mapper/provider/base/FilesetMetaBaseSQLProvider.java:
##########
@@ -320,25 +330,37 @@ public String
insertFilesetMetaOnDuplicateKeyUpdate(@Param("filesetMeta") Filese
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";
+ // Null POs only reach this method from SQL-text probes. Keep the stricter
guard then,
+ // rather than emit a statement that skips a check the caller may have
needed.
+ boolean allocatesSnapshot =
+ newFilesetPO == null
+ || oldFilesetPO == null
+ || !Objects.equals(newFilesetPO.getCurrentVersion(),
oldFilesetPO.getCurrentVersion());
Review Comment:
[Nit] The null handling exists only because two tests call this method with
`null`, which lets a test shape production SQL.
MyBatis always binds both `@Param` POs here, so the null arms are
unreachable in production; the only callers that pass null are
`TestFilesetMetaBaseSQLProvider.testUpdateUsesVersionCasAndRejectsAnOccupiedSnapshotVersion`
(`TestFilesetMetaBaseSQLProvider.java:45`) and its PostgreSQL sibling. The new
`testUpdateDropsTheSnapshotCheckWhenNoVersionIsAllocated` already shows the
better pattern — two mocked POs. Passing mocks in the existing test too would
let this method take `newFilesetPO.getCurrentVersion()` /
`oldFilesetPO.getCurrentVersion()` directly and emit exactly one shape per
input, rather than a third shape no caller asks for.
Verified by: reading `updateFilesetMeta` and both provider test classes in
this run, and confirming `FilesetMetaService.tryUpdateFileset`
(`FilesetMetaService.java:496-500`) is the only production caller.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]