anoopj commented on code in PR #18108:
URL: https://github.com/apache/iceberg/pull/18108#discussion_r4150631673
##########
core/src/main/java/org/apache/iceberg/TrackedFileStruct.java:
##########
@@ -80,14 +82,16 @@ 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());
+ this.partition = new PartitionData(partType.asStructType());
Review Comment:
You're right - It was unnecessary because reader already registers
`setCustomType` for `PartitionData`. Ive removed it.
##########
core/src/test/java/org/apache/iceberg/TestTrackedFileStruct.java:
##########
@@ -340,6 +342,78 @@ void projectedStructLike() {
assertThat(file.get(1, Long.class)).isEqualTo(1024L);
}
+ @Test
+ void partitionIsProjectedOntoResolvedSpec() {
+ // 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()));
+
+ // 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().get(0,
CharSequence.class)).hasToString("books");
+
+ StructLike copyPartition = file.copy().partition();
+ unionPartition.set(categoryUnionPos, "changed");
+ assertThat(copyPartition.get(0, CharSequence.class)).hasToString("books");
+ }
+
+ @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);
+ }
+
+ private static TrackedFileStruct trackedFile(int specId, PartitionData
partition) {
Review Comment:
Moved it to below.
##########
core/src/main/java/org/apache/iceberg/TrackedFileStruct.java:
##########
@@ -329,7 +341,7 @@ public String toString() {
.add("record_count", recordCount)
.add("file_size_in_bytes", fileSizeInBytes)
.add("spec_id", specId())
- .add("partition", partitionData)
+ .add("partition", partition)
Review Comment:
Will do.
##########
core/src/test/java/org/apache/iceberg/TestV4ManifestReader.java:
##########
@@ -1495,40 +1495,109 @@ public void
partitionFilterWithMultipleSpecs(FileFormat format) throws IOExcepti
.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 unknownSpecPartitionIsNotProjected(FileFormat format) throws
IOException {
Review Comment:
I agree. I went with `null` so that an unknown spec can't be projected to
its output type, so `partition()` returns null instead of the raw union tuple.
Also renamed the tests.
##########
core/src/test/java/org/apache/iceberg/TestTrackedFileStruct.java:
##########
@@ -387,6 +461,38 @@ void
serializationRoundTrip(RoundTripSerializer<TrackedFileStruct> serializer) t
assertThat(deserialized.equalityIds()).containsExactly(1, 2, 3);
}
+ @ParameterizedTest
+ @MethodSource("org.apache.iceberg.TestHelpers#serializers")
+ void
partitionProjectionSurvivesSerialization(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()));
+
+ TrackedFileStruct deserialized = serializer.apply((TrackedFileStruct)
file.copy());
Review Comment:
Done.
##########
core/src/test/java/org/apache/iceberg/TestTrackedFileStruct.java:
##########
@@ -340,6 +342,78 @@ void projectedStructLike() {
assertThat(file.get(1, Long.class)).isEqualTo(1024L);
}
+ @Test
+ void partitionIsProjectedOntoResolvedSpec() {
Review Comment:
Done.
##########
core/src/test/java/org/apache/iceberg/TestTrackedFileStruct.java:
##########
@@ -340,6 +342,78 @@ void projectedStructLike() {
assertThat(file.get(1, Long.class)).isEqualTo(1024L);
}
+ @Test
+ void partitionIsProjectedOntoResolvedSpec() {
+ // 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()));
+
+ // 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().get(0,
CharSequence.class)).hasToString("books");
Review Comment:
Thanks, done.
--
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]