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(); 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..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,6 +39,7 @@ 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; @@ -310,10 +311,11 @@ 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 = FileGenerationUtil.generateDataFile(table, null); 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..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,6 +45,7 @@ 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; @@ -318,10 +319,11 @@ 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 = FileGenerationUtil.generateDataFile(table, null); 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..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,6 +45,7 @@ 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; @@ -318,10 +319,11 @@ 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 = FileGenerationUtil.generateDataFile(table, null); 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