yuqi1129 commented on code in PR #13484:
URL: https://github.com/apache/gravitino/pull/13484#discussion_r4123900777


##########
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:
   Good catch, thanks. Fixed in 09d345cc77.
   
   `filesetSnapshotUnchanged` now compares properties by value, the same way as 
`policyContentUnchanged`: if the JSON strings are equal we stop there; if not, 
we read the stored JSON into a map and compare it with 
`newFileset.properties()`.
   
   New tests:
   - 
`TestFilesetMetaService.testRenameWithReorderedPropertiesWritesNoSnapshot`: 
renames a fileset with properties passed through `new HashMap<>(...)`, like 
`FilesetCatalogOperations` does, and checks that no snapshot is written.
   - `TestPOConverters.testUpdateFilesetPOVersionComparesPropertiesByValue`: 
the same check at the converter level.
   
   Both tests first assert that the HashMap order really differs from the 
stored order, so they can't pass by luck. Both fail without the fix (`expected: 
<1> but was: <2>`), on H2, MySQL and PostgreSQL.



##########
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:
   Agreed, removed in 09d345cc77. The overwrite path always writes a snapshot, 
because `storedPO` has no snapshot rows. I put back the plain 
`insertFilesetVersions` call and added a short comment at the call site to say 
this. Skipping the snapshot for an overwrite with the same content can be a 
follow-up.



##########
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:
   Done in 09d345cc77. `updateFilesetMeta` no longer handles null. It just 
compares the two current versions. 
`testUpdateUsesVersionCasAndRejectsAnOccupiedSnapshotVersion` now passes two 
mocked POs (current version 3 -> 4), the same way as 
`testUpdateDropsTheSnapshotCheckWhenNoVersionIsAllocated`.



##########
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:
   Rolling upgrade is not supported for 1.3.0 -> 2.0.0. 
`docs/how-to-upgrade.md` Step 1 says to shut down Gravitino before running the 
upgrade scripts, so 1.3.0 and 2.0.0 servers never write to the same database at 
the same time. I added a line about this to the PR description.



-- 
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]

Reply via email to