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