This is an automated email from the ASF dual-hosted git repository.
rdblue pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/iceberg.git
The following commit(s) were added to refs/heads/main by this push:
new c81a261852 Core: Project tracked-file partitions onto the resolved
spec (#18108)
c81a261852 is described below
commit c81a26185257cca14fdfaf7e9ff9d6db881a5939
Author: Anoop Johnson <[email protected]>
AuthorDate: Fri Oct 2 09:35:16 2026 -0700
Core: Project tracked-file partitions onto the resolved spec (#18108)
TrackedFileAdapters exposed the manifest's union-type partition tuple
unchanged while reporting the file's own spec ID. Residual evaluation and
partition-constant injection read the tuple by the spec's field ordinals, so
on partition-evolved tables a later spec's field (at a different position in
the union than in its own spec) was read from the wrong slot, silently
dropping rows or injecting null partition constants.
Project each file's partition onto its resolved spec's partition type before
exposing it, so partition() is always in the file's own spec order.
Co-authored-by: Isaac <[email protected]>
---
.../main/java/org/apache/iceberg/TrackedFile.java | 4 +-
.../java/org/apache/iceberg/TrackedFileStruct.java | 40 ++++--
.../java/org/apache/iceberg/V4ManifestReader.java | 71 +++++++---
.../org/apache/iceberg/TestTrackedFileStruct.java | 156 ++++++++++++++++++++-
.../org/apache/iceberg/TestV4ManifestReader.java | 132 +++++++++++++----
5 files changed, 338 insertions(+), 65 deletions(-)
diff --git a/core/src/main/java/org/apache/iceberg/TrackedFile.java
b/core/src/main/java/org/apache/iceberg/TrackedFile.java
index a02d50cf35..e2db837c48 100644
--- a/core/src/main/java/org/apache/iceberg/TrackedFile.java
+++ b/core/src/main/java/org/apache/iceberg/TrackedFile.java
@@ -169,7 +169,9 @@ interface TrackedFile {
/** Returns the ID of the partition spec used to partition this file, or
null. */
Integer specId();
- /** Returns partition for this file as a {@link StructLike}, or null. */
+ /**
+ * Returns the partition for this file as a struct with the partition spec's
output type, or null.
+ */
StructLike partition();
/** Returns the content stats for this entry. */
diff --git a/core/src/main/java/org/apache/iceberg/TrackedFileStruct.java
b/core/src/main/java/org/apache/iceberg/TrackedFileStruct.java
index 7a33396824..fffd822a56 100644
--- a/core/src/main/java/org/apache/iceberg/TrackedFileStruct.java
+++ b/core/src/main/java/org/apache/iceberg/TrackedFileStruct.java
@@ -24,11 +24,13 @@ import java.util.Arrays;
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.types.Type;
import org.apache.iceberg.types.Types;
import org.apache.iceberg.util.ArrayUtil;
import org.apache.iceberg.util.ByteBuffers;
+import org.apache.iceberg.util.StructLikeUtil;
+import org.apache.iceberg.util.StructProjection;
/** Mutable {@link StructLike} implementation of {@link TrackedFile}. */
class TrackedFileStruct extends SupportsIndexProjection implements
TrackedFile, Serializable {
@@ -68,7 +70,7 @@ class TrackedFileStruct extends SupportsIndexProjection
implements TrackedFile,
private Tracking tracking = null;
private long recordCount = -1L;
private long fileSizeInBytes = -1L;
- private PartitionData partitionData = null;
+ private StructLike partition = null;
// optional fields
private Integer specId = null;
@@ -80,15 +82,11 @@ class TrackedFileStruct extends SupportsIndexProjection
implements TrackedFile,
private long[] splitOffsets = null;
private int[] equalityIds = null;
+ private transient StructProjection partitionProjection = null;
+
/** Used by internal readers to instantiate this class with a projection
schema. */
TrackedFileStruct(Types.StructType projection) {
super(BASE_TYPE, projection);
- // partition type may be null if the field was not projected, or unknown
for unpartitioned
- // manifests
- Type partType = projection.fieldType(TrackedFile.PARTITION_NAME);
- if (partType != null && partType.isStructType()) {
- this.partitionData = new PartitionData(partType.asStructType());
- }
}
/** No-projection constructor for direct construction. */
@@ -122,7 +120,7 @@ class TrackedFileStruct extends SupportsIndexProjection
implements TrackedFile,
this.recordCount = recordCount;
this.fileSizeInBytes = fileSizeInBytes;
this.specId = specId;
- this.partitionData = partition;
+ this.partition = partition;
this.contentStats = contentStats;
this.sortOrderId = sortOrderId;
this.deletionVector = deletionVector;
@@ -142,7 +140,7 @@ class TrackedFileStruct extends SupportsIndexProjection
implements TrackedFile,
this.recordCount = toCopy.recordCount;
this.fileSizeInBytes = toCopy.fileSizeInBytes;
this.specId = toCopy.specId;
- this.partitionData = toCopy.partitionData != null ?
toCopy.partitionData.copy() : null;
+ this.partition = toCopy.partition == null ? null :
StructLikeUtil.copy(toCopy.partition());
this.tracking = toCopy.tracking != null ? toCopy.tracking.copy() : null;
this.sortOrderId = toCopy.sortOrderId;
this.deletionVector = toCopy.deletionVector != null ?
toCopy.deletionVector.copy() : null;
@@ -210,6 +208,15 @@ class TrackedFileStruct extends SupportsIndexProjection
implements TrackedFile,
return fileSizeInBytes;
}
+ void setPartitionProjection(StructProjection projection) {
+ this.partitionProjection = projection;
+ }
+
+ void clearPartition() {
+ this.partition = null;
+ this.partitionProjection = null;
+ }
+
@Override
public Integer specId() {
return specId;
@@ -217,7 +224,12 @@ class TrackedFileStruct extends SupportsIndexProjection
implements TrackedFile,
@Override
public StructLike partition() {
- return partitionData;
+ if (partition == null) {
+ ValidationException.check(specId == null, "Missing partition for spec
%s", specId);
+ return null;
+ }
+
+ return partitionProjection != null ? partitionProjection.wrap(partition) :
partition;
}
@Override
@@ -280,7 +292,7 @@ class TrackedFileStruct extends SupportsIndexProjection
implements TrackedFile,
case 5 -> recordCount;
case 6 -> fileSizeInBytes;
case 7 -> specId;
- case 8 -> partitionData;
+ case 8 -> partition;
case 9 -> contentStats;
case 10 -> sortOrderId;
case 11 -> deletionVector;
@@ -305,7 +317,7 @@ class TrackedFileStruct extends SupportsIndexProjection
implements TrackedFile,
case 5 -> this.recordCount = (long) value;
case 6 -> this.fileSizeInBytes = (long) value;
case 7 -> this.specId = (Integer) value;
- case 8 -> this.partitionData = (PartitionData) value;
+ case 8 -> this.partition = (StructLike) value;
case 9 -> this.contentStats = (ContentStats) value;
case 10 -> this.sortOrderId = (Integer) value;
case 11 -> this.deletionVector = (DeletionVector) value;
@@ -329,7 +341,7 @@ class TrackedFileStruct extends SupportsIndexProjection
implements TrackedFile,
.add("record_count", recordCount)
.add("file_size_in_bytes", fileSizeInBytes)
.add("spec_id", specId())
- .add("partition", partitionData)
+ .add("partition", partition)
.add("tracking", tracking)
.add("content_stats", contentStats)
.add("sort_order_id", sortOrderId)
diff --git a/core/src/main/java/org/apache/iceberg/V4ManifestReader.java
b/core/src/main/java/org/apache/iceberg/V4ManifestReader.java
index 09515b3f11..4a207dacd9 100644
--- a/core/src/main/java/org/apache/iceberg/V4ManifestReader.java
+++ b/core/src/main/java/org/apache/iceberg/V4ManifestReader.java
@@ -44,7 +44,6 @@ import org.apache.iceberg.types.TypeUtil;
import org.apache.iceberg.types.Types;
import org.apache.iceberg.util.ArrayUtil;
import org.apache.iceberg.util.LocationUtil;
-import org.apache.iceberg.util.Pair;
import org.apache.iceberg.util.StructProjection;
/** Reader that reads a v4+ manifest file as {@link TrackedFile}s. */
@@ -73,7 +72,8 @@ class V4ManifestReader extends CloseableGroup implements
CloseableIterable<Track
private final Schema readSchema;
private final String tableLocation;
private final InclusiveStatsEvaluator statsFilter;
- private final Map<Integer, Pair<Evaluator, StructProjection>>
partitionFilters; // by spec ID
+ private final Map<Integer, Evaluator> partitionFilters; // by spec ID
+ private final Map<Integer, StructProjection> partitionProjections; // by
spec ID
private final Set<Integer> requestedStatsFieldIds;
private final boolean isUncommitted;
private final ScanMetrics scanMetrics;
@@ -86,7 +86,8 @@ class V4ManifestReader extends CloseableGroup implements
CloseableIterable<Track
Schema readSchema,
String tableLocation,
InclusiveStatsEvaluator statsFilter,
- Map<Integer, Pair<Evaluator, StructProjection>> partitionFilters,
+ Map<Integer, Evaluator> partitionFilters,
+ Map<Integer, StructProjection> partitionProjections,
Set<Integer> requestedStatsFieldIds,
boolean isUncommitted,
ScanMetrics scanMetrics) {
@@ -97,6 +98,7 @@ class V4ManifestReader extends CloseableGroup implements
CloseableIterable<Track
this.tableLocation = tableLocation;
this.statsFilter = statsFilter;
this.partitionFilters = partitionFilters;
+ this.partitionProjections = partitionProjections;
this.requestedStatsFieldIds = requestedStatsFieldIds;
this.isUncommitted = isUncommitted;
this.scanMetrics = scanMetrics;
@@ -117,6 +119,8 @@ class V4ManifestReader extends CloseableGroup implements
CloseableIterable<Track
CloseableIterable<TrackedFile> files =
CloseableIterable.transform(open(), this::applyInheritance);
+ files = CloseableIterable.transform(files, this::projectPartition);
+
if (dv != null) {
files =
CloseableIterable.filter(
@@ -161,6 +165,18 @@ class V4ManifestReader extends CloseableGroup implements
CloseableIterable<Track
return file;
}
+ private TrackedFile projectPartition(TrackedFile file) {
+ Integer specId = file.specId();
+ TrackedFileStruct struct = (TrackedFileStruct) file;
+ if (specId != null && !partitionProjections.containsKey(specId)) {
+ struct.clearPartition();
+ } else {
+ struct.setPartitionProjection(partitionProjections.get(specId));
+ }
+
+ return file;
+ }
+
private boolean isDeletedByMDV(TrackedFile file) {
return dv.isSet(Math.toIntExact(file.tracking().manifestPos()));
}
@@ -172,15 +188,13 @@ class V4ManifestReader extends CloseableGroup implements
CloseableIterable<Track
return true;
}
- Pair<Evaluator, StructProjection> partitionFilter =
partitionFilters.get(specId);
- if (partitionFilter == null) {
+ Evaluator evaluator = partitionFilters.get(specId);
+ if (evaluator == null) {
// the row filter does not project to a partition filter for this spec
return true;
}
- Evaluator evaluator = partitionFilter.first();
- StructProjection projection = partitionFilter.second();
- boolean matches = evaluator.eval(projection.wrap(trackedFile.partition()));
+ boolean matches = evaluator.eval(trackedFile.partition());
if (!matches) {
incrementSkipCount(trackedFile);
}
@@ -408,7 +422,7 @@ class V4ManifestReader extends CloseableGroup implements
CloseableIterable<Track
requestedStatsFieldIds = ImmutableSet.of();
}
- Map<Integer, Pair<Evaluator, StructProjection>> partitionFilters =
projectFilters();
+ Map<Integer, Evaluator> partitionFilters = projectFilters();
return new V4ManifestReader(
manifest,
@@ -417,6 +431,7 @@ class V4ManifestReader extends CloseableGroup implements
CloseableIterable<Track
tableLocation,
statsFilter(),
partitionFilters,
+ partitionProjections(),
requestedStatsFieldIds,
isUncommitted,
scanMetrics);
@@ -430,21 +445,33 @@ class V4ManifestReader extends CloseableGroup implements
CloseableIterable<Track
}
}
- private Map<Integer, Pair<Evaluator, StructProjection>> projectFilters() {
- Map<Integer, Pair<Evaluator, StructProjection>> evaluatorAndProjections
= Maps.newHashMap();
+ private Map<Integer, Evaluator> projectFilters() {
+ Map<Integer, Evaluator> evaluators = Maps.newHashMap();
if (rowFilter != Expressions.alwaysTrue() &&
!unionPartitionType.fields().isEmpty()) {
for (PartitionSpec spec : specsById.values()) {
Expression partFilter = Projections.inclusive(spec,
caseSensitive).project(rowFilter);
if (partFilter != Expressions.alwaysTrue()) {
- Evaluator evaluator = new Evaluator(spec.partitionType(),
partFilter, caseSensitive);
- StructProjection projection =
- StructProjection.create(unionPartitionType,
spec.partitionType());
- evaluatorAndProjections.put(spec.specId(), Pair.of(evaluator,
projection));
+ evaluators.put(
+ spec.specId(), new Evaluator(spec.partitionType(), partFilter,
caseSensitive));
}
}
}
- return evaluatorAndProjections;
+ return evaluators;
+ }
+
+ private Map<Integer, StructProjection> partitionProjections() {
+ Map<Integer, StructProjection> projections = Maps.newHashMap();
+ for (PartitionSpec spec : specsById.values()) {
+ Types.StructType partitionType = spec.partitionType();
+ StructProjection projection =
+ partitionType.equals(unionPartitionType)
+ ? null
+ : StructProjection.create(unionPartitionType, partitionType);
+ projections.put(spec.specId(), projection);
+ }
+
+ return projections;
}
private Schema readSchema(boolean includePartition) {
@@ -465,7 +492,9 @@ class V4ManifestReader extends CloseableGroup implements
CloseableIterable<Track
if (requestedProjection != null) {
return RestoreColumns.restore(
- tableManifestSchema, requestedProjection,
idsToRestore(includePartition));
+ tableManifestSchema,
+ requestedProjection,
+ idsToRestore(projectsPartition(includePartition,
requestedProjection)));
}
if (requestedColumns != null) {
@@ -475,12 +504,18 @@ class V4ManifestReader extends CloseableGroup implements
CloseableIterable<Track
: tableManifestSchema.caseInsensitiveSelect(requestedColumns);
return RestoreColumns.restore(
- tableManifestSchema, projection, idsToRestore(includePartition));
+ tableManifestSchema,
+ projection,
+ idsToRestore(projectsPartition(includePartition, projection)));
}
return tableManifestSchema;
}
+ private boolean projectsPartition(boolean includePartition, Schema
projection) {
+ return includePartition ||
projection.findField(TrackedFile.PARTITION_ID) != null;
+ }
+
/** Return a set of manifest field IDs that should be projected. */
private Set<Integer> idsToRestore(boolean includePartition) {
Set<Integer> ids = Sets.newHashSet(REQUIRED_COLUMN_IDS);
diff --git a/core/src/test/java/org/apache/iceberg/TestTrackedFileStruct.java
b/core/src/test/java/org/apache/iceberg/TestTrackedFileStruct.java
index 14265c91f6..f5edda5567 100644
--- a/core/src/test/java/org/apache/iceberg/TestTrackedFileStruct.java
+++ b/core/src/test/java/org/apache/iceberg/TestTrackedFileStruct.java
@@ -19,13 +19,21 @@
package org.apache.iceberg;
import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
import java.nio.ByteBuffer;
+import java.util.Comparator;
import java.util.List;
+import java.util.Map;
import org.apache.iceberg.TestHelpers.RoundTripSerializer;
+import org.apache.iceberg.exceptions.ValidationException;
import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList;
+import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap;
import org.apache.iceberg.relocated.com.google.common.collect.ImmutableSet;
+import org.apache.iceberg.transforms.Transforms;
+import org.apache.iceberg.types.Comparators;
import org.apache.iceberg.types.Types;
+import org.apache.iceberg.util.StructProjection;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.MethodSource;
@@ -41,7 +49,6 @@ class TestTrackedFileStruct {
private static final Tracking TRACKING_COPY = Mockito.mock(Tracking.class);
private static final PartitionData PARTITION =
Mockito.mock(PartitionData.class);
- private static final PartitionData PARTITION_COPY =
Mockito.mock(PartitionData.class);
private static final ContentStats CONTENT_STATS =
Mockito.mock(ContentStats.class);
private static final ContentStats CONTENT_STATS_COPY =
Mockito.mock(ContentStats.class);
@@ -54,7 +61,6 @@ class TestTrackedFileStruct {
static {
Mockito.when(TRACKING.copy()).thenReturn(TRACKING_COPY);
- Mockito.when(PARTITION.copy()).thenReturn(PARTITION_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);
@@ -217,7 +223,7 @@ class TestTrackedFileStruct {
assertThat(copy.keyMetadata()).isEqualTo(ByteBuffer.wrap(new byte[] {1, 2,
3}));
assertThat(copy.splitOffsets()).containsExactly(100L, 200L);
assertThat(copy.equalityIds()).containsExactly(1, 2, 3);
- assertThat(copy.partition()).isSameAs(PARTITION_COPY);
+ assertThat(copy.partition()).isNotSameAs(PARTITION);
// mutable fields are deep-copied, not shared with the original
assertThat(copy.keyMetadata()).isNotSameAs(file.keyMetadata());
@@ -266,7 +272,7 @@ class TestTrackedFileStruct {
assertThat(copy.keyMetadata()).isEqualTo(ByteBuffer.wrap(new byte[] {1, 2,
3}));
assertThat(copy.splitOffsets()).containsExactly(100L, 200L);
assertThat(copy.equalityIds()).containsExactly(1, 2, 3);
- assertThat(copy.partition()).isSameAs(PARTITION_COPY);
+ assertThat(copy.partition()).isNotSameAs(PARTITION);
// mutable fields are deep-copied, not shared with the original
assertThat(copy.keyMetadata()).isNotSameAs(file.keyMetadata());
@@ -316,7 +322,7 @@ class TestTrackedFileStruct {
assertThat(copy.keyMetadata()).isEqualTo(ByteBuffer.wrap(new byte[] {1, 2,
3}));
assertThat(copy.splitOffsets()).containsExactly(100L, 200L);
assertThat(copy.equalityIds()).containsExactly(1, 2, 3);
- assertThat(copy.partition()).isSameAs(PARTITION_COPY);
+ assertThat(copy.partition()).isNotSameAs(PARTITION);
// mutable fields are deep-copied, not shared with the original
assertThat(copy.keyMetadata()).isNotSameAs(file.keyMetadata());
@@ -340,6 +346,78 @@ class TestTrackedFileStruct {
assertThat(file.get(1, Long.class)).isEqualTo(1024L);
}
+ @Test
+ void partitionIsProjectedToResolvedSpec() {
+ // a table whose partitioning evolved from id to category
+ Schema schema =
+ new Schema(
+ Types.NestedField.required(1, "id", Types.IntegerType.get()),
+ Types.NestedField.required(2, "category", Types.StringType.get()));
+ PartitionSpec idSpec =
+ PartitionSpec.builderFor(schema)
+ .withSpecId(0)
+ .add(1, 1000, "id", Transforms.identity())
+ .build();
+ PartitionSpec categorySpec =
+ PartitionSpec.builderFor(schema)
+ .withSpecId(1)
+ .add(2, 1001, "category", Transforms.identity())
+ .build();
+ Map<Integer, PartitionSpec> specsById =
+ ImmutableMap.of(idSpec.specId(), idSpec, categorySpec.specId(),
categorySpec);
+
+ // the manifest stores partitions in the union type, where category sits
after id
+ Types.StructType unionType =
Partitioning.unionPartitionTypes(specsById.values());
+ int categoryUnionPos =
unionType.fields().indexOf(unionType.field("category"));
+ PartitionData unionPartition = new PartitionData(unionType);
+ unionPartition.set(categoryUnionPos, "books");
+
+ TrackedFileStruct file = trackedFile(categorySpec.specId(),
unionPartition);
+ file.setPartitionProjection(StructProjection.create(unionType,
categorySpec.partitionType()));
+
+ Comparator<StructLike> comparator =
Comparators.forType(categorySpec.partitionType());
+ PartitionData expected = new PartitionData(categorySpec.partitionType());
+ expected.set(0, "books");
+
+ // category is at position 1 in the union but position 0 in categorySpec;
reading by the spec's
+ // ordinal must return category, not id (null)
+
assertThat(file.partition()).usingComparator(comparator).isEqualTo(expected);
+
+ StructLike copyPartition = file.copy().partition();
+ unionPartition.set(categoryUnionPos, "changed");
+ assertThat(copyPartition).usingComparator(comparator).isEqualTo(expected);
+ }
+
+ @Test
+ void partitionReturnedAsIsWhenNoProjection() {
+ PartitionData partition =
+ new PartitionData(
+ Types.StructType.of(
+ Types.NestedField.required(1000, "category",
Types.StringType.get())));
+ partition.set(0, "music");
+
+ TrackedFileStruct file = trackedFile(1, partition);
+
+ // no projection is set, so partition() returns the stored tuple unchanged
+ assertThat(file.partition()).isSameAs(partition);
+ }
+
+ @Test
+ void partitionFailsWhenSpecIsSetButPartitionIsMissing() {
+ TrackedFileStruct file = trackedFile(1, null);
+
+ assertThatThrownBy(file::partition)
+ .isInstanceOf(ValidationException.class)
+ .hasMessage("Missing partition for spec 1");
+ }
+
+ @Test
+ void partitionIsNullWhenThereIsNoSpec() {
+ TrackedFileStruct file = trackedFile(null, null);
+
+ assertThat(file.partition()).isNull();
+ }
+
@Test
void structLikeSize() {
TrackedFileStruct file = new TrackedFileStruct();
@@ -349,6 +427,11 @@ class TestTrackedFileStruct {
@ParameterizedTest
@MethodSource("org.apache.iceberg.TestHelpers#serializers")
void serializationRoundTrip(RoundTripSerializer<TrackedFileStruct>
serializer) throws Exception {
+ Types.StructType partitionType =
+ Types.StructType.of(Types.NestedField.required(1000, "id",
Types.IntegerType.get()));
+ PartitionData partition = new PartitionData(partitionType);
+ partition.set(0, 7);
+
TrackedFileStruct file =
new TrackedFileStruct(
null, // TrackingStruct has its own serialization tests
@@ -359,7 +442,7 @@ class TestTrackedFileStruct {
100L,
1024L,
7,
- null, // PartitionData has its own serialization tests
+ partition,
null,
1,
null, // DeletionVector has its own serialization tests
@@ -375,7 +458,9 @@ class TestTrackedFileStruct {
assertThat(deserialized.formatVersion()).isEqualTo(FORMAT_VERSION_V4);
assertThat(deserialized.location()).isEqualTo("s3://bucket/data/file.parquet");
assertThat(deserialized.fileFormat()).isEqualTo(FileFormat.PARQUET);
- assertThat(deserialized.partition()).isNull();
+ assertThat(deserialized.partition())
+ .usingComparator(Comparators.forType(partitionType))
+ .isEqualTo(partition);
assertThat(deserialized.recordCount()).isEqualTo(100L);
assertThat(deserialized.fileSizeInBytes()).isEqualTo(1024L);
assertThat(deserialized.specId()).isEqualTo(7);
@@ -387,6 +472,63 @@ class TestTrackedFileStruct {
assertThat(deserialized.equalityIds()).containsExactly(1, 2, 3);
}
+ @ParameterizedTest
+ @MethodSource("org.apache.iceberg.TestHelpers#serializers")
+ void partitionProjectionSurvivesSerializationAfterCopy(
+ RoundTripSerializer<TrackedFileStruct> serializer) throws Exception {
+ Schema schema =
+ new Schema(
+ Types.NestedField.required(1, "id", Types.IntegerType.get()),
+ Types.NestedField.required(2, "category", Types.StringType.get()));
+ PartitionSpec idSpec =
+ PartitionSpec.builderFor(schema)
+ .withSpecId(0)
+ .add(1, 1000, "id", Transforms.identity())
+ .build();
+ PartitionSpec categorySpec =
+ PartitionSpec.builderFor(schema)
+ .withSpecId(1)
+ .add(2, 1001, "category", Transforms.identity())
+ .build();
+ Map<Integer, PartitionSpec> specsById =
+ ImmutableMap.of(idSpec.specId(), idSpec, categorySpec.specId(),
categorySpec);
+
+ Types.StructType unionType =
Partitioning.unionPartitionTypes(specsById.values());
+ PartitionData unionPartition = new PartitionData(unionType);
+
unionPartition.set(unionType.fields().indexOf(unionType.field("category")),
"books");
+
+ TrackedFileStruct file = trackedFile(categorySpec.specId(),
unionPartition);
+ file.setPartitionProjection(StructProjection.create(unionType,
categorySpec.partitionType()));
+
+ PartitionData expected = new PartitionData(categorySpec.partitionType());
+ expected.set(0, "books");
+
+ TrackedFileStruct deserialized = serializer.apply((TrackedFileStruct)
file.copy());
+ assertThat(deserialized.partition())
+ .usingComparator(Comparators.forType(categorySpec.partitionType()))
+ .isEqualTo(expected);
+ }
+
+ private static TrackedFileStruct trackedFile(Integer specId, PartitionData
partition) {
+ return new TrackedFileStruct(
+ null, // tracking
+ FileContent.DATA,
+ FORMAT_VERSION_V4,
+ "s3://bucket/file.parquet",
+ FileFormat.PARQUET,
+ 100L, // recordCount
+ 1024L, // fileSizeInBytes
+ specId,
+ partition,
+ null, // contentStats
+ null, // sortOrderId
+ null, // deletionVector
+ null, // manifestInfo
+ null, // keyMetadata
+ null, // splitOffsets
+ null); // equalityIds
+ }
+
private static int pos(String fieldName) {
for (int i = 0; i < DEFAULT_FIELDS.size(); i += 1) {
if (DEFAULT_FIELDS.get(i).name().equals(fieldName)) {
diff --git a/core/src/test/java/org/apache/iceberg/TestV4ManifestReader.java
b/core/src/test/java/org/apache/iceberg/TestV4ManifestReader.java
index 7b8e41a534..9d1d81d3c6 100644
--- a/core/src/test/java/org/apache/iceberg/TestV4ManifestReader.java
+++ b/core/src/test/java/org/apache/iceberg/TestV4ManifestReader.java
@@ -1495,40 +1495,118 @@ class TestV4ManifestReader {
.build();
Map<Integer, PartitionSpec> specsById =
ImmutableMap.of(idSpec.specId(), idSpec, dataSpec.specId(), dataSpec);
-
- PartitionData dataPartitionX = partition(dataSpec, "x");
- TrackedFile dataPartitionedFile =
- dataFileWithoutStats(
- "s3://bucket/table/data=x/file-c.parquet", dataSpec.specId(),
dataPartitionX);
-
- ManifestFile idPartitionedManifest =
- writeManifest(format, idSpec.partitionType(), ImmutableList.of(FILE_A,
FILE_B));
- ManifestFile dataPartitionedManifest =
- writeManifest(format, dataSpec.partitionType(), dataPartitionedFile);
-
- List<TrackedFile> files = Lists.newArrayList();
- files.addAll(
- read(
- V4ManifestReader.builder(idPartitionedManifest, IO, TABLE_SCHEMA,
specsById)
- .metricsConfig(METRICS_CONFIG)));
- files.addAll(
- read(
- V4ManifestReader.builder(dataPartitionedManifest, IO,
TABLE_SCHEMA, specsById)
- .metricsConfig(METRICS_CONFIG)));
-
Types.StructType unionType =
Partitioning.unionPartitionTypes(specsById.values());
+ int idPos = unionType.fields().indexOf(unionType.field("id"));
+ int dataPos = unionType.fields().indexOf(unionType.field("data"));
+
+ // a mixed-spec manifest stores every file's partition in the union type
+ PartitionData idOne = new PartitionData(unionType);
+ idOne.set(idPos, 1);
+ PartitionData idTwo = new PartitionData(unionType);
+ idTwo.set(idPos, 2);
+ PartitionData dataX = new PartitionData(unionType);
+ dataX.set(dataPos, "x");
+
+ TrackedFile idFileKept =
+ dataFileWithoutStats("s3://bucket/table/id=1/file-a.parquet",
idSpec.specId(), idOne);
+ TrackedFile idFilePruned =
+ dataFileWithoutStats("s3://bucket/table/id=2/file-b.parquet",
idSpec.specId(), idTwo);
+ TrackedFile dataFile =
+ dataFileWithoutStats("s3://bucket/table/data=x/file-c.parquet",
dataSpec.specId(), dataX);
- ManifestFile manifest = writeManifest(format, unionType, files);
+ ManifestFile manifest =
+ writeManifest(format, unionType, ImmutableList.of(idFileKept,
idFilePruned, dataFile));
V4ManifestReader.Builder builder =
V4ManifestReader.builder(manifest, IO, TABLE_SCHEMA, specsById)
.filter(Expressions.equal("id", 1))
.metricsConfig(METRICS_CONFIG);
- // the comparator is built for ID partitioning, so only check the location
assertThat(read(builder))
.extracting(TrackedFile::location)
- .containsExactlyInAnyOrder(FILE_A.location(),
dataPartitionedFile.location());
+ .containsExactlyInAnyOrder(idFileKept.location(), dataFile.location());
+ }
+
+ @ParameterizedTest
+ @FieldSource("MANIFEST_FORMATS")
+ public void narrowPartitionProjectionReadsFullUnionTuple(FileFormat format)
throws IOException {
+ PartitionSpec idSpec =
+ PartitionSpec.builderFor(TABLE_SCHEMA)
+ .withSpecId(1)
+ .add(1, 1000, "id", Transforms.identity())
+ .build();
+ PartitionSpec dataSpec =
+ PartitionSpec.builderFor(TABLE_SCHEMA)
+ .withSpecId(2)
+ .add(2, 1001, "data", Transforms.identity())
+ .build();
+ Map<Integer, PartitionSpec> specsById =
+ ImmutableMap.of(idSpec.specId(), idSpec, dataSpec.specId(), dataSpec);
+ Types.StructType unionType =
Partitioning.unionPartitionTypes(specsById.values());
+
+ PartitionData unionPartition = new PartitionData(unionType);
+ unionPartition.set(unionType.fields().indexOf(unionType.field("data")),
"x");
+ TrackedFile file =
+ dataFileWithoutStats(
+ "s3://bucket/table/data=x/file.parquet", dataSpec.specId(),
unionPartition);
+ ManifestFile manifest = writeManifest(format, unionType,
ImmutableList.of(file));
+
+ TrackedFile actual =
+ readOne(
+ V4ManifestReader.builder(manifest, IO, TABLE_SCHEMA, specsById)
+ .metricsConfig(METRICS_CONFIG)
+ .select("partition.id"));
+ assertThat(actual.partition().get(0, CharSequence.class)).hasToString("x");
+ }
+
+ @ParameterizedTest
+ @FieldSource("MANIFEST_FORMATS")
+ public void unknownSpecPartitionFails(FileFormat format) throws IOException {
+ PartitionSpec idSpec =
+ PartitionSpec.builderFor(TABLE_SCHEMA)
+ .withSpecId(1)
+ .add(1, 1000, "id", Transforms.identity())
+ .build();
+ PartitionSpec dataSpec =
+ PartitionSpec.builderFor(TABLE_SCHEMA)
+ .withSpecId(2)
+ .add(2, 1001, "data", Transforms.identity())
+ .build();
+ Map<Integer, PartitionSpec> specsById =
+ ImmutableMap.of(idSpec.specId(), idSpec, dataSpec.specId(), dataSpec);
+ Types.StructType unionType =
Partitioning.unionPartitionTypes(specsById.values());
+ int dataPos = unionType.fields().indexOf(unionType.field("data"));
+
+ PartitionData knownPartition = new PartitionData(unionType);
+ knownPartition.set(unionType.fields().indexOf(unionType.field("id")), 7);
+ TrackedFile known =
+ dataFileWithoutStats("s3://bucket/table/known.parquet",
idSpec.specId(), knownPartition);
+
+ // spec id 5 is not in specsById; it follows a known-spec file in the same
(reused) reader
+ PartitionData unknownPartition = new PartitionData(unionType);
+ unknownPartition.set(dataPos, "x");
+ TrackedFile unknown =
+ dataFileWithoutStats("s3://bucket/table/unknown.parquet", 5,
unknownPartition);
+
+ ManifestFile manifest = writeManifest(format, unionType,
ImmutableList.of(known, unknown));
+
+ List<TrackedFile> files =
+ read(
+ V4ManifestReader.builder(manifest, IO, TABLE_SCHEMA, specsById)
+ .metricsConfig(METRICS_CONFIG));
+
+ TrackedFile knownActual =
+ files.stream().filter(f ->
Integer.valueOf(1).equals(f.specId())).findFirst().orElseThrow();
+ assertThat(knownActual.partition().get(0, Integer.class))
+ .as("known spec's partition is projected to its output type")
+ .isEqualTo(7);
+
+ TrackedFile unknownActual =
+ files.stream().filter(f ->
Integer.valueOf(5).equals(f.specId())).findFirst().orElseThrow();
+ assertThatThrownBy(unknownActual::partition)
+ .as("partition cannot be projected to an unknown spec's output type")
+ .isInstanceOf(ValidationException.class)
+ .hasMessage("Missing partition for spec 5");
}
@ParameterizedTest
@@ -1548,7 +1626,11 @@ class TestV4ManifestReader {
TrackedFile actual = readOne(builder);
- assertThat(actual).usingComparator(FILE_COMPARATOR).isEqualTo(file);
+ assertThat(actual.location()).isEqualTo(file.location());
+ assertThatThrownBy(actual::partition)
+ .as("unknown spec's partition cannot be projected to its output type")
+ .isInstanceOf(ValidationException.class)
+ .hasMessage("Missing partition for spec 5");
}
@ParameterizedTest