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


##########
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:
   That's a fair concern about the semantics here. Even after my new change to 
allow list, this concern still holds. My thought here is that the custom 
properties are usually the ones that hold state of a table (for example 
(`ingested-by: service`). With the allow list, the table owner opt into these 
properties to be carried forward by replace. So what the parent's property said 
about the table is still true since replace operations don't change the data in 
a table. 



##########
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,

Review Comment:
   Agreed, I switched to an allow list. The new table property 
`write.summary.carry-forward-keys` lists the custom properties to copy, and 
nothing is copied unless it's set. That also removed the reserved-key list 
entirely. Anything the commit itself sets, including Iceberg's computed 
metrics, is never overwritten.



##########
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:
   Good catch, I updated that Javadoc to mention the properties carried forward 
for replace commits.



##########
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:
   This PR Is scoped only to replace commits. A replace commit doesn't change 
the table's data, so whatever the parent's properties had just before the 
replace commit will be carried forward. 
   
   Writers that change data are expected to set the property themselves.



##########
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:
   Thanks for the suggestion. This code is gone now that the property is an 
allow list, so there's no reserved-key predicate left to filter on.



##########
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:
   Thanks, renamed it to `write.summary.carry-forward-keys`, next to 
`WRITE_PARTITION_SUMMARY_LIMIT`. It now takes a list of keys, since it became 
an allow list.



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