diff --git a/core/src/main/java/org/apache/iceberg/DeleteFileIndex.java b/core/src/main/java/org/apache/iceberg/DeleteFileIndex.java index 51d0cf5d9410..416d24e4a97b 100644 --- a/core/src/main/java/org/apache/iceberg/DeleteFileIndex.java +++ b/core/src/main/java/org/apache/iceberg/DeleteFileIndex.java @@ -476,7 +476,6 @@ private Iterable>> deleteManifestRea matchingManifests, manifest -> ManifestFiles.readDeleteManifest(manifest, io, specsById) - .filterRows(dataFilter) .filterPartitions(partitionFilter) .filterPartitions(partitionSet) .caseSensitive(caseSensitive) diff --git a/core/src/test/java/org/apache/iceberg/TestOverwriteWithValidation.java b/core/src/test/java/org/apache/iceberg/TestOverwriteWithValidation.java index 1938ddddc649..102155430172 100644 --- a/core/src/test/java/org/apache/iceberg/TestOverwriteWithValidation.java +++ b/core/src/test/java/org/apache/iceberg/TestOverwriteWithValidation.java @@ -893,37 +893,6 @@ public void testConcurrentConflictingEqualityDeletes() { overwrite::commit); } - @Test - public void testConcurrentNonConflictingEqualityDeletes() { - Assume.assumeTrue(formatVersion == 2); - - Assert.assertNull("Should be empty table", table.currentSnapshot()); - - table.newAppend() - .appendFile(FILE_DAY_2) - .appendFile(FILE_DAY_2_ANOTHER_RANGE) - .commit(); - - Snapshot firstSnapshot = table.currentSnapshot(); - - OverwriteFiles overwrite = table.newOverwrite() - .deleteFile(FILE_DAY_2) - .addFile(FILE_DAY_2_MODIFIED) - .validateFromSnapshot(firstSnapshot.snapshotId()) - .conflictDetectionFilter(EXPRESSION_DAY_2_ID_RANGE) - .validateNoConflictingData() - .validateNoConflictingDeletes(); - - table.newRowDelta() - .addDeletes(FILE_DAY_2_ANOTHER_RANGE_EQ_DELETES) - .commit(); - - overwrite.commit(); - - validateTableFiles(table, FILE_DAY_2_ANOTHER_RANGE, FILE_DAY_2_MODIFIED); - validateTableDeleteFiles(table, FILE_DAY_2_ANOTHER_RANGE_EQ_DELETES); - } - @Test public void testOverwriteByFilterInheritsConflictDetectionFilter() { Assume.assumeTrue(formatVersion == 2); diff --git a/flink/v1.14/flink/src/test/java/org/apache/iceberg/flink/TestFlinkTableSink.java b/flink/v1.14/flink/src/test/java/org/apache/iceberg/flink/TestFlinkTableSink.java index 0c30b09166fc..0c6474cf39cd 100644 --- a/flink/v1.14/flink/src/test/java/org/apache/iceberg/flink/TestFlinkTableSink.java +++ b/flink/v1.14/flink/src/test/java/org/apache/iceberg/flink/TestFlinkTableSink.java @@ -253,6 +253,34 @@ public void testInsertIntoPartition() throws Exception { } } + @Test + public void testInsertWithUpsertAndScanFilterWithNonEqualityField() { + Assume.assumeTrue(format == FileFormat.PARQUET); + String tableName = "test_insert"; + + Map tableProps = ImmutableMap.of( + "write.format.default", format.name(), + TableProperties.FORMAT_VERSION, "2", + TableProperties.UPSERT_ENABLED, "true" + ); + + sql("CREATE TABLE %s(id INT NOT NULL, dt DATE, PRIMARY KEY (id) NOT ENFORCED) PARTITIONED BY (id) WITH %s", + tableName, toWithClause(tableProps)); + + // insert data set + sql("INSERT INTO %s VALUES " + + "(1, to_date('2021-01-01','yyyy-MM-dd'))", tableName); + + sql("INSERT INTO %s VALUES " + + "(1, to_date('2022-01-01','yyyy-MM-dd'))", tableName); + + List rows = sql("SELECT * FROM %s WHERE dt < '2022-01-01'", tableName); + + Assert.assertEquals("rows should be 0", 0, rows.size()); + + sql("DROP TABLE IF EXISTS %s.%s", flinkDatabase, tableName); + } + @Test public void testHashDistributeMode() throws Exception { String tableName = "test_hash_distribution_mode";