rdblue commented on code in PR #16936:
URL: https://github.com/apache/iceberg/pull/16936#discussion_r4200083191
##########
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:
Minor: This error message is long and not very direct. I'd prefer a message
like `Non-zero replaced file count: %s`
--
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]