dramaticlly commented on code in PR #18401:
URL: https://github.com/apache/iceberg/pull/18401#discussion_r4212881898


##########
core/src/main/java/org/apache/iceberg/SnapshotProducer.java:
##########
@@ -418,8 +418,18 @@ private Map<String, String> summary(TableMetadata 
previous) {
       previousSummary = summaryBuilder.build();
     }
 
+    Map<String, String> environment = EnvironmentContext.get();
     ImmutableMap.Builder<String, String> builder = ImmutableMap.builder();
 
+    if (DataOperations.REPLACE.equals(operation())

Review Comment:
   I think this expand the current scope, does it warrant a update to existing 
comment on LINE 391
   `  /** Returns the snapshot summary from the implementation and updates 
totals. */`?



##########
core/src/main/java/org/apache/iceberg/SnapshotSummary.java:
##########
@@ -69,8 +70,73 @@ public class SnapshotSummary {
 
   public static final MapJoiner MAP_JOINER = 
Joiner.on(",").withKeyValueSeparator("=");
 
+  // computed by each commit or describing the operation and engine that 
produced a snapshot
+  private static final Set<String> RESERVED_PROPERTIES =
+      ImmutableSet.of(
+          ADDED_FILES_PROP,
+          DELETED_FILES_PROP,
+          TOTAL_DATA_FILES_PROP,
+          ADDED_DELETE_FILES_PROP,
+          ADD_EQ_DELETE_FILES_PROP,
+          REMOVED_EQ_DELETE_FILES_PROP,
+          ADD_POS_DELETE_FILES_PROP,
+          REMOVED_POS_DELETE_FILES_PROP,
+          ADDED_DVS_PROP,
+          REMOVED_DVS_PROP,
+          REMOVED_DELETE_FILES_PROP,
+          TOTAL_DELETE_FILES_PROP,
+          ADDED_RECORDS_PROP,
+          DELETED_RECORDS_PROP,
+          TOTAL_RECORDS_PROP,
+          ADDED_FILE_SIZE_PROP,
+          REMOVED_FILE_SIZE_PROP,
+          TOTAL_FILE_SIZE_PROP,
+          ADDED_POS_DELETES_PROP,
+          REMOVED_POS_DELETES_PROP,
+          TOTAL_POS_DELETES_PROP,
+          ADDED_EQ_DELETES_PROP,
+          REMOVED_EQ_DELETES_PROP,
+          TOTAL_EQ_DELETES_PROP,
+          DELETED_DUPLICATE_FILES,
+          CHANGED_PARTITION_COUNT_PROP,
+          PARTITION_SUMMARY_PROP,
+          STAGED_WAP_ID_PROP,
+          PUBLISHED_WAP_ID_PROP,
+          SOURCE_SNAPSHOT_ID_PROP,
+          REPLACE_PARTITIONS_PROP,
+          CREATED_MANIFESTS_COUNT,
+          REPLACED_MANIFESTS_COUNT,
+          KEPT_MANIFESTS_COUNT,
+          PROCESSED_MANIFEST_ENTRY_COUNT,
+          EnvironmentContext.ENGINE_NAME,
+          EnvironmentContext.ENGINE_VERSION,
+          CatalogProperties.APP_ID,
+          CatalogProperties.APP_NAME);
+
   private SnapshotSummary() {}
 
+  /**
+   * Returns the entries of a previous snapshot summary that were supplied by 
a caller rather than
+   * computed by Iceberg, excluding any key already set in {@code summary} or 
{@code environment}.
+   */
+  static Map<String, String> carriedForwardProperties(
+      Map<String, String> previousSummary,
+      Map<String, String> summary,
+      Map<String, String> environment) {
+    Map<String, String> carried = Maps.newHashMap();
+    for (Map.Entry<String, String> entry : previousSummary.entrySet()) {
+      String key = entry.getKey();
+      if (!RESERVED_PROPERTIES.contains(key)
+          && !key.startsWith(CHANGED_PARTITION_PREFIX)
+          && !environment.containsKey(key)
+          && !summary.containsKey(key)) {
+        carried.put(key, entry.getValue());
+      }
+    }
+

Review Comment:
   I think we can use built-in, also easier to see the predicate this way on 
which key to carry forward with 
   ```suggestion
       return PropertyUtil.filterProperties(
           previousSummary,
           key ->
               !RESERVED_PROPERTIES.contains(key)
                   && !key.startsWith(CHANGED_PARTITION_PREFIX)
                   && !environment.containsKey(key)
                   && !summary.containsKey(key));
   
   ```



##########
core/src/main/java/org/apache/iceberg/TableProperties.java:
##########
@@ -121,6 +121,10 @@ private TableProperties() {}
   public static final String MANIFEST_MERGE_ENABLED = 
"commit.manifest-merge.enabled";
   public static final boolean MANIFEST_MERGE_ENABLED_DEFAULT = true;
 
+  public static final String SNAPSHOT_PROPERTY_CARRY_FORWARD_ENABLED =
+      "commit.snapshot-property-carry-forward.enabled";

Review Comment:
   suggestion: I believe the terminology is `snapshot-summary`  instead of 
snapshot-property. so maybe consider `write.summary.carry-forward`, similar to 
existing `WRITE_PARTITION_SUMMARY_LIMIT`



##########
core/src/test/java/org/apache/iceberg/TestSnapshotSummary.java:
##########
@@ -528,4 +533,136 @@ public void rewriteManifestsWithDuplicateFiles() {
         .containsEntry(SnapshotSummary.KEPT_MANIFESTS_COUNT, "0")
         .containsEntry(SnapshotSummary.REPLACED_MANIFESTS_COUNT, "3");
   }
+
+  @TestTemplate
+  void rewriteCarriesForwardCustomProperties() {
+    enableSnapshotPropertyCarryForward();
+    table.newAppend().appendFile(FILE_A).set("custom-key", 
"custom-value").commit();
+
+    table.newRewrite().deleteFile(FILE_A).addFile(FILE_A2).commit();
+
+    assertThat(table.currentSnapshot().summary()).containsEntry("custom-key", 
"custom-value");
+  }
+
+  @TestTemplate
+  void rewriteDoesNotCarryForwardPropertiesByDefault() {
+    table.newAppend().appendFile(FILE_A).set("custom-key", 
"custom-value").commit();
+
+    table.newRewrite().deleteFile(FILE_A).addFile(FILE_A2).commit();
+
+    
assertThat(table.currentSnapshot().summary()).doesNotContainKey("custom-key");
+  }
+
+  @TestTemplate
+  void rewritePropertyOverridesCarriedForwardValue() {
+    enableSnapshotPropertyCarryForward();
+    table.newAppend().appendFile(FILE_A).set("custom-key", 
"previous-value").commit();
+
+    table
+        .newRewrite()
+        .deleteFile(FILE_A)
+        .addFile(FILE_A2)
+        .set("custom-key", "explicit-value")
+        .commit();
+
+    assertThat(table.currentSnapshot().summary()).containsEntry("custom-key", 
"explicit-value");
+  }
+
+  @TestTemplate
+  void rewriteManifestsDoesNotCarryForwardReservedProperties() {
+    enableSnapshotPropertyCarryForward();
+    table
+        .newAppend()
+        .appendFile(FILE_A)
+        .set("custom-key", "custom-value")
+        .set(SnapshotSummary.PUBLISHED_WAP_ID_PROP, "wap-1")
+        .set(SnapshotSummary.SOURCE_SNAPSHOT_ID_PROP, "1")
+        .set(EnvironmentContext.ENGINE_NAME, "engine")
+        .set(EnvironmentContext.ENGINE_VERSION, "1.0")
+        .set(CatalogProperties.APP_ID, "app")
+        .set(CatalogProperties.APP_NAME, "app-name")
+        .commit();
+    Map<String, String> appendSummary = table.currentSnapshot().summary();
+    assertThat(appendSummary)
+        .containsKeys(SnapshotSummary.ADDED_FILES_PROP, 
SnapshotSummary.ADDED_RECORDS_PROP);
+
+    table.rewriteManifests().clusterBy(file -> "file").rewriteIf(ignored -> 
true).commit();
+
+    assertThat(table.currentSnapshot().summary())
+        .containsEntry("custom-key", "custom-value")
+        .doesNotContainKeys(
+            SnapshotSummary.ADDED_FILES_PROP,
+            SnapshotSummary.ADDED_RECORDS_PROP,
+            SnapshotSummary.ADDED_FILE_SIZE_PROP,
+            SnapshotSummary.PUBLISHED_WAP_ID_PROP,
+            SnapshotSummary.SOURCE_SNAPSHOT_ID_PROP,
+            EnvironmentContext.ENGINE_NAME,
+            EnvironmentContext.ENGINE_VERSION,
+            CatalogProperties.APP_ID,
+            CatalogProperties.APP_NAME);
+  }
+
+  @Test
+  void carriedForwardPropertiesExcludesSummaryConstants() throws 
IllegalAccessException {
+    Map<String, String> previousSummary = Maps.newHashMap();
+    for (Field field : SnapshotSummary.class.getFields()) {
+      if (Modifier.isStatic(field.getModifiers())
+          && field.getType() == String.class
+          && !field.getName().endsWith("_PREFIX")) {
+        previousSummary.put((String) field.get(null), "value");
+      }
+    }
+
+    assertThat(previousSummary).isNotEmpty();
+    assertThat(
+            SnapshotSummary.carriedForwardProperties(
+                previousSummary, ImmutableMap.of(), ImmutableMap.of()))
+        .isEmpty();
+  }
+
+  @TestTemplate
+  void rewriteCarriesForwardFromParentAtCommitTime() {
+    enableSnapshotPropertyCarryForward();
+    table.newAppend().appendFile(FILE_A).set("custom-key", 
"previous-value").commit();
+    RewriteFiles rewrite = 
table.newRewrite().deleteFile(FILE_A).addFile(FILE_A2);
+
+    table.newAppend().appendFile(FILE_B).set("custom-key", 
"concurrent-value").commit();
+    rewrite.commit();
+
+    assertThat(table.currentSnapshot().summary()).containsEntry("custom-key", 
"concurrent-value");
+  }
+
+  @TestTemplate
+  void rewriteOnBranchCarriesForwardFromBranchHead() {
+    enableSnapshotPropertyCarryForward();
+    table.newAppend().appendFile(FILE_A).set("custom-key", 
"main-value").commit();
+    table.manageSnapshots().createBranch("branch", 
table.currentSnapshot().snapshotId()).commit();
+    table
+        .newAppend()
+        .appendFile(FILE_B)
+        .set("custom-key", "branch-value")
+        .toBranch("branch")
+        .commit();
+
+    
table.newRewrite().deleteFile(FILE_B).addFile(FILE_A2).toBranch("branch").commit();
+
+    assertThat(table.snapshot("branch").summary()).containsEntry("custom-key", 
"branch-value");
+  }
+
+  @TestTemplate
+  void appendDoesNotCarryForwardProperties() {

Review Comment:
   seems a bit inconsistent if we do conditional carry forward based on 
operation type



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


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to