From 5f67fbd518bbcc66f0710864b39fac644ded2eeb Mon Sep 17 00:00:00 2001 From: Amogh Jahagirdar Date: Fri, 24 Jul 2026 13:15:09 -0600 Subject: [PATCH 1/3] Core: Disallow adding and removing the same file in a rewrite files commit --- .../org/apache/iceberg/BaseRewriteFiles.java | 7 +++++++ .../org/apache/iceberg/TestRewriteFiles.java | 16 ++++++++++++++++ 2 files changed, 23 insertions(+) diff --git a/core/src/main/java/org/apache/iceberg/BaseRewriteFiles.java b/core/src/main/java/org/apache/iceberg/BaseRewriteFiles.java index b25681de4238..60961ca72f7f 100644 --- a/core/src/main/java/org/apache/iceberg/BaseRewriteFiles.java +++ b/core/src/main/java/org/apache/iceberg/BaseRewriteFiles.java @@ -152,5 +152,12 @@ private void validateReplacedAndAddedFiles() { Preconditions.checkArgument( deletesDeleteFiles() || !addsDeleteFiles(), "Delete files to add must be empty because there's no delete file to be rewritten"); + + for (DataFile added : addedDataFiles()) { + Preconditions.checkArgument( + !replacedDataFiles.contains(added), + "Cannot add and delete the same file in the same rewrite: %s", + added.location()); + } } } diff --git a/core/src/test/java/org/apache/iceberg/TestRewriteFiles.java b/core/src/test/java/org/apache/iceberg/TestRewriteFiles.java index 72a3c89b74d5..5eba212a9bdd 100644 --- a/core/src/test/java/org/apache/iceberg/TestRewriteFiles.java +++ b/core/src/test/java/org/apache/iceberg/TestRewriteFiles.java @@ -172,6 +172,22 @@ public void testDeleteOnly() { .hasMessage("Files to delete cannot be empty"); } + @TestTemplate + public void addingAndDeletingSameFileDisallowed() { + assertThat(listManifestFiles()).isEmpty(); + + commit(table, table.newAppend().appendFile(FILE_A).appendFile(FILE_B), branch); + + assertThatThrownBy( + () -> + apply( + table.newRewrite().rewriteFiles(Sets.newSet(FILE_A), Sets.newSet(FILE_A)), + branch)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage( + "Cannot add and delete the same file in the same rewrite: " + FILE_A.location()); + } + @TestTemplate public void testDeleteWithDuplicateEntriesInManifest() { assertThat(listManifestFiles()).isEmpty(); From 189f37bdb47abefe19d6c59a349708aa5d30b875 Mon Sep 17 00:00:00 2001 From: Amogh Jahagirdar Date: Sun, 26 Jul 2026 12:08:46 -0600 Subject: [PATCH 2/3] Update flink test which uses same file in rewrite --- .../flink/maintenance/operator/TestMonitorSource.java | 9 +++++++-- .../flink/maintenance/operator/TestMonitorSource.java | 9 +++++++-- .../flink/maintenance/operator/TestMonitorSource.java | 9 +++++++-- 3 files changed, 21 insertions(+), 6 deletions(-) diff --git a/flink/v1.20/flink/src/test/java/org/apache/iceberg/flink/maintenance/operator/TestMonitorSource.java b/flink/v1.20/flink/src/test/java/org/apache/iceberg/flink/maintenance/operator/TestMonitorSource.java index 34bb14336511..1fb63020427a 100644 --- a/flink/v1.20/flink/src/test/java/org/apache/iceberg/flink/maintenance/operator/TestMonitorSource.java +++ b/flink/v1.20/flink/src/test/java/org/apache/iceberg/flink/maintenance/operator/TestMonitorSource.java @@ -45,6 +45,7 @@ import org.apache.iceberg.data.GenericAppenderHelper; import org.apache.iceberg.data.RandomGenericData; import org.apache.iceberg.data.Record; +import org.apache.iceberg.flink.SimpleDataUtil; import org.apache.iceberg.flink.TableLoader; import org.apache.iceberg.relocated.com.google.common.collect.Lists; import org.awaitility.Awaitility; @@ -59,6 +60,7 @@ class TestMonitorSource extends OperatorTestBase { private static final RateLimiterStrategy LOW_RATE = RateLimiterStrategy.perSecond(1.0 / 10000.0); @TempDir private File checkpointDir; + @TempDir private java.nio.file.Path dataDir; @ParameterizedTest @ValueSource(booleans = {true, false}) @@ -310,10 +312,13 @@ void testSkipReplace() throws IOException { // Create a DataOperations.REPLACE snapshot DataFile dataFile = SnapshotChanges.builderFor(table).build().addedDataFiles().iterator().next(); + // Replace the file with a new file to produce a REPLACE snapshot + DataFile replacement = + new GenericAppenderHelper(table, FileFormat.PARQUET, dataDir) + .writeFile(Lists.newArrayList(SimpleDataUtil.createRecord(2, "b"))); RewriteFiles rewrite = tableLoader.loadTable().newRewrite(); - // Replace the file with itself for testing purposes rewrite.deleteFile(dataFile); - rewrite.addFile(dataFile); + rewrite.addFile(replacement); rewrite.commit(); // Check that the rewrite is ignored diff --git a/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/maintenance/operator/TestMonitorSource.java b/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/maintenance/operator/TestMonitorSource.java index 3dca6c421c76..38ead2e532e1 100644 --- a/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/maintenance/operator/TestMonitorSource.java +++ b/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/maintenance/operator/TestMonitorSource.java @@ -51,6 +51,7 @@ import org.apache.iceberg.data.GenericAppenderHelper; import org.apache.iceberg.data.RandomGenericData; import org.apache.iceberg.data.Record; +import org.apache.iceberg.flink.SimpleDataUtil; import org.apache.iceberg.flink.TableLoader; import org.apache.iceberg.relocated.com.google.common.collect.Lists; import org.awaitility.Awaitility; @@ -65,6 +66,7 @@ class TestMonitorSource extends OperatorTestBase { private static final RateLimiterStrategy LOW_RATE = RateLimiterStrategy.perSecond(1.0 / 10000.0); @TempDir private File checkpointDir; + @TempDir private java.nio.file.Path dataDir; @ParameterizedTest @ValueSource(booleans = {true, false}) @@ -318,10 +320,13 @@ void testSkipReplace() throws IOException { // Create a DataOperations.REPLACE snapshot DataFile dataFile = SnapshotChanges.builderFor(table).build().addedDataFiles().iterator().next(); + // Replace the file with a new file to produce a REPLACE snapshot + DataFile replacement = + new GenericAppenderHelper(table, FileFormat.PARQUET, dataDir) + .writeFile(Lists.newArrayList(SimpleDataUtil.createRecord(2, "b"))); RewriteFiles rewrite = tableLoader.loadTable().newRewrite(); - // Replace the file with itself for testing purposes rewrite.deleteFile(dataFile); - rewrite.addFile(dataFile); + rewrite.addFile(replacement); rewrite.commit(); // Check that the rewrite is ignored diff --git a/flink/v2.1/flink/src/test/java/org/apache/iceberg/flink/maintenance/operator/TestMonitorSource.java b/flink/v2.1/flink/src/test/java/org/apache/iceberg/flink/maintenance/operator/TestMonitorSource.java index 3dca6c421c76..38ead2e532e1 100644 --- a/flink/v2.1/flink/src/test/java/org/apache/iceberg/flink/maintenance/operator/TestMonitorSource.java +++ b/flink/v2.1/flink/src/test/java/org/apache/iceberg/flink/maintenance/operator/TestMonitorSource.java @@ -51,6 +51,7 @@ import org.apache.iceberg.data.GenericAppenderHelper; import org.apache.iceberg.data.RandomGenericData; import org.apache.iceberg.data.Record; +import org.apache.iceberg.flink.SimpleDataUtil; import org.apache.iceberg.flink.TableLoader; import org.apache.iceberg.relocated.com.google.common.collect.Lists; import org.awaitility.Awaitility; @@ -65,6 +66,7 @@ class TestMonitorSource extends OperatorTestBase { private static final RateLimiterStrategy LOW_RATE = RateLimiterStrategy.perSecond(1.0 / 10000.0); @TempDir private File checkpointDir; + @TempDir private java.nio.file.Path dataDir; @ParameterizedTest @ValueSource(booleans = {true, false}) @@ -318,10 +320,13 @@ void testSkipReplace() throws IOException { // Create a DataOperations.REPLACE snapshot DataFile dataFile = SnapshotChanges.builderFor(table).build().addedDataFiles().iterator().next(); + // Replace the file with a new file to produce a REPLACE snapshot + DataFile replacement = + new GenericAppenderHelper(table, FileFormat.PARQUET, dataDir) + .writeFile(Lists.newArrayList(SimpleDataUtil.createRecord(2, "b"))); RewriteFiles rewrite = tableLoader.loadTable().newRewrite(); - // Replace the file with itself for testing purposes rewrite.deleteFile(dataFile); - rewrite.addFile(dataFile); + rewrite.addFile(replacement); rewrite.commit(); // Check that the rewrite is ignored From 852caa48f0171d9e025a90d27beff53d6384b6d9 Mon Sep 17 00:00:00 2001 From: Amogh Jahagirdar Date: Sun, 26 Jul 2026 12:17:09 -0600 Subject: [PATCH 3/3] use existing utils --- .../flink/maintenance/operator/TestMonitorSource.java | 7 ++----- .../flink/maintenance/operator/TestMonitorSource.java | 7 ++----- .../flink/maintenance/operator/TestMonitorSource.java | 7 ++----- 3 files changed, 6 insertions(+), 15 deletions(-) diff --git a/flink/v1.20/flink/src/test/java/org/apache/iceberg/flink/maintenance/operator/TestMonitorSource.java b/flink/v1.20/flink/src/test/java/org/apache/iceberg/flink/maintenance/operator/TestMonitorSource.java index 1fb63020427a..56c89874b111 100644 --- a/flink/v1.20/flink/src/test/java/org/apache/iceberg/flink/maintenance/operator/TestMonitorSource.java +++ b/flink/v1.20/flink/src/test/java/org/apache/iceberg/flink/maintenance/operator/TestMonitorSource.java @@ -39,13 +39,13 @@ import org.apache.iceberg.DataFile; import org.apache.iceberg.DeleteFile; import org.apache.iceberg.FileFormat; +import org.apache.iceberg.FileGenerationUtil; import org.apache.iceberg.RewriteFiles; import org.apache.iceberg.SnapshotChanges; import org.apache.iceberg.Table; import org.apache.iceberg.data.GenericAppenderHelper; import org.apache.iceberg.data.RandomGenericData; import org.apache.iceberg.data.Record; -import org.apache.iceberg.flink.SimpleDataUtil; import org.apache.iceberg.flink.TableLoader; import org.apache.iceberg.relocated.com.google.common.collect.Lists; import org.awaitility.Awaitility; @@ -60,7 +60,6 @@ class TestMonitorSource extends OperatorTestBase { private static final RateLimiterStrategy LOW_RATE = RateLimiterStrategy.perSecond(1.0 / 10000.0); @TempDir private File checkpointDir; - @TempDir private java.nio.file.Path dataDir; @ParameterizedTest @ValueSource(booleans = {true, false}) @@ -313,9 +312,7 @@ void testSkipReplace() throws IOException { DataFile dataFile = SnapshotChanges.builderFor(table).build().addedDataFiles().iterator().next(); // Replace the file with a new file to produce a REPLACE snapshot - DataFile replacement = - new GenericAppenderHelper(table, FileFormat.PARQUET, dataDir) - .writeFile(Lists.newArrayList(SimpleDataUtil.createRecord(2, "b"))); + DataFile replacement = FileGenerationUtil.generateDataFile(table, null); RewriteFiles rewrite = tableLoader.loadTable().newRewrite(); rewrite.deleteFile(dataFile); rewrite.addFile(replacement); diff --git a/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/maintenance/operator/TestMonitorSource.java b/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/maintenance/operator/TestMonitorSource.java index 38ead2e532e1..af73178ed56d 100644 --- a/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/maintenance/operator/TestMonitorSource.java +++ b/flink/v2.0/flink/src/test/java/org/apache/iceberg/flink/maintenance/operator/TestMonitorSource.java @@ -45,13 +45,13 @@ import org.apache.iceberg.DataFile; import org.apache.iceberg.DeleteFile; import org.apache.iceberg.FileFormat; +import org.apache.iceberg.FileGenerationUtil; import org.apache.iceberg.RewriteFiles; import org.apache.iceberg.SnapshotChanges; import org.apache.iceberg.Table; import org.apache.iceberg.data.GenericAppenderHelper; import org.apache.iceberg.data.RandomGenericData; import org.apache.iceberg.data.Record; -import org.apache.iceberg.flink.SimpleDataUtil; import org.apache.iceberg.flink.TableLoader; import org.apache.iceberg.relocated.com.google.common.collect.Lists; import org.awaitility.Awaitility; @@ -66,7 +66,6 @@ class TestMonitorSource extends OperatorTestBase { private static final RateLimiterStrategy LOW_RATE = RateLimiterStrategy.perSecond(1.0 / 10000.0); @TempDir private File checkpointDir; - @TempDir private java.nio.file.Path dataDir; @ParameterizedTest @ValueSource(booleans = {true, false}) @@ -321,9 +320,7 @@ void testSkipReplace() throws IOException { DataFile dataFile = SnapshotChanges.builderFor(table).build().addedDataFiles().iterator().next(); // Replace the file with a new file to produce a REPLACE snapshot - DataFile replacement = - new GenericAppenderHelper(table, FileFormat.PARQUET, dataDir) - .writeFile(Lists.newArrayList(SimpleDataUtil.createRecord(2, "b"))); + DataFile replacement = FileGenerationUtil.generateDataFile(table, null); RewriteFiles rewrite = tableLoader.loadTable().newRewrite(); rewrite.deleteFile(dataFile); rewrite.addFile(replacement); diff --git a/flink/v2.1/flink/src/test/java/org/apache/iceberg/flink/maintenance/operator/TestMonitorSource.java b/flink/v2.1/flink/src/test/java/org/apache/iceberg/flink/maintenance/operator/TestMonitorSource.java index 38ead2e532e1..af73178ed56d 100644 --- a/flink/v2.1/flink/src/test/java/org/apache/iceberg/flink/maintenance/operator/TestMonitorSource.java +++ b/flink/v2.1/flink/src/test/java/org/apache/iceberg/flink/maintenance/operator/TestMonitorSource.java @@ -45,13 +45,13 @@ import org.apache.iceberg.DataFile; import org.apache.iceberg.DeleteFile; import org.apache.iceberg.FileFormat; +import org.apache.iceberg.FileGenerationUtil; import org.apache.iceberg.RewriteFiles; import org.apache.iceberg.SnapshotChanges; import org.apache.iceberg.Table; import org.apache.iceberg.data.GenericAppenderHelper; import org.apache.iceberg.data.RandomGenericData; import org.apache.iceberg.data.Record; -import org.apache.iceberg.flink.SimpleDataUtil; import org.apache.iceberg.flink.TableLoader; import org.apache.iceberg.relocated.com.google.common.collect.Lists; import org.awaitility.Awaitility; @@ -66,7 +66,6 @@ class TestMonitorSource extends OperatorTestBase { private static final RateLimiterStrategy LOW_RATE = RateLimiterStrategy.perSecond(1.0 / 10000.0); @TempDir private File checkpointDir; - @TempDir private java.nio.file.Path dataDir; @ParameterizedTest @ValueSource(booleans = {true, false}) @@ -321,9 +320,7 @@ void testSkipReplace() throws IOException { DataFile dataFile = SnapshotChanges.builderFor(table).build().addedDataFiles().iterator().next(); // Replace the file with a new file to produce a REPLACE snapshot - DataFile replacement = - new GenericAppenderHelper(table, FileFormat.PARQUET, dataDir) - .writeFile(Lists.newArrayList(SimpleDataUtil.createRecord(2, "b"))); + DataFile replacement = FileGenerationUtil.generateDataFile(table, null); RewriteFiles rewrite = tableLoader.loadTable().newRewrite(); rewrite.deleteFile(dataFile); rewrite.addFile(replacement);