diff --git a/api/src/main/java/org/apache/iceberg/ManifestFile.java b/api/src/main/java/org/apache/iceberg/ManifestFile.java index cfedb78a7cee..32c7adae679f 100644 --- a/api/src/main/java/org/apache/iceberg/ManifestFile.java +++ b/api/src/main/java/org/apache/iceberg/ManifestFile.java @@ -186,6 +186,26 @@ default boolean hasDeletedFiles() { /** Returns the total number of rows in all files with status DELETED in the manifest file. */ Long deletedRowsCount(); + /** Returns the number of files with status REPLACED in the manifest file. */ + default Integer replacedFilesCount() { + return 0; + } + + /** Returns the total number of rows in all files with status REPLACED in the manifest file. */ + default Long replacedRowsCount() { + return 0L; + } + + /** Returns the number of files with status MODIFIED in the manifest file. */ + default Integer modifiedFilesCount() { + return 0; + } + + /** Returns the total number of rows in all files with status MODIFIED in the manifest file. */ + default Long modifiedRowsCount() { + return 0L; + } + /** * Returns a list of {@link PartitionFieldSummary partition field summaries}. * diff --git a/api/src/main/java/org/apache/iceberg/geospatial/GeospatialBound.java b/api/src/main/java/org/apache/iceberg/geospatial/GeospatialBound.java index 3425a4663801..cdb76ae3c46e 100644 --- a/api/src/main/java/org/apache/iceberg/geospatial/GeospatialBound.java +++ b/api/src/main/java/org/apache/iceberg/geospatial/GeospatialBound.java @@ -21,6 +21,7 @@ import java.nio.ByteBuffer; import java.nio.ByteOrder; import java.util.Objects; +import org.apache.iceberg.StructLike; import org.apache.iceberg.relocated.com.google.common.base.Preconditions; /** @@ -44,7 +45,7 @@ *

This class represents a lower or upper geospatial bound and handles serialization and * deserialization of these bounds to/from byte arrays, conforming to the Iceberg specification. */ -public class GeospatialBound { +public class GeospatialBound implements StructLike { /** * Parses a geospatial bound from a byte buffer according to Iceberg spec. * @@ -282,6 +283,27 @@ public boolean hasM() { return !Double.isNaN(m); } + @Override + public int size() { + return 4; + } + + @Override + public T get(int pos, Class javaClass) { + return switch (pos) { + case 0 -> javaClass.cast(x); + case 1 -> javaClass.cast(y); + case 2 -> hasZ() ? javaClass.cast(z) : null; + case 3 -> hasM() ? javaClass.cast(m) : null; + default -> throw new IllegalArgumentException("Invalid position: " + pos); + }; + } + + @Override + public void set(int pos, T value) { + throw new UnsupportedOperationException("GeospatialBound is read only"); + } + @Override public String toString() { return "GeospatialBound(" + simpleString() + ")"; diff --git a/api/src/test/java/org/apache/iceberg/geospatial/TestGeospatialBound.java b/api/src/test/java/org/apache/iceberg/geospatial/TestGeospatialBound.java index 8a6a6325cc85..7b60ddc38c1d 100644 --- a/api/src/test/java/org/apache/iceberg/geospatial/TestGeospatialBound.java +++ b/api/src/test/java/org/apache/iceberg/geospatial/TestGeospatialBound.java @@ -19,9 +19,11 @@ package org.apache.iceberg.geospatial; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; import java.nio.ByteBuffer; import java.nio.ByteOrder; +import org.apache.iceberg.StructLike; import org.apache.iceberg.util.ByteBuffers; import org.junit.jupiter.api.Test; @@ -71,6 +73,34 @@ public void testCreateXYZM() { assertThat(bound.hasM()).isTrue(); } + @Test + public void testStructLikeXY() { + StructLike bound = GeospatialBound.createXY(1.0, 2.0); + assertThat(bound.size()).isEqualTo(4); + assertThat(bound.get(0, Double.class)).isEqualTo(1.0); + assertThat(bound.get(1, Double.class)).isEqualTo(2.0); + assertThat(bound.get(2, Double.class)).isNull(); + assertThat(bound.get(3, Double.class)).isNull(); + } + + @Test + public void testStructLikeXYZM() { + StructLike bound = GeospatialBound.createXYZM(3.0, 4.0, 5.0, 6.0); + assertThat(bound.size()).isEqualTo(4); + assertThat(bound.get(0, Double.class)).isEqualTo(3.0); + assertThat(bound.get(1, Double.class)).isEqualTo(4.0); + assertThat(bound.get(2, Double.class)).isEqualTo(5.0); + assertThat(bound.get(3, Double.class)).isEqualTo(6.0); + } + + @Test + public void testSet() { + GeospatialBound bound = GeospatialBound.createXY(1.0, 2.0); + assertThatThrownBy(() -> bound.set(0, 9.0)) + .isInstanceOf(UnsupportedOperationException.class) + .hasMessage("GeospatialBound is read only"); + } + @Test public void testEqualsAndHashCode() { GeospatialBound xy1 = GeospatialBound.createXY(1.0, 2.0); diff --git a/core/src/main/java/org/apache/iceberg/MapBackedContentStats.java b/core/src/main/java/org/apache/iceberg/MapBackedContentStats.java new file mode 100644 index 000000000000..66e9b6694f5f --- /dev/null +++ b/core/src/main/java/org/apache/iceberg/MapBackedContentStats.java @@ -0,0 +1,244 @@ +/* + * 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; + +import java.nio.ByteBuffer; +import java.util.Map; +import java.util.Objects; +import java.util.Set; +import org.apache.iceberg.relocated.com.google.common.base.Preconditions; +import org.apache.iceberg.relocated.com.google.common.collect.Iterables; +import org.apache.iceberg.relocated.com.google.common.collect.Maps; +import org.apache.iceberg.relocated.com.google.common.collect.Sets; +import org.apache.iceberg.types.Conversions; +import org.apache.iceberg.types.Type; +import org.apache.iceberg.types.Types; + +/** Reusable {@link ContentStats} view over a {@link ContentFile}'s stat maps. */ +class MapBackedContentStats implements ContentStats { + private final Schema tableSchema; + private final Map> statsById = Maps.newHashMap(); + + private Types.StructType type; + private Map valueCounts; + private Map nullValueCounts; + private Map nanValueCounts; + private Map avgValueSizes; + private Map lowerBounds; + private Map upperBounds; + + MapBackedContentStats(Schema tableSchema) { + Preconditions.checkArgument(tableSchema != null, "Invalid table schema: null"); + this.tableSchema = tableSchema; + } + + MapBackedContentStats wrap(ContentFile file) { + this.valueCounts = file.valueCounts(); + this.nullValueCounts = file.nullValueCounts(); + this.nanValueCounts = file.nanValueCounts(); + this.avgValueSizes = file.avgValueSizes(); + this.lowerBounds = file.lowerBounds(); + this.upperBounds = file.upperBounds(); + this.type = null; + return this; + } + + @Override + public Iterable> fieldStats() { + return Iterables.filter(Iterables.transform(statsFieldIds(), this::statsFor), Objects::nonNull); + } + + @Override + @SuppressWarnings("unchecked") + public FieldStats statsFor(int fieldId) { + // Schema is fixed for this instance; wrap() rebinds the metric maps. An id absent + // from the current maps is left uncached so a later file can still surface it. + if (!containsFieldInMaps(fieldId)) { + return null; + } + + if (!statsById.containsKey(fieldId)) { + statsById.put(fieldId, createFieldStats(fieldId)); + } + + return (FieldStats) statsById.get(fieldId); + } + + private FieldStats createFieldStats(int fieldId) { + Types.NestedField field = tableSchema.findField(fieldId); + // A file can carry metrics for an id this schema does not have. Skip it. + if (field == null) { + return null; + } + + Type fieldType = field.type(); + Types.StructType struct = + StatsUtil.fieldStatsStruct(fieldType, StatsUtil.toBaseId(fieldId), MetricsModes.Full.get()); + // null means a struct, list, or map, or an id outside the stats window. The schema is bound + // once during wrapper creation, so the cached null stays valid across wrap(). + if (struct == null) { + return null; + } + + return new MapBackedFieldStats<>(fieldId, fieldType, struct); + } + + @Override + public Types.StructType type() { + if (type == null) { + this.type = StatsUtil.statsReadSchema(tableSchema, statsFieldIds()); + } + + return type; + } + + @Override + public ContentStats copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + + @Override + public ContentStats copy(Set fieldIds) { + throw new UnsupportedOperationException("copy is not implemented"); + } + + private boolean containsFieldInMaps(int fieldId) { + return containsId(valueCounts, fieldId) + || containsId(nullValueCounts, fieldId) + || containsId(nanValueCounts, fieldId) + || containsId(avgValueSizes, fieldId) + || containsId(lowerBounds, fieldId) + || containsId(upperBounds, fieldId); + } + + private Set statsFieldIds() { + Set ids = + Sets.newHashSetWithExpectedSize(valueCounts == null ? 0 : valueCounts.size()); + addKeys(ids, valueCounts); + addKeys(ids, nullValueCounts); + addKeys(ids, nanValueCounts); + addKeys(ids, avgValueSizes); + addKeys(ids, lowerBounds); + addKeys(ids, upperBounds); + return ids; + } + + private static void addKeys(Set ids, Map map) { + if (map != null) { + ids.addAll(map.keySet()); + } + } + + private static boolean containsId(Map map, int id) { + return map != null && map.containsKey(id); + } + + /** Reusable {@link FieldStats} view over one field's entries in a {@link ContentFile}'s maps. */ + private class MapBackedFieldStats implements FieldStats { + private final int fieldId; + private final Types.StructType struct; + private final Type fieldType; + private final Type boundType; + + private MapBackedFieldStats(int fieldId, Type fieldType, Types.StructType struct) { + this.fieldId = fieldId; + this.fieldType = fieldType; + this.struct = struct; + this.boundType = struct.fieldType(StatsUtil.LOWER_BOUND_NAME); + } + + @Override + public int fieldId() { + return fieldId; + } + + @Override + public Types.StructType type() { + return struct; + } + + @Override + public T lowerBound() { + return decodeBound(lowerBounds); + } + + @Override + public T upperBound() { + return decodeBound(upperBounds); + } + + @SuppressWarnings("unchecked") + private T decodeBound(Map bounds) { + if (boundType == null || bounds == null) { + return null; + } + + return (T) Conversions.fromByteBuffer(fieldType, bounds.get(fieldId)); + } + + @Override + public boolean tightBounds() { + return false; + } + + @Override + public boolean hasValueCount() { + return count(valueCounts) != null; + } + + @Override + public long valueCount() { + return count(valueCounts); + } + + @Override + public boolean hasNullValueCount() { + return count(nullValueCounts) != null; + } + + @Override + public long nullValueCount() { + return count(nullValueCounts); + } + + @Override + public boolean hasNanValueCount() { + return count(nanValueCounts) != null; + } + + @Override + public long nanValueCount() { + return count(nanValueCounts); + } + + @Override + public Integer avgValueSizeInBytes() { + return avgValueSizes == null ? null : avgValueSizes.get(fieldId); + } + + @Override + public FieldStats copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + + private Long count(Map counts) { + return counts == null ? null : counts.get(fieldId); + } + } +} diff --git a/core/src/main/java/org/apache/iceberg/TrackedFileAdapters.java b/core/src/main/java/org/apache/iceberg/TrackedFileAdapters.java index 5acb03b9089c..dc5fec4879b0 100644 --- a/core/src/main/java/org/apache/iceberg/TrackedFileAdapters.java +++ b/core/src/main/java/org/apache/iceberg/TrackedFileAdapters.java @@ -22,10 +22,14 @@ import java.nio.ByteBuffer; import java.util.List; import java.util.Map; +import java.util.Objects; import java.util.Set; import org.apache.iceberg.relocated.com.google.common.base.Preconditions; -/** Adapts {@link TrackedFile} entries to their read APIs, for example {@link DataFile}. */ +/** + * Adapts between the {@link TrackedFile} model and the {@link DataFile} / {@link DeleteFile} / + * {@link ManifestFile} APIs in both directions. + */ class TrackedFileAdapters { private TrackedFileAdapters() {} @@ -63,6 +67,16 @@ static ManifestFile asManifestFile(TrackedFile file) { return new TrackedManifestFile(file); } + /** Returns a reusable adapter from {@link DataFile} to {@link TrackedFile}. */ + static DataTrackedFile forDataFile(Schema tableSchema) { + return new DataTrackedFile(tableSchema); + } + + /** Returns a reusable adapter from {@link ManifestFile} to {@link TrackedFile}. */ + static ManifestTrackedFile forManifestFile() { + return new ManifestTrackedFile(); + } + /** Shared base for data and delete file adapters. */ private abstract static class TrackedFileAdapter> implements ContentFile, Serializable { @@ -424,6 +438,10 @@ private TrackedManifestFile(TrackedFile file) { this.file = file; } + private TrackedFile file() { + return file; + } + @Override public String path() { return file.location(); @@ -498,6 +516,26 @@ public Long deletedRowsCount() { return file.manifestInfo().deletedRowsCount(); } + @Override + public Integer replacedFilesCount() { + return file.manifestInfo().replacedFilesCount(); + } + + @Override + public Long replacedRowsCount() { + return file.manifestInfo().replacedRowsCount(); + } + + @Override + public Integer modifiedFilesCount() { + return file.manifestInfo().modifiedFilesCount(); + } + + @Override + public Long modifiedRowsCount() { + return file.manifestInfo().modifiedRowsCount(); + } + @Override public List partitions() { return null; @@ -529,6 +567,547 @@ public ManifestFile copy() { } } + /** Adapts a {@link DataFile} to {@link TrackedFile}. */ + static class DataTrackedFile implements TrackedFile { + private final MapBackedContentStats statsWrapper; + private final WrappedEntryTracking trackingWrapper = new WrappedEntryTracking(); + + private Tracking tracking; + private DataFile file; + private ContentStats stats; + + DataTrackedFile(Schema tableSchema) { + this.statsWrapper = new MapBackedContentStats(tableSchema); + } + + /** + * Re-points this adapter at a {@link DataFile} from the public API. Tracking is unset. + * + *

Returns the inner {@link TrackedFile} when {@code newFile} is a {@link TrackedDataFile}. + */ + public TrackedFile wrap(DataFile newFile) { + if (newFile instanceof TrackedDataFile tracked) { + return tracked.file(); + } + + wrapInternal(newFile); + this.tracking = null; + return this; + } + + /** + * Re-points this adapter at a {@link ManifestEntry}. Converts the contained data file and the + * entry's tracking fields. + * + *

Returns the inner {@link TrackedFile} when the entry's file is a {@link TrackedDataFile}. + */ + public TrackedFile wrap(ManifestEntry entry) { + Preconditions.checkArgument(entry != null, "Invalid entry: null"); + if (entry.file() instanceof TrackedDataFile tracked) { + return tracked.file(); + } + + wrapInternal(entry.file()); + this.tracking = trackingWrapper.wrap(entry); + return this; + } + + private void wrapInternal(DataFile newFile) { + Preconditions.checkArgument(newFile != null, "Invalid file: null"); + Preconditions.checkArgument( + newFile.content() == FileContent.DATA, + "Invalid content for data file: %s", + newFile.content()); + + this.file = newFile; + this.stats = hasContentStats(newFile) ? statsWrapper.wrap(newFile) : null; + } + + @Override + public Tracking tracking() { + return tracking; + } + + @Override + public FileContent contentType() { + return FileContent.DATA; + } + + @Override + public int formatVersion() { + throw new IllegalStateException("Format version is assigned at write time"); + } + + @Override + public String location() { + return file.location(); + } + + @Override + public FileFormat fileFormat() { + return file.format(); + } + + @Override + public long recordCount() { + return file.recordCount(); + } + + @Override + public long fileSizeInBytes() { + return file.fileSizeInBytes(); + } + + @Override + public Integer specId() { + // Files in one manifest may use different specs; this is the spec for this data file only. + return file.specId(); + } + + @Override + public StructLike partition() { + return file.partition(); + } + + @Override + public ContentStats contentStats() { + return stats; + } + + @Override + public Integer sortOrderId() { + return file.sortOrderId(); + } + + @Override + public DeletionVector deletionVector() { + return file.deletionVector(); + } + + @Override + public ManifestInfo manifestInfo() { + return null; + } + + @Override + public ByteBuffer keyMetadata() { + return file.keyMetadata(); + } + + @Override + public List splitOffsets() { + return file.splitOffsets(); + } + + @Override + public List equalityIds() { + return null; + } + + @Override + public TrackedFile copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + + @Override + public TrackedFile copyWithStats(Set requestedColumnIds) { + throw new UnsupportedOperationException("copy is not implemented"); + } + } + + /** Adapts a {@link ManifestFile} to {@link TrackedFile}. */ + static class ManifestTrackedFile implements TrackedFile { + private final WrappedManifestInfo manifestInfo = new WrappedManifestInfo(); + private final WrappedManifestTracking trackingWrapper = new WrappedManifestTracking(); + private Tracking tracking; + private ManifestFile manifest; + private long recordCount; + private FileContent contentType; + + ManifestTrackedFile() {} + + /** + * Re-points this adapter at {@code newManifest}. Converts the manifest's own fields only; + * write-time tracking updates are applied by the versioned writer. + * + *

Returns the inner {@link TrackedFile} when {@code newManifest} is a {@link + * TrackedManifestFile}. + */ + public TrackedFile wrap(ManifestFile newManifest) { + if (newManifest instanceof TrackedManifestFile tracked) { + return tracked.file(); + } + + Preconditions.checkArgument(newManifest != null, "Invalid manifest file: null"); + + this.manifest = newManifest; + this.contentType = + newManifest.content() == ManifestContent.DATA + ? FileContent.DATA_MANIFEST + : FileContent.DELETE_MANIFEST; + this.recordCount = manifestRecordCount(newManifest); + this.tracking = trackingWrapper.wrap(newManifest); + this.manifestInfo.wrap(newManifest); + return this; + } + + @Override + public Tracking tracking() { + return tracking; + } + + @Override + public FileContent contentType() { + return contentType; + } + + @Override + public int formatVersion() { + return manifest.formatVersion(); + } + + @Override + public String location() { + return manifest.path(); + } + + @Override + public FileFormat fileFormat() { + // Manifest files before v4 are always Avro. + return FileFormat.AVRO; + } + + @Override + public long recordCount() { + // Number of TrackedFile rows stored in the manifest. + return recordCount; + } + + @Override + public long fileSizeInBytes() { + return manifest.length(); + } + + @Override + public Integer specId() { + // Spec the wrapped manifest was written with. Data file entries in a v4 manifest file may + // use different specs. + return manifest.partitionSpecId(); + } + + @Override + public StructLike partition() { + return null; + } + + @Override + public ContentStats contentStats() { + return null; + } + + @Override + public Integer sortOrderId() { + // Manifests have no table sort order. + return null; + } + + @Override + public DeletionVector deletionVector() { + return null; + } + + @Override + public ManifestInfo manifestInfo() { + return manifestInfo; + } + + @Override + public ByteBuffer keyMetadata() { + return manifest.keyMetadata(); + } + + @Override + public List splitOffsets() { + return null; + } + + @Override + public List equalityIds() { + return null; + } + + @Override + public TrackedFile copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + + @Override + public TrackedFile copyWithStats(Set requestedColumnIds) { + throw new UnsupportedOperationException("copy is not implemented"); + } + } + + /** Reusable {@link ManifestInfo} view over a {@link ManifestFile}'s counts. */ + private static class WrappedManifestInfo implements ManifestInfo { + private ManifestFile manifest = null; + + WrappedManifestInfo wrap(ManifestFile newManifest) { + this.manifest = newManifest; + return this; + } + + @Override + public int addedFilesCount() { + return manifest.addedFilesCount(); + } + + @Override + public int existingFilesCount() { + return manifest.existingFilesCount(); + } + + @Override + public int deletedFilesCount() { + return manifest.deletedFilesCount(); + } + + @Override + public int replacedFilesCount() { + return manifest.replacedFilesCount(); + } + + @Override + public int modifiedFilesCount() { + return manifest.modifiedFilesCount(); + } + + @Override + public long addedRowsCount() { + return manifest.addedRowsCount(); + } + + @Override + public long existingRowsCount() { + return manifest.existingRowsCount(); + } + + @Override + public long deletedRowsCount() { + return manifest.deletedRowsCount(); + } + + @Override + public long replacedRowsCount() { + return manifest.replacedRowsCount(); + } + + @Override + public long modifiedRowsCount() { + return manifest.modifiedRowsCount(); + } + + @Override + public long minSequenceNumber() { + return manifest.minSequenceNumber(); + } + + @Override + public ManifestBitmap manifestDeletionVector() { + return manifest.manifestDeletionVector(); + } + + @Override + public ManifestInfo copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + } + + private static class WrappedEntryTracking implements Tracking { + private ManifestEntry entry = null; + + WrappedEntryTracking wrap(ManifestEntry newEntry) { + this.entry = newEntry; + return this; + } + + @Override + public EntryStatus status() { + return entryStatus(entry.status()); + } + + @Override + public Long snapshotId() { + return entry.snapshotId(); + } + + @Override + public Long dataSequenceNumber() { + return entry.dataSequenceNumber(); + } + + @Override + public Long fileSequenceNumber() { + return entry.fileSequenceNumber(); + } + + @Override + public Long dvSnapshotId() { + // dvSnapshotId is null because this wrapper has no DV entry + return null; + } + + @Override + public Long firstRowId() { + return entry.file().firstRowId(); + } + + @Override + public ByteBuffer deletedPositions() { + return null; + } + + @Override + public ByteBuffer replacedPositions() { + return null; + } + + @Override + public String manifestLocation() { + return entry.file().manifestLocation(); + } + + @Override + public long manifestPos() { + Long pos = entry.file().pos(); + return pos != null ? pos : -1L; + } + + @Override + public Tracking copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + } + + private static class WrappedManifestTracking implements Tracking { + private ManifestFile manifest = null; + + WrappedManifestTracking wrap(ManifestFile newManifest) { + this.manifest = newManifest; + return this; + } + + @Override + public EntryStatus status() { + // Pre-v4 manifests have no status and are live; the writer sets EXISTING or MODIFIED. + return EntryStatus.EXISTING; + } + + @Override + public Long snapshotId() { + return manifest.snapshotId(); + } + + @Override + public Long dataSequenceNumber() { + return manifest.sequenceNumber(); + } + + @Override + public Long fileSequenceNumber() { + return manifest.sequenceNumber(); + } + + @Override + public Long dvSnapshotId() { + return null; + } + + @Override + public Long firstRowId() { + return manifest.firstRowId(); + } + + @Override + public ByteBuffer deletedPositions() { + return null; + } + + @Override + public ByteBuffer replacedPositions() { + return null; + } + + @Override + public String manifestLocation() { + return null; + } + + @Override + public long manifestPos() { + return -1L; + } + + @Override + public Tracking copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + } + + private static EntryStatus entryStatus(ManifestEntry.Status status) { + return switch (status) { + case EXISTING -> EntryStatus.EXISTING; + case ADDED -> EntryStatus.ADDED; + case DELETED -> EntryStatus.DELETED; + }; + } + + private static boolean hasContentStats(ContentFile file) { + return isPresent(file.valueCounts()) + || isPresent(file.nullValueCounts()) + || isPresent(file.nanValueCounts()) + || isPresent(file.avgValueSizes()) + || isPresent(file.lowerBounds()) + || isPresent(file.upperBounds()); + } + + private static boolean isPresent(Map map) { + return map != null && !map.isEmpty(); + } + + /** + * Record count of a manifest is the number of TrackedFile rows it stores: the sum of per-status + * file counts. Missing counts fail rather than producing an incorrect total. + */ + private static long manifestRecordCount(ManifestFile manifest) { + Preconditions.checkNotNull( + manifest.addedFilesCount(), + "Cannot convert manifest %s: missing added files count", + manifest.path()); + Preconditions.checkNotNull( + manifest.existingFilesCount(), + "Cannot convert manifest %s: missing existing files count", + manifest.path()); + Preconditions.checkNotNull( + manifest.deletedFilesCount(), + "Cannot convert manifest %s: missing deleted files count", + manifest.path()); + Preconditions.checkArgument( + Objects.equals(manifest.replacedFilesCount(), 0), + "Cannot convert manifest %s: Invalid replaced file count: %s", + manifest.path(), + manifest.replacedFilesCount()); + Preconditions.checkArgument( + Objects.equals(manifest.modifiedFilesCount(), 0), + "Cannot convert manifest %s: Invalid modified file count: %s", + manifest.path(), + manifest.modifiedFilesCount()); + return (long) manifest.addedFilesCount() + + manifest.existingFilesCount() + + manifest.deletedFilesCount() + + manifest.replacedFilesCount() + + manifest.modifiedFilesCount(); + } + private static int resolveSpecId(TrackedFile file, Map specsById) { Integer specId = file.specId(); if (specId != null) { diff --git a/core/src/test/java/org/apache/iceberg/TestMapBackedContentStats.java b/core/src/test/java/org/apache/iceberg/TestMapBackedContentStats.java new file mode 100644 index 000000000000..ca8e92dbcae0 --- /dev/null +++ b/core/src/test/java/org/apache/iceberg/TestMapBackedContentStats.java @@ -0,0 +1,382 @@ +/* + * 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; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +import java.math.BigDecimal; +import java.nio.ByteBuffer; +import java.util.Comparator; +import java.util.List; +import java.util.Map; +import java.util.UUID; +import java.util.stream.Stream; +import org.apache.iceberg.geospatial.GeospatialBound; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; +import org.apache.iceberg.types.Comparators; +import org.apache.iceberg.types.Conversions; +import org.apache.iceberg.types.Type; +import org.apache.iceberg.types.Types; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.MethodSource; + +class TestMapBackedContentStats { + + private static final Schema SCHEMA = + new Schema( + Types.NestedField.required(1, "id", Types.IntegerType.get()), + Types.NestedField.optional(2, "score", Types.FloatType.get()), + Types.NestedField.optional(3, "ts", Types.LongType.get()), + Types.NestedField.optional(4, "name", Types.StringType.get()), + Types.NestedField.optional(5, "flag", Types.BooleanType.get())); + + private static final PartitionData EMPTY_PARTITION = + new PartitionData(PartitionSpec.unpartitioned().partitionType()); + + /** + * Stats on fields 1-4; field 5 is absent. Field 3 has no value-count entry, so valueCount() + * throws on unboxing. + */ + private static final DataFile FILE_WITH_STATS = + dataFile( + 100L, + ImmutableMap.of(1, 100L, 2, 100L, 4, 100L), + ImmutableMap.of(2, 5L, 3, 1L, 4, 2L), + ImmutableMap.of(2, 3L), + ImmutableMap.of( + 1, buf(Types.IntegerType.get(), 1), + 2, buf(Types.FloatType.get(), 1.5f), + 3, buf(Types.LongType.get(), 100L), + 4, buf(Types.StringType.get(), "aaa")), + ImmutableMap.of( + 1, buf(Types.IntegerType.get(), 1000), + 2, buf(Types.FloatType.get(), 9.5f), + 3, buf(Types.LongType.get(), 999L), + 4, buf(Types.StringType.get(), "zzz"))); + + @Test + void typeMatchesIdsPresentInMaps() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA); + + assertThat(stats.type().fields()).isEmpty(); + + stats.wrap(FILE_WITH_STATS); + + assertThat(stats.type().fields()) + .containsExactlyInAnyOrderElementsOf( + StatsUtil.statsReadSchema(SCHEMA, List.of(1, 2, 3, 4)).fields()); + } + + @Test + void wrapInvalidatesType() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + Types.StructType firstType = stats.type(); + assertThat(firstType.fields()) + .containsExactlyInAnyOrderElementsOf( + StatsUtil.statsReadSchema(SCHEMA, List.of(1, 2, 3, 4)).fields()); + + DataFile file2 = + dataFile( + 100L, + ImmutableMap.of(1, 50L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(1, buf(Types.IntegerType.get(), 500)), + ImmutableMap.of(1, buf(Types.IntegerType.get(), 5000))); + stats.wrap(file2); + + assertThat(stats.type().fields()) + .containsExactlyInAnyOrderElementsOf( + StatsUtil.statsReadSchema(SCHEMA, List.of(1)).fields()); + assertThat(stats.type()).isNotEqualTo(firstType); + } + + @ParameterizedTest + @MethodSource("typesAndBounds") + void boundDecodingPerType(Type type, Object lower, Object upper) { + int fieldId = 1; + Schema schema = new Schema(Types.NestedField.optional(fieldId, "col", type)); + DataFile file = + dataFile( + 100L, + ImmutableMap.of(fieldId, 1L), + ImmutableMap.of(), + ImmutableMap.of(), + bound(fieldId, Conversions.toByteBuffer(type, lower)), + bound(fieldId, Conversions.toByteBuffer(type, upper))); + FieldStats stats = new MapBackedContentStats(schema).wrap(file).statsFor(fieldId); + + Comparator comparator = Comparators.forType(type.asPrimitiveType()); + assertThat(comparator.compare(stats.lowerBound(), lower)).isZero(); + assertThat(comparator.compare(stats.upperBound(), upper)).isZero(); + } + + private static Stream typesAndBounds() { + return Stream.of( + Arguments.of(Types.BooleanType.get(), false, true), + Arguments.of(Types.IntegerType.get(), -5, 100), + Arguments.of(Types.LongType.get(), 0L, 1_000L), + Arguments.of(Types.FloatType.get(), 1.5f, 9.5f), + Arguments.of(Types.DoubleType.get(), 0.0d, 25.0d), + Arguments.of(Types.DateType.get(), 100, 200), + Arguments.of(Types.TimeType.get(), 1_000L, 2_000L), + Arguments.of(Types.TimestampType.withZone(), 111L, 222L), + Arguments.of(Types.TimestampNanoType.withZone(), 111L, 222L), + Arguments.of(Types.StringType.get(), "a", "z"), + Arguments.of( + Types.UUIDType.get(), + UUID.fromString("07ceab48-62b2-4219-9172-856e687c92ad"), + UUID.fromString("a9a7c24d-2869-4c66-8803-b1f6a36256ed")), + Arguments.of( + Types.FixedType.ofLength(4), + ByteBuffer.wrap(new byte[] {0, 1, 2, 3}), + ByteBuffer.wrap(new byte[] {4, 5, 6, 7})), + Arguments.of( + Types.BinaryType.get(), + ByteBuffer.wrap(new byte[] {1, 2}), + ByteBuffer.wrap(new byte[] {3, 4, 5})), + Arguments.of(Types.DecimalType.of(9, 2), new BigDecimal("1.23"), new BigDecimal("9.99")), + Arguments.of(Types.UnknownType.get(), null, null)); + } + + @Test + void geoBoundsDecode() { + GeospatialBound lower = GeospatialBound.createXY(1.0, 2.0); + GeospatialBound upper = GeospatialBound.createXYZM(3.0, 4.0, 5.0, 6.0); + Schema schema = new Schema(Types.NestedField.optional(10, "geom", Types.GeometryType.crs84())); + DataFile file = + dataFile( + 100L, + ImmutableMap.of(10, 26L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(10, lower.toByteBuffer()), + ImmutableMap.of(10, upper.toByteBuffer())); + MapBackedContentStats stats = new MapBackedContentStats(schema).wrap(file); + + FieldStats geom = stats.statsFor(10); + assertThat(geom.lowerBound()).isInstanceOf(GeospatialBound.class).isEqualTo(lower); + assertThat(geom.upperBound()).isInstanceOf(GeospatialBound.class).isEqualTo(upper); + } + + @Test + void missingBoundsDecodeToNull() { + DataFile file = + dataFile( + 100L, + ImmutableMap.of(1, 100L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of()); + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(file); + + FieldStats id = stats.statsFor(1); + assertThat(id.lowerBound()).isNull(); + assertThat(id.upperBound()).isNull(); + assertThat(id.valueCount()).isEqualTo(100L); + } + + @Test + void countsAndTightBounds() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + + FieldStats id = stats.statsFor(1); + assertThat(id.hasValueCount()).isTrue(); + assertThat(id.valueCount()).isEqualTo(100L); + assertThat(id.hasNullValueCount()).isFalse(); + assertThat(id.hasNanValueCount()).isFalse(); + assertThatThrownBy(id::nullValueCount) + .isInstanceOf(NullPointerException.class) + .hasMessageContaining("Long.longValue()"); + assertThatThrownBy(id::nanValueCount) + .isInstanceOf(NullPointerException.class) + .hasMessageContaining("Long.longValue()"); + // ContentFile maps have no tight-bounds flag so the view always reports false. + assertThat(id.tightBounds()).isFalse(); + // FILE_WITH_STATS has no avg-size map. + assertThat(id.avgValueSizeInBytes()).isNull(); + + FieldStats score = stats.statsFor(2); + assertThat(score.hasNullValueCount()).isTrue(); + assertThat(score.nullValueCount()).isEqualTo(5L); + assertThat(score.hasNanValueCount()).isTrue(); + assertThat(score.nanValueCount()).isEqualTo(3L); + + FieldStats ts = stats.statsFor(3); + assertThat(ts.hasValueCount()).isFalse(); + assertThatThrownBy(ts::valueCount) + .isInstanceOf(NullPointerException.class) + .hasMessageContaining("Long.longValue()"); + assertThat(ts.hasNullValueCount()).isTrue(); + assertThat(ts.nullValueCount()).isEqualTo(1L); + + FieldStats name = stats.statsFor(4); + assertThat(name.hasValueCount()).isTrue(); + assertThat(name.valueCount()).isEqualTo(100L); + assertThat(name.hasNullValueCount()).isTrue(); + assertThat(name.nullValueCount()).isEqualTo(2L); + assertThat(name.hasNanValueCount()).isFalse(); + assertThatThrownBy(name::nanValueCount) + .isInstanceOf(NullPointerException.class) + .hasMessageContaining("Long.longValue()"); + } + + @Test + void fieldWithoutStatsIsExcluded() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + assertThat(stats.type().field(StatsUtil.toBaseId(5))).isNull(); + assertThat(stats.statsFor(5)).isNull(); + } + + @Test + void unknownFieldIdInMaps() { + DataFile file = + dataFile( + 100L, + ImmutableMap.of(99, 1L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of()); + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(file); + + assertThat(stats.statsFor(99)).isNull(); + assertThat(stats.type().fields()).isEmpty(); + assertThat(stats.fieldStats()).isEmpty(); + } + + @Test + void reuseRebindsFields() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA); + + stats.wrap(FILE_WITH_STATS); + assertThat(stats.statsFor(1).lowerBound()).isEqualTo(1); + assertThat(stats.statsFor(1).upperBound()).isEqualTo(1000); + assertThat(stats.statsFor(1).valueCount()).isEqualTo(100L); + + DataFile file2 = + dataFile( + 100L, + ImmutableMap.of(1, 50L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(1, buf(Types.IntegerType.get(), 500)), + ImmutableMap.of(1, buf(Types.IntegerType.get(), 5000))); + stats.wrap(file2); + + assertThat(stats.statsFor(1).lowerBound()).isEqualTo(500); + assertThat(stats.statsFor(1).upperBound()).isEqualTo(5000); + assertThat(stats.statsFor(1).valueCount()).isEqualTo(50L); + assertThat(stats.statsFor(2)).isNull(); + } + + @Test + void absentIdIsRereadWhenPresentAgain() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA); + stats.wrap(FILE_WITH_STATS); + assertThat(stats.statsFor(1)).isNotNull(); + + stats.wrap( + dataFile( + 100L, + ImmutableMap.of(2, 10L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(2, buf(Types.FloatType.get(), 0.0f)), + ImmutableMap.of(2, buf(Types.FloatType.get(), 1.0f)))); + assertThat(stats.statsFor(1)).isNull(); + + stats.wrap( + dataFile( + 100L, + ImmutableMap.of(1, 7L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(1, buf(Types.IntegerType.get(), 42)), + ImmutableMap.of(1, buf(Types.IntegerType.get(), 43)))); + FieldStats id = stats.statsFor(1); + assertThat(id.lowerBound()).isEqualTo(42); + assertThat(id.upperBound()).isEqualTo(43); + assertThat(id.valueCount()).isEqualTo(7L); + } + + @Test + void listElementFieldStats() { + Schema schema = + new Schema( + Types.NestedField.required( + 1, "nums", Types.ListType.ofRequired(2, Types.IntegerType.get()))); + DataFile file = + dataFile( + 100L, + ImmutableMap.of(2, 4L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(2, buf(Types.IntegerType.get(), 1)), + ImmutableMap.of(2, buf(Types.IntegerType.get(), 9))); + MapBackedContentStats stats = new MapBackedContentStats(schema).wrap(file); + + FieldStats element = stats.statsFor(2); + assertThat(element.lowerBound()).isEqualTo(1); + assertThat(element.upperBound()).isEqualTo(9); + assertThat(stats.fieldStats()).extracting(FieldStats::fieldId).containsExactly(2); + } + + private static ByteBuffer buf(Type type, Object value) { + return Conversions.toByteBuffer(type, value); + } + + private static Map bound(int fieldId, ByteBuffer value) { + return value == null ? ImmutableMap.of() : ImmutableMap.of(fieldId, value); + } + + private static DataFile dataFile( + long recordCount, + Map valueCounts, + Map nullValueCounts, + Map nanValueCounts, + Map lowerBounds, + Map upperBounds) { + Metrics metrics = + new Metrics( + recordCount, + null, + valueCounts, + nullValueCounts, + nanValueCounts, + lowerBounds, + upperBounds); + return new GenericDataFile( + PartitionSpec.unpartitioned().specId(), + "s3://bucket/data/file.parquet", + FileFormat.PARQUET, + EMPTY_PARTITION, + 1024L, + metrics, + null, + ImmutableList.of(0L), + null, + null); + } +} diff --git a/core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java b/core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java index a7bfe777c147..3c6799989766 100644 --- a/core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java +++ b/core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java @@ -25,10 +25,12 @@ import static org.mockito.Mockito.when; import java.nio.ByteBuffer; +import java.util.List; import java.util.Map; import org.apache.iceberg.mumbling.MumblingTestUtil; import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList; import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; +import org.apache.iceberg.types.Comparators; import org.apache.iceberg.types.Conversions; import org.apache.iceberg.types.Types; import org.junit.jupiter.api.Test; @@ -50,6 +52,16 @@ class TestTrackedFileAdapters { private static final long DATA_SEQUENCE_NUMBER = 10L; private static final long FILE_SEQUENCE_NUMBER = 11L; private static final long FIRST_ROW_ID = 1000L; + private static final int SORT_ORDER_ID = 3; + + private static final long MANIFEST_SEQUENCE_NUMBER = 5L; + private static final long MANIFEST_MIN_SEQUENCE_NUMBER = 4L; + private static final int ADDED_FILES_COUNT = 2; + private static final long ADDED_ROWS_COUNT = 200L; + private static final int EXISTING_FILES_COUNT = 3; + private static final long EXISTING_ROWS_COUNT = 300L; + private static final int DELETED_FILES_COUNT = 1; + private static final long DELETED_ROWS_COUNT = 100L; private static final int UNPARTITIONED_SPEC_ID = PartitionSpec.unpartitioned().specId(); private static final Map UNPARTITIONED = @@ -70,6 +82,9 @@ class TestTrackedFileAdapters { // these are populated by readers using the setter with the position of the field. private static final int MANIFEST_LOCATION_ORDINAL = Tracking.schema().fields().size(); private static final int MANIFEST_POSITION_ORDINAL = Tracking.schema().fields().size() + 1; + // row_position follows the data file schema and stores the manifest position. + private static final int DATA_FILE_POS_ORDINAL = + DataFile.getType(Types.StructType.of()).fields().size(); private static final Schema TABLE_SCHEMA = new Schema( @@ -96,6 +111,9 @@ class TestTrackedFileAdapters { CONTENT_STATS.setStats(3, GEOM_STATS); } + private static final PartitionSpec UNPARTITIONED_SPEC = PartitionSpec.unpartitioned(); + private static final Types.StructType PARTITION_TYPE = UNPARTITIONED_SPEC.partitionType(); + private static final Tracking MANIFEST_TRACKING = new TrackingStruct( EntryStatus.ADDED, @@ -114,17 +132,57 @@ class TestTrackedFileAdapters { .addedFilesCount(3) .existingFilesCount(5) .deletedFilesCount(2) - .replacedFilesCount(0) - .modifiedFilesCount(0) + .replacedFilesCount(4) + .modifiedFilesCount(1) .addedRowsCount(300L) .existingRowsCount(500L) .deletedRowsCount(200L) - .replacedRowsCount(0L) - .modifiedRowsCount(0L) + .replacedRowsCount(40L) + .modifiedRowsCount(10L) .minSequenceNumber(7L) .dv(ByteBuffer.wrap(MumblingTestUtil.onlyFirstBitSetBytes())) .build(); + private static final Metrics METRICS_WITH_BOUNDS = + new Metrics( + 100L, + ImmutableMap.of(1, 16L, 2, 64L), + ImmutableMap.of(1, 100L, 2, 100L), + ImmutableMap.of(1, 0L, 2, 5L), + ImmutableMap.of(), + ImmutableMap.of(1, Conversions.toByteBuffer(Types.IntegerType.get(), 1)), + ImmutableMap.of(1, Conversions.toByteBuffer(Types.IntegerType.get(), 1000))); + + private static final DataFile DATA_FILE = + new GenericDataFile( + UNPARTITIONED_SPEC.specId(), + DATA_FILE_LOCATION, + FileFormat.PARQUET, + PartitionData.EMPTY, + 1024L, + new Metrics(100L), + KEY_METADATA, + ImmutableList.of(0L), + SORT_ORDER_ID, + FIRST_ROW_ID); + private static final DataFile DATA_FILE_WITH_METRICS = + new GenericDataFile( + UNPARTITIONED_SPEC.specId(), + DATA_FILE_LOCATION, + FileFormat.PARQUET, + PartitionData.EMPTY, + 1024L, + METRICS_WITH_BOUNDS, + null, + ImmutableList.of(0L), + null, + null); + + static { + assignManifestPosition(DATA_FILE, MANIFEST_LOCATION, MANIFEST_POS); + assignManifestPosition(DATA_FILE_WITH_METRICS, MANIFEST_LOCATION, MANIFEST_POS); + } + @Test void dataFileAdapterDelegation() { TrackingStruct tracking = @@ -199,7 +257,7 @@ void dataFileAdapterDelegation() { @ParameterizedTest @EnumSource(value = FileContent.class, mode = EnumSource.Mode.EXCLUDE, names = "DATA") void dataFileAdapterRejectsNonDataContent(FileContent contentType) { - TrackedFileStruct file = dummyTrackedFile(contentType); + TrackedFileStruct file = trackedFile(contentType); assertThatThrownBy(() -> TrackedFileAdapters.asDataFile(file, UNPARTITIONED)) .isInstanceOf(IllegalArgumentException.class) @@ -279,7 +337,7 @@ void equalityDeleteFileAdapterDelegation() { @ParameterizedTest @EnumSource(value = FileContent.class, mode = EnumSource.Mode.EXCLUDE, names = "EQUALITY_DELETES") void equalityDeleteFileAdapterRejectsNonEqualityContent(FileContent contentType) { - TrackedFileStruct file = dummyTrackedFile(contentType); + TrackedFileStruct file = trackedFile(contentType); assertThatThrownBy(() -> TrackedFileAdapters.asEqualityDeleteFile(file, UNPARTITIONED)) .isInstanceOf(IllegalArgumentException.class) @@ -465,7 +523,7 @@ private static TrackedFile serializableDataEntry() { @ParameterizedTest @EnumSource(value = FileContent.class, mode = EnumSource.Mode.EXCLUDE, names = "DATA") void dvDeleteFileAdapterRejectsNonDataContent(FileContent contentType) { - TrackedFileStruct file = dummyTrackedFile(contentType); + TrackedFileStruct file = trackedFile(contentType); assertThatThrownBy(() -> TrackedFileAdapters.asDVDeleteFile(file, UNPARTITIONED)) .isInstanceOf(IllegalArgumentException.class) @@ -474,7 +532,7 @@ void dvDeleteFileAdapterRejectsNonDataContent(FileContent contentType) { @Test void dvDeleteFileAdapterRejectsNullDeletionVector() { - TrackedFileStruct file = dummyTrackedFile(FileContent.DATA); + TrackedFileStruct file = trackedFile(FileContent.DATA); assertThatThrownBy(() -> TrackedFileAdapters.asDVDeleteFile(file, UNPARTITIONED)) .isInstanceOf(IllegalArgumentException.class) @@ -521,6 +579,10 @@ void manifestFileAdapterDelegation(FileContent contentType) { assertThat(manifest.existingRowsCount()).isEqualTo(MANIFEST_INFO.existingRowsCount()); assertThat(manifest.deletedFilesCount()).isEqualTo(MANIFEST_INFO.deletedFilesCount()); assertThat(manifest.deletedRowsCount()).isEqualTo(MANIFEST_INFO.deletedRowsCount()); + assertThat(manifest.replacedFilesCount()).isEqualTo(MANIFEST_INFO.replacedFilesCount()); + assertThat(manifest.replacedRowsCount()).isEqualTo(MANIFEST_INFO.replacedRowsCount()); + assertThat(manifest.modifiedFilesCount()).isEqualTo(MANIFEST_INFO.modifiedFilesCount()); + assertThat(manifest.modifiedRowsCount()).isEqualTo(MANIFEST_INFO.modifiedRowsCount()); assertThat(manifest.firstRowId()).isEqualTo(FIRST_ROW_ID); assertThat(manifest.keyMetadata()).isEqualTo(MANIFEST_KEY_METADATA); assertThat(manifest.manifestDeletionVector().buffer()) @@ -553,7 +615,7 @@ void manifestFileAdapterCopy() { mode = EnumSource.Mode.EXCLUDE, names = {"DATA_MANIFEST", "DELETE_MANIFEST"}) void manifestFileAdapterRejectsNonManifestContent(FileContent contentType) { - TrackedFileStruct file = dummyTrackedFile(contentType); + TrackedFileStruct file = trackedFile(contentType); assertThatThrownBy(() -> TrackedFileAdapters.asManifestFile(file)) .isInstanceOf(IllegalArgumentException.class) @@ -572,7 +634,7 @@ void dataFileWithoutDeletionVectorReturnsNull() { @Test void nullContentStatsReturnsNullStats() { - TrackedFileStruct file = dummyTrackedFile(FileContent.DATA); + TrackedFileStruct file = trackedFile(FileContent.DATA); DataFile dataFile = TrackedFileAdapters.asDataFile(file, UNPARTITIONED); @@ -588,10 +650,10 @@ void nullTrackingReturnsNullTrackingFields() { // Files read before manifest inheritance have no tracking; tracking-derived fields must be // null rather than throwing. assertNullTrackingFields( - TrackedFileAdapters.asDataFile(dummyTrackedFile(FileContent.DATA), UNPARTITIONED)); + TrackedFileAdapters.asDataFile(trackedFile(FileContent.DATA), UNPARTITIONED)); assertNullTrackingFields( TrackedFileAdapters.asEqualityDeleteFile( - dummyTrackedFile(FileContent.EQUALITY_DELETES), UNPARTITIONED)); + trackedFile(FileContent.EQUALITY_DELETES), UNPARTITIONED)); TrackedFileStruct fileWithDV = new TrackedFileStruct( @@ -616,7 +678,7 @@ void nullTrackingReturnsNullTrackingFields() { @Test void unpartitionedFilePartitionIsEmpty() { - TrackedFileStruct file = dummyTrackedFile(FileContent.DATA); + TrackedFileStruct file = trackedFile(FileContent.DATA); DataFile dataFile = TrackedFileAdapters.asDataFile(file, UNPARTITIONED); @@ -627,7 +689,7 @@ void unpartitionedFilePartitionIsEmpty() { @Test void nullSpecIdResolvesToUnpartitionedSpec() { PartitionSpec unpartitioned = PartitionSpec.builderFor(new Schema()).withSpecId(5).build(); - TrackedFileStruct file = dummyTrackedFile(FileContent.DATA); + TrackedFileStruct file = trackedFile(FileContent.DATA); DataFile dataFile = TrackedFileAdapters.asDataFile(file, specsById(unpartitioned)); @@ -638,7 +700,7 @@ void nullSpecIdResolvesToUnpartitionedSpec() { void nullSpecIdThrowsWhenNoUnpartitionedSpec() { Schema schema = new Schema(Types.NestedField.required(1, "id", Types.IntegerType.get())); PartitionSpec partitioned = PartitionSpec.builderFor(schema).identity("id").build(); - TrackedFileStruct file = dummyTrackedFile(FileContent.DATA); + TrackedFileStruct file = trackedFile(FileContent.DATA); assertThatThrownBy(() -> TrackedFileAdapters.asDataFile(file, specsById(partitioned))) .isInstanceOf(IllegalArgumentException.class) @@ -671,6 +733,312 @@ void unknownSpecIdThrows() { .hasMessageContaining("Cannot find partition spec for spec ID"); } + @Test + void dataTrackedFileAdapterFromDataFile() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE); + + assertThat(result.tracking()).isNull(); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + assertThatThrownBy(result::formatVersion) + .isInstanceOf(IllegalStateException.class) + .hasMessage("Format version is assigned at write time"); + } + + @Test + void dataTrackedFileAdapterFromExistingManifestEntry() { + TrackedFile result = + TrackedFileAdapters.forDataFile(TABLE_SCHEMA) + .wrap( + newEntry() + .wrapExisting( + SNAPSHOT_ID, DATA_SEQUENCE_NUMBER, FILE_SEQUENCE_NUMBER, DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().dataSequenceNumber()).isEqualTo(DATA_SEQUENCE_NUMBER); + assertThat(result.tracking().fileSequenceNumber()).isEqualTo(FILE_SEQUENCE_NUMBER); + assertThat(result.tracking().firstRowId()).isEqualTo(FIRST_ROW_ID); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterReuse() { + TrackedFileAdapters.DataTrackedFile adapter = TrackedFileAdapters.forDataFile(TABLE_SCHEMA); + + adapter.wrap(DATA_FILE); + assertWrappedDataFileMatchesFileFields(adapter, DATA_FILE); + assertThat(adapter.tracking()).isNull(); + + DataFile file2 = + new GenericDataFile( + UNPARTITIONED_SPEC.specId(), + "s3://bucket/data/file2.parquet", + FileFormat.PARQUET, + PartitionData.EMPTY, + 2048L, + new Metrics(200L, null, null, null, null), + null, + ImmutableList.of(0L), + null, + null); + assignManifestPosition(file2, "s3://bucket/table/manifest-2.parquet", 8L); + adapter.wrap( + newEntry().wrapExisting(SNAPSHOT_ID, DATA_SEQUENCE_NUMBER, FILE_SEQUENCE_NUMBER, file2)); + assertWrappedDataFileMatchesFileFields(adapter, file2); + assertThat(adapter.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertManifestPosition(adapter.tracking(), file2); + } + + @Test + void dataTrackedFileAdapterRejectsNullFile() { + TrackedFileAdapters.DataTrackedFile adapter = TrackedFileAdapters.forDataFile(TABLE_SCHEMA); + assertThatThrownBy(() -> adapter.wrap((DataFile) null)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("Invalid file: null"); + } + + @Test + void dataTrackedFileAdapterContentStats() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE_WITH_METRICS); + + ContentStats stats = result.contentStats(); + assertThat(stats).isNotNull(); + assertThat(stats.fieldStats()).extracting(FieldStats::fieldId).containsExactlyInAnyOrder(1, 2); + + FieldStats idStats = stats.statsFor(1); + assertThat(idStats.valueCount()).isEqualTo(100L); + assertThat(idStats.lowerBound()).isEqualTo(1); + assertThat(idStats.upperBound()).isEqualTo(1000); + } + + @Test + void dataTrackedFileAdapterWithoutMetricsHasNoContentStats() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE); + + assertThat(result.contentStats()).isNull(); + } + + @Test + void dataTrackedFileAdapterKeepsPartitionTuple() { + DataFile partitioned = + new GenericDataFile( + PARTITIONED_SPEC.specId(), + DATA_FILE_LOCATION, + FileFormat.PARQUET, + PARTITION, + 1024L, + new Metrics(100L, null, null, null, null), + null, + ImmutableList.of(0L), + null, + null); + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(partitioned); + + assertThat(result.partition()) + .usingComparator(Comparators.forType(PARTITIONED_SPEC.partitionType())) + .isEqualTo(PARTITION); + assertThat(TrackedFileAdapters.asDataFile(result, specsById(PARTITIONED_SPEC)).partition()) + .usingComparator(Comparators.forType(PARTITIONED_SPEC.partitionType())) + .isEqualTo(PARTITION); + } + + @Test + void dataTrackedFileAdapterUnwrapsToOriginalTrackedFile() { + TrackedFile source = trackedFile(FileContent.DATA); + DataFile dataFile = TrackedFileAdapters.asDataFile(source, UNPARTITIONED); + TrackedFile roundTripped = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(dataFile); + assertThat(roundTripped).isSameAs(source); + } + + @Test + void manifestTrackedFileAdapterUnwrapsToOriginalTrackedFile() { + TrackedFile original = trackedFile(FileContent.DATA_MANIFEST, 0); + ManifestFile adapted = TrackedFileAdapters.asManifestFile(original); + TrackedFile result = TrackedFileAdapters.forManifestFile().wrap(adapted); + assertThat(result).isSameAs(original); + assertThat(result.formatVersion()).isZero(); + } + + @Test + void dataManifestTrackedFileAdapter() { + ManifestFile manifest = writeManifestFile(ManifestContent.DATA); + TrackedFile result = TrackedFileAdapters.forManifestFile().wrap(manifest); + + assertThat(result.contentType()).isEqualTo(FileContent.DATA_MANIFEST); + assertThat(result.formatVersion()).isZero(); + assertThat(result.location()).isEqualTo(MANIFEST_LOCATION); + assertThat(result.fileFormat()).isEqualTo(FileFormat.AVRO); + assertThat(result.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().firstRowId()).isNull(); + assertThat(result.recordCount()) + .isEqualTo(ADDED_FILES_COUNT + EXISTING_FILES_COUNT + DELETED_FILES_COUNT); + assertThat(result.manifestInfo()).isNotNull(); + assertThat(result.manifestInfo().addedFilesCount()).isEqualTo(ADDED_FILES_COUNT); + assertThat(result.manifestInfo().existingFilesCount()).isEqualTo(EXISTING_FILES_COUNT); + assertThat(result.manifestInfo().deletedFilesCount()).isEqualTo(DELETED_FILES_COUNT); + assertThat(result.manifestInfo().addedRowsCount()).isEqualTo(ADDED_ROWS_COUNT); + assertThat(result.manifestInfo().replacedFilesCount()).isEqualTo(0); + assertThat(result.manifestInfo().replacedRowsCount()).isEqualTo(0L); + assertThat(result.manifestInfo().modifiedFilesCount()).isEqualTo(0); + assertThat(result.manifestInfo().modifiedRowsCount()).isEqualTo(0L); + } + + @Test + void deleteManifestTrackedFileAdapter() { + ManifestFile manifest = writeManifestFile(ManifestContent.DELETES); + TrackedFile result = TrackedFileAdapters.forManifestFile().wrap(manifest); + + assertThat(result.contentType()).isEqualTo(FileContent.DELETE_MANIFEST); + assertThat(result.formatVersion()).isZero(); + assertThat(result.recordCount()) + .isEqualTo(ADDED_FILES_COUNT + EXISTING_FILES_COUNT + DELETED_FILES_COUNT); + assertThat(result.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertThat(result.tracking().firstRowId()).isNull(); + } + + @Test + void manifestTrackedFileAdapterFailsWhenAddedFilesCountMissing() { + ManifestFile manifest = mock(ManifestFile.class); + when(manifest.path()).thenReturn(MANIFEST_LOCATION); + when(manifest.content()).thenReturn(ManifestContent.DATA); + when(manifest.addedFilesCount()).thenReturn(null); + when(manifest.existingFilesCount()).thenReturn(1); + when(manifest.deletedFilesCount()).thenReturn(0); + when(manifest.replacedFilesCount()).thenReturn(0); + + assertThatThrownBy(() -> TrackedFileAdapters.forManifestFile().wrap(manifest)) + .isInstanceOf(NullPointerException.class) + .hasMessageContaining("missing added files count"); + } + + @Test + void manifestTrackedFileAdapterRejectsNullReplacedFilesCount() { + ManifestFile manifest = manifestWithCounts(1, 1, 0, null, 0); + + assertThatThrownBy(() -> TrackedFileAdapters.forManifestFile().wrap(manifest)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage( + "Cannot convert manifest %s: Invalid replaced file count: null", MANIFEST_LOCATION); + } + + @Test + void manifestTrackedFileAdapterRejectsNonZeroReplacedFilesCount() { + ManifestFile manifest = manifestWithCounts(1, 1, 0, 1, 0); + + assertThatThrownBy(() -> TrackedFileAdapters.forManifestFile().wrap(manifest)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage( + "Cannot convert manifest %s: Invalid replaced file count: 1", MANIFEST_LOCATION); + } + + @Test + void manifestTrackedFileAdapterRejectsNullModifiedFilesCount() { + ManifestFile manifest = manifestWithCounts(1, 1, 0, 0, null); + + assertThatThrownBy(() -> TrackedFileAdapters.forManifestFile().wrap(manifest)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage( + "Cannot convert manifest %s: Invalid modified file count: null", MANIFEST_LOCATION); + } + + @Test + void manifestTrackedFileAdapterRejectsNonZeroModifiedFilesCount() { + ManifestFile manifest = manifestWithCounts(1, 1, 0, 0, 1); + + assertThatThrownBy(() -> TrackedFileAdapters.forManifestFile().wrap(manifest)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage( + "Cannot convert manifest %s: Invalid modified file count: 1", MANIFEST_LOCATION); + } + + @Test + void manifestTrackedFileAdapterFailsWhenAddedRowsCountMissing() { + ManifestFile manifest = + writeManifestFile( + ManifestContent.DATA, MANIFEST_SEQUENCE_NUMBER, MANIFEST_MIN_SEQUENCE_NUMBER, null); + TrackedFile tracked = TrackedFileAdapters.forManifestFile().wrap(manifest); + + assertThatThrownBy(() -> tracked.manifestInfo().addedRowsCount()) + .isInstanceOf(NullPointerException.class) + .hasMessageContaining("null"); + } + + private static void assertWrappedDataFileMatchesFileFields(TrackedFile result, DataFile file) { + assertThat(result.contentType()).isEqualTo(FileContent.DATA); + assertThat(result.location()).isEqualTo(file.location()); + assertThat(result.fileFormat()).isEqualTo(file.format()); + assertThat(result.recordCount()).isEqualTo(file.recordCount()); + assertThat(result.fileSizeInBytes()).isEqualTo(file.fileSizeInBytes()); + assertThat(result.specId()).isEqualTo(file.specId()); + assertThat(result.sortOrderId()).isEqualTo(file.sortOrderId()); + assertThat(result.keyMetadata()).isEqualTo(file.keyMetadata()); + assertThat(result.splitOffsets()).isEqualTo(file.splitOffsets()); + assertThat(result.manifestInfo()).isNull(); + assertThat(result.deletionVector()).isNull(); + assertThat(result.equalityIds()).isNull(); + } + + private static ManifestFile manifestWithCounts( + Integer addedFilesCount, + Integer existingFilesCount, + Integer deletedFilesCount, + Integer replacedFilesCount, + Integer modifiedFilesCount) { + ManifestFile manifest = mock(ManifestFile.class); + when(manifest.path()).thenReturn(MANIFEST_LOCATION); + when(manifest.content()).thenReturn(ManifestContent.DATA); + when(manifest.addedFilesCount()).thenReturn(addedFilesCount); + when(manifest.existingFilesCount()).thenReturn(existingFilesCount); + when(manifest.deletedFilesCount()).thenReturn(deletedFilesCount); + when(manifest.replacedFilesCount()).thenReturn(replacedFilesCount); + when(manifest.modifiedFilesCount()).thenReturn(modifiedFilesCount); + return manifest; + } + + private static ManifestFile writeManifestFile(ManifestContent content) { + return writeManifestFile( + content, MANIFEST_SEQUENCE_NUMBER, MANIFEST_MIN_SEQUENCE_NUMBER, ADDED_ROWS_COUNT); + } + + private static ManifestFile writeManifestFile( + ManifestContent content, long sequenceNumber, long minSequenceNumber, Long addedRowsCount) { + List partitions = ImmutableList.of(); + return new GenericManifestFile( + MANIFEST_LOCATION, + MANIFEST_FILE_SIZE, + UNPARTITIONED_SPEC.specId(), + content, + sequenceNumber, + minSequenceNumber, + SNAPSHOT_ID, + partitions, + null, + ADDED_FILES_COUNT, + addedRowsCount, + EXISTING_FILES_COUNT, + EXISTING_ROWS_COUNT, + DELETED_FILES_COUNT, + DELETED_ROWS_COUNT, + null); + } + + private static GenericManifestEntry newEntry() { + return new GenericManifestEntry<>(ManifestEntry.getSchema(PARTITION_TYPE).asStruct()); + } + + private static void assignManifestPosition(DataFile file, String location, long manifestPos) { + GenericDataFile dataFile = (GenericDataFile) file; + dataFile.setManifestLocation(location); + dataFile.set(DATA_FILE_POS_ORDINAL, manifestPos); + } + + private static void assertManifestPosition(Tracking tracking, DataFile file) { + assertThat(tracking.manifestLocation()).isEqualTo(file.manifestLocation()); + assertThat(tracking.manifestPos()).isEqualTo(file.pos()); + } + private static void assertNullTrackingFields(ContentFile file) { assertThat(file.pos()).isNull(); assertThat(file.manifestLocation()).isNull(); @@ -690,12 +1058,16 @@ private static PartitionData partition(String category) { return partition; } - /** Minimal file with no tracking, used by the rejection and null-tracking tests. */ - private static TrackedFileStruct dummyTrackedFile(FileContent contentType) { + /** Minimal file for the rejection and null-tracking tests. */ + private static TrackedFileStruct trackedFile(FileContent contentType) { + return trackedFile(contentType, FORMAT_VERSION_V4); + } + + private static TrackedFileStruct trackedFile(FileContent contentType, int formatVersion) { return new TrackedFileStruct( null, contentType, - FORMAT_VERSION_V4, + formatVersion, DATA_FILE_LOCATION, FileFormat.PARQUET, 1L,