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]