stevenzwu commented on code in PR #16936: URL: https://github.com/apache/iceberg/pull/16936#discussion_r4170363705
########## core/src/main/java/org/apache/iceberg/MapBackedContentStats.java: ########## @@ -0,0 +1,216 @@ +/* + * 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.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 MetricsConfig metricsConfig; + private final Types.StructType type; + private final Map<Integer, FieldStats<?>> statsById = Maps.newHashMap(); + + private Map<Integer, Long> valueCounts; + private Map<Integer, Long> nullValueCounts; + private Map<Integer, Long> nanValueCounts; + private Map<Integer, Integer> avgValueSizes; + private Map<Integer, ByteBuffer> lowerBounds; + private Map<Integer, ByteBuffer> upperBounds; + + MapBackedContentStats(Schema tableSchema, MetricsConfig metricsConfig) { + Preconditions.checkArgument(tableSchema != null, "Invalid table schema: null"); + Preconditions.checkArgument(metricsConfig != null, "Invalid metrics config: null"); + this.tableSchema = tableSchema; + this.metricsConfig = metricsConfig; + this.type = StatsUtil.statsWriteSchema(tableSchema, metricsConfig); Review Comment: the type is built once from the table schema and MetricsConfig. map IDs absent from that type are ignored: `hasStats` is false and `statsFor` returns null. building once is 3.6–5.4× faster than rebuilding on every wrap. throughput is ops/ms, JDK 21, higher is better. | Columns | Rebuilt each wrap | Built once | Speedup | | ---: | ---: | ---: | ---: | | 2 | 2284 | 9527 | 4.17× | | 10 | 559 | 2009 | 3.59× | | 50 | 83 | 445 | 5.36× | | 200 | 18 | 69 | 3.83× | ########## core/src/main/java/org/apache/iceberg/TrackedFileAdapters.java: ########## @@ -529,6 +566,427 @@ public ManifestFile copy() { } } + /** Adapts a {@link DataFile} to {@link TrackedFile}. */ + static class DataTrackedFile implements TrackedFile { + private final MapBackedContentStats statsWrapper; + private Tracking tracking; + private DataFile file; + private ContentStats stats; + + DataTrackedFile(Schema tableSchema, MetricsConfig metricsConfig) { + this.statsWrapper = new MapBackedContentStats(tableSchema, metricsConfig); + } + + /** Re-points this adapter at a {@link DataFile} from the public API. Tracking is unset. */ + public TrackedFile wrap(DataFile newFile) { + return wrapFile(newFile, null); + } + + /** + * Re-points this adapter at a {@link ManifestEntry}. Converts the contained data file and the + * entry's tracking fields. + */ + public TrackedFile wrap(ManifestEntry<DataFile> entry) { + Preconditions.checkArgument(entry != null, "Invalid entry: null"); + return wrapFile(entry.file(), trackingFrom(entry, entry.file())); + } + + private TrackedFile wrapFile(DataFile newFile, Tracking newTracking) { + if (newFile instanceof TrackedDataFile tracked) { + return tracked.file(); + } + + 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; + this.tracking = newTracking; + return this; + } + + @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<Long> splitOffsets() { + return file.splitOffsets(); + } + + @Override + public List<Integer> equalityIds() { + return null; + } + + @Override + public TrackedFile copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + + @Override + public TrackedFile copyWithStats(Set<Integer> 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 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. + */ + 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 = trackingFrom(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() { + return FileFormat.fromFileName(manifest.path()); + } + + @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() { + return null; + } + + @Override + public DeletionVector deletionVector() { + return null; + } + + @Override + public ManifestInfo manifestInfo() { + return manifestInfo; + } + + @Override + public ByteBuffer keyMetadata() { + return manifest.keyMetadata(); + } + + @Override + public List<Long> splitOffsets() { + return null; + } + + @Override + public List<Integer> equalityIds() { + return null; + } + + @Override + public TrackedFile copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + + @Override + public TrackedFile copyWithStats(Set<Integer> 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; + + void wrap(ManifestFile newManifest) { + this.manifest = newManifest; + } + + @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 Tracking trackingFrom(ManifestEntry<?> entry, ContentFile<?> file) { + return new TrackingStruct( + entryStatus(entry.status()), + entry.snapshotId(), + entry.dataSequenceNumber(), + entry.fileSequenceNumber(), + null, + file.firstRowId(), + null, + null); + } + + private static Tracking trackingFrom(ManifestFile manifest) { + return new TrackingStruct( + null, Review Comment: status() returns EntryStatus.EXISTING. the comment is: `Pre-v4 manifests have no status and are live; the writer sets EXISTING or MODIFIED.` ########## core/src/main/java/org/apache/iceberg/MapBackedContentStats.java: ########## @@ -0,0 +1,206 @@ +/* + * 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.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 Types.StructType type; + private final Map<Integer, FieldStats<?>> statsById = Maps.newHashMap(); + + private Map<Integer, Long> valueCounts; + private Map<Integer, Long> nullValueCounts; + private Map<Integer, Long> nanValueCounts; + private Map<Integer, Integer> avgValueSizes; + private Map<Integer, ByteBuffer> lowerBounds; + private Map<Integer, ByteBuffer> upperBounds; + + MapBackedContentStats(Schema tableSchema, MetricsConfig metricsConfig) { + Preconditions.checkArgument(tableSchema != null, "Invalid table schema: null"); + Preconditions.checkArgument(metricsConfig != null, "Invalid metrics config: null"); + this.type = StatsUtil.statsWriteSchema(tableSchema, metricsConfig); + } + + 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(); + return this; + } + + @Override + public Iterable<FieldStats<?>> fieldStats() { + return Iterables.filter( Review Comment: this is based on the option 2 for content stats type. `hasStats` checks the precomputed type first. if the id is not in that type, hasStats is false and statsFor returns null even when the id is in a map. fieldStats walks that type and drops ids with no map entry. ########## core/src/main/java/org/apache/iceberg/TrackedFileAdapters.java: ########## @@ -529,6 +566,427 @@ public ManifestFile copy() { } } + /** Adapts a {@link DataFile} to {@link TrackedFile}. */ + static class DataTrackedFile implements TrackedFile { + private final MapBackedContentStats statsWrapper; + private Tracking tracking; + private DataFile file; + private ContentStats stats; + + DataTrackedFile(Schema tableSchema, MetricsConfig metricsConfig) { + this.statsWrapper = new MapBackedContentStats(tableSchema, metricsConfig); + } + + /** Re-points this adapter at a {@link DataFile} from the public API. Tracking is unset. */ + public TrackedFile wrap(DataFile newFile) { + return wrapFile(newFile, null); + } + + /** + * Re-points this adapter at a {@link ManifestEntry}. Converts the contained data file and the + * entry's tracking fields. + */ + public TrackedFile wrap(ManifestEntry<DataFile> entry) { + Preconditions.checkArgument(entry != null, "Invalid entry: null"); + return wrapFile(entry.file(), trackingFrom(entry, entry.file())); + } + + private TrackedFile wrapFile(DataFile newFile, Tracking newTracking) { + if (newFile instanceof TrackedDataFile tracked) { + return tracked.file(); + } + + 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; + this.tracking = newTracking; + return this; + } + + @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<Long> splitOffsets() { + return file.splitOffsets(); + } + + @Override + public List<Integer> equalityIds() { + return null; + } + + @Override + public TrackedFile copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + + @Override + public TrackedFile copyWithStats(Set<Integer> 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 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. + */ + 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 = trackingFrom(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() { + return FileFormat.fromFileName(manifest.path()); + } + + @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() { + return null; + } + + @Override + public DeletionVector deletionVector() { + return null; + } + + @Override + public ManifestInfo manifestInfo() { + return manifestInfo; + } + + @Override + public ByteBuffer keyMetadata() { + return manifest.keyMetadata(); + } + + @Override + public List<Long> splitOffsets() { + return null; + } + + @Override + public List<Integer> equalityIds() { + return null; + } + + @Override + public TrackedFile copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + + @Override + public TrackedFile copyWithStats(Set<Integer> 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; + + void wrap(ManifestFile newManifest) { Review Comment: done ########## core/src/main/java/org/apache/iceberg/TrackedFileAdapters.java: ########## @@ -529,6 +566,427 @@ public ManifestFile copy() { } } + /** Adapts a {@link DataFile} to {@link TrackedFile}. */ + static class DataTrackedFile implements TrackedFile { + private final MapBackedContentStats statsWrapper; + private Tracking tracking; + private DataFile file; + private ContentStats stats; + + DataTrackedFile(Schema tableSchema, MetricsConfig metricsConfig) { + this.statsWrapper = new MapBackedContentStats(tableSchema, metricsConfig); + } + + /** Re-points this adapter at a {@link DataFile} from the public API. Tracking is unset. */ + public TrackedFile wrap(DataFile newFile) { + return wrapFile(newFile, null); + } + + /** + * Re-points this adapter at a {@link ManifestEntry}. Converts the contained data file and the + * entry's tracking fields. + */ + public TrackedFile wrap(ManifestEntry<DataFile> entry) { + Preconditions.checkArgument(entry != null, "Invalid entry: null"); + return wrapFile(entry.file(), trackingFrom(entry, entry.file())); + } + + private TrackedFile wrapFile(DataFile newFile, Tracking newTracking) { + if (newFile instanceof TrackedDataFile tracked) { + return tracked.file(); + } + + 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; + this.tracking = newTracking; + return this; + } + + @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<Long> splitOffsets() { + return file.splitOffsets(); + } + + @Override + public List<Integer> equalityIds() { + return null; + } + + @Override + public TrackedFile copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + + @Override + public TrackedFile copyWithStats(Set<Integer> 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 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. + */ + 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 = trackingFrom(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() { + return FileFormat.fromFileName(manifest.path()); + } + + @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() { + return null; + } + + @Override + public DeletionVector deletionVector() { + return null; + } + + @Override + public ManifestInfo manifestInfo() { + return manifestInfo; + } + + @Override + public ByteBuffer keyMetadata() { + return manifest.keyMetadata(); + } + + @Override + public List<Long> splitOffsets() { + return null; + } + + @Override + public List<Integer> equalityIds() { + return null; + } + + @Override + public TrackedFile copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + + @Override + public TrackedFile copyWithStats(Set<Integer> 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; + + void wrap(ManifestFile newManifest) { + this.manifest = newManifest; + } + + @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 Tracking trackingFrom(ManifestEntry<?> entry, ContentFile<?> file) { + return new TrackingStruct( + entryStatus(entry.status()), + entry.snapshotId(), + entry.dataSequenceNumber(), + entry.fileSequenceNumber(), + null, Review Comment: dvSnapshotId stays null. the comment is: `dvSnapshotId is null because this wrapper has no DV entry`. ########## core/src/main/java/org/apache/iceberg/MapBackedContentStats.java: ########## @@ -0,0 +1,206 @@ +/* + * 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.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 Types.StructType type; + private final Map<Integer, FieldStats<?>> statsById = Maps.newHashMap(); + + private Map<Integer, Long> valueCounts; + private Map<Integer, Long> nullValueCounts; + private Map<Integer, Long> nanValueCounts; + private Map<Integer, Integer> avgValueSizes; + private Map<Integer, ByteBuffer> lowerBounds; + private Map<Integer, ByteBuffer> upperBounds; + + MapBackedContentStats(Schema tableSchema, MetricsConfig metricsConfig) { + Preconditions.checkArgument(tableSchema != null, "Invalid table schema: null"); + Preconditions.checkArgument(metricsConfig != null, "Invalid metrics config: null"); + this.type = StatsUtil.statsWriteSchema(tableSchema, metricsConfig); + } + + 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(); + return this; + } + + @Override + public Iterable<FieldStats<?>> fieldStats() { + return Iterables.filter( + Iterables.transform(type.fields(), field -> statsFor(StatsUtil.toFieldId(field.fieldId()))), + Objects::nonNull); + } + + @Override + @SuppressWarnings("unchecked") + public <T> FieldStats<T> statsFor(int fieldId) { + if (!hasStats(fieldId)) { + return null; + } + + FieldStats<?> fieldStats = statsById.get(fieldId); + if (fieldStats == null) { + fieldStats = new MapBackedFieldStats<>(fieldId); + statsById.put(fieldId, fieldStats); + } + + return (FieldStats<T>) fieldStats; + } + + @Override + public Types.StructType type() { Review Comment: see the reply above. ########## core/src/main/java/org/apache/iceberg/TrackedFileAdapters.java: ########## @@ -529,6 +566,427 @@ public ManifestFile copy() { } } + /** Adapts a {@link DataFile} to {@link TrackedFile}. */ + static class DataTrackedFile implements TrackedFile { + private final MapBackedContentStats statsWrapper; + private Tracking tracking; + private DataFile file; + private ContentStats stats; + + DataTrackedFile(Schema tableSchema, MetricsConfig metricsConfig) { + this.statsWrapper = new MapBackedContentStats(tableSchema, metricsConfig); + } + + /** Re-points this adapter at a {@link DataFile} from the public API. Tracking is unset. */ + public TrackedFile wrap(DataFile newFile) { + return wrapFile(newFile, null); + } + + /** + * Re-points this adapter at a {@link ManifestEntry}. Converts the contained data file and the + * entry's tracking fields. + */ + public TrackedFile wrap(ManifestEntry<DataFile> entry) { + Preconditions.checkArgument(entry != null, "Invalid entry: null"); + return wrapFile(entry.file(), trackingFrom(entry, entry.file())); + } + + private TrackedFile wrapFile(DataFile newFile, Tracking newTracking) { Review Comment: renamed it to wrapInternal and narrowed it to storing the file and stats. wrap(DataFile) clears tracking, while wrap(ManifestEntry) rebinds entry tracking. ########## core/src/main/java/org/apache/iceberg/TrackedFileAdapters.java: ########## @@ -529,6 +566,427 @@ public ManifestFile copy() { } } + /** Adapts a {@link DataFile} to {@link TrackedFile}. */ + static class DataTrackedFile implements TrackedFile { + private final MapBackedContentStats statsWrapper; + private Tracking tracking; + private DataFile file; + private ContentStats stats; + + DataTrackedFile(Schema tableSchema, MetricsConfig metricsConfig) { + this.statsWrapper = new MapBackedContentStats(tableSchema, metricsConfig); + } + + /** Re-points this adapter at a {@link DataFile} from the public API. Tracking is unset. */ + public TrackedFile wrap(DataFile newFile) { + return wrapFile(newFile, null); + } + + /** + * Re-points this adapter at a {@link ManifestEntry}. Converts the contained data file and the + * entry's tracking fields. + */ + public TrackedFile wrap(ManifestEntry<DataFile> entry) { + Preconditions.checkArgument(entry != null, "Invalid entry: null"); + return wrapFile(entry.file(), trackingFrom(entry, entry.file())); + } + + private TrackedFile wrapFile(DataFile newFile, Tracking newTracking) { + if (newFile instanceof TrackedDataFile tracked) { + return tracked.file(); + } + + 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; + this.tracking = newTracking; + return this; + } + + @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<Long> splitOffsets() { + return file.splitOffsets(); + } + + @Override + public List<Integer> equalityIds() { + return null; + } + + @Override + public TrackedFile copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + + @Override + public TrackedFile copyWithStats(Set<Integer> 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 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. + */ + 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 = trackingFrom(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() { + return FileFormat.fromFileName(manifest.path()); + } + + @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() { + return null; + } + + @Override + public DeletionVector deletionVector() { + return null; + } + + @Override + public ManifestInfo manifestInfo() { + return manifestInfo; + } + + @Override + public ByteBuffer keyMetadata() { + return manifest.keyMetadata(); + } + + @Override + public List<Long> splitOffsets() { + return null; + } + + @Override + public List<Integer> equalityIds() { + return null; + } + + @Override + public TrackedFile copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + + @Override + public TrackedFile copyWithStats(Set<Integer> 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; + + void wrap(ManifestFile newManifest) { + this.manifest = newManifest; + } + + @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 Tracking trackingFrom(ManifestEntry<?> entry, ContentFile<?> file) { Review Comment: wrap(ManifestEntry) takes only the entry, typed ManifestEntry<DataFile>. manifestLocation() and manifestPos() come from DataFile.manifestLocation() and DataFile.pos(). ########## core/src/main/java/org/apache/iceberg/TrackedFileAdapters.java: ########## @@ -529,6 +566,427 @@ public ManifestFile copy() { } } + /** Adapts a {@link DataFile} to {@link TrackedFile}. */ + static class DataTrackedFile implements TrackedFile { + private final MapBackedContentStats statsWrapper; + private Tracking tracking; + private DataFile file; + private ContentStats stats; + + DataTrackedFile(Schema tableSchema, MetricsConfig metricsConfig) { + this.statsWrapper = new MapBackedContentStats(tableSchema, metricsConfig); + } + + /** Re-points this adapter at a {@link DataFile} from the public API. Tracking is unset. */ + public TrackedFile wrap(DataFile newFile) { + return wrapFile(newFile, null); + } + + /** + * Re-points this adapter at a {@link ManifestEntry}. Converts the contained data file and the + * entry's tracking fields. + */ + public TrackedFile wrap(ManifestEntry<DataFile> entry) { + Preconditions.checkArgument(entry != null, "Invalid entry: null"); + return wrapFile(entry.file(), trackingFrom(entry, entry.file())); + } + + private TrackedFile wrapFile(DataFile newFile, Tracking newTracking) { + if (newFile instanceof TrackedDataFile tracked) { + return tracked.file(); + } + + 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; + this.tracking = newTracking; + return this; + } + + @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<Long> splitOffsets() { + return file.splitOffsets(); + } + + @Override + public List<Integer> equalityIds() { + return null; + } + + @Override + public TrackedFile copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + + @Override + public TrackedFile copyWithStats(Set<Integer> 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 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. + */ + 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 = trackingFrom(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() { + return FileFormat.fromFileName(manifest.path()); + } + + @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() { + return null; + } + + @Override + public DeletionVector deletionVector() { + return null; + } + + @Override + public ManifestInfo manifestInfo() { + return manifestInfo; + } + + @Override + public ByteBuffer keyMetadata() { + return manifest.keyMetadata(); + } + + @Override + public List<Long> splitOffsets() { + return null; + } + + @Override + public List<Integer> equalityIds() { + return null; + } + + @Override + public TrackedFile copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + + @Override + public TrackedFile copyWithStats(Set<Integer> 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; + + void wrap(ManifestFile newManifest) { + this.manifest = newManifest; + } + + @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 Tracking trackingFrom(ManifestEntry<?> entry, ContentFile<?> file) { + return new TrackingStruct( Review Comment: Makes sense to avoid unnecessary allocation even though TrackingStruct is relatively compact. Implemented the wrapping. a Tracking reference is only valid until the next wrap. ########## core/src/main/java/org/apache/iceberg/TrackedFileAdapters.java: ########## @@ -529,6 +566,427 @@ public ManifestFile copy() { } } + /** Adapts a {@link DataFile} to {@link TrackedFile}. */ + static class DataTrackedFile implements TrackedFile { + private final MapBackedContentStats statsWrapper; + private Tracking tracking; + private DataFile file; + private ContentStats stats; + + DataTrackedFile(Schema tableSchema, MetricsConfig metricsConfig) { + this.statsWrapper = new MapBackedContentStats(tableSchema, metricsConfig); + } + + /** Re-points this adapter at a {@link DataFile} from the public API. Tracking is unset. */ + public TrackedFile wrap(DataFile newFile) { + return wrapFile(newFile, null); + } + + /** + * Re-points this adapter at a {@link ManifestEntry}. Converts the contained data file and the + * entry's tracking fields. + */ + public TrackedFile wrap(ManifestEntry<DataFile> entry) { + Preconditions.checkArgument(entry != null, "Invalid entry: null"); + return wrapFile(entry.file(), trackingFrom(entry, entry.file())); + } + + private TrackedFile wrapFile(DataFile newFile, Tracking newTracking) { + if (newFile instanceof TrackedDataFile tracked) { + return tracked.file(); + } + + 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; + this.tracking = newTracking; + return this; + } + + @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<Long> splitOffsets() { + return file.splitOffsets(); + } + + @Override + public List<Integer> equalityIds() { + return null; + } + + @Override + public TrackedFile copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + + @Override + public TrackedFile copyWithStats(Set<Integer> 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 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. + */ + 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 = trackingFrom(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() { + return FileFormat.fromFileName(manifest.path()); + } + + @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() { + return null; + } + + @Override + public DeletionVector deletionVector() { + return null; + } + + @Override + public ManifestInfo manifestInfo() { + return manifestInfo; + } + + @Override + public ByteBuffer keyMetadata() { + return manifest.keyMetadata(); + } + + @Override + public List<Long> splitOffsets() { + return null; + } + + @Override + public List<Integer> equalityIds() { + return null; + } + + @Override + public TrackedFile copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + + @Override + public TrackedFile copyWithStats(Set<Integer> 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; + + void wrap(ManifestFile newManifest) { + this.manifest = newManifest; + } + + @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 Tracking trackingFrom(ManifestEntry<?> entry, ContentFile<?> file) { + return new TrackingStruct( + entryStatus(entry.status()), + entry.snapshotId(), + entry.dataSequenceNumber(), + entry.fileSequenceNumber(), + null, + file.firstRowId(), + null, + null); + } + + private static Tracking trackingFrom(ManifestFile manifest) { + return new TrackingStruct( + null, + manifest.snapshotId(), + manifest.sequenceNumber(), + manifest.sequenceNumber(), + null, + manifest.firstRowId(), + null, + null); + } + + 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.checkNotNull( + manifest.replacedFilesCount(), + "Cannot convert manifest %s: missing replaced files count", + manifest.path()); + Preconditions.checkNotNull( + manifest.modifiedFilesCount(), + "Cannot convert manifest %s: missing modified files count", + manifest.path()); + Preconditions.checkArgument( + manifest.replacedFilesCount() == 0, + "Cannot convert manifest %s: replaced files count must be 0 for v3 or earlier manifests, was %s", Review Comment: done ########## core/src/main/java/org/apache/iceberg/MapBackedContentStats.java: ########## @@ -0,0 +1,206 @@ +/* + * 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.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 Types.StructType type; + private final Map<Integer, FieldStats<?>> statsById = Maps.newHashMap(); + + private Map<Integer, Long> valueCounts; + private Map<Integer, Long> nullValueCounts; + private Map<Integer, Long> nanValueCounts; + private Map<Integer, Integer> avgValueSizes; + private Map<Integer, ByteBuffer> lowerBounds; + private Map<Integer, ByteBuffer> upperBounds; + + MapBackedContentStats(Schema tableSchema, MetricsConfig metricsConfig) { + Preconditions.checkArgument(tableSchema != null, "Invalid table schema: null"); + Preconditions.checkArgument(metricsConfig != null, "Invalid metrics config: null"); + this.type = StatsUtil.statsWriteSchema(tableSchema, metricsConfig); + } + + 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(); + return this; + } + + @Override + public Iterable<FieldStats<?>> fieldStats() { + return Iterables.filter( + Iterables.transform(type.fields(), field -> statsFor(StatsUtil.toFieldId(field.fieldId()))), + Objects::nonNull); + } + + @Override + @SuppressWarnings("unchecked") + public <T> FieldStats<T> statsFor(int fieldId) { + if (!hasStats(fieldId)) { + return null; + } + + FieldStats<?> fieldStats = statsById.get(fieldId); + if (fieldStats == null) { + fieldStats = new MapBackedFieldStats<>(fieldId); + statsById.put(fieldId, fieldStats); + } + + return (FieldStats<T>) fieldStats; + } + + @Override + public Types.StructType type() { + return type; + } + + @Override + public ContentStats copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + + @Override + public ContentStats copy(Set<Integer> fieldIds) { + throw new UnsupportedOperationException("copy is not implemented"); + } + + private boolean hasStats(int id) { + return containsId(valueCounts, id) + || containsId(nullValueCounts, id) + || containsId(nanValueCounts, id) + || containsId(avgValueSizes, id) + || containsId(lowerBounds, id) + || containsId(upperBounds, id); + } + + private static boolean containsId(Map<Integer, ?> 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<T> implements FieldStats<T> { + private final int fieldId; + private final Types.StructType struct; + private final Type boundType; + + MapBackedFieldStats(int fieldId) { + Types.NestedField field = type.field(StatsUtil.toBaseId(fieldId)); Review Comment: with option 2, full content stats schema is only constructed once for the wrapper object (not per wrap action). ########## core/src/main/java/org/apache/iceberg/TrackedFileAdapters.java: ########## @@ -529,6 +566,427 @@ public ManifestFile copy() { } } + /** Adapts a {@link DataFile} to {@link TrackedFile}. */ + static class DataTrackedFile implements TrackedFile { + private final MapBackedContentStats statsWrapper; + private Tracking tracking; + private DataFile file; + private ContentStats stats; + + DataTrackedFile(Schema tableSchema, MetricsConfig metricsConfig) { + this.statsWrapper = new MapBackedContentStats(tableSchema, metricsConfig); + } + + /** Re-points this adapter at a {@link DataFile} from the public API. Tracking is unset. */ + public TrackedFile wrap(DataFile newFile) { + return wrapFile(newFile, null); + } + + /** + * Re-points this adapter at a {@link ManifestEntry}. Converts the contained data file and the + * entry's tracking fields. + */ + public TrackedFile wrap(ManifestEntry<DataFile> entry) { + Preconditions.checkArgument(entry != null, "Invalid entry: null"); + return wrapFile(entry.file(), trackingFrom(entry, entry.file())); + } + + private TrackedFile wrapFile(DataFile newFile, Tracking newTracking) { + if (newFile instanceof TrackedDataFile tracked) { + return tracked.file(); + } + + 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; + this.tracking = newTracking; + return this; + } + + @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<Long> splitOffsets() { + return file.splitOffsets(); + } + + @Override + public List<Integer> equalityIds() { + return null; + } + + @Override + public TrackedFile copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + + @Override + public TrackedFile copyWithStats(Set<Integer> 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 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. Review Comment: updated Javadoc ########## core/src/main/java/org/apache/iceberg/TrackedFileAdapters.java: ########## @@ -529,6 +566,427 @@ public ManifestFile copy() { } } + /** Adapts a {@link DataFile} to {@link TrackedFile}. */ + static class DataTrackedFile implements TrackedFile { + private final MapBackedContentStats statsWrapper; + private Tracking tracking; + private DataFile file; + private ContentStats stats; + + DataTrackedFile(Schema tableSchema, MetricsConfig metricsConfig) { + this.statsWrapper = new MapBackedContentStats(tableSchema, metricsConfig); + } + + /** Re-points this adapter at a {@link DataFile} from the public API. Tracking is unset. */ + public TrackedFile wrap(DataFile newFile) { + return wrapFile(newFile, null); + } + + /** + * Re-points this adapter at a {@link ManifestEntry}. Converts the contained data file and the + * entry's tracking fields. + */ + public TrackedFile wrap(ManifestEntry<DataFile> entry) { + Preconditions.checkArgument(entry != null, "Invalid entry: null"); + return wrapFile(entry.file(), trackingFrom(entry, entry.file())); + } + + private TrackedFile wrapFile(DataFile newFile, Tracking newTracking) { + if (newFile instanceof TrackedDataFile tracked) { + return tracked.file(); + } + + 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; + this.tracking = newTracking; + return this; + } + + @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<Long> splitOffsets() { + return file.splitOffsets(); + } + + @Override + public List<Integer> equalityIds() { + return null; + } + + @Override + public TrackedFile copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + + @Override + public TrackedFile copyWithStats(Set<Integer> 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 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. + */ + 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 = trackingFrom(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() { + return FileFormat.fromFileName(manifest.path()); + } + + @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() { + return null; + } + + @Override + public DeletionVector deletionVector() { + return null; + } + + @Override + public ManifestInfo manifestInfo() { + return manifestInfo; + } + + @Override + public ByteBuffer keyMetadata() { + return manifest.keyMetadata(); + } + + @Override + public List<Long> splitOffsets() { + return null; + } + + @Override + public List<Integer> equalityIds() { + return null; + } + + @Override + public TrackedFile copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + + @Override + public TrackedFile copyWithStats(Set<Integer> 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; Review Comment: done ########## core/src/main/java/org/apache/iceberg/TrackedFileAdapters.java: ########## @@ -529,6 +566,427 @@ public ManifestFile copy() { } } + /** Adapts a {@link DataFile} to {@link TrackedFile}. */ + static class DataTrackedFile implements TrackedFile { + private final MapBackedContentStats statsWrapper; + private Tracking tracking; + private DataFile file; + private ContentStats stats; + + DataTrackedFile(Schema tableSchema, MetricsConfig metricsConfig) { Review Comment: MetricsConfig is needed if we stay with the current option 2. ########## core/src/main/java/org/apache/iceberg/MapBackedContentStats.java: ########## @@ -0,0 +1,206 @@ +/* + * 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.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 Types.StructType type; + private final Map<Integer, FieldStats<?>> statsById = Maps.newHashMap(); + + private Map<Integer, Long> valueCounts; + private Map<Integer, Long> nullValueCounts; + private Map<Integer, Long> nanValueCounts; + private Map<Integer, Integer> avgValueSizes; + private Map<Integer, ByteBuffer> lowerBounds; + private Map<Integer, ByteBuffer> upperBounds; + + MapBackedContentStats(Schema tableSchema, MetricsConfig metricsConfig) { + Preconditions.checkArgument(tableSchema != null, "Invalid table schema: null"); + Preconditions.checkArgument(metricsConfig != null, "Invalid metrics config: null"); + this.type = StatsUtil.statsWriteSchema(tableSchema, metricsConfig); + } + + 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(); + return this; + } + + @Override + public Iterable<FieldStats<?>> fieldStats() { + return Iterables.filter( + Iterables.transform(type.fields(), field -> statsFor(StatsUtil.toFieldId(field.fieldId()))), + Objects::nonNull); + } + + @Override + @SuppressWarnings("unchecked") + public <T> FieldStats<T> statsFor(int fieldId) { + if (!hasStats(fieldId)) { + return null; + } + + FieldStats<?> fieldStats = statsById.get(fieldId); + if (fieldStats == null) { + fieldStats = new MapBackedFieldStats<>(fieldId); + statsById.put(fieldId, fieldStats); + } + + return (FieldStats<T>) fieldStats; + } + + @Override + public Types.StructType type() { + return type; + } + + @Override + public ContentStats copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + + @Override + public ContentStats copy(Set<Integer> fieldIds) { + throw new UnsupportedOperationException("copy is not implemented"); + } + + private boolean hasStats(int id) { + return containsId(valueCounts, id) + || containsId(nullValueCounts, id) + || containsId(nanValueCounts, id) + || containsId(avgValueSizes, id) + || containsId(lowerBounds, id) + || containsId(upperBounds, id); + } + + private static boolean containsId(Map<Integer, ?> 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<T> implements FieldStats<T> { + private final int fieldId; + private final Types.StructType struct; + private final Type boundType; + + MapBackedFieldStats(int fieldId) { + Types.NestedField field = type.field(StatsUtil.toBaseId(fieldId)); + Preconditions.checkArgument( + field != null, + "Cannot convert stats for field ID %s: unknown, not a scalar, or not in metrics config", + fieldId); + this.fieldId = fieldId; + this.struct = field.type().asStructType(); + 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); Review Comment: GeospatialBound implements StructLike in x, y, z, m order. set throws because the fields are final. the write path reads those positions through StructLike.get and re-encodes with toByteBuffer. ########## core/src/main/java/org/apache/iceberg/TrackedFileAdapters.java: ########## @@ -529,6 +566,427 @@ public ManifestFile copy() { } } + /** Adapts a {@link DataFile} to {@link TrackedFile}. */ + static class DataTrackedFile implements TrackedFile { + private final MapBackedContentStats statsWrapper; + private Tracking tracking; + private DataFile file; + private ContentStats stats; + + DataTrackedFile(Schema tableSchema, MetricsConfig metricsConfig) { + this.statsWrapper = new MapBackedContentStats(tableSchema, metricsConfig); + } + + /** Re-points this adapter at a {@link DataFile} from the public API. Tracking is unset. */ + public TrackedFile wrap(DataFile newFile) { + return wrapFile(newFile, null); + } + + /** + * Re-points this adapter at a {@link ManifestEntry}. Converts the contained data file and the + * entry's tracking fields. + */ + public TrackedFile wrap(ManifestEntry<DataFile> entry) { + Preconditions.checkArgument(entry != null, "Invalid entry: null"); + return wrapFile(entry.file(), trackingFrom(entry, entry.file())); + } + + private TrackedFile wrapFile(DataFile newFile, Tracking newTracking) { + if (newFile instanceof TrackedDataFile tracked) { + return tracked.file(); + } + + 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; + this.tracking = newTracking; + return this; + } + + @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<Long> splitOffsets() { + return file.splitOffsets(); + } + + @Override + public List<Integer> equalityIds() { + return null; + } + + @Override + public TrackedFile copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + + @Override + public TrackedFile copyWithStats(Set<Integer> 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 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. + */ + 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 = trackingFrom(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() { + return FileFormat.fromFileName(manifest.path()); Review Comment: fileFormat() returns FileFormat.AVRO. the comment is: manifest files before v4 are always Avro. ########## core/src/main/java/org/apache/iceberg/TrackedFileAdapters.java: ########## @@ -529,6 +566,427 @@ public ManifestFile copy() { } } + /** Adapts a {@link DataFile} to {@link TrackedFile}. */ + static class DataTrackedFile implements TrackedFile { + private final MapBackedContentStats statsWrapper; + private Tracking tracking; + private DataFile file; + private ContentStats stats; + + DataTrackedFile(Schema tableSchema, MetricsConfig metricsConfig) { + this.statsWrapper = new MapBackedContentStats(tableSchema, metricsConfig); + } + + /** Re-points this adapter at a {@link DataFile} from the public API. Tracking is unset. */ + public TrackedFile wrap(DataFile newFile) { + return wrapFile(newFile, null); + } + + /** + * Re-points this adapter at a {@link ManifestEntry}. Converts the contained data file and the + * entry's tracking fields. Review Comment: updated Javadoc ########## core/src/main/java/org/apache/iceberg/TrackedFileAdapters.java: ########## @@ -529,6 +566,427 @@ public ManifestFile copy() { } } + /** Adapts a {@link DataFile} to {@link TrackedFile}. */ + static class DataTrackedFile implements TrackedFile { + private final MapBackedContentStats statsWrapper; + private Tracking tracking; + private DataFile file; + private ContentStats stats; + + DataTrackedFile(Schema tableSchema, MetricsConfig metricsConfig) { + this.statsWrapper = new MapBackedContentStats(tableSchema, metricsConfig); + } + + /** Re-points this adapter at a {@link DataFile} from the public API. Tracking is unset. */ + public TrackedFile wrap(DataFile newFile) { + return wrapFile(newFile, null); + } + + /** + * Re-points this adapter at a {@link ManifestEntry}. Converts the contained data file and the + * entry's tracking fields. + */ + public TrackedFile wrap(ManifestEntry<DataFile> entry) { + Preconditions.checkArgument(entry != null, "Invalid entry: null"); + return wrapFile(entry.file(), trackingFrom(entry, entry.file())); + } + + private TrackedFile wrapFile(DataFile newFile, Tracking newTracking) { + if (newFile instanceof TrackedDataFile tracked) { + return tracked.file(); + } + + 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; + this.tracking = newTracking; + return this; + } + + @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<Long> splitOffsets() { + return file.splitOffsets(); + } + + @Override + public List<Integer> equalityIds() { + return null; + } + + @Override + public TrackedFile copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + + @Override + public TrackedFile copyWithStats(Set<Integer> 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 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. + */ + 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 = trackingFrom(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() { + return FileFormat.fromFileName(manifest.path()); + } + + @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() { + return null; Review Comment: added inline comment: `manifests have no table sort order` -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
