stevenzwu commented on code in PR #16285:
URL: https://github.com/apache/iceberg/pull/16285#discussion_r3738431274
##########
core/src/main/java/org/apache/iceberg/TrackingStruct.java:
##########
@@ -130,11 +135,13 @@ void inheritFrom(Tracking manifestTracking) {
manifestTracking.dataSequenceNumber(),
manifestTracking.fileSequenceNumber());
- if (status == EntryStatus.ADDED) {
+ if (status == EntryStatus.ADDED || status == EntryStatus.MODIFIED) {
if (dataSequenceNumber == null) {
this.dataSequenceNumber = manifestTracking.fileSequenceNumber();
}
+ }
+ if (status == EntryStatus.ADDED) {
Review Comment:
wondering if this guardrail is needed.
if we want it, it maybe be better to add a `Preconditions.check(status ==
EntryStatus.ADDED)` inside the `if (fileSequenceNumber == null)` block as that
is the only valid scenario.
##########
core/src/main/java/org/apache/iceberg/ColumnFileStruct.java:
##########
@@ -0,0 +1,271 @@
+/*
+ * 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.FORMAT_VERSION,
+ ColumnFile.FIELD_IDS,
+ ColumnFile.LOCATION,
+ ColumnFile.FILE_FORMAT,
+ ColumnFile.FILE_SIZE_IN_BYTES,
+ ColumnFile.KEY_METADATA,
+ ColumnFile.SPLIT_OFFSETS);
+
+ private int formatVersion = -1;
+ private int[] fieldIds = null;
+ private String location = null;
+ private FileFormat fileFormat = null;
+ private long fileSizeInBytes = -1L;
+ private byte[] keyMetadata = null;
+ private long[] splitOffsets = null;
+
+ /** Used by internal readers to instantiate this class with a projection
schema. */
+ ColumnFileStruct(Types.StructType projection) {
+ super(BASE_TYPE, projection);
+ }
+
+ ColumnFileStruct(
+ int formatVersion,
+ List<Integer> fieldIds,
+ String location,
+ FileFormat fileFormat,
+ long fileSizeInBytes,
+ ByteBuffer keyMetadata,
+ List<Long> splitOffsets) {
+ super(BASE_TYPE.fields().size());
+ this.formatVersion = formatVersion;
+ this.fieldIds = ArrayUtil.toIntArray(fieldIds);
+ this.location = location;
+ this.fileFormat = fileFormat;
+ this.fileSizeInBytes = fileSizeInBytes;
+ this.keyMetadata = ByteBuffers.toByteArray(keyMetadata);
+ this.splitOffsets = ArrayUtil.toLongArray(splitOffsets);
+ }
+
+ /** Copy constructor. */
+ private ColumnFileStruct(ColumnFileStruct toCopy) {
+ super(toCopy);
+ this.formatVersion = toCopy.formatVersion;
+ this.fieldIds =
+ toCopy.fieldIds != null ? Arrays.copyOf(toCopy.fieldIds,
toCopy.fieldIds.length) : null;
+ this.location = toCopy.location;
+ this.fileFormat = toCopy.fileFormat;
+ this.fileSizeInBytes = toCopy.fileSizeInBytes;
+ this.keyMetadata =
+ toCopy.keyMetadata != null
+ ? Arrays.copyOf(toCopy.keyMetadata, toCopy.keyMetadata.length)
+ : null;
+ this.splitOffsets =
+ toCopy.splitOffsets != null
+ ? Arrays.copyOf(toCopy.splitOffsets, toCopy.splitOffsets.length)
+ : null;
+ }
+
+ /** Constructor for Java serialization. */
+ ColumnFileStruct() {
+ super(BASE_TYPE.fields().size());
+ }
+
+ @Override
+ public int formatVersion() {
+ return formatVersion;
+ }
+
+ @Override
+ public List<Integer> fieldIds() {
+ return fieldIds != null ? ArrayUtil.toUnmodifiableIntList(fieldIds) : null;
+ }
+
+ @Override
+ public String location() {
+ return location;
+ }
+
+ @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 List<Long> splitOffsets() {
+ return splitOffsets != null ?
ArrayUtil.toUnmodifiableLongList(splitOffsets) : null;
+ }
+
+ @Override
+ public ColumnFile copy() {
+ return new ColumnFileStruct(this);
+ }
+
+ @Override
+ protected <T> T internalGet(int pos, Class<T> javaClass) {
+ return javaClass.cast(getByPos(pos));
+ }
+
+ private Object getByPos(int pos) {
+ return switch (pos) {
+ case 0 -> formatVersion;
+ case 1 -> fieldIds();
+ case 2 -> location;
+ case 3 -> fileFormat != null ? fileFormat.toString() : null;
+ case 4 -> fileSizeInBytes;
+ case 5 -> keyMetadata();
+ case 6 -> splitOffsets();
+ default -> throw new UnsupportedOperationException("Unknown field
ordinal: " + pos);
+ };
+ }
+
+ @Override
+ @SuppressWarnings("unchecked")
+ protected <T> void internalSet(int pos, T value) {
+ switch (pos) {
+ case 0 -> this.formatVersion = (int) value;
+ case 1 -> this.fieldIds = ArrayUtil.toIntArray((List<Integer>) value);
+ // always coerce to String for Serializable
+ case 2 -> this.location = value.toString();
+ case 3 -> this.fileFormat = FileFormat.fromString(value.toString());
+ case 4 -> this.fileSizeInBytes = (long) value;
+ case 5 -> this.keyMetadata = ByteBuffers.toByteArray((ByteBuffer) value);
+ case 6 -> this.splitOffsets = ArrayUtil.toLongArray((List<Long>) value);
+ default -> {
+ // ignore the object, it must be from a newer version of the format
Review Comment:
nit: should the comment say `ignore the unknown positions, as they must come
from a newer version of the format"
##########
core/src/main/java/org/apache/iceberg/TrackingBuilder.java:
##########
@@ -111,10 +115,27 @@ TrackingBuilder dvUpdated() {
return this;
}
+ /** Indicates that the column files list has been updated for the new
Tracking. */
+ TrackingBuilder columnFilesUpdated() {
+ Preconditions.checkState(
+ deletedPositions == null && replacedPositions == null,
+ "Cannot mark column files updated on a manifest entry
(deleted/replaced positions are set)");
+ this.latestColumnFileSnapshotId = newSnapshotId;
+ if (status == EntryStatus.EXISTING) {
+ this.status = EntryStatus.MODIFIED;
+ }
+ // Bumping 'dataSequenceNumber' to avoid having both equality deletes and
column files.
+ this.dataSequenceNumber = null;
Review Comment:
In a previous [google doc
discussion](https://docs.google.com/document/d/160-FizR6zOASMb86NycfgCm7cZbh6HK7FLcUW8_Xp-0/edit?disco=AAAB7BXsORk),
@pvary raised the question if we should just bump up the `dataSequenceNumber`
which captures the logical age of the row. Column file should materialize the
`_last_updated_sequence_number` for unmodified rows and leave the modified rows
with null value for inheritance. From row lineage perspective, bumping up
`dataSequenceNumber` is correct and simpler semantically.
data sequence number is only used for v2 equality and position delete
matching. it seems that we might be able to forbid writing new equality deletes
for v4 tables. I also remember some previous discussion on rewriting equality
delete and v2 position delete files when adding a new column file. With the
writer requirement, it is safe to just bump up the `dataSequenceNumber` here.
But the comment line is a bit confusing. I would write as following: `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.`
##########
core/src/main/java/org/apache/iceberg/ColumnFileStruct.java:
##########
@@ -0,0 +1,271 @@
+/*
+ * 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.FORMAT_VERSION,
+ ColumnFile.FIELD_IDS,
+ ColumnFile.LOCATION,
+ ColumnFile.FILE_FORMAT,
+ ColumnFile.FILE_SIZE_IN_BYTES,
+ ColumnFile.KEY_METADATA,
+ ColumnFile.SPLIT_OFFSETS);
+
+ private int formatVersion = -1;
+ private int[] fieldIds = null;
+ private String location = null;
+ private FileFormat fileFormat = null;
+ private long fileSizeInBytes = -1L;
+ private byte[] keyMetadata = null;
+ private long[] splitOffsets = null;
+
+ /** Used by internal readers to instantiate this class with a projection
schema. */
+ ColumnFileStruct(Types.StructType projection) {
+ super(BASE_TYPE, projection);
+ }
+
+ ColumnFileStruct(
+ int formatVersion,
+ List<Integer> fieldIds,
+ String location,
+ FileFormat fileFormat,
+ long fileSizeInBytes,
+ ByteBuffer keyMetadata,
+ List<Long> splitOffsets) {
+ super(BASE_TYPE.fields().size());
+ this.formatVersion = formatVersion;
+ this.fieldIds = ArrayUtil.toIntArray(fieldIds);
+ this.location = location;
+ this.fileFormat = fileFormat;
+ this.fileSizeInBytes = fileSizeInBytes;
+ this.keyMetadata = ByteBuffers.toByteArray(keyMetadata);
+ this.splitOffsets = ArrayUtil.toLongArray(splitOffsets);
+ }
+
+ /** Copy constructor. */
+ private ColumnFileStruct(ColumnFileStruct toCopy) {
+ super(toCopy);
+ this.formatVersion = toCopy.formatVersion;
+ this.fieldIds =
+ toCopy.fieldIds != null ? Arrays.copyOf(toCopy.fieldIds,
toCopy.fieldIds.length) : null;
+ this.location = toCopy.location;
+ this.fileFormat = toCopy.fileFormat;
+ this.fileSizeInBytes = toCopy.fileSizeInBytes;
+ this.keyMetadata =
+ toCopy.keyMetadata != null
+ ? Arrays.copyOf(toCopy.keyMetadata, toCopy.keyMetadata.length)
+ : null;
+ this.splitOffsets =
+ toCopy.splitOffsets != null
+ ? Arrays.copyOf(toCopy.splitOffsets, toCopy.splitOffsets.length)
+ : null;
+ }
+
+ /** Constructor for Java serialization. */
+ ColumnFileStruct() {
+ super(BASE_TYPE.fields().size());
+ }
+
+ @Override
+ public int formatVersion() {
+ return formatVersion;
+ }
+
+ @Override
+ public List<Integer> fieldIds() {
+ return fieldIds != null ? ArrayUtil.toUnmodifiableIntList(fieldIds) : null;
+ }
+
+ @Override
+ public String location() {
+ return location;
+ }
+
+ @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 List<Long> splitOffsets() {
+ return splitOffsets != null ?
ArrayUtil.toUnmodifiableLongList(splitOffsets) : null;
+ }
+
+ @Override
+ public ColumnFile copy() {
+ return new ColumnFileStruct(this);
+ }
+
+ @Override
+ protected <T> T internalGet(int pos, Class<T> javaClass) {
+ return javaClass.cast(getByPos(pos));
+ }
+
+ private Object getByPos(int pos) {
+ return switch (pos) {
+ case 0 -> formatVersion;
+ case 1 -> fieldIds();
+ case 2 -> location;
+ case 3 -> fileFormat != null ? fileFormat.toString() : null;
+ case 4 -> fileSizeInBytes;
+ case 5 -> keyMetadata();
+ case 6 -> splitOffsets();
+ default -> throw new UnsupportedOperationException("Unknown field
ordinal: " + pos);
+ };
+ }
+
+ @Override
+ @SuppressWarnings("unchecked")
+ protected <T> void internalSet(int pos, T value) {
+ switch (pos) {
+ case 0 -> this.formatVersion = (int) value;
+ case 1 -> this.fieldIds = ArrayUtil.toIntArray((List<Integer>) value);
+ // always coerce to String for Serializable
Review Comment:
nit: indention should move 2 spaces left.
--
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]