diff --git a/core/src/main/java/org/apache/iceberg/ColumnFile.java b/core/src/main/java/org/apache/iceberg/ColumnFile.java new file mode 100644 index 000000000000..dde090666a91 --- /dev/null +++ b/core/src/main/java/org/apache/iceberg/ColumnFile.java @@ -0,0 +1,72 @@ +/* + * 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.List; +import org.apache.iceberg.types.Types; + +interface ColumnFile { + Types.NestedField LOCATION = + Types.NestedField.required( + 160, "location", Types.StringType.get(), "Location of the column file"); + Types.NestedField FIELD_IDS = + Types.NestedField.required( + 161, + "field_ids", + Types.ListType.ofRequired(162, Types.IntegerType.get()), + "Live field IDs in this column file"); + Types.NestedField FILE_FORMAT = + Types.NestedField.required( + 163, + "file_format", + Types.StringType.get(), + "String file format name for this column file"); + Types.NestedField FILE_SIZE_IN_BYTES = + Types.NestedField.required( + 164, "file_size_in_bytes", Types.LongType.get(), "Total column file size in bytes"); + Types.NestedField KEY_METADATA = + Types.NestedField.optional( + 165, + "key_metadata", + Types.BinaryType.get(), + "Key metadata for encryption; specific to the encryption scheme."); + + static Types.StructType schema() { + return Types.StructType.of(LOCATION, FIELD_IDS, FILE_FORMAT, FILE_SIZE_IN_BYTES, KEY_METADATA); + } + + /** Returns the location of this column file. */ + String location(); + + /** Returns the field IDs contained in this column file. */ + List fieldIds(); + + /** Returns the format of this column file. */ + FileFormat fileFormat(); + + /** Returns the total size of this column file in bytes. */ + long fileSizeInBytes(); + + /** Returns encryption key metadata, or null if this column file is not encrypted. */ + ByteBuffer keyMetadata(); + + /** Copies this column file. */ + ColumnFile copy(); +} diff --git a/core/src/main/java/org/apache/iceberg/ColumnFileStruct.java b/core/src/main/java/org/apache/iceberg/ColumnFileStruct.java new file mode 100644 index 000000000000..a51f0c0cc320 --- /dev/null +++ b/core/src/main/java/org/apache/iceberg/ColumnFileStruct.java @@ -0,0 +1,218 @@ +/* + * 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.io.Serializable; +import java.nio.ByteBuffer; +import java.util.Arrays; +import java.util.List; +import org.apache.iceberg.avro.SupportsIndexProjection; +import org.apache.iceberg.relocated.com.google.common.base.MoreObjects; +import org.apache.iceberg.relocated.com.google.common.base.Preconditions; +import org.apache.iceberg.relocated.com.google.common.collect.Sets; +import org.apache.iceberg.types.Types; +import org.apache.iceberg.util.ArrayUtil; +import org.apache.iceberg.util.ByteBuffers; + +/** Mutable {@link StructLike} implementation of {@link ColumnFile}. */ +class ColumnFileStruct extends SupportsIndexProjection implements ColumnFile, Serializable { + private static final Types.StructType BASE_TYPE = + Types.StructType.of( + ColumnFile.LOCATION, + ColumnFile.FIELD_IDS, + ColumnFile.FILE_FORMAT, + ColumnFile.FILE_SIZE_IN_BYTES, + ColumnFile.KEY_METADATA); + + private String location = null; + private int[] fieldIds = null; + private FileFormat fileFormat = null; + private long fileSizeInBytes = -1L; + private byte[] keyMetadata = null; + + /** Used by internal readers to instantiate this class with a projection schema. */ + ColumnFileStruct(Types.StructType projection) { + super(BASE_TYPE, projection); + } + + ColumnFileStruct( + String location, + List fieldIds, + FileFormat fileFormat, + long fileSizeInBytes, + ByteBuffer keyMetadata) { + super(BASE_TYPE.fields().size()); + this.location = location; + this.fieldIds = ArrayUtil.toIntArray(fieldIds); + this.fileFormat = fileFormat; + this.fileSizeInBytes = fileSizeInBytes; + this.keyMetadata = ByteBuffers.toByteArray(keyMetadata); + } + + /** Copy constructor. */ + private ColumnFileStruct(ColumnFileStruct toCopy) { + super(toCopy); + this.location = toCopy.location; + this.fieldIds = + toCopy.fieldIds != null ? Arrays.copyOf(toCopy.fieldIds, toCopy.fieldIds.length) : null; + this.fileFormat = toCopy.fileFormat; + this.fileSizeInBytes = toCopy.fileSizeInBytes; + this.keyMetadata = + toCopy.keyMetadata != null + ? Arrays.copyOf(toCopy.keyMetadata, toCopy.keyMetadata.length) + : null; + } + + /** Constructor for Java serialization. */ + ColumnFileStruct() { + super(BASE_TYPE.fields().size()); + } + + @Override + public String location() { + return location; + } + + @Override + public List fieldIds() { + return fieldIds != null ? ArrayUtil.toUnmodifiableIntList(fieldIds) : null; + } + + @Override + public FileFormat fileFormat() { + return fileFormat; + } + + @Override + public long fileSizeInBytes() { + return fileSizeInBytes; + } + + @Override + public ByteBuffer keyMetadata() { + return keyMetadata != null ? ByteBuffer.wrap(keyMetadata) : null; + } + + @Override + public ColumnFile copy() { + return new ColumnFileStruct(this); + } + + @Override + protected T internalGet(int pos, Class javaClass) { + return javaClass.cast(getByPos(pos)); + } + + private Object getByPos(int pos) { + return switch (pos) { + case 0 -> location; + case 1 -> fieldIds(); + case 2 -> fileFormat != null ? fileFormat.toString() : null; + case 3 -> fileSizeInBytes; + case 4 -> keyMetadata(); + default -> throw new UnsupportedOperationException("Unknown field ordinal: " + pos); + }; + } + + @Override + @SuppressWarnings("unchecked") + protected void internalSet(int pos, T value) { + switch (pos) { + // always coerce to String for Serializable + case 0 -> this.location = value.toString(); + case 1 -> this.fieldIds = ArrayUtil.toIntArray((List) value); + case 2 -> this.fileFormat = FileFormat.fromString(value.toString()); + case 3 -> this.fileSizeInBytes = (long) value; + case 4 -> this.keyMetadata = ByteBuffers.toByteArray((ByteBuffer) value); + default -> { + // ignore the object, it must be from a newer version of the format + } + } + } + + static Builder builder() { + return new Builder(); + } + + @Override + public String toString() { + return MoreObjects.toStringHelper(this) + .add("location", location) + .add("field_ids", fieldIds) + .add("file_format", fileFormat) + .add("file_size_in_bytes", fileSizeInBytes) + .add("key_metadata", keyMetadata == null ? "null" : "(redacted)") + .toString(); + } + + static class Builder { + private String location = null; + private List fieldIds = null; + private FileFormat fileFormat = null; + private Long fileSizeInBytes = null; + private ByteBuffer keyMetadata = null; + + Builder location(String newLocation) { + Preconditions.checkArgument(newLocation != null, "Invalid location: null"); + Preconditions.checkArgument(!newLocation.isEmpty(), "Invalid location: empty"); + this.location = newLocation; + return this; + } + + Builder fieldIds(List newFieldIds) { + Preconditions.checkArgument(newFieldIds != null, "Invalid field IDs: null"); + Preconditions.checkArgument(!newFieldIds.isEmpty(), "Invalid field IDs: empty"); + Preconditions.checkArgument( + Sets.newHashSet(newFieldIds).size() == newFieldIds.size(), + "Invalid field IDs: duplicated IDs found in: %s", + newFieldIds); + this.fieldIds = newFieldIds; + return this; + } + + Builder fileFormat(FileFormat newFileFormat) { + Preconditions.checkArgument(newFileFormat != null, "Invalid file format: null"); + this.fileFormat = newFileFormat; + return this; + } + + Builder fileSizeInBytes(long newFileSizeInBytes) { + Preconditions.checkArgument( + newFileSizeInBytes >= 0, + "Invalid file size in bytes: %s (must be >= 0)", + newFileSizeInBytes); + this.fileSizeInBytes = newFileSizeInBytes; + return this; + } + + Builder keyMetadata(ByteBuffer newKeyMetadata) { + this.keyMetadata = newKeyMetadata; + return this; + } + + ColumnFile build() { + Preconditions.checkArgument(location != null, "Missing required value: location"); + Preconditions.checkArgument(fieldIds != null, "Missing required value: field IDs"); + Preconditions.checkArgument(fileFormat != null, "Missing required value: file format"); + Preconditions.checkArgument( + fileSizeInBytes != null, "Missing required value: file size in bytes"); + return new ColumnFileStruct(location, fieldIds, fileFormat, fileSizeInBytes, keyMetadata); + } + } +} diff --git a/core/src/main/java/org/apache/iceberg/TrackedFile.java b/core/src/main/java/org/apache/iceberg/TrackedFile.java index e2db837c4804..e77f8c55cdab 100644 --- a/core/src/main/java/org/apache/iceberg/TrackedFile.java +++ b/core/src/main/java/org/apache/iceberg/TrackedFile.java @@ -96,6 +96,9 @@ interface TrackedFile { "equality_ids", Types.ListType.ofRequired(136, Types.IntegerType.get()), "Field ids used to determine row equality in equality delete files"); + Types.NestedField COLUMN_FILES = + Types.NestedField.optional( + 158, "column_files", Types.ListType.ofRequired(159, ColumnFile.schema()), "Column files"); private static List fields( Types.StructType partitionType, Types.StructType contentStatsType) { @@ -120,7 +123,8 @@ PARTITION_ID, PARTITION_NAME, typeOrUnknown(partitionType), PARTITION_DOC), MANIFEST_INFO, KEY_METADATA, SPLIT_OFFSETS, - EQUALITY_IDS); + EQUALITY_IDS, + COLUMN_FILES); } private static Type typeOrUnknown(Types.StructType structType) { @@ -195,6 +199,9 @@ static Schema readSchema(Types.StructType partitionType, Types.StructType conten /** Returns the set of field IDs used for equality comparison in equality delete files. */ List equalityIds(); + /** Returns the column files for this file. */ + List columnFiles(); + /** Copies this tracked file. */ TrackedFile copy(); diff --git a/core/src/main/java/org/apache/iceberg/TrackedFileStruct.java b/core/src/main/java/org/apache/iceberg/TrackedFileStruct.java index fffd822a56a2..273f31147381 100644 --- a/core/src/main/java/org/apache/iceberg/TrackedFileStruct.java +++ b/core/src/main/java/org/apache/iceberg/TrackedFileStruct.java @@ -21,11 +21,13 @@ import java.io.Serializable; import java.nio.ByteBuffer; import java.util.Arrays; +import java.util.Collections; import java.util.List; import java.util.Set; import org.apache.iceberg.avro.SupportsIndexProjection; import org.apache.iceberg.exceptions.ValidationException; import org.apache.iceberg.relocated.com.google.common.base.MoreObjects; +import org.apache.iceberg.relocated.com.google.common.collect.Lists; import org.apache.iceberg.types.Types; import org.apache.iceberg.util.ArrayUtil; import org.apache.iceberg.util.ByteBuffers; @@ -61,7 +63,8 @@ class TrackedFileStruct extends SupportsIndexProjection implements TrackedFile, TrackedFile.MANIFEST_INFO, TrackedFile.KEY_METADATA, TrackedFile.SPLIT_OFFSETS, - TrackedFile.EQUALITY_IDS); + TrackedFile.EQUALITY_IDS, + TrackedFile.COLUMN_FILES); private FileContent contentType = null; private int formatVersion = -1; @@ -81,6 +84,7 @@ class TrackedFileStruct extends SupportsIndexProjection implements TrackedFile, private byte[] keyMetadata = null; private long[] splitOffsets = null; private int[] equalityIds = null; + private List columnFiles = null; private transient StructProjection partitionProjection = null; @@ -110,7 +114,8 @@ class TrackedFileStruct extends SupportsIndexProjection implements TrackedFile, ManifestInfo manifestInfo, ByteBuffer keyMetadata, List splitOffsets, - List equalityIds) { + List equalityIds, + List columnFiles) { super(BASE_TYPE.fields().size()); this.tracking = tracking; this.contentType = contentType; @@ -128,9 +133,11 @@ class TrackedFileStruct extends SupportsIndexProjection implements TrackedFile, this.keyMetadata = ByteBuffers.toByteArray(keyMetadata); this.splitOffsets = ArrayUtil.toLongArray(splitOffsets); this.equalityIds = ArrayUtil.toIntArray(equalityIds); + this.columnFiles = columnFiles != null ? Lists.newArrayList(columnFiles) : null; } /** Copy constructor. */ + @SuppressWarnings("CyclomaticComplexity") private TrackedFileStruct(TrackedFileStruct toCopy, Set statsIds) { super(toCopy); this.contentType = toCopy.contentType; @@ -165,6 +172,15 @@ private TrackedFileStruct(TrackedFileStruct toCopy, Set statsIds) { toCopy.equalityIds != null ? Arrays.copyOf(toCopy.equalityIds, toCopy.equalityIds.length) : null; + + if (toCopy.columnFiles != null) { + this.columnFiles = Lists.newArrayListWithCapacity(toCopy.columnFiles.size()); + for (ColumnFile columnFile : toCopy.columnFiles) { + this.columnFiles.add(columnFile != null ? columnFile.copy() : null); + } + } else { + this.columnFiles = null; + } } @Override @@ -267,6 +283,11 @@ public List equalityIds() { return equalityIds != null ? ArrayUtil.toUnmodifiableIntList(equalityIds) : null; } + @Override + public List columnFiles() { + return columnFiles != null ? Collections.unmodifiableList(columnFiles) : null; + } + @Override public TrackedFile copy() { return new TrackedFileStruct(this, null); @@ -300,6 +321,7 @@ private Object getByPos(int pos) { case 13 -> keyMetadata(); case 14 -> splitOffsets(); case 15 -> equalityIds(); + case 16 -> columnFiles(); default -> throw new UnsupportedOperationException("Unknown field ordinal: " + pos); }; } @@ -325,6 +347,7 @@ protected void internalSet(int pos, T value) { case 13 -> this.keyMetadata = ByteBuffers.toByteArray((ByteBuffer) value); case 14 -> this.splitOffsets = ArrayUtil.toLongArray((List) value); case 15 -> this.equalityIds = ArrayUtil.toIntArray((List) value); + case 16 -> this.columnFiles = (List) value; default -> { // ignore the object, it must be from a newer version of the format } @@ -350,6 +373,7 @@ public String toString() { .add("key_metadata", keyMetadata == null ? "null" : "(redacted)") .add("split_offsets", splitOffsets == null ? "null" : splitOffsets()) .add("equality_ids", equalityIds == null ? "null" : equalityIds()) + .add("column_files", columnFiles) .toString(); } } diff --git a/core/src/main/java/org/apache/iceberg/Tracking.java b/core/src/main/java/org/apache/iceberg/Tracking.java index fcdc4e50b236..48ff5da5f876 100644 --- a/core/src/main/java/org/apache/iceberg/Tracking.java +++ b/core/src/main/java/org/apache/iceberg/Tracking.java @@ -44,12 +44,12 @@ interface Tracking { "file_sequence_number", Types.LongType.get(), "File sequence number indicating when the file was added"); - Types.NestedField DV_SNAPSHOT_ID = + Types.NestedField MODIFIED_SNAPSHOT_ID = Types.NestedField.optional( 5, - "dv_snapshot_id", + "modified_snapshot_id", Types.LongType.get(), - "Snapshot ID where the DV was added; null if there is no DV"); + "Snapshot ID where the file was last modified."); Types.NestedField FIRST_ROW_ID = Types.NestedField.optional( 142, "first_row_id", Types.LongType.get(), "ID of the first row in the data file"); @@ -72,7 +72,7 @@ static Types.StructType schema() { SNAPSHOT_ID, SEQUENCE_NUMBER, FILE_SEQUENCE_NUMBER, - DV_SNAPSHOT_ID, + MODIFIED_SNAPSHOT_ID, FIRST_ROW_ID, DELETED_POSITIONS, REPLACED_POSITIONS); @@ -95,8 +95,8 @@ default boolean isLive() { /** Returns the file sequence number indicating when the file was added. */ Long fileSequenceNumber(); - /** Returns the snapshot ID where the DV was added; null if there is no DV. */ - Long dvSnapshotId(); + /** Returns the snapshot ID where the file was last modified. */ + Long modifiedSnapshotId(); /** Returns the ID of the first row in the data file. */ Long firstRowId(); diff --git a/core/src/main/java/org/apache/iceberg/TrackingBuilder.java b/core/src/main/java/org/apache/iceberg/TrackingBuilder.java index 3c733825e262..64e7cae043b6 100644 --- a/core/src/main/java/org/apache/iceberg/TrackingBuilder.java +++ b/core/src/main/java/org/apache/iceberg/TrackingBuilder.java @@ -25,11 +25,11 @@ class TrackingBuilder { private final long newSnapshotId; private final Long snapshotId; - private final Long dataSequenceNumber; private final Long fileSequenceNumber; private final Long firstRowId; private EntryStatus status; - private Long dvSnapshotId; + private Long dataSequenceNumber; + private Long modifiedSnapshotId; private byte[] deletedPositions; private byte[] replacedPositions; @@ -83,7 +83,7 @@ private TrackingBuilder(long newSnapshotId) { this.dataSequenceNumber = null; this.fileSequenceNumber = null; this.firstRowId = null; - this.dvSnapshotId = null; + this.modifiedSnapshotId = null; this.deletedPositions = null; this.replacedPositions = null; } @@ -97,7 +97,7 @@ private TrackingBuilder(Tracking source, long newSnapshotId) { this.dataSequenceNumber = source.dataSequenceNumber(); this.fileSequenceNumber = source.fileSequenceNumber(); this.firstRowId = source.firstRowId(); - this.dvSnapshotId = source.dvSnapshotId(); + this.modifiedSnapshotId = source.modifiedSnapshotId(); this.deletedPositions = null; this.replacedPositions = null; } @@ -107,7 +107,7 @@ TrackingBuilder dvUpdated() { Preconditions.checkState( deletedPositions == null && replacedPositions == null, "Cannot mark DV updated on a manifest entry (deleted/replaced positions are set)"); - this.dvSnapshotId = newSnapshotId; + this.modifiedSnapshotId = newSnapshotId; if (status == EntryStatus.EXISTING) { this.status = EntryStatus.MODIFIED; } @@ -115,12 +115,25 @@ TrackingBuilder dvUpdated() { return this; } + /** Indicates that the column files list has been updated for the new Tracking. */ + TrackingBuilder columnFilesUpdated() { + this.modifiedSnapshotId = newSnapshotId; + if (status == EntryStatus.EXISTING) { + this.status = EntryStatus.MODIFIED; + } + // Reset to null to inherit from the new snapshot sequence number. It is safe to bump up the + // dataSequenceNumber as writers are required to rewrite v2 equality and position deletes to DVs + // when applying column update. + this.dataSequenceNumber = null; + return this; + } + /** Sets the positions deleted by this commit for a manifest entry. */ TrackingBuilder deletedPositions(ByteBuffer positions) { Preconditions.checkState( status != EntryStatus.ADDED, "Cannot set deleted positions on ADDED entry"); this.deletedPositions = ByteBuffers.toByteArray(positions); - this.dvSnapshotId = newSnapshotId; + this.modifiedSnapshotId = newSnapshotId; this.status = EntryStatus.MODIFIED; return this; } @@ -130,7 +143,7 @@ TrackingBuilder replacedPositions(ByteBuffer positions) { Preconditions.checkState( status != EntryStatus.ADDED, "Cannot set replaced positions on ADDED entry"); this.replacedPositions = ByteBuffers.toByteArray(positions); - this.dvSnapshotId = newSnapshotId; + this.modifiedSnapshotId = newSnapshotId; this.status = EntryStatus.MODIFIED; return this; } @@ -141,7 +154,7 @@ Tracking build() { snapshotId, dataSequenceNumber, fileSequenceNumber, - dvSnapshotId, + modifiedSnapshotId, firstRowId, deletedPositions, replacedPositions); @@ -155,7 +168,7 @@ private static Tracking terminal(EntryStatus to, Tracking source, long newSnapsh newSnapshotId, source.dataSequenceNumber(), source.fileSequenceNumber(), - source.dvSnapshotId(), + source.modifiedSnapshotId(), source.firstRowId(), null, null); diff --git a/core/src/main/java/org/apache/iceberg/TrackingStruct.java b/core/src/main/java/org/apache/iceberg/TrackingStruct.java index e5b351711113..01196b08c656 100644 --- a/core/src/main/java/org/apache/iceberg/TrackingStruct.java +++ b/core/src/main/java/org/apache/iceberg/TrackingStruct.java @@ -35,7 +35,7 @@ class TrackingStruct extends SupportsIndexProjection implements Tracking, Serial Tracking.SNAPSHOT_ID, Tracking.SEQUENCE_NUMBER, Tracking.FILE_SEQUENCE_NUMBER, - Tracking.DV_SNAPSHOT_ID, + Tracking.MODIFIED_SNAPSHOT_ID, Tracking.FIRST_ROW_ID, Tracking.DELETED_POSITIONS, Tracking.REPLACED_POSITIONS, @@ -58,7 +58,7 @@ class TrackingStruct extends SupportsIndexProjection implements Tracking, Serial private Long snapshotId = null; private Long dataSequenceNumber = null; private Long fileSequenceNumber = null; - private Long dvSnapshotId = null; + private Long modifiedSnapshotId = null; private Long firstRowId = null; private byte[] deletedPositions = null; private byte[] replacedPositions = null; @@ -82,7 +82,7 @@ private TrackingStruct(TrackingStruct toCopy) { this.snapshotId = toCopy.snapshotId; this.dataSequenceNumber = toCopy.dataSequenceNumber; this.fileSequenceNumber = toCopy.fileSequenceNumber; - this.dvSnapshotId = toCopy.dvSnapshotId; + this.modifiedSnapshotId = toCopy.modifiedSnapshotId; this.firstRowId = toCopy.firstRowId; this.deletedPositions = toCopy.deletedPositions != null @@ -101,7 +101,7 @@ private TrackingStruct(TrackingStruct toCopy) { Long snapshotId, Long dataSequenceNumber, Long fileSequenceNumber, - Long dvSnapshotId, + Long modifiedSnapshotId, Long firstRowId, byte[] deletedPositions, byte[] replacedPositions) { @@ -110,7 +110,7 @@ private TrackingStruct(TrackingStruct toCopy) { this.snapshotId = snapshotId; this.dataSequenceNumber = dataSequenceNumber; this.fileSequenceNumber = fileSequenceNumber; - this.dvSnapshotId = dvSnapshotId; + this.modifiedSnapshotId = modifiedSnapshotId; this.firstRowId = firstRowId; this.deletedPositions = deletedPositions; this.replacedPositions = replacedPositions; @@ -141,8 +141,9 @@ void inherit(long manifestSnapshotId, long manifestSeqNumber) { } boolean isAdded = status == EntryStatus.ADDED; + boolean isModified = status == EntryStatus.MODIFIED; - if (null == dataSequenceNumber && (isAdded || manifestSeqNumber == 0)) { + if (null == dataSequenceNumber && (isAdded || isModified || manifestSeqNumber == 0)) { this.dataSequenceNumber = manifestSeqNumber; } @@ -200,8 +201,8 @@ public Long fileSequenceNumber() { } @Override - public Long dvSnapshotId() { - return dvSnapshotId; + public Long modifiedSnapshotId() { + return modifiedSnapshotId; } @Override @@ -240,67 +241,37 @@ protected T internalGet(int pos, Class javaClass) { } private Object getByPos(int pos) { - switch (pos) { - case 0: - return status != null ? status.id() : null; - case 1: - return snapshotId(); - case 2: - return dataSequenceNumber(); - case 3: - return fileSequenceNumber(); - case 4: - return dvSnapshotId; - case 5: - return firstRowId; - case 6: - return deletedPositions(); - case 7: - return replacedPositions(); - case 8: - return manifestLocation; - case 9: - return manifestPos; - default: - throw new UnsupportedOperationException("Unknown field ordinal: " + pos); - } + return switch (pos) { + case 0 -> status != null ? status.id() : null; + case 1 -> snapshotId(); + case 2 -> dataSequenceNumber(); + case 3 -> fileSequenceNumber(); + case 4 -> modifiedSnapshotId; + case 5 -> firstRowId; + case 6 -> deletedPositions(); + case 7 -> replacedPositions(); + case 8 -> manifestLocation; + case 9 -> manifestPos; + default -> throw new UnsupportedOperationException("Unknown field ordinal: " + pos); + }; } @Override protected void internalSet(int pos, T value) { switch (pos) { - case 0: - this.status = EntryStatus.fromId((Integer) value); - break; - case 1: - this.snapshotId = (Long) value; - break; - case 2: - this.dataSequenceNumber = (Long) value; - break; - case 3: - this.fileSequenceNumber = (Long) value; - break; - case 4: - this.dvSnapshotId = (Long) value; - break; - case 5: - this.firstRowId = (Long) value; - break; - case 6: - this.deletedPositions = ByteBuffers.toByteArray((ByteBuffer) value); - break; - case 7: - this.replacedPositions = ByteBuffers.toByteArray((ByteBuffer) value); - break; - case 8: - this.manifestLocation = (String) value; - break; - case 9: - this.manifestPos = (long) value; - break; - default: + case 0 -> this.status = EntryStatus.fromId((Integer) value); + case 1 -> this.snapshotId = (Long) value; + case 2 -> this.dataSequenceNumber = (Long) value; + case 3 -> this.fileSequenceNumber = (Long) value; + case 4 -> this.modifiedSnapshotId = (Long) value; + case 5 -> this.firstRowId = (Long) value; + case 6 -> this.deletedPositions = ByteBuffers.toByteArray((ByteBuffer) value); + case 7 -> this.replacedPositions = ByteBuffers.toByteArray((ByteBuffer) value); + case 8 -> this.manifestLocation = (String) value; + case 9 -> this.manifestPos = (long) value; + default -> { // ignore the object, it must be from a newer version of the format + } } } @@ -311,7 +282,7 @@ public String toString() { .add("snapshot_id", snapshotId) .add("data_sequence_number", dataSequenceNumber) .add("file_sequence_number", fileSequenceNumber) - .add("dv_snapshot_id", dvSnapshotId) + .add("modified_snapshot_id", modifiedSnapshotId) .add("first_row_id", firstRowId) .add("deleted_positions", deletedPositions == null ? "null" : "(binary)") .add("replaced_positions", replacedPositions == null ? "null" : "(binary)") diff --git a/core/src/test/java/org/apache/iceberg/StatsTestUtil.java b/core/src/test/java/org/apache/iceberg/StatsTestUtil.java index 93309ceb9b49..b5c50cf211ee 100644 --- a/core/src/test/java/org/apache/iceberg/StatsTestUtil.java +++ b/core/src/test/java/org/apache/iceberg/StatsTestUtil.java @@ -42,6 +42,7 @@ static TrackedFile trackedFile(String location, long recordCount, ContentStats s null, null, null, + null, null); } diff --git a/core/src/test/java/org/apache/iceberg/TestColumnFileStruct.java b/core/src/test/java/org/apache/iceberg/TestColumnFileStruct.java new file mode 100644 index 000000000000..acae6aa23906 --- /dev/null +++ b/core/src/test/java/org/apache/iceberg/TestColumnFileStruct.java @@ -0,0 +1,210 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.iceberg; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +import java.io.IOException; +import java.nio.ByteBuffer; +import java.util.List; +import org.apache.iceberg.relocated.com.google.common.collect.Lists; +import org.apache.iceberg.types.Comparators; +import org.apache.iceberg.types.Types; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.MethodSource; + +class TestColumnFileStruct { + + private static final String LOCATION = "s3://bucket/data/column.parquet"; + private static final List FIELD_IDS = Lists.newArrayList(1, 2, 3); + private static final FileFormat FILE_FORMAT = FileFormat.PARQUET; + private static final long FILE_SIZE_IN_BYTES = 1024L; + private static final ByteBuffer KEY_METADATA = ByteBuffer.wrap(new byte[] {1, 2, 3}); + + @Test + void fieldAccess() { + ColumnFile columnFile = + new ColumnFileStruct(LOCATION, FIELD_IDS, FILE_FORMAT, FILE_SIZE_IN_BYTES, KEY_METADATA); + + assertThat(columnFile.location()).isEqualTo(LOCATION); + assertThat(columnFile.fieldIds()).containsExactlyElementsOf(FIELD_IDS); + assertThat(columnFile.fileFormat()).isEqualTo(FILE_FORMAT); + assertThat(columnFile.fileSizeInBytes()).isEqualTo(FILE_SIZE_IN_BYTES); + assertThat(columnFile.keyMetadata()).isEqualTo(KEY_METADATA); + } + + @Test + void copy() { + ColumnFile columnFile = + new ColumnFileStruct(LOCATION, FIELD_IDS, FILE_FORMAT, FILE_SIZE_IN_BYTES, KEY_METADATA); + + ColumnFile copy = columnFile.copy(); + + assertThat(copy.location()).isEqualTo(LOCATION); + assertThat(copy.fieldIds()).containsExactlyElementsOf(FIELD_IDS); + assertThat(copy.fileFormat()).isEqualTo(FILE_FORMAT); + assertThat(copy.fileSizeInBytes()).isEqualTo(FILE_SIZE_IN_BYTES); + assertThat(copy.keyMetadata()).isEqualTo(KEY_METADATA); + } + + @Test + void structLikeSize() { + ColumnFileStruct columnFile = new ColumnFileStruct(); + assertThat(columnFile.size()).isEqualTo(5); + } + + @Test + void setFieldsByOrdinals() { + ColumnFileStruct columnFile = new ColumnFileStruct(); + + columnFile.set(0, LOCATION); + columnFile.set(1, FIELD_IDS); + columnFile.set(2, FILE_FORMAT.toString()); + columnFile.set(3, FILE_SIZE_IN_BYTES); + columnFile.set(4, KEY_METADATA); + + assertThat(columnFile.location()).isEqualTo(LOCATION); + assertThat(columnFile.fieldIds()).containsExactlyElementsOf(FIELD_IDS); + assertThat(columnFile.fileFormat()).isEqualTo(FILE_FORMAT); + assertThat(columnFile.fileSizeInBytes()).isEqualTo(FILE_SIZE_IN_BYTES); + assertThat(columnFile.keyMetadata()).isEqualTo(KEY_METADATA); + } + + @Test + void getFieldsByOrdinals() { + ColumnFileStruct columnFile = + new ColumnFileStruct(LOCATION, FIELD_IDS, FILE_FORMAT, FILE_SIZE_IN_BYTES, KEY_METADATA); + + assertThat(columnFile.get(0, String.class)).isEqualTo(LOCATION); + assertThat(columnFile.get(1, List.class)).containsExactlyElementsOf(FIELD_IDS); + assertThat(columnFile.get(2, String.class)).isEqualTo(FILE_FORMAT.toString()); + assertThat(columnFile.get(3, Long.class)).isEqualTo(FILE_SIZE_IN_BYTES); + assertThat(columnFile.get(4, ByteBuffer.class)).isEqualTo(KEY_METADATA); + } + + @Test + void projectedStructLike() { + Types.StructType projection = + Types.StructType.of(ColumnFile.LOCATION, ColumnFile.FILE_SIZE_IN_BYTES); + + ColumnFileStruct columnFile = new ColumnFileStruct(projection); + assertThat(columnFile.size()).isEqualTo(2); + + // projected position 0 maps to internal position of location + // projected position 1 maps to internal position of file_size_in_bytes + columnFile.set(0, LOCATION); + columnFile.set(1, FILE_SIZE_IN_BYTES); + + assertThat(columnFile.location()).isEqualTo(LOCATION); + assertThat(columnFile.fileSizeInBytes()).isEqualTo(FILE_SIZE_IN_BYTES); + assertThat(columnFile.get(0, String.class)).isEqualTo(LOCATION); + assertThat(columnFile.get(1, Long.class)).isEqualTo(FILE_SIZE_IN_BYTES); + } + + @ParameterizedTest + @MethodSource("org.apache.iceberg.TestHelpers#serializers") + void serializationRoundTrip(TestHelpers.RoundTripSerializer roundTripSerializer) + throws IOException, ClassNotFoundException { + ColumnFile columnFile = + new ColumnFileStruct(LOCATION, FIELD_IDS, FILE_FORMAT, FILE_SIZE_IN_BYTES, KEY_METADATA); + + ColumnFile deserialized = roundTripSerializer.apply(columnFile); + + assertThat((StructLike) deserialized) + .usingComparator(Comparators.forType(ColumnFile.schema())) + .isEqualTo(columnFile); + } + + @Test + void invalidBuilderValues() { + assertThatThrownBy(() -> ColumnFileStruct.builder().location(null).build()) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("Invalid location: null"); + + assertThatThrownBy(() -> ColumnFileStruct.builder().location("").build()) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("Invalid location: empty"); + + assertThatThrownBy(() -> ColumnFileStruct.builder().fieldIds(null).build()) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("Invalid field IDs: null"); + + assertThatThrownBy(() -> ColumnFileStruct.builder().fieldIds(Lists.newArrayList()).build()) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("Invalid field IDs: empty"); + + assertThatThrownBy( + () -> ColumnFileStruct.builder().fieldIds(Lists.newArrayList(1, 2, 1)).build()) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("Invalid field IDs: duplicated IDs found in: [1, 2, 1]"); + + assertThatThrownBy(() -> ColumnFileStruct.builder().fileFormat(null)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("Invalid file format: null"); + + assertThatThrownBy(() -> ColumnFileStruct.builder().fileSizeInBytes(-1).build()) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("Invalid file size in bytes: -1 (must be >= 0)"); + } + + @Test + void missingBuilderValues() { + assertThatThrownBy( + () -> + ColumnFileStruct.builder() + .fieldIds(FIELD_IDS) + .fileFormat(FILE_FORMAT) + .fileSizeInBytes(FILE_SIZE_IN_BYTES) + .build()) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("Missing required value: location"); + + assertThatThrownBy( + () -> + ColumnFileStruct.builder() + .location(LOCATION) + .fileFormat(FILE_FORMAT) + .fileSizeInBytes(FILE_SIZE_IN_BYTES) + .build()) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("Missing required value: field IDs"); + + assertThatThrownBy( + () -> + ColumnFileStruct.builder() + .location(LOCATION) + .fieldIds(FIELD_IDS) + .fileSizeInBytes(FILE_SIZE_IN_BYTES) + .build()) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("Missing required value: file format"); + + assertThatThrownBy( + () -> + ColumnFileStruct.builder() + .location(LOCATION) + .fieldIds(FIELD_IDS) + .fileFormat(FILE_FORMAT) + .build()) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("Missing required value: file size in bytes"); + } +} diff --git a/core/src/test/java/org/apache/iceberg/TestTrackedFile.java b/core/src/test/java/org/apache/iceberg/TestTrackedFile.java index 9a32243cc12f..554737e412b9 100644 --- a/core/src/test/java/org/apache/iceberg/TestTrackedFile.java +++ b/core/src/test/java/org/apache/iceberg/TestTrackedFile.java @@ -61,7 +61,8 @@ public void schemaFieldOrder() { "manifest_info", "key_metadata", "split_offsets", - "equality_ids"); + "equality_ids", + "column_files"); } @Test @@ -72,7 +73,7 @@ public void schemaFieldIds() { assertThat(fields) .extracting(Types.NestedField::fieldId) .containsExactly( - 147, 134, 157, 100, 101, 103, 104, 141, 102, 146, 140, 148, 150, 131, 132, 135); + 147, 134, 157, 100, 101, 103, 104, 141, 102, 146, 140, 148, 150, 131, 132, 135, 158); } @Test diff --git a/core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java b/core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java index a7bfe777c147..7dff7d265d7b 100644 --- a/core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java +++ b/core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java @@ -102,7 +102,7 @@ class TestTrackedFileAdapters { SNAPSHOT_ID, DATA_SEQUENCE_NUMBER, FILE_SEQUENCE_NUMBER, - null, // dvSnapshotId + null, // modifiedSnapshotId FIRST_ROW_ID, null, // deletedPositions null); // replacedPositions @@ -158,6 +158,7 @@ void dataFileAdapterDelegation() { null, ByteBuffer.wrap(new byte[] {1, 2, 3}), ImmutableList.of(50L, 100L), + null, null); DataFile dataFile = TrackedFileAdapters.asDataFile(file, specsById(PARTITIONED_SPEC)); @@ -238,7 +239,8 @@ void equalityDeleteFileAdapterDelegation() { null, ByteBuffer.wrap(new byte[] {4, 5}), ImmutableList.of(200L), - ImmutableList.of(1, 2, 3)); + ImmutableList.of(1, 2, 3), + null); DeleteFile deleteFile = TrackedFileAdapters.asEqualityDeleteFile(file, specsById(PARTITIONED_SPEC)); @@ -327,6 +329,7 @@ void dvDeleteFileAdapterDelegation() { null, null, null, + null, null); DeleteFile dvFile = TrackedFileAdapters.asDVDeleteFile(file, specsById(PARTITIONED_SPEC)); @@ -459,7 +462,8 @@ private static TrackedFile serializableDataEntry() { null, // manifestInfo KEY_METADATA, ImmutableList.of(50L, 100L), // splitOffsets - null); // equalityIds + null, // equalityIds + null); // columnFiles } @ParameterizedTest @@ -503,7 +507,8 @@ void manifestFileAdapterDelegation(FileContent contentType) { MANIFEST_INFO, MANIFEST_KEY_METADATA, null, // splitOffsets - null); // equalityIds + null, // equalityIds + null); // columnFiles ManifestFile manifest = TrackedFileAdapters.asManifestFile(file); @@ -610,6 +615,7 @@ void nullTrackingReturnsNullTrackingFields() { null, null, null, + null, null); assertNullTrackingFields(TrackedFileAdapters.asDVDeleteFile(fileWithDV, UNPARTITIONED)); } @@ -664,6 +670,7 @@ void unknownSpecIdThrows() { null, null, null, + null, null); assertThatThrownBy(() -> TrackedFileAdapters.asDataFile(file, ImmutableMap.of())) @@ -708,6 +715,7 @@ private static TrackedFileStruct dummyTrackedFile(FileContent contentType) { null, null, null, + null, null); } diff --git a/core/src/test/java/org/apache/iceberg/TestTrackedFileStruct.java b/core/src/test/java/org/apache/iceberg/TestTrackedFileStruct.java index f5edda55678f..d3062e865b73 100644 --- a/core/src/test/java/org/apache/iceberg/TestTrackedFileStruct.java +++ b/core/src/test/java/org/apache/iceberg/TestTrackedFileStruct.java @@ -59,11 +59,18 @@ class TestTrackedFileStruct { private static final ManifestInfo MANIFEST_INFO = Mockito.mock(ManifestInfo.class); private static final ManifestInfo MANIFEST_INFO_COPY = Mockito.mock(ManifestInfo.class); + private static final ColumnFile COLUMN_FILE_1 = Mockito.mock(ColumnFile.class); + private static final ColumnFile COLUMN_FILE_2 = Mockito.mock(ColumnFile.class); + private static final ColumnFile COLUMN_FILE_1_COPY = Mockito.mock(ColumnFile.class); + private static final ColumnFile COLUMN_FILE_2_COPY = Mockito.mock(ColumnFile.class); + static { Mockito.when(TRACKING.copy()).thenReturn(TRACKING_COPY); Mockito.when(CONTENT_STATS.copy()).thenReturn(CONTENT_STATS_COPY); Mockito.when(DELETION_VECTOR.copy()).thenReturn(DELETION_VECTOR_COPY); Mockito.when(MANIFEST_INFO.copy()).thenReturn(MANIFEST_INFO_COPY); + Mockito.when(COLUMN_FILE_1.copy()).thenReturn(COLUMN_FILE_1_COPY); + Mockito.when(COLUMN_FILE_2.copy()).thenReturn(COLUMN_FILE_2_COPY); } @Test @@ -85,7 +92,8 @@ void fieldAccess() { MANIFEST_INFO, ByteBuffer.wrap(new byte[] {1, 2, 3}), ImmutableList.of(100L, 200L), - ImmutableList.of(1, 2, 3)); + ImmutableList.of(1, 2, 3), + ImmutableList.of(COLUMN_FILE_1, COLUMN_FILE_2)); assertThat(file.tracking()).isSameAs(TRACKING); assertThat(file.contentType()).isEqualTo(FileContent.DATA); @@ -103,6 +111,7 @@ void fieldAccess() { assertThat(file.keyMetadata()).isEqualTo(ByteBuffer.wrap(new byte[] {1, 2, 3})); assertThat(file.splitOffsets()).containsExactly(100L, 200L); assertThat(file.equalityIds()).containsExactly(1, 2, 3); + assertThat(file.columnFiles()).containsExactly(COLUMN_FILE_1, COLUMN_FILE_2); } @Test @@ -162,7 +171,8 @@ void getByPosition() { MANIFEST_INFO, ByteBuffer.wrap(new byte[] {1, 2, 3}), ImmutableList.of(100L, 200L), - ImmutableList.of(1, 2, 3)); + ImmutableList.of(1, 2, 3), + ImmutableList.of(COLUMN_FILE_1, COLUMN_FILE_2)); assertThat(file.get(pos("tracking"), Tracking.class)).isSameAs(TRACKING); assertThat(file.get(pos("content_type"), Integer.class)).isEqualTo(FileContent.DATA.id()); @@ -203,7 +213,8 @@ void copy() { MANIFEST_INFO, ByteBuffer.wrap(new byte[] {1, 2, 3}), ImmutableList.of(100L, 200L), - ImmutableList.of(1, 2, 3)); + ImmutableList.of(1, 2, 3), + ImmutableList.of(COLUMN_FILE_1, COLUMN_FILE_2)); TrackedFile copy = file.copy(); @@ -224,6 +235,7 @@ void copy() { assertThat(copy.splitOffsets()).containsExactly(100L, 200L); assertThat(copy.equalityIds()).containsExactly(1, 2, 3); assertThat(copy.partition()).isNotSameAs(PARTITION); + assertThat(copy.columnFiles()).containsExactly(COLUMN_FILE_1_COPY, COLUMN_FILE_2_COPY); // mutable fields are deep-copied, not shared with the original assertThat(copy.keyMetadata()).isNotSameAs(file.keyMetadata()); @@ -252,7 +264,8 @@ void copyWithStats() { MANIFEST_INFO, ByteBuffer.wrap(new byte[] {1, 2, 3}), ImmutableList.of(100L, 200L), - ImmutableList.of(1, 2, 3)); + ImmutableList.of(1, 2, 3), + ImmutableList.of(COLUMN_FILE_1, COLUMN_FILE_2)); TrackedFile copy = file.copyWithStats(ImmutableSet.of(1)); @@ -273,6 +286,7 @@ void copyWithStats() { assertThat(copy.splitOffsets()).containsExactly(100L, 200L); assertThat(copy.equalityIds()).containsExactly(1, 2, 3); assertThat(copy.partition()).isNotSameAs(PARTITION); + assertThat(copy.columnFiles()).containsExactly(COLUMN_FILE_1_COPY, COLUMN_FILE_2_COPY); // mutable fields are deep-copied, not shared with the original assertThat(copy.keyMetadata()).isNotSameAs(file.keyMetadata()); @@ -299,7 +313,8 @@ void copyWithoutStats() { MANIFEST_INFO, ByteBuffer.wrap(new byte[] {1, 2, 3}), ImmutableList.of(100L, 200L), - ImmutableList.of(1, 2, 3)); + ImmutableList.of(1, 2, 3), + ImmutableList.of(COLUMN_FILE_1, COLUMN_FILE_2)); TrackedFile copy = file.copyWithoutStats(); @@ -323,9 +338,11 @@ void copyWithoutStats() { assertThat(copy.splitOffsets()).containsExactly(100L, 200L); assertThat(copy.equalityIds()).containsExactly(1, 2, 3); assertThat(copy.partition()).isNotSameAs(PARTITION); + assertThat(copy.columnFiles()).containsExactly(COLUMN_FILE_1_COPY, COLUMN_FILE_2_COPY); // mutable fields are deep-copied, not shared with the original assertThat(copy.keyMetadata()).isNotSameAs(file.keyMetadata()); + assertThat(copy.columnFiles()).isNotSameAs(file.columnFiles()); } @Test @@ -449,7 +466,8 @@ void serializationRoundTrip(RoundTripSerializer serializer) t null, // ManifestInfo has its own serialization tests ByteBuffer.wrap(new byte[] {1, 2, 3}), ImmutableList.of(50L), - ImmutableList.of(1, 2, 3)); + ImmutableList.of(1, 2, 3), + null); // ColumnFile has its own serialization tests TrackedFileStruct deserialized = serializer.apply(file); @@ -470,6 +488,7 @@ void serializationRoundTrip(RoundTripSerializer serializer) t assertThat(deserialized.keyMetadata()).isEqualTo(ByteBuffer.wrap(new byte[] {1, 2, 3})); assertThat(deserialized.splitOffsets()).containsExactly(50L); assertThat(deserialized.equalityIds()).containsExactly(1, 2, 3); + assertThat(deserialized.columnFiles()).isNull(); } @ParameterizedTest @@ -526,7 +545,8 @@ private static TrackedFileStruct trackedFile(Integer specId, PartitionData parti null, // manifestInfo null, // keyMetadata null, // splitOffsets - null); // equalityIds + null, // equalityIds + null); // column files } private static int pos(String fieldName) { diff --git a/core/src/test/java/org/apache/iceberg/TestTrackingBuilder.java b/core/src/test/java/org/apache/iceberg/TestTrackingBuilder.java index f7e62271c76b..0813e5561ae6 100644 --- a/core/src/test/java/org/apache/iceberg/TestTrackingBuilder.java +++ b/core/src/test/java/org/apache/iceberg/TestTrackingBuilder.java @@ -45,7 +45,7 @@ void addedWithSameCommitDvStaysAdded() { assertThat(tracking.status()).isEqualTo(EntryStatus.ADDED); assertThat(tracking.snapshotId()).isEqualTo(42L); - assertThat(tracking.dvSnapshotId()).isEqualTo(42L); + assertThat(tracking.modifiedSnapshotId()).isEqualTo(42L); assertThat(tracking.deletedPositions()).isNull(); assertThat(tracking.replacedPositions()).isNull(); // sequence numbers and firstRowId remain null; populated by inheritance @@ -62,7 +62,7 @@ void existingBuilderPreservesSourceFields() { assertThat(existing.snapshotId()).isEqualTo(SOURCE_TRACKING_ADDED.snapshotId()); assertThat(existing.dataSequenceNumber()).isEqualTo(SOURCE_TRACKING_ADDED.dataSequenceNumber()); assertThat(existing.fileSequenceNumber()).isEqualTo(SOURCE_TRACKING_ADDED.fileSequenceNumber()); - assertThat(existing.dvSnapshotId()).isEqualTo(SOURCE_TRACKING_ADDED.dvSnapshotId()); + assertThat(existing.modifiedSnapshotId()).isEqualTo(SOURCE_TRACKING_ADDED.modifiedSnapshotId()); assertThat(existing.firstRowId()).isEqualTo(SOURCE_TRACKING_ADDED.firstRowId()); } @@ -74,7 +74,7 @@ void deleteUpdatesSnapshotIdAndPreservesRest() { assertThat(deleted.snapshotId()).isEqualTo(999L); assertThat(deleted.dataSequenceNumber()).isEqualTo(SOURCE_TRACKING_ADDED.dataSequenceNumber()); assertThat(deleted.fileSequenceNumber()).isEqualTo(SOURCE_TRACKING_ADDED.fileSequenceNumber()); - assertThat(deleted.dvSnapshotId()).isEqualTo(SOURCE_TRACKING_ADDED.dvSnapshotId()); + assertThat(deleted.modifiedSnapshotId()).isEqualTo(SOURCE_TRACKING_ADDED.modifiedSnapshotId()); assertThat(deleted.firstRowId()).isEqualTo(SOURCE_TRACKING_ADDED.firstRowId()); } @@ -86,7 +86,7 @@ void replaceUpdatesSnapshotIdAndPreservesRest() { assertThat(replaced.snapshotId()).isEqualTo(999L); assertThat(replaced.dataSequenceNumber()).isEqualTo(SOURCE_TRACKING_ADDED.dataSequenceNumber()); assertThat(replaced.fileSequenceNumber()).isEqualTo(SOURCE_TRACKING_ADDED.fileSequenceNumber()); - assertThat(replaced.dvSnapshotId()).isEqualTo(SOURCE_TRACKING_ADDED.dvSnapshotId()); + assertThat(replaced.modifiedSnapshotId()).isEqualTo(SOURCE_TRACKING_ADDED.modifiedSnapshotId()); assertThat(replaced.firstRowId()).isEqualTo(SOURCE_TRACKING_ADDED.firstRowId()); } @@ -110,7 +110,7 @@ void sourceDvPositionsAreNotCarriedForward() { } @Test - void dvUpdatedProducesModifiedAndAdvancesDvSnapshotId() { + void dvUpdatedProducesModifiedAndAdvancesModifiedSnapshotId() { Tracking modified = TrackingBuilder.from(SOURCE_TRACKING_ADDED, 999L).dvUpdated().build(); assertThat(modified.status()).isEqualTo(EntryStatus.MODIFIED); @@ -118,8 +118,8 @@ void dvUpdatedProducesModifiedAndAdvancesDvSnapshotId() { assertThat(modified.snapshotId()) .isEqualTo(SOURCE_TRACKING_ADDED.snapshotId()) .isNotEqualTo(999L); - // only the DV snapshot id advances to the commit snapshot - assertThat(modified.dvSnapshotId()).isEqualTo(999L); + // only the modified snapshot id advances to the commit snapshot + assertThat(modified.modifiedSnapshotId()).isEqualTo(999L); } @Test @@ -255,8 +255,8 @@ void carryForwardFromModifiedSourceChangesToExisting() { assertThat(carried.snapshotId()) .isEqualTo(SOURCE_TRACKING_MODIFIED.snapshotId()) .isNotEqualTo(999L); - assertThat(carried.dvSnapshotId()) - .isEqualTo(SOURCE_TRACKING_MODIFIED.dvSnapshotId()) + assertThat(carried.modifiedSnapshotId()) + .isEqualTo(SOURCE_TRACKING_MODIFIED.modifiedSnapshotId()) .isNotEqualTo(999L); assertThat(carried.dataSequenceNumber()) .isEqualTo(SOURCE_TRACKING_MODIFIED.dataSequenceNumber()); @@ -272,11 +272,39 @@ void manifestDVPositionsProduceModified() { Tracking modified = TrackingBuilder.from(SOURCE_TRACKING_ADDED, 999L).deletedPositions(deletedBytes).build(); assertThat(modified.status()).isEqualTo(EntryStatus.MODIFIED); - // the entry snapshot id is preserved; only the DV snapshot id advances to the commit snapshot + // the entry snapshot id is preserved; only the modified snapshot id advances to the commit + // snapshot assertThat(modified.snapshotId()) .isEqualTo(SOURCE_TRACKING_ADDED.snapshotId()) .isNotEqualTo(999L); - assertThat(modified.dvSnapshotId()).isEqualTo(999L); + assertThat(modified.modifiedSnapshotId()).isEqualTo(999L); assertThat(modified.deletedPositions()).isEqualTo(deletedBytes); } + + @Test + void manifestPositionsWithColumnFilesUpdated() { + ByteBuffer deletedBytes = ByteBuffer.wrap(new byte[] {1}); + Tracking withDeletedPositions = + TrackingBuilder.from(SOURCE_TRACKING_ADDED, 999L) + .columnFilesUpdated() + .deletedPositions(deletedBytes) + .build(); + + assertThat(withDeletedPositions.status()).isEqualTo(EntryStatus.MODIFIED); + assertThat(withDeletedPositions.modifiedSnapshotId()).isEqualTo(999L); + assertThat(withDeletedPositions.deletedPositions()).isEqualTo(deletedBytes); + assertThat(withDeletedPositions.dataSequenceNumber()).isNull(); + + ByteBuffer replacedBytes = ByteBuffer.wrap(new byte[] {2}); + Tracking withReplacedPositions = + TrackingBuilder.from(SOURCE_TRACKING_ADDED, 999L) + .columnFilesUpdated() + .replacedPositions(replacedBytes) + .build(); + + assertThat(withReplacedPositions.status()).isEqualTo(EntryStatus.MODIFIED); + assertThat(withReplacedPositions.modifiedSnapshotId()).isEqualTo(999L); + assertThat(withReplacedPositions.replacedPositions()).isEqualTo(replacedBytes); + assertThat(withReplacedPositions.dataSequenceNumber()).isNull(); + } } diff --git a/core/src/test/java/org/apache/iceberg/TestTrackingStruct.java b/core/src/test/java/org/apache/iceberg/TestTrackingStruct.java index 253384a6834c..eea6e4e0a78d 100644 --- a/core/src/test/java/org/apache/iceberg/TestTrackingStruct.java +++ b/core/src/test/java/org/apache/iceberg/TestTrackingStruct.java @@ -49,7 +49,7 @@ void fieldAccess() { assertThat(tracking.snapshotId()).isEqualTo(42L); assertThat(tracking.dataSequenceNumber()).isEqualTo(10L); assertThat(tracking.fileSequenceNumber()).isEqualTo(11L); - assertThat(tracking.dvSnapshotId()).isEqualTo(43L); + assertThat(tracking.modifiedSnapshotId()).isEqualTo(43L); assertThat(tracking.firstRowId()).isEqualTo(1000L); assertThat(tracking.deletedPositions()).isEqualTo(ByteBuffer.wrap(DELETED_POSITIONS)); assertThat(tracking.replacedPositions()).isEqualTo(ByteBuffer.wrap(REPLACED_POSITIONS)); @@ -65,7 +65,7 @@ void setByPosition() { tracking.set(pos("snapshot_id"), 42L); tracking.set(pos("sequence_number"), 10L); tracking.set(pos("file_sequence_number"), 11L); - tracking.set(pos("dv_snapshot_id"), 43L); + tracking.set(pos("modified_snapshot_id"), 43L); tracking.set(pos("first_row_id"), 1000L); tracking.set(pos("deleted_positions"), ByteBuffer.wrap(DELETED_POSITIONS)); tracking.set(pos("replaced_positions"), ByteBuffer.wrap(REPLACED_POSITIONS)); @@ -76,7 +76,7 @@ void setByPosition() { assertThat(tracking.snapshotId()).isEqualTo(42L); assertThat(tracking.dataSequenceNumber()).isEqualTo(10L); assertThat(tracking.fileSequenceNumber()).isEqualTo(11L); - assertThat(tracking.dvSnapshotId()).isEqualTo(43L); + assertThat(tracking.modifiedSnapshotId()).isEqualTo(43L); assertThat(tracking.firstRowId()).isEqualTo(1000L); assertThat(tracking.deletedPositions()).isEqualTo(ByteBuffer.wrap(DELETED_POSITIONS)); assertThat(tracking.replacedPositions()).isEqualTo(ByteBuffer.wrap(REPLACED_POSITIONS)); @@ -96,7 +96,7 @@ void getByPosition() { assertThat(tracking.get(pos("snapshot_id"), Long.class)).isEqualTo(42L); assertThat(tracking.get(pos("sequence_number"), Long.class)).isEqualTo(10L); assertThat(tracking.get(pos("file_sequence_number"), Long.class)).isEqualTo(11L); - assertThat(tracking.get(pos("dv_snapshot_id"), Long.class)).isEqualTo(43L); + assertThat(tracking.get(pos("modified_snapshot_id"), Long.class)).isEqualTo(43L); assertThat(tracking.get(pos("first_row_id"), Long.class)).isEqualTo(1000L); assertThat(tracking.get(pos("deleted_positions"), ByteBuffer.class)) .isEqualTo(ByteBuffer.wrap(DELETED_POSITIONS)); @@ -121,7 +121,7 @@ void copy() { assertThat(copy.snapshotId()).isEqualTo(tracking.snapshotId()); assertThat(copy.dataSequenceNumber()).isEqualTo(tracking.dataSequenceNumber()); assertThat(copy.fileSequenceNumber()).isEqualTo(tracking.fileSequenceNumber()); - assertThat(copy.dvSnapshotId()).isEqualTo(tracking.dvSnapshotId()); + assertThat(copy.modifiedSnapshotId()).isEqualTo(tracking.modifiedSnapshotId()); assertThat(copy.firstRowId()).isEqualTo(tracking.firstRowId()); assertThat(copy.deletedPositions()).isEqualTo(tracking.deletedPositions()); assertThat(copy.replacedPositions()).isEqualTo(tracking.replacedPositions()); @@ -189,9 +189,19 @@ void inheritanceAddedEntriesInheritSequenceNumber() { assertThat(tracking.fileSequenceNumber()).isEqualTo(60L); } + @Test + void inheritanceModifiedEntriesInheritDataSequenceNumber() { + TrackingStruct tracking = + new TrackingStruct(EntryStatus.MODIFIED, 42L, null, null, null, null, null, null); + + tracking.inherit(100L, 60L); + + assertThat(tracking.dataSequenceNumber()).isEqualTo(60L); + assertThat(tracking.fileSequenceNumber()).isNull(); + } + private static final List NON_INHERITING_STATUSES = - List.of( - EntryStatus.EXISTING, EntryStatus.MODIFIED, EntryStatus.DELETED, EntryStatus.REPLACED); + List.of(EntryStatus.EXISTING, EntryStatus.DELETED, EntryStatus.REPLACED); @ParameterizedTest @FieldSource("NON_INHERITING_STATUSES") @@ -280,7 +290,7 @@ void internalSetIgnoresUnknownOrdinal() { assertThat(tracking.snapshotId()).isEqualTo(42L); assertThat(tracking.dataSequenceNumber()).isEqualTo(10L); assertThat(tracking.fileSequenceNumber()).isEqualTo(11L); - assertThat(tracking.dvSnapshotId()).isEqualTo(43L); + assertThat(tracking.modifiedSnapshotId()).isEqualTo(43L); assertThat(tracking.firstRowId()).isEqualTo(1000L); assertThat(tracking.deletedPositions()).isEqualTo(ByteBuffer.wrap(DELETED_POSITIONS)); assertThat(tracking.replacedPositions()).isEqualTo(ByteBuffer.wrap(REPLACED_POSITIONS)); @@ -321,7 +331,7 @@ void serializationRoundTrip(TestHelpers.RoundTripSerializer roundTripS assertThat(deserialized.snapshotId()).isEqualTo(42L); assertThat(deserialized.dataSequenceNumber()).isEqualTo(10L); assertThat(deserialized.fileSequenceNumber()).isEqualTo(11L); - assertThat(deserialized.dvSnapshotId()).isEqualTo(43L); + assertThat(deserialized.modifiedSnapshotId()).isEqualTo(43L); assertThat(deserialized.firstRowId()).isEqualTo(1000L); assertThat(deserialized.deletedPositions()).isEqualTo(ByteBuffer.wrap(DELETED_POSITIONS)); assertThat(deserialized.replacedPositions()).isEqualTo(ByteBuffer.wrap(REPLACED_POSITIONS)); diff --git a/core/src/test/java/org/apache/iceberg/TestV4ManifestReader.java b/core/src/test/java/org/apache/iceberg/TestV4ManifestReader.java index 9d1d81d3c6e5..f2a41bca28cb 100644 --- a/core/src/test/java/org/apache/iceberg/TestV4ManifestReader.java +++ b/core/src/test/java/org/apache/iceberg/TestV4ManifestReader.java @@ -171,7 +171,8 @@ public void readDataFile(FileFormat format) throws IOException { null, // manifest info ByteBuffer.wrap(new byte[] {1, 2, 3}), // key metadata ImmutableList.of(50L, 100L), - null); // equality field IDs + null, // equality field IDs + null); // column files ManifestFile manifest = writeManifest(format, ID_PARTITIONED_TYPE, file); @@ -203,7 +204,8 @@ public void readForScanPlanningDoesNotCopyStats(FileFormat format) throws IOExce null, // manifest info ByteBuffer.wrap(new byte[] {1, 2, 3}), // key metadata ImmutableList.of(50L, 100L), - null); // equality field IDs + null, // equality field IDs + null); // column files ManifestFile manifest = writeManifest(format, ID_PARTITIONED_TYPE, file); @@ -237,7 +239,8 @@ public void readForScanPlanningCopiesRequestedStats(FileFormat format) throws IO null, // manifest info ByteBuffer.wrap(new byte[] {1, 2, 3}), // key metadata ImmutableList.of(50L, 100L), - null); // equality field IDs + null, // equality field IDs + null); // column files ManifestFile manifest = writeManifest(format, ID_PARTITIONED_TYPE, file); @@ -273,7 +276,8 @@ public void readEqualityDelete(FileFormat format) throws IOException { null, // manifest info ByteBuffer.wrap(new byte[] {1, 2, 3}), // key metadata null, // split offsets - ImmutableList.of(1, 2)); + ImmutableList.of(1, 2), + null); // column files ManifestFile manifest = writeManifest(format, ID_PARTITIONED_TYPE, delete); @@ -305,7 +309,8 @@ public void readManifestFile(FileFormat format) throws IOException { MANIFEST_INFO, ByteBuffer.wrap(new byte[] {1, 2, 3}), // key metadata null, // split offsets - ImmutableList.of(1, 2)); + ImmutableList.of(1, 2), + null); // column files ManifestFile manifest = writeManifest(format, ID_PARTITIONED_TYPE, manifestRef); @@ -517,20 +522,22 @@ public void inheritanceSnapshotId(FileFormat format) throws IOException { @ParameterizedTest @FieldSource("MANIFEST_FORMATS") - public void inheritanceDVSnapshotIdNotInherited(FileFormat format) throws IOException { - TrackedFile withDVSnapshotId = + public void inheritanceModifiedSnapshotIdNotInherited(FileFormat format) throws IOException { + TrackedFile withModifiedSnapshotId = unpartitionedDataFile( new TrackingStruct( EntryStatus.ADDED, SNAPSHOT_ID, 5L, 5L, 1234567L, 5_000L, null, null), "s3://bucket/table/file-b.parquet"); - TrackedFile withoutDVSnapshotId = + TrackedFile withoutModifiedSnapshotId = unpartitionedDataFile( new TrackingStruct(EntryStatus.ADDED, SNAPSHOT_ID, 5L, 5L, null, 5_100L, null, null), "s3://bucket/table/file-a.parquet"); ManifestFile manifest = writeManifest( - format, UNPARTITIONED_TYPE, ImmutableList.of(withDVSnapshotId, withoutDVSnapshotId)); + format, + UNPARTITIONED_TYPE, + ImmutableList.of(withModifiedSnapshotId, withoutModifiedSnapshotId)); when(manifest.firstRowId()).thenReturn(10_000L); when(manifest.snapshotId()).thenReturn(34L); @@ -1169,7 +1176,8 @@ public void statsFilterRecordCountFiltering(FileFormat format) throws IOExceptio null, null, List.of(4L), - null); + null, + null); // column files ManifestFile manifest = writeManifest(format, UNPARTITIONED_TYPE, ImmutableList.of(emptyTrackedFile, FILE_D)); @@ -1206,7 +1214,8 @@ public void statsFilterInvalidRecordCountNotFiltered(FileFormat format) throws I null, null, List.of(4L), - null); + null, + null); // column files ManifestFile manifest = writeManifest(format, UNPARTITIONED_TYPE, ImmutableList.of(invalidRecordCountFile, FILE_D)); @@ -1391,7 +1400,8 @@ public void partitionFilterBucketPartitionID(FileFormat format) throws IOExcepti null, // manifest info null, // key metadata List.of(4L), - null); // eq delete ids + null, // eq delete ids + null); // column files ManifestFile manifest = writeManifest( @@ -1841,7 +1851,7 @@ private static TrackedFile unpartitionedFileWithoutStats(String location) { SNAPSHOT_ID, 3L, // data sequence number 3L, // file sequence number - null, // dv snapshot id + null, // modified snapshot id null, // first row id null, // deleted positions null); // replaced positions @@ -1871,7 +1881,8 @@ private static TrackedFile dataFileWithoutStats( null, // manifest_info null, // key_metadata ImmutableList.of(4L), // split offsets - null); // equality_ids + null, // equality field IDs + null); // column files } private static TrackedFile idPartitionedDeleteFileWithoutStats( @@ -1892,7 +1903,8 @@ private static TrackedFile idPartitionedDeleteFileWithoutStats( null, // manifest_info null, // key_metadata ImmutableList.of(4L), // split offsets - ImmutableList.of(1)); // equality_ids + ImmutableList.of(1), // equality_ids + null); // column files } private static TrackedFile unpartitionedDataWithDVFile(String location, String dvLocation) { @@ -1924,7 +1936,8 @@ private static TrackedFile manifestRef(FileContent content, String location, Con MANIFEST_INFO, null, // key_metadata ImmutableList.of(4L), // split_offsets - null); // equality_ids + null, // equality_ids + null); // column files } private static TrackedFile unpartitionedDataFile(Tracking tracking, String location) { @@ -1949,7 +1962,8 @@ private static TrackedFile unpartitionedDataFile( null, // manifest info null, // key metadata ImmutableList.of(4L), // split offsets - null); // equality ids + null, // equality ids + null); // column files } private static PartitionData idPartition(int id) { diff --git a/core/src/test/java/org/apache/iceberg/V4TestComparators.java b/core/src/test/java/org/apache/iceberg/V4TestComparators.java index 9c3227afd7a3..08ea44ed353f 100644 --- a/core/src/test/java/org/apache/iceberg/V4TestComparators.java +++ b/core/src/test/java/org/apache/iceberg/V4TestComparators.java @@ -91,7 +91,7 @@ private static > Comparator natural() { .thenComparing(Tracking::snapshotId, natural()) .thenComparing(Tracking::dataSequenceNumber, natural()) .thenComparing(Tracking::fileSequenceNumber, natural()) - .thenComparing(Tracking::dvSnapshotId, natural()) + .thenComparing(Tracking::modifiedSnapshotId, natural()) .thenComparing(Tracking::firstRowId, natural()) .thenComparing(Tracking::deletedPositions, BYTES) .thenComparing(Tracking::replacedPositions, BYTES));