Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,51 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance
* with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing,
* software distributed under the License is distributed on an
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
* KIND, either express or implied. See the License for the
* specific language governing permissions and limitations
* under the License.
*/
package org.apache.iceberg.expressions;

import java.util.List;
import java.util.Map;
import java.util.function.Function;
import org.apache.iceberg.ContentFile;
import org.apache.iceberg.FileContent;
import org.apache.iceberg.relocated.com.google.common.collect.Maps;

final class ContentFileStats {
private ContentFileStats() {}

/**
* Returns stats narrowed to the columns relevant for pruning.
*
* <p>An equality delete file may carry stats for extra columns that are not part of its match
* condition; only its equality fields' stats are relevant for pruning.
*/
static <V> Map<Integer, V> forColumns(
ContentFile<?> file, Function<ContentFile<?>, Map<Integer, V>> stats) {
Map<Integer, V> columnStats = stats.apply(file);
if (columnStats == null || file.content() != FileContent.EQUALITY_DELETES) {
return columnStats;
}

List<Integer> equalityFieldIds = file.equalityFieldIds();
if (equalityFieldIds == null || equalityFieldIds.isEmpty()) {
return columnStats;
}

return Maps.filterKeys(columnStats, equalityFieldIds::contains);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -87,11 +87,11 @@ private boolean eval(ContentFile<?> file) {
return ROWS_MIGHT_MATCH;
}

this.valueCounts = file.valueCounts();
this.nullCounts = file.nullValueCounts();
this.nanCounts = file.nanValueCounts();
this.lowerBounds = file.lowerBounds();
this.upperBounds = file.upperBounds();
this.valueCounts = ContentFileStats.forColumns(file, ContentFile::valueCounts);
this.nullCounts = ContentFileStats.forColumns(file, ContentFile::nullValueCounts);
this.nanCounts = ContentFileStats.forColumns(file, ContentFile::nanValueCounts);
this.lowerBounds = ContentFileStats.forColumns(file, ContentFile::lowerBounds);
this.upperBounds = ContentFileStats.forColumns(file, ContentFile::upperBounds);

return ExpressionVisitors.visitEvaluator(expr, this);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -84,11 +84,11 @@ private boolean eval(ContentFile<?> file) {
return ROWS_MUST_MATCH;
}

this.valueCounts = file.valueCounts();
this.nullCounts = file.nullValueCounts();
this.nanCounts = file.nanValueCounts();
this.lowerBounds = file.lowerBounds();
this.upperBounds = file.upperBounds();
this.valueCounts = ContentFileStats.forColumns(file, ContentFile::valueCounts);
this.nullCounts = ContentFileStats.forColumns(file, ContentFile::nullValueCounts);
this.nanCounts = ContentFileStats.forColumns(file, ContentFile::nanValueCounts);
this.lowerBounds = ContentFileStats.forColumns(file, ContentFile::lowerBounds);
this.upperBounds = ContentFileStats.forColumns(file, ContentFile::upperBounds);

return ExpressionVisitors.visitEvaluator(expr, this);
}
Expand Down
17 changes: 16 additions & 1 deletion core/src/main/java/org/apache/iceberg/V4ManifestReader.java
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
package org.apache.iceberg;

import java.util.Arrays;
import java.util.List;
import java.util.Map;
import java.util.Set;
import org.apache.iceberg.expressions.Binder;
Expand Down Expand Up @@ -134,7 +135,7 @@ public CloseableIterator<TrackedFile> iterator() {
CloseableIterable.filter(
this::incrementSkipCount,
files,
file -> statsFilter.eval(file.contentStats(), file.recordCount()));
file -> statsFilter.eval(statsForFiltering(file), file.recordCount()));
} else {
files =
CloseableIterable.filter(
Expand Down Expand Up @@ -181,6 +182,20 @@ private boolean isDeletedByMDV(TrackedFile file) {
return dv.isSet(Math.toIntExact(file.tracking().manifestPos()));
}

private static ContentStats statsForFiltering(TrackedFile file) {
ContentStats stats = file.contentStats();
if (stats == null || file.contentType() != FileContent.EQUALITY_DELETES) {
return stats;
}

List<Integer> equalityIds = file.equalityIds();
if (equalityIds == null || equalityIds.isEmpty()) {
return stats;
}

return stats.copy(Sets.newHashSet(equalityIds));
}

private boolean matchesPartition(TrackedFile trackedFile) {
Integer specId = trackedFile.specId();
if (specId == null) {
Expand Down
55 changes: 55 additions & 0 deletions core/src/test/java/org/apache/iceberg/DeleteFileIndexTestBase.java
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,9 @@

import static org.apache.iceberg.expressions.Expressions.bucket;
import static org.apache.iceberg.expressions.Expressions.equal;
import static org.apache.iceberg.expressions.Expressions.greaterThanOrEqual;
import static org.apache.iceberg.types.Types.NestedField.optional;
import static org.apache.iceberg.types.Types.NestedField.required;
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.assertj.core.api.Assumptions.assumeThat;
Expand Down Expand Up @@ -674,6 +677,58 @@ public void testEqualityDeleteDiscardMetrics() {
.containsExactly(Map.entry(fieldId, ByteBuffer.wrap(new byte[20])));
}

@TestTemplate
public void testEqualityDeleteNotPrunedByNonKeyColumnStats() {
// v4 tables can still inherit an equality delete written before an upgrade to v4
int createFormatVersion = Math.min(formatVersion, 3);

Schema schema =
new Schema(
required(1, "id", Types.IntegerType.get()),
required(2, "data", Types.StringType.get()),
optional(3, "extra", Types.IntegerType.get()));
PartitionSpec spec = PartitionSpec.builderFor(schema).withSpecId(0).build();
Table table =
TestTables.create(tableDir, "eq-delete-full-row", schema, spec, createFormatVersion);

DataFile dataFile =
DataFiles.builder(spec)
.withPath("/path/to/data-a.parquet")
.withFileSizeInBytes(10)
.withRecordCount(1)
.build();
table.newAppend().appendFile(dataFile).commit();

// equality field is "id" (1); "extra" (3) is also in the file but is not part of the match
// condition, so its stats must not be used to prune the delete against the scan's filter
DeleteFile eqDeletes =
FileMetadata.deleteFileBuilder(spec)
.ofEqualityDeletes(1)
.withPath(UUID.randomUUID() + "/path/to/eq-delete-full-row.parquet")
.withFileSizeInBytes(10)
.withRecordCount(1)
.withMetrics(
new Metrics(
1L, null, ImmutableMap.of(1, 1L, 3, 1L), ImmutableMap.of(1, 0L, 3, 1L), null))
.build();
table.newRowDelta().addDeletes(eqDeletes).commit();

if (formatVersion > createFormatVersion) {
table
.updateProperties()
.set(TableProperties.FORMAT_VERSION, String.valueOf(formatVersion))
.commit();
}

List<T> tasks =
Lists.newArrayList(
newScan(table).filter(greaterThanOrEqual("extra", 10)).planFiles().iterator());
assertThat(tasks).as("Should have one task").hasSize(1);

FileScanTask task = (FileScanTask) tasks.get(0);
assertThat(task.deletes()).as("Should have one associated delete file").hasSize(1);
}

@TestTemplate
public void testPositionDeletesGroup() {
DeleteFile file1 = withDataSequenceNumber(1, partitionedPosDeletes(SPEC, FILE_A.partition()));
Expand Down
54 changes: 54 additions & 0 deletions core/src/test/java/org/apache/iceberg/TestDeleteFiles.java
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@
import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap;
import org.apache.iceberg.relocated.com.google.common.collect.Iterables;
import org.apache.iceberg.relocated.com.google.common.collect.Lists;
import org.apache.iceberg.types.Conversions;
import org.apache.iceberg.types.Types;
import org.apache.iceberg.util.StructLikeWrapper;
import org.junit.jupiter.api.Test;
Expand Down Expand Up @@ -470,6 +471,59 @@ public void testDeleteWithCollision() {
assertThat(afterDeletePartitions).containsExactly(partitionOne);
}

@TestTemplate
public void deleteFromRowFilterKeepsEqualityDeleteWithMisleadingNonKeyColumnStats() {
Schema schema =
new Schema(
Types.NestedField.required(3, "id", Types.IntegerType.get()),
Types.NestedField.required(4, "data", Types.StringType.get()));
int eqDeleteFormatVersion = Math.min(Math.max(formatVersion, 2), 3);
Table eqDeleteTable =
TestTables.create(
tableDir,
"eq-delete-row-filter",
schema,
PartitionSpec.unpartitioned(),
eqDeleteFormatVersion);

int idFieldId = eqDeleteTable.schema().findField("id").fieldId();
int dataFieldId = eqDeleteTable.schema().findField("data").fieldId();

DeleteFile eqDeletes =
FileMetadata.deleteFileBuilder(PartitionSpec.unpartitioned())
.ofEqualityDeletes(idFieldId)
.withPath("/path/to/eq-delete.parquet")
.withFileSizeInBytes(10)
.withRecordCount(1)
.withMetrics(
new Metrics(
1L,
null, // no column sizes
ImmutableMap.of(idFieldId, 1L, dataFieldId, 1L), // value counts
ImmutableMap.of(idFieldId, 0L, dataFieldId, 0L), // null counts
null, // no nan value counts
ImmutableMap.of(
dataFieldId, Conversions.toByteBuffer(Types.StringType.get(), "zzz")),
ImmutableMap.of(
dataFieldId, Conversions.toByteBuffer(Types.StringType.get(), "zzz"))))
.build();
eqDeleteTable.newRowDelta().addDeletes(eqDeletes).commit();
Snapshot addSnapshot = eqDeleteTable.currentSnapshot();

eqDeleteTable.newDelete().deleteFromRowFilter(Expressions.equal("data", "zzz")).commit();

// the delete manifest is untouched: the filter can only reach the non-key "data" stats,
Snapshot afterFilter = eqDeleteTable.currentSnapshot();
assertThat(afterFilter.deleteManifests(eqDeleteTable.io())).hasSize(1);
validateDeleteManifest(
afterFilter.deleteManifests(eqDeleteTable.io()).get(0),
null,
null,
ids(addSnapshot.snapshotId()),
files(eqDeletes),
statuses(Status.ADDED));
}

@TestTemplate
public void testDeleteValidateFileExistence() {
Snapshot append = commit(table, table.newFastAppend().appendFile(FILE_B), branch);
Expand Down
44 changes: 44 additions & 0 deletions core/src/test/java/org/apache/iceberg/TestV4ManifestReader.java
Original file line number Diff line number Diff line change
Expand Up @@ -1248,6 +1248,29 @@ public void statsFilterDataFileBoundsFiltering(FileFormat format) throws IOExcep
.containsExactly(FILE_D);
}

@ParameterizedTest
@FieldSource("MANIFEST_FORMATS")
public void statsFilterIgnoresEqualityDeleteNonKeyColumnStats(FileFormat format)
throws IOException {
// CONTENT_STATS has id in [0, 99] and data in [a, z], but only id is an equality field, so the
// data stats must not be used to filter out the delete file
TrackedFile equalityDelete =
unpartitionedEqualityDeleteFileWithStats(
"s3://bucket/table/eq-delete.parquet", CONTENT_STATS, ImmutableList.of(1));
ManifestFile manifest = writeManifest(format, UNPARTITIONED_TYPE, equalityDelete);

V4ManifestReader.Builder builder =
V4ManifestReader.builder(manifest, IO, TABLE_SCHEMA, UNPARTITIONED_SPECS)
.filter(Expressions.equal("data", "zzz")) // outside data stats bounds
.metricsConfig(METRICS_CONFIG);

List<TrackedFile> actualFiles = read(builder);

assertThat(actualFiles)
.usingComparatorForType(FILE_COMPARATOR, TrackedFile.class)
.containsExactly(equalityDelete);
}

@ParameterizedTest
@FieldSource("MANIFEST_FORMATS")
public void statsFilterManifestBoundsFiltering(FileFormat format) throws IOException {
Expand Down Expand Up @@ -1826,6 +1849,27 @@ public void resolutionResolvesOnlyRelativePaths(FileFormat format) throws IOExce
"s3://bucket/db/table/data/rel.parquet", "s3://other/abs-dv.puffin"));
}

private static TrackedFile unpartitionedEqualityDeleteFileWithStats(
String location, ContentStats stats, List<Integer> equalityIds) {
return new TrackedFileStruct(
ADDED_TRACKING,
FileContent.EQUALITY_DELETES,
FORMAT_VERSION_V4,
location,
FileFormat.PARQUET,
RECORD_COUNT,
FILE_SIZE_IN_BYTES,
null, // unpartitioned
null, // null partition data
stats,
SortOrder.unsorted().orderId(),
null, // deletion vector
null, // manifest info
null, // key metadata
ImmutableList.of(4L), // split offsets
equalityIds);
}

private static TrackedFile unpartitionedFileWithoutStats(String location) {
Tracking tracking =
new TrackingStruct(
Expand Down
Loading