stevenzwu commented on code in PR #16936: URL: https://github.com/apache/iceberg/pull/16936#discussion_r4226328932
########## core/src/test/java/org/apache/iceberg/TestMapBackedContentStats.java: ########## @@ -0,0 +1,382 @@ +/* + * 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.math.BigDecimal; +import java.nio.ByteBuffer; +import java.util.Comparator; +import java.util.List; +import java.util.Map; +import java.util.UUID; +import java.util.stream.Stream; +import org.apache.iceberg.geospatial.GeospatialBound; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; +import org.apache.iceberg.types.Comparators; +import org.apache.iceberg.types.Conversions; +import org.apache.iceberg.types.Type; +import org.apache.iceberg.types.Types; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.MethodSource; + +class TestMapBackedContentStats { + + private static final Schema SCHEMA = + new Schema( + Types.NestedField.required(1, "id", Types.IntegerType.get()), + Types.NestedField.optional(2, "score", Types.FloatType.get()), + Types.NestedField.optional(3, "ts", Types.LongType.get()), + Types.NestedField.optional(4, "name", Types.StringType.get()), + Types.NestedField.optional(5, "flag", Types.BooleanType.get())); + + private static final PartitionData EMPTY_PARTITION = + new PartitionData(PartitionSpec.unpartitioned().partitionType()); + + /** + * Stats on fields 1-4; field 5 is absent. Field 3 has no value-count entry, so valueCount() + * throws on unboxing. + */ + private static final DataFile FILE_WITH_STATS = + dataFile( + 100L, + ImmutableMap.of(1, 100L, 2, 100L, 4, 100L), + ImmutableMap.of(2, 5L, 3, 1L, 4, 2L), + ImmutableMap.of(2, 3L), + ImmutableMap.of( + 1, buf(Types.IntegerType.get(), 1), + 2, buf(Types.FloatType.get(), 1.5f), + 3, buf(Types.LongType.get(), 100L), + 4, buf(Types.StringType.get(), "aaa")), + ImmutableMap.of( + 1, buf(Types.IntegerType.get(), 1000), + 2, buf(Types.FloatType.get(), 9.5f), + 3, buf(Types.LongType.get(), 999L), + 4, buf(Types.StringType.get(), "zzz"))); + + @Test + void typeMatchesIdsPresentInMaps() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA); + + assertThat(stats.type().fields()).isEmpty(); + + stats.wrap(FILE_WITH_STATS); + + assertThat(stats.type().fields()) + .containsExactlyInAnyOrderElementsOf( + StatsUtil.statsReadSchema(SCHEMA, List.of(1, 2, 3, 4)).fields()); + } + + @Test + void wrapInvalidatesType() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + Types.StructType firstType = stats.type(); + assertThat(firstType.fields()) + .containsExactlyInAnyOrderElementsOf( + StatsUtil.statsReadSchema(SCHEMA, List.of(1, 2, 3, 4)).fields()); + + DataFile file2 = + dataFile( + 100L, + ImmutableMap.of(1, 50L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(1, buf(Types.IntegerType.get(), 500)), + ImmutableMap.of(1, buf(Types.IntegerType.get(), 5000))); + stats.wrap(file2); + + assertThat(stats.type().fields()) + .containsExactlyInAnyOrderElementsOf( + StatsUtil.statsReadSchema(SCHEMA, List.of(1)).fields()); + assertThat(stats.type()).isNotEqualTo(firstType); + } + + @ParameterizedTest + @MethodSource("typesAndBounds") + void boundDecodingPerType(Type type, Object lower, Object upper) { + int fieldId = 1; + Schema schema = new Schema(Types.NestedField.optional(fieldId, "col", type)); + DataFile file = + dataFile( + 100L, + ImmutableMap.of(fieldId, 1L), + ImmutableMap.of(), + ImmutableMap.of(), + bound(fieldId, Conversions.toByteBuffer(type, lower)), + bound(fieldId, Conversions.toByteBuffer(type, upper))); + FieldStats<?> stats = new MapBackedContentStats(schema).wrap(file).statsFor(fieldId); + + Comparator<Object> comparator = Comparators.forType(type.asPrimitiveType()); + assertThat(comparator.compare(stats.lowerBound(), lower)).isZero(); Review Comment: Done. `boundDecodingPerType` uses `usingComparator(comparator).isEqualTo(...)`. ########## core/src/test/java/org/apache/iceberg/TestMapBackedContentStats.java: ########## @@ -0,0 +1,382 @@ +/* + * 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.math.BigDecimal; +import java.nio.ByteBuffer; +import java.util.Comparator; +import java.util.List; +import java.util.Map; +import java.util.UUID; +import java.util.stream.Stream; +import org.apache.iceberg.geospatial.GeospatialBound; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; +import org.apache.iceberg.types.Comparators; +import org.apache.iceberg.types.Conversions; +import org.apache.iceberg.types.Type; +import org.apache.iceberg.types.Types; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.MethodSource; + +class TestMapBackedContentStats { + + private static final Schema SCHEMA = + new Schema( + Types.NestedField.required(1, "id", Types.IntegerType.get()), + Types.NestedField.optional(2, "score", Types.FloatType.get()), + Types.NestedField.optional(3, "ts", Types.LongType.get()), + Types.NestedField.optional(4, "name", Types.StringType.get()), + Types.NestedField.optional(5, "flag", Types.BooleanType.get())); + + private static final PartitionData EMPTY_PARTITION = + new PartitionData(PartitionSpec.unpartitioned().partitionType()); + + /** + * Stats on fields 1-4; field 5 is absent. Field 3 has no value-count entry, so valueCount() + * throws on unboxing. + */ + private static final DataFile FILE_WITH_STATS = + dataFile( + 100L, + ImmutableMap.of(1, 100L, 2, 100L, 4, 100L), + ImmutableMap.of(2, 5L, 3, 1L, 4, 2L), + ImmutableMap.of(2, 3L), + ImmutableMap.of( + 1, buf(Types.IntegerType.get(), 1), + 2, buf(Types.FloatType.get(), 1.5f), + 3, buf(Types.LongType.get(), 100L), + 4, buf(Types.StringType.get(), "aaa")), + ImmutableMap.of( + 1, buf(Types.IntegerType.get(), 1000), + 2, buf(Types.FloatType.get(), 9.5f), + 3, buf(Types.LongType.get(), 999L), + 4, buf(Types.StringType.get(), "zzz"))); + + @Test + void typeMatchesIdsPresentInMaps() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA); + + assertThat(stats.type().fields()).isEmpty(); + + stats.wrap(FILE_WITH_STATS); + + assertThat(stats.type().fields()) + .containsExactlyInAnyOrderElementsOf( + StatsUtil.statsReadSchema(SCHEMA, List.of(1, 2, 3, 4)).fields()); + } + + @Test + void wrapInvalidatesType() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + Types.StructType firstType = stats.type(); + assertThat(firstType.fields()) + .containsExactlyInAnyOrderElementsOf( + StatsUtil.statsReadSchema(SCHEMA, List.of(1, 2, 3, 4)).fields()); + + DataFile file2 = + dataFile( + 100L, + ImmutableMap.of(1, 50L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(1, buf(Types.IntegerType.get(), 500)), + ImmutableMap.of(1, buf(Types.IntegerType.get(), 5000))); + stats.wrap(file2); + + assertThat(stats.type().fields()) + .containsExactlyInAnyOrderElementsOf( + StatsUtil.statsReadSchema(SCHEMA, List.of(1)).fields()); + assertThat(stats.type()).isNotEqualTo(firstType); + } + + @ParameterizedTest + @MethodSource("typesAndBounds") + void boundDecodingPerType(Type type, Object lower, Object upper) { + int fieldId = 1; + Schema schema = new Schema(Types.NestedField.optional(fieldId, "col", type)); + DataFile file = + dataFile( + 100L, + ImmutableMap.of(fieldId, 1L), + ImmutableMap.of(), + ImmutableMap.of(), + bound(fieldId, Conversions.toByteBuffer(type, lower)), + bound(fieldId, Conversions.toByteBuffer(type, upper))); + FieldStats<?> stats = new MapBackedContentStats(schema).wrap(file).statsFor(fieldId); + + Comparator<Object> comparator = Comparators.forType(type.asPrimitiveType()); + assertThat(comparator.compare(stats.lowerBound(), lower)).isZero(); + assertThat(comparator.compare(stats.upperBound(), upper)).isZero(); + } + + private static Stream<Arguments> typesAndBounds() { + return Stream.of( + Arguments.of(Types.BooleanType.get(), false, true), + Arguments.of(Types.IntegerType.get(), -5, 100), + Arguments.of(Types.LongType.get(), 0L, 1_000L), + Arguments.of(Types.FloatType.get(), 1.5f, 9.5f), + Arguments.of(Types.DoubleType.get(), 0.0d, 25.0d), + Arguments.of(Types.DateType.get(), 100, 200), + Arguments.of(Types.TimeType.get(), 1_000L, 2_000L), + Arguments.of(Types.TimestampType.withZone(), 111L, 222L), + Arguments.of(Types.TimestampNanoType.withZone(), 111L, 222L), + Arguments.of(Types.StringType.get(), "a", "z"), + Arguments.of( + Types.UUIDType.get(), + UUID.fromString("07ceab48-62b2-4219-9172-856e687c92ad"), + UUID.fromString("a9a7c24d-2869-4c66-8803-b1f6a36256ed")), + Arguments.of( + Types.FixedType.ofLength(4), + ByteBuffer.wrap(new byte[] {0, 1, 2, 3}), + ByteBuffer.wrap(new byte[] {4, 5, 6, 7})), + Arguments.of( + Types.BinaryType.get(), + ByteBuffer.wrap(new byte[] {1, 2}), + ByteBuffer.wrap(new byte[] {3, 4, 5})), + Arguments.of(Types.DecimalType.of(9, 2), new BigDecimal("1.23"), new BigDecimal("9.99")), + Arguments.of(Types.UnknownType.get(), null, null)); + } + + @Test + void geoBoundsDecode() { + GeospatialBound lower = GeospatialBound.createXY(1.0, 2.0); + GeospatialBound upper = GeospatialBound.createXYZM(3.0, 4.0, 5.0, 6.0); + Schema schema = new Schema(Types.NestedField.optional(10, "geom", Types.GeometryType.crs84())); + DataFile file = + dataFile( + 100L, + ImmutableMap.of(10, 26L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(10, lower.toByteBuffer()), + ImmutableMap.of(10, upper.toByteBuffer())); + MapBackedContentStats stats = new MapBackedContentStats(schema).wrap(file); + + FieldStats<?> geom = stats.statsFor(10); + assertThat(geom.lowerBound()).isInstanceOf(GeospatialBound.class).isEqualTo(lower); + assertThat(geom.upperBound()).isInstanceOf(GeospatialBound.class).isEqualTo(upper); + } + + @Test + void missingBoundsDecodeToNull() { + DataFile file = + dataFile( + 100L, + ImmutableMap.of(1, 100L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of()); + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(file); + + FieldStats<?> id = stats.statsFor(1); + assertThat(id.lowerBound()).isNull(); + assertThat(id.upperBound()).isNull(); + assertThat(id.valueCount()).isEqualTo(100L); + } + + @Test + void countsAndTightBounds() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + + FieldStats<?> id = stats.statsFor(1); + assertThat(id.hasValueCount()).isTrue(); + assertThat(id.valueCount()).isEqualTo(100L); + assertThat(id.hasNullValueCount()).isFalse(); + assertThat(id.hasNanValueCount()).isFalse(); + assertThatThrownBy(id::nullValueCount) + .isInstanceOf(NullPointerException.class) + .hasMessageContaining("Long.longValue()"); + assertThatThrownBy(id::nanValueCount) + .isInstanceOf(NullPointerException.class) + .hasMessageContaining("Long.longValue()"); + // ContentFile maps have no tight-bounds flag so the view always reports false. + assertThat(id.tightBounds()).isFalse(); + // FILE_WITH_STATS has no avg-size map. + assertThat(id.avgValueSizeInBytes()).isNull(); + + FieldStats<?> score = stats.statsFor(2); Review Comment: Bound checks stay in `boundDecodingPerType`. `countsAndTightBounds` stays on the count maps, so I didn't add lower/upper asserts for fields 2–4. Field 1 still checks `tightBounds()` is false because the content-file maps have no tight-bounds flag. ########## core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java: ########## @@ -114,17 +132,57 @@ class TestTrackedFileAdapters { .addedFilesCount(3) .existingFilesCount(5) .deletedFilesCount(2) - .replacedFilesCount(0) - .modifiedFilesCount(0) + .replacedFilesCount(4) + .modifiedFilesCount(1) .addedRowsCount(300L) .existingRowsCount(500L) .deletedRowsCount(200L) - .replacedRowsCount(0L) - .modifiedRowsCount(0L) + .replacedRowsCount(40L) + .modifiedRowsCount(10L) .minSequenceNumber(7L) .dv(ByteBuffer.wrap(MumblingTestUtil.onlyFirstBitSetBytes())) .build(); + private static final Metrics METRICS_WITH_BOUNDS = + new Metrics( + 100L, + ImmutableMap.of(1, 16L, 2, 64L), + ImmutableMap.of(1, 100L, 2, 100L), + ImmutableMap.of(1, 0L, 2, 5L), + ImmutableMap.of(), + ImmutableMap.of(1, Conversions.toByteBuffer(Types.IntegerType.get(), 1)), + ImmutableMap.of(1, Conversions.toByteBuffer(Types.IntegerType.get(), 1000))); + + private static final DataFile DATA_FILE = + new GenericDataFile( + UNPARTITIONED_SPEC.specId(), + DATA_FILE_LOCATION, + FileFormat.PARQUET, + PartitionData.EMPTY, + 1024L, + new Metrics(100L), + KEY_METADATA, + ImmutableList.of(0L), + SORT_ORDER_ID, + FIRST_ROW_ID); + private static final DataFile DATA_FILE_WITH_METRICS = + new GenericDataFile( + UNPARTITIONED_SPEC.specId(), + DATA_FILE_LOCATION, + FileFormat.PARQUET, + PartitionData.EMPTY, + 1024L, + METRICS_WITH_BOUNDS, + null, Review Comment: Done. `DATA_FILE_WITH_METRICS` now uses `KEY_METADATA`, `SORT_ORDER_ID`, and `FIRST_ROW_ID`. ########## core/src/test/java/org/apache/iceberg/TestMapBackedContentStats.java: ########## @@ -0,0 +1,382 @@ +/* + * 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.math.BigDecimal; +import java.nio.ByteBuffer; +import java.util.Comparator; +import java.util.List; +import java.util.Map; +import java.util.UUID; +import java.util.stream.Stream; +import org.apache.iceberg.geospatial.GeospatialBound; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; +import org.apache.iceberg.types.Comparators; +import org.apache.iceberg.types.Conversions; +import org.apache.iceberg.types.Type; +import org.apache.iceberg.types.Types; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.MethodSource; + +class TestMapBackedContentStats { + + private static final Schema SCHEMA = + new Schema( + Types.NestedField.required(1, "id", Types.IntegerType.get()), + Types.NestedField.optional(2, "score", Types.FloatType.get()), + Types.NestedField.optional(3, "ts", Types.LongType.get()), + Types.NestedField.optional(4, "name", Types.StringType.get()), + Types.NestedField.optional(5, "flag", Types.BooleanType.get())); + + private static final PartitionData EMPTY_PARTITION = + new PartitionData(PartitionSpec.unpartitioned().partitionType()); + + /** + * Stats on fields 1-4; field 5 is absent. Field 3 has no value-count entry, so valueCount() + * throws on unboxing. + */ + private static final DataFile FILE_WITH_STATS = + dataFile( + 100L, + ImmutableMap.of(1, 100L, 2, 100L, 4, 100L), + ImmutableMap.of(2, 5L, 3, 1L, 4, 2L), + ImmutableMap.of(2, 3L), + ImmutableMap.of( + 1, buf(Types.IntegerType.get(), 1), + 2, buf(Types.FloatType.get(), 1.5f), + 3, buf(Types.LongType.get(), 100L), + 4, buf(Types.StringType.get(), "aaa")), + ImmutableMap.of( + 1, buf(Types.IntegerType.get(), 1000), + 2, buf(Types.FloatType.get(), 9.5f), + 3, buf(Types.LongType.get(), 999L), + 4, buf(Types.StringType.get(), "zzz"))); + + @Test + void typeMatchesIdsPresentInMaps() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA); + + assertThat(stats.type().fields()).isEmpty(); + + stats.wrap(FILE_WITH_STATS); + + assertThat(stats.type().fields()) + .containsExactlyInAnyOrderElementsOf( + StatsUtil.statsReadSchema(SCHEMA, List.of(1, 2, 3, 4)).fields()); + } + + @Test + void wrapInvalidatesType() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + Types.StructType firstType = stats.type(); + assertThat(firstType.fields()) + .containsExactlyInAnyOrderElementsOf( + StatsUtil.statsReadSchema(SCHEMA, List.of(1, 2, 3, 4)).fields()); + + DataFile file2 = + dataFile( + 100L, + ImmutableMap.of(1, 50L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(1, buf(Types.IntegerType.get(), 500)), + ImmutableMap.of(1, buf(Types.IntegerType.get(), 5000))); + stats.wrap(file2); + + assertThat(stats.type().fields()) + .containsExactlyInAnyOrderElementsOf( + StatsUtil.statsReadSchema(SCHEMA, List.of(1)).fields()); + assertThat(stats.type()).isNotEqualTo(firstType); + } + + @ParameterizedTest + @MethodSource("typesAndBounds") + void boundDecodingPerType(Type type, Object lower, Object upper) { + int fieldId = 1; + Schema schema = new Schema(Types.NestedField.optional(fieldId, "col", type)); + DataFile file = + dataFile( + 100L, + ImmutableMap.of(fieldId, 1L), + ImmutableMap.of(), + ImmutableMap.of(), + bound(fieldId, Conversions.toByteBuffer(type, lower)), + bound(fieldId, Conversions.toByteBuffer(type, upper))); + FieldStats<?> stats = new MapBackedContentStats(schema).wrap(file).statsFor(fieldId); + + Comparator<Object> comparator = Comparators.forType(type.asPrimitiveType()); + assertThat(comparator.compare(stats.lowerBound(), lower)).isZero(); + assertThat(comparator.compare(stats.upperBound(), upper)).isZero(); + } + + private static Stream<Arguments> typesAndBounds() { + return Stream.of( + Arguments.of(Types.BooleanType.get(), false, true), + Arguments.of(Types.IntegerType.get(), -5, 100), + Arguments.of(Types.LongType.get(), 0L, 1_000L), + Arguments.of(Types.FloatType.get(), 1.5f, 9.5f), + Arguments.of(Types.DoubleType.get(), 0.0d, 25.0d), + Arguments.of(Types.DateType.get(), 100, 200), + Arguments.of(Types.TimeType.get(), 1_000L, 2_000L), + Arguments.of(Types.TimestampType.withZone(), 111L, 222L), + Arguments.of(Types.TimestampNanoType.withZone(), 111L, 222L), + Arguments.of(Types.StringType.get(), "a", "z"), + Arguments.of( + Types.UUIDType.get(), + UUID.fromString("07ceab48-62b2-4219-9172-856e687c92ad"), + UUID.fromString("a9a7c24d-2869-4c66-8803-b1f6a36256ed")), + Arguments.of( + Types.FixedType.ofLength(4), + ByteBuffer.wrap(new byte[] {0, 1, 2, 3}), + ByteBuffer.wrap(new byte[] {4, 5, 6, 7})), + Arguments.of( + Types.BinaryType.get(), + ByteBuffer.wrap(new byte[] {1, 2}), + ByteBuffer.wrap(new byte[] {3, 4, 5})), + Arguments.of(Types.DecimalType.of(9, 2), new BigDecimal("1.23"), new BigDecimal("9.99")), + Arguments.of(Types.UnknownType.get(), null, null)); + } + + @Test + void geoBoundsDecode() { + GeospatialBound lower = GeospatialBound.createXY(1.0, 2.0); + GeospatialBound upper = GeospatialBound.createXYZM(3.0, 4.0, 5.0, 6.0); + Schema schema = new Schema(Types.NestedField.optional(10, "geom", Types.GeometryType.crs84())); + DataFile file = + dataFile( + 100L, + ImmutableMap.of(10, 26L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(10, lower.toByteBuffer()), + ImmutableMap.of(10, upper.toByteBuffer())); + MapBackedContentStats stats = new MapBackedContentStats(schema).wrap(file); + + FieldStats<?> geom = stats.statsFor(10); + assertThat(geom.lowerBound()).isInstanceOf(GeospatialBound.class).isEqualTo(lower); + assertThat(geom.upperBound()).isInstanceOf(GeospatialBound.class).isEqualTo(upper); + } + + @Test + void missingBoundsDecodeToNull() { + DataFile file = + dataFile( + 100L, + ImmutableMap.of(1, 100L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of()); + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(file); + + FieldStats<?> id = stats.statsFor(1); + assertThat(id.lowerBound()).isNull(); + assertThat(id.upperBound()).isNull(); + assertThat(id.valueCount()).isEqualTo(100L); + } + + @Test + void countsAndTightBounds() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + + FieldStats<?> id = stats.statsFor(1); + assertThat(id.hasValueCount()).isTrue(); + assertThat(id.valueCount()).isEqualTo(100L); + assertThat(id.hasNullValueCount()).isFalse(); + assertThat(id.hasNanValueCount()).isFalse(); + assertThatThrownBy(id::nullValueCount) + .isInstanceOf(NullPointerException.class) + .hasMessageContaining("Long.longValue()"); + assertThatThrownBy(id::nanValueCount) + .isInstanceOf(NullPointerException.class) + .hasMessageContaining("Long.longValue()"); + // ContentFile maps have no tight-bounds flag so the view always reports false. + assertThat(id.tightBounds()).isFalse(); + // FILE_WITH_STATS has no avg-size map. + assertThat(id.avgValueSizeInBytes()).isNull(); + + FieldStats<?> score = stats.statsFor(2); + assertThat(score.hasNullValueCount()).isTrue(); + assertThat(score.nullValueCount()).isEqualTo(5L); + assertThat(score.hasNanValueCount()).isTrue(); + assertThat(score.nanValueCount()).isEqualTo(3L); + + FieldStats<?> ts = stats.statsFor(3); + assertThat(ts.hasValueCount()).isFalse(); + assertThatThrownBy(ts::valueCount) + .isInstanceOf(NullPointerException.class) + .hasMessageContaining("Long.longValue()"); + assertThat(ts.hasNullValueCount()).isTrue(); + assertThat(ts.nullValueCount()).isEqualTo(1L); + + FieldStats<?> name = stats.statsFor(4); + assertThat(name.hasValueCount()).isTrue(); + assertThat(name.valueCount()).isEqualTo(100L); + assertThat(name.hasNullValueCount()).isTrue(); + assertThat(name.nullValueCount()).isEqualTo(2L); + assertThat(name.hasNanValueCount()).isFalse(); + assertThatThrownBy(name::nanValueCount) + .isInstanceOf(NullPointerException.class) + .hasMessageContaining("Long.longValue()"); + } + + @Test + void fieldWithoutStatsIsExcluded() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + assertThat(stats.type().field(StatsUtil.toBaseId(5))).isNull(); + assertThat(stats.statsFor(5)).isNull(); + } + + @Test + void unknownFieldIdInMaps() { + DataFile file = + dataFile( + 100L, + ImmutableMap.of(99, 1L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of()); + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(file); + + assertThat(stats.statsFor(99)).isNull(); + assertThat(stats.type().fields()).isEmpty(); + assertThat(stats.fieldStats()).isEmpty(); + } + + @Test + void reuseRebindsFields() { Review Comment: Done. `wrapInvalidatesType` now sits next to `reuseRebindsFields`. ########## core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java: ########## @@ -671,6 +733,312 @@ void unknownSpecIdThrows() { .hasMessageContaining("Cannot find partition spec for spec ID"); } + @Test + void dataTrackedFileAdapterFromDataFile() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE); + + assertThat(result.tracking()).isNull(); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + assertThatThrownBy(result::formatVersion) + .isInstanceOf(IllegalStateException.class) + .hasMessage("Format version is assigned at write time"); + } + + @Test + void dataTrackedFileAdapterFromExistingManifestEntry() { + TrackedFile result = + TrackedFileAdapters.forDataFile(TABLE_SCHEMA) + .wrap( + newEntry() + .wrapExisting( + SNAPSHOT_ID, DATA_SEQUENCE_NUMBER, FILE_SEQUENCE_NUMBER, DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().dataSequenceNumber()).isEqualTo(DATA_SEQUENCE_NUMBER); + assertThat(result.tracking().fileSequenceNumber()).isEqualTo(FILE_SEQUENCE_NUMBER); + assertThat(result.tracking().firstRowId()).isEqualTo(FIRST_ROW_ID); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterReuse() { + TrackedFileAdapters.DataTrackedFile adapter = TrackedFileAdapters.forDataFile(TABLE_SCHEMA); + + adapter.wrap(DATA_FILE); + assertWrappedDataFileMatchesFileFields(adapter, DATA_FILE); + assertThat(adapter.tracking()).isNull(); + + DataFile file2 = + new GenericDataFile( + UNPARTITIONED_SPEC.specId(), + "s3://bucket/data/file2.parquet", + FileFormat.PARQUET, + PartitionData.EMPTY, + 2048L, + new Metrics(200L, null, null, null, null), + null, + ImmutableList.of(0L), + null, + null); + assignManifestPosition(file2, "s3://bucket/table/manifest-2.parquet", 8L); + adapter.wrap( + newEntry().wrapExisting(SNAPSHOT_ID, DATA_SEQUENCE_NUMBER, FILE_SEQUENCE_NUMBER, file2)); + assertWrappedDataFileMatchesFileFields(adapter, file2); + assertThat(adapter.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertManifestPosition(adapter.tracking(), file2); + } + + @Test + void dataTrackedFileAdapterRejectsNullFile() { + TrackedFileAdapters.DataTrackedFile adapter = TrackedFileAdapters.forDataFile(TABLE_SCHEMA); + assertThatThrownBy(() -> adapter.wrap((DataFile) null)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("Invalid file: null"); + } + + @Test + void dataTrackedFileAdapterContentStats() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE_WITH_METRICS); + + ContentStats stats = result.contentStats(); + assertThat(stats).isNotNull(); + assertThat(stats.fieldStats()).extracting(FieldStats::fieldId).containsExactlyInAnyOrder(1, 2); + + FieldStats<?> idStats = stats.statsFor(1); + assertThat(idStats.valueCount()).isEqualTo(100L); + assertThat(idStats.lowerBound()).isEqualTo(1); + assertThat(idStats.upperBound()).isEqualTo(1000); + } + + @Test + void dataTrackedFileAdapterWithoutMetricsHasNoContentStats() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE); + + assertThat(result.contentStats()).isNull(); + } + + @Test + void dataTrackedFileAdapterKeepsPartitionTuple() { + DataFile partitioned = + new GenericDataFile( + PARTITIONED_SPEC.specId(), + DATA_FILE_LOCATION, + FileFormat.PARQUET, + PARTITION, + 1024L, + new Metrics(100L, null, null, null, null), + null, + ImmutableList.of(0L), + null, + null); + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(partitioned); + + assertThat(result.partition()) + .usingComparator(Comparators.forType(PARTITIONED_SPEC.partitionType())) + .isEqualTo(PARTITION); + assertThat(TrackedFileAdapters.asDataFile(result, specsById(PARTITIONED_SPEC)).partition()) + .usingComparator(Comparators.forType(PARTITIONED_SPEC.partitionType())) + .isEqualTo(PARTITION); + } + + @Test + void dataTrackedFileAdapterUnwrapsToOriginalTrackedFile() { + TrackedFile source = trackedFile(FileContent.DATA); + DataFile dataFile = TrackedFileAdapters.asDataFile(source, UNPARTITIONED); + TrackedFile roundTripped = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(dataFile); + assertThat(roundTripped).isSameAs(source); + } + + @Test + void manifestTrackedFileAdapterUnwrapsToOriginalTrackedFile() { + TrackedFile original = trackedFile(FileContent.DATA_MANIFEST, 0); + ManifestFile adapted = TrackedFileAdapters.asManifestFile(original); + TrackedFile result = TrackedFileAdapters.forManifestFile().wrap(adapted); + assertThat(result).isSameAs(original); + assertThat(result.formatVersion()).isZero(); + } + + @Test + void dataManifestTrackedFileAdapter() { + ManifestFile manifest = writeManifestFile(ManifestContent.DATA); + TrackedFile result = TrackedFileAdapters.forManifestFile().wrap(manifest); + + assertThat(result.contentType()).isEqualTo(FileContent.DATA_MANIFEST); + assertThat(result.formatVersion()).isZero(); + assertThat(result.location()).isEqualTo(MANIFEST_LOCATION); + assertThat(result.fileFormat()).isEqualTo(FileFormat.AVRO); + assertThat(result.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().firstRowId()).isNull(); + assertThat(result.recordCount()) + .isEqualTo(ADDED_FILES_COUNT + EXISTING_FILES_COUNT + DELETED_FILES_COUNT); + assertThat(result.manifestInfo()).isNotNull(); + assertThat(result.manifestInfo().addedFilesCount()).isEqualTo(ADDED_FILES_COUNT); + assertThat(result.manifestInfo().existingFilesCount()).isEqualTo(EXISTING_FILES_COUNT); + assertThat(result.manifestInfo().deletedFilesCount()).isEqualTo(DELETED_FILES_COUNT); + assertThat(result.manifestInfo().addedRowsCount()).isEqualTo(ADDED_ROWS_COUNT); + assertThat(result.manifestInfo().replacedFilesCount()).isEqualTo(0); + assertThat(result.manifestInfo().replacedRowsCount()).isEqualTo(0L); + assertThat(result.manifestInfo().modifiedFilesCount()).isEqualTo(0); + assertThat(result.manifestInfo().modifiedRowsCount()).isEqualTo(0L); + } + + @Test + void deleteManifestTrackedFileAdapter() { + ManifestFile manifest = writeManifestFile(ManifestContent.DELETES); + TrackedFile result = TrackedFileAdapters.forManifestFile().wrap(manifest); + + assertThat(result.contentType()).isEqualTo(FileContent.DELETE_MANIFEST); + assertThat(result.formatVersion()).isZero(); + assertThat(result.recordCount()) + .isEqualTo(ADDED_FILES_COUNT + EXISTING_FILES_COUNT + DELETED_FILES_COUNT); + assertThat(result.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertThat(result.tracking().firstRowId()).isNull(); + } + + @Test + void manifestTrackedFileAdapterFailsWhenAddedFilesCountMissing() { + ManifestFile manifest = mock(ManifestFile.class); + when(manifest.path()).thenReturn(MANIFEST_LOCATION); + when(manifest.content()).thenReturn(ManifestContent.DATA); + when(manifest.addedFilesCount()).thenReturn(null); + when(manifest.existingFilesCount()).thenReturn(1); + when(manifest.deletedFilesCount()).thenReturn(0); + when(manifest.replacedFilesCount()).thenReturn(0); + + assertThatThrownBy(() -> TrackedFileAdapters.forManifestFile().wrap(manifest)) + .isInstanceOf(NullPointerException.class) + .hasMessageContaining("missing added files count"); + } + + @Test + void manifestTrackedFileAdapterRejectsNullReplacedFilesCount() { + ManifestFile manifest = manifestWithCounts(1, 1, 0, null, 0); + + assertThatThrownBy(() -> TrackedFileAdapters.forManifestFile().wrap(manifest)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage( + "Cannot convert manifest %s: Invalid replaced file count: null", MANIFEST_LOCATION); + } + + @Test + void manifestTrackedFileAdapterRejectsNonZeroReplacedFilesCount() { + ManifestFile manifest = manifestWithCounts(1, 1, 0, 1, 0); + + assertThatThrownBy(() -> TrackedFileAdapters.forManifestFile().wrap(manifest)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage( + "Cannot convert manifest %s: Invalid replaced file count: 1", MANIFEST_LOCATION); + } + + @Test + void manifestTrackedFileAdapterRejectsNullModifiedFilesCount() { + ManifestFile manifest = manifestWithCounts(1, 1, 0, 0, null); + + assertThatThrownBy(() -> TrackedFileAdapters.forManifestFile().wrap(manifest)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage( + "Cannot convert manifest %s: Invalid modified file count: null", MANIFEST_LOCATION); + } + + @Test + void manifestTrackedFileAdapterRejectsNonZeroModifiedFilesCount() { + ManifestFile manifest = manifestWithCounts(1, 1, 0, 0, 1); + + assertThatThrownBy(() -> TrackedFileAdapters.forManifestFile().wrap(manifest)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage( + "Cannot convert manifest %s: Invalid modified file count: 1", MANIFEST_LOCATION); + } + + @Test + void manifestTrackedFileAdapterFailsWhenAddedRowsCountMissing() { + ManifestFile manifest = + writeManifestFile( + ManifestContent.DATA, MANIFEST_SEQUENCE_NUMBER, MANIFEST_MIN_SEQUENCE_NUMBER, null); + TrackedFile tracked = TrackedFileAdapters.forManifestFile().wrap(manifest); + + assertThatThrownBy(() -> tracked.manifestInfo().addedRowsCount()) + .isInstanceOf(NullPointerException.class) + .hasMessageContaining("null"); + } + + private static void assertWrappedDataFileMatchesFileFields(TrackedFile result, DataFile file) { Review Comment: Leaving this assert helper where it is (after all test methods). `assertSameDataFile` and `assertSameDeleteFile` predate this PR and were placed before some test methods. Didn't relocate them to avoid code churning. ########## core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java: ########## @@ -671,6 +733,312 @@ void unknownSpecIdThrows() { .hasMessageContaining("Cannot find partition spec for spec ID"); } + @Test + void dataTrackedFileAdapterFromDataFile() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE); + + assertThat(result.tracking()).isNull(); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + assertThatThrownBy(result::formatVersion) + .isInstanceOf(IllegalStateException.class) + .hasMessage("Format version is assigned at write time"); + } + + @Test + void dataTrackedFileAdapterFromExistingManifestEntry() { + TrackedFile result = + TrackedFileAdapters.forDataFile(TABLE_SCHEMA) + .wrap( + newEntry() + .wrapExisting( + SNAPSHOT_ID, DATA_SEQUENCE_NUMBER, FILE_SEQUENCE_NUMBER, DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().dataSequenceNumber()).isEqualTo(DATA_SEQUENCE_NUMBER); + assertThat(result.tracking().fileSequenceNumber()).isEqualTo(FILE_SEQUENCE_NUMBER); + assertThat(result.tracking().firstRowId()).isEqualTo(FIRST_ROW_ID); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterReuse() { + TrackedFileAdapters.DataTrackedFile adapter = TrackedFileAdapters.forDataFile(TABLE_SCHEMA); + + adapter.wrap(DATA_FILE); + assertWrappedDataFileMatchesFileFields(adapter, DATA_FILE); + assertThat(adapter.tracking()).isNull(); + + DataFile file2 = + new GenericDataFile( + UNPARTITIONED_SPEC.specId(), + "s3://bucket/data/file2.parquet", + FileFormat.PARQUET, + PartitionData.EMPTY, + 2048L, + new Metrics(200L, null, null, null, null), + null, + ImmutableList.of(0L), + null, + null); + assignManifestPosition(file2, "s3://bucket/table/manifest-2.parquet", 8L); + adapter.wrap( + newEntry().wrapExisting(SNAPSHOT_ID, DATA_SEQUENCE_NUMBER, FILE_SEQUENCE_NUMBER, file2)); + assertWrappedDataFileMatchesFileFields(adapter, file2); + assertThat(adapter.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertManifestPosition(adapter.tracking(), file2); + } + + @Test + void dataTrackedFileAdapterRejectsNullFile() { + TrackedFileAdapters.DataTrackedFile adapter = TrackedFileAdapters.forDataFile(TABLE_SCHEMA); + assertThatThrownBy(() -> adapter.wrap((DataFile) null)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("Invalid file: null"); + } + + @Test + void dataTrackedFileAdapterContentStats() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE_WITH_METRICS); + + ContentStats stats = result.contentStats(); + assertThat(stats).isNotNull(); + assertThat(stats.fieldStats()).extracting(FieldStats::fieldId).containsExactlyInAnyOrder(1, 2); + + FieldStats<?> idStats = stats.statsFor(1); + assertThat(idStats.valueCount()).isEqualTo(100L); + assertThat(idStats.lowerBound()).isEqualTo(1); + assertThat(idStats.upperBound()).isEqualTo(1000); + } + + @Test + void dataTrackedFileAdapterWithoutMetricsHasNoContentStats() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE); + + assertThat(result.contentStats()).isNull(); + } + + @Test + void dataTrackedFileAdapterKeepsPartitionTuple() { + DataFile partitioned = + new GenericDataFile( + PARTITIONED_SPEC.specId(), + DATA_FILE_LOCATION, + FileFormat.PARQUET, + PARTITION, + 1024L, + new Metrics(100L, null, null, null, null), + null, + ImmutableList.of(0L), + null, + null); + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(partitioned); + + assertThat(result.partition()) + .usingComparator(Comparators.forType(PARTITIONED_SPEC.partitionType())) + .isEqualTo(PARTITION); + assertThat(TrackedFileAdapters.asDataFile(result, specsById(PARTITIONED_SPEC)).partition()) + .usingComparator(Comparators.forType(PARTITIONED_SPEC.partitionType())) + .isEqualTo(PARTITION); + } + + @Test + void dataTrackedFileAdapterUnwrapsToOriginalTrackedFile() { + TrackedFile source = trackedFile(FileContent.DATA); + DataFile dataFile = TrackedFileAdapters.asDataFile(source, UNPARTITIONED); + TrackedFile roundTripped = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(dataFile); + assertThat(roundTripped).isSameAs(source); + } + + @Test + void manifestTrackedFileAdapterUnwrapsToOriginalTrackedFile() { + TrackedFile original = trackedFile(FileContent.DATA_MANIFEST, 0); + ManifestFile adapted = TrackedFileAdapters.asManifestFile(original); + TrackedFile result = TrackedFileAdapters.forManifestFile().wrap(adapted); + assertThat(result).isSameAs(original); + assertThat(result.formatVersion()).isZero(); Review Comment: Dropped the format-version 0 check on the unwrap test, and that test now constructs the tracked file without passing version 0. `isZero()` remains on `manifestTrackedFileAdapter` because `GenericManifestFile` does not override `formatVersion()` (the interface default is 0). ########## core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java: ########## @@ -671,6 +733,312 @@ void unknownSpecIdThrows() { .hasMessageContaining("Cannot find partition spec for spec ID"); } + @Test + void dataTrackedFileAdapterFromDataFile() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE); + + assertThat(result.tracking()).isNull(); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + assertThatThrownBy(result::formatVersion) + .isInstanceOf(IllegalStateException.class) + .hasMessage("Format version is assigned at write time"); + } + + @Test + void dataTrackedFileAdapterFromExistingManifestEntry() { + TrackedFile result = + TrackedFileAdapters.forDataFile(TABLE_SCHEMA) + .wrap( + newEntry() + .wrapExisting( + SNAPSHOT_ID, DATA_SEQUENCE_NUMBER, FILE_SEQUENCE_NUMBER, DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().dataSequenceNumber()).isEqualTo(DATA_SEQUENCE_NUMBER); + assertThat(result.tracking().fileSequenceNumber()).isEqualTo(FILE_SEQUENCE_NUMBER); + assertThat(result.tracking().firstRowId()).isEqualTo(FIRST_ROW_ID); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterReuse() { + TrackedFileAdapters.DataTrackedFile adapter = TrackedFileAdapters.forDataFile(TABLE_SCHEMA); + + adapter.wrap(DATA_FILE); + assertWrappedDataFileMatchesFileFields(adapter, DATA_FILE); + assertThat(adapter.tracking()).isNull(); + + DataFile file2 = + new GenericDataFile( + UNPARTITIONED_SPEC.specId(), + "s3://bucket/data/file2.parquet", + FileFormat.PARQUET, + PartitionData.EMPTY, + 2048L, + new Metrics(200L, null, null, null, null), + null, + ImmutableList.of(0L), + null, + null); + assignManifestPosition(file2, "s3://bucket/table/manifest-2.parquet", 8L); + adapter.wrap( + newEntry().wrapExisting(SNAPSHOT_ID, DATA_SEQUENCE_NUMBER, FILE_SEQUENCE_NUMBER, file2)); + assertWrappedDataFileMatchesFileFields(adapter, file2); + assertThat(adapter.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertManifestPosition(adapter.tracking(), file2); + } + + @Test + void dataTrackedFileAdapterRejectsNullFile() { + TrackedFileAdapters.DataTrackedFile adapter = TrackedFileAdapters.forDataFile(TABLE_SCHEMA); + assertThatThrownBy(() -> adapter.wrap((DataFile) null)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("Invalid file: null"); + } + + @Test + void dataTrackedFileAdapterContentStats() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE_WITH_METRICS); + + ContentStats stats = result.contentStats(); + assertThat(stats).isNotNull(); + assertThat(stats.fieldStats()).extracting(FieldStats::fieldId).containsExactlyInAnyOrder(1, 2); + + FieldStats<?> idStats = stats.statsFor(1); + assertThat(idStats.valueCount()).isEqualTo(100L); + assertThat(idStats.lowerBound()).isEqualTo(1); + assertThat(idStats.upperBound()).isEqualTo(1000); + } + + @Test + void dataTrackedFileAdapterWithoutMetricsHasNoContentStats() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE); + + assertThat(result.contentStats()).isNull(); + } + + @Test + void dataTrackedFileAdapterKeepsPartitionTuple() { + DataFile partitioned = + new GenericDataFile( + PARTITIONED_SPEC.specId(), + DATA_FILE_LOCATION, + FileFormat.PARQUET, + PARTITION, + 1024L, + new Metrics(100L, null, null, null, null), + null, + ImmutableList.of(0L), + null, + null); + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(partitioned); + + assertThat(result.partition()) + .usingComparator(Comparators.forType(PARTITIONED_SPEC.partitionType())) + .isEqualTo(PARTITION); + assertThat(TrackedFileAdapters.asDataFile(result, specsById(PARTITIONED_SPEC)).partition()) + .usingComparator(Comparators.forType(PARTITIONED_SPEC.partitionType())) + .isEqualTo(PARTITION); + } + + @Test + void dataTrackedFileAdapterUnwrapsToOriginalTrackedFile() { + TrackedFile source = trackedFile(FileContent.DATA); + DataFile dataFile = TrackedFileAdapters.asDataFile(source, UNPARTITIONED); + TrackedFile roundTripped = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(dataFile); + assertThat(roundTripped).isSameAs(source); + } + + @Test + void manifestTrackedFileAdapterUnwrapsToOriginalTrackedFile() { + TrackedFile original = trackedFile(FileContent.DATA_MANIFEST, 0); + ManifestFile adapted = TrackedFileAdapters.asManifestFile(original); + TrackedFile result = TrackedFileAdapters.forManifestFile().wrap(adapted); + assertThat(result).isSameAs(original); + assertThat(result.formatVersion()).isZero(); + } + + @Test + void dataManifestTrackedFileAdapter() { Review Comment: Combined into `manifestTrackedFileAdapter(ManifestContent)`. DATA maps to `DATA_MANIFEST` and DELETES maps to `DELETE_MANIFEST`. ########## core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java: ########## @@ -671,6 +733,312 @@ void unknownSpecIdThrows() { .hasMessageContaining("Cannot find partition spec for spec ID"); } + @Test + void dataTrackedFileAdapterFromDataFile() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE); + + assertThat(result.tracking()).isNull(); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + assertThatThrownBy(result::formatVersion) + .isInstanceOf(IllegalStateException.class) + .hasMessage("Format version is assigned at write time"); + } + + @Test + void dataTrackedFileAdapterFromExistingManifestEntry() { + TrackedFile result = + TrackedFileAdapters.forDataFile(TABLE_SCHEMA) + .wrap( + newEntry() + .wrapExisting( + SNAPSHOT_ID, DATA_SEQUENCE_NUMBER, FILE_SEQUENCE_NUMBER, DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().dataSequenceNumber()).isEqualTo(DATA_SEQUENCE_NUMBER); + assertThat(result.tracking().fileSequenceNumber()).isEqualTo(FILE_SEQUENCE_NUMBER); + assertThat(result.tracking().firstRowId()).isEqualTo(FIRST_ROW_ID); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterReuse() { + TrackedFileAdapters.DataTrackedFile adapter = TrackedFileAdapters.forDataFile(TABLE_SCHEMA); + + adapter.wrap(DATA_FILE); + assertWrappedDataFileMatchesFileFields(adapter, DATA_FILE); + assertThat(adapter.tracking()).isNull(); + + DataFile file2 = + new GenericDataFile( + UNPARTITIONED_SPEC.specId(), + "s3://bucket/data/file2.parquet", + FileFormat.PARQUET, + PartitionData.EMPTY, + 2048L, + new Metrics(200L, null, null, null, null), + null, + ImmutableList.of(0L), + null, + null); + assignManifestPosition(file2, "s3://bucket/table/manifest-2.parquet", 8L); + adapter.wrap( + newEntry().wrapExisting(SNAPSHOT_ID, DATA_SEQUENCE_NUMBER, FILE_SEQUENCE_NUMBER, file2)); + assertWrappedDataFileMatchesFileFields(adapter, file2); + assertThat(adapter.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertManifestPosition(adapter.tracking(), file2); + } + + @Test + void dataTrackedFileAdapterRejectsNullFile() { + TrackedFileAdapters.DataTrackedFile adapter = TrackedFileAdapters.forDataFile(TABLE_SCHEMA); + assertThatThrownBy(() -> adapter.wrap((DataFile) null)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("Invalid file: null"); + } + + @Test + void dataTrackedFileAdapterContentStats() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE_WITH_METRICS); + + ContentStats stats = result.contentStats(); + assertThat(stats).isNotNull(); + assertThat(stats.fieldStats()).extracting(FieldStats::fieldId).containsExactlyInAnyOrder(1, 2); + + FieldStats<?> idStats = stats.statsFor(1); + assertThat(idStats.valueCount()).isEqualTo(100L); + assertThat(idStats.lowerBound()).isEqualTo(1); + assertThat(idStats.upperBound()).isEqualTo(1000); + } + + @Test + void dataTrackedFileAdapterWithoutMetricsHasNoContentStats() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE); + + assertThat(result.contentStats()).isNull(); + } + + @Test + void dataTrackedFileAdapterKeepsPartitionTuple() { + DataFile partitioned = + new GenericDataFile( + PARTITIONED_SPEC.specId(), + DATA_FILE_LOCATION, + FileFormat.PARQUET, + PARTITION, + 1024L, + new Metrics(100L, null, null, null, null), + null, + ImmutableList.of(0L), + null, + null); + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(partitioned); + + assertThat(result.partition()) + .usingComparator(Comparators.forType(PARTITIONED_SPEC.partitionType())) + .isEqualTo(PARTITION); + assertThat(TrackedFileAdapters.asDataFile(result, specsById(PARTITIONED_SPEC)).partition()) + .usingComparator(Comparators.forType(PARTITIONED_SPEC.partitionType())) + .isEqualTo(PARTITION); + } + + @Test + void dataTrackedFileAdapterUnwrapsToOriginalTrackedFile() { + TrackedFile source = trackedFile(FileContent.DATA); + DataFile dataFile = TrackedFileAdapters.asDataFile(source, UNPARTITIONED); + TrackedFile roundTripped = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(dataFile); + assertThat(roundTripped).isSameAs(source); Review Comment: Both round trips now use `original` -> `adapted` -> `roundTripped`. ########## core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java: ########## @@ -671,6 +733,312 @@ void unknownSpecIdThrows() { .hasMessageContaining("Cannot find partition spec for spec ID"); } + @Test + void dataTrackedFileAdapterFromDataFile() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE); + + assertThat(result.tracking()).isNull(); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + assertThatThrownBy(result::formatVersion) + .isInstanceOf(IllegalStateException.class) + .hasMessage("Format version is assigned at write time"); + } + + @Test + void dataTrackedFileAdapterFromExistingManifestEntry() { + TrackedFile result = + TrackedFileAdapters.forDataFile(TABLE_SCHEMA) + .wrap( + newEntry() + .wrapExisting( + SNAPSHOT_ID, DATA_SEQUENCE_NUMBER, FILE_SEQUENCE_NUMBER, DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().dataSequenceNumber()).isEqualTo(DATA_SEQUENCE_NUMBER); + assertThat(result.tracking().fileSequenceNumber()).isEqualTo(FILE_SEQUENCE_NUMBER); + assertThat(result.tracking().firstRowId()).isEqualTo(FIRST_ROW_ID); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterReuse() { + TrackedFileAdapters.DataTrackedFile adapter = TrackedFileAdapters.forDataFile(TABLE_SCHEMA); + + adapter.wrap(DATA_FILE); + assertWrappedDataFileMatchesFileFields(adapter, DATA_FILE); + assertThat(adapter.tracking()).isNull(); + + DataFile file2 = + new GenericDataFile( + UNPARTITIONED_SPEC.specId(), + "s3://bucket/data/file2.parquet", + FileFormat.PARQUET, + PartitionData.EMPTY, + 2048L, + new Metrics(200L, null, null, null, null), + null, + ImmutableList.of(0L), + null, + null); + assignManifestPosition(file2, "s3://bucket/table/manifest-2.parquet", 8L); + adapter.wrap( + newEntry().wrapExisting(SNAPSHOT_ID, DATA_SEQUENCE_NUMBER, FILE_SEQUENCE_NUMBER, file2)); + assertWrappedDataFileMatchesFileFields(adapter, file2); + assertThat(adapter.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertManifestPosition(adapter.tracking(), file2); + } + + @Test + void dataTrackedFileAdapterRejectsNullFile() { + TrackedFileAdapters.DataTrackedFile adapter = TrackedFileAdapters.forDataFile(TABLE_SCHEMA); + assertThatThrownBy(() -> adapter.wrap((DataFile) null)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("Invalid file: null"); + } + + @Test + void dataTrackedFileAdapterContentStats() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE_WITH_METRICS); + + ContentStats stats = result.contentStats(); + assertThat(stats).isNotNull(); + assertThat(stats.fieldStats()).extracting(FieldStats::fieldId).containsExactlyInAnyOrder(1, 2); + + FieldStats<?> idStats = stats.statsFor(1); + assertThat(idStats.valueCount()).isEqualTo(100L); + assertThat(idStats.lowerBound()).isEqualTo(1); + assertThat(idStats.upperBound()).isEqualTo(1000); + } + + @Test + void dataTrackedFileAdapterWithoutMetricsHasNoContentStats() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE); + + assertThat(result.contentStats()).isNull(); + } + + @Test + void dataTrackedFileAdapterKeepsPartitionTuple() { + DataFile partitioned = + new GenericDataFile( + PARTITIONED_SPEC.specId(), + DATA_FILE_LOCATION, + FileFormat.PARQUET, + PARTITION, + 1024L, + new Metrics(100L, null, null, null, null), + null, + ImmutableList.of(0L), + null, + null); + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(partitioned); + + assertThat(result.partition()) + .usingComparator(Comparators.forType(PARTITIONED_SPEC.partitionType())) + .isEqualTo(PARTITION); + assertThat(TrackedFileAdapters.asDataFile(result, specsById(PARTITIONED_SPEC)).partition()) + .usingComparator(Comparators.forType(PARTITIONED_SPEC.partitionType())) + .isEqualTo(PARTITION); + } + + @Test + void dataTrackedFileAdapterUnwrapsToOriginalTrackedFile() { + TrackedFile source = trackedFile(FileContent.DATA); + DataFile dataFile = TrackedFileAdapters.asDataFile(source, UNPARTITIONED); + TrackedFile roundTripped = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(dataFile); + assertThat(roundTripped).isSameAs(source); + } + + @Test + void manifestTrackedFileAdapterUnwrapsToOriginalTrackedFile() { + TrackedFile original = trackedFile(FileContent.DATA_MANIFEST, 0); + ManifestFile adapted = TrackedFileAdapters.asManifestFile(original); + TrackedFile result = TrackedFileAdapters.forManifestFile().wrap(adapted); + assertThat(result).isSameAs(original); + assertThat(result.formatVersion()).isZero(); + } + + @Test + void dataManifestTrackedFileAdapter() { + ManifestFile manifest = writeManifestFile(ManifestContent.DATA); + TrackedFile result = TrackedFileAdapters.forManifestFile().wrap(manifest); + + assertThat(result.contentType()).isEqualTo(FileContent.DATA_MANIFEST); + assertThat(result.formatVersion()).isZero(); + assertThat(result.location()).isEqualTo(MANIFEST_LOCATION); + assertThat(result.fileFormat()).isEqualTo(FileFormat.AVRO); + assertThat(result.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().firstRowId()).isNull(); + assertThat(result.recordCount()) + .isEqualTo(ADDED_FILES_COUNT + EXISTING_FILES_COUNT + DELETED_FILES_COUNT); + assertThat(result.manifestInfo()).isNotNull(); + assertThat(result.manifestInfo().addedFilesCount()).isEqualTo(ADDED_FILES_COUNT); + assertThat(result.manifestInfo().existingFilesCount()).isEqualTo(EXISTING_FILES_COUNT); + assertThat(result.manifestInfo().deletedFilesCount()).isEqualTo(DELETED_FILES_COUNT); + assertThat(result.manifestInfo().addedRowsCount()).isEqualTo(ADDED_ROWS_COUNT); + assertThat(result.manifestInfo().replacedFilesCount()).isEqualTo(0); + assertThat(result.manifestInfo().replacedRowsCount()).isEqualTo(0L); + assertThat(result.manifestInfo().modifiedFilesCount()).isEqualTo(0); + assertThat(result.manifestInfo().modifiedRowsCount()).isEqualTo(0L); + } + + @Test + void deleteManifestTrackedFileAdapter() { + ManifestFile manifest = writeManifestFile(ManifestContent.DELETES); + TrackedFile result = TrackedFileAdapters.forManifestFile().wrap(manifest); + + assertThat(result.contentType()).isEqualTo(FileContent.DELETE_MANIFEST); + assertThat(result.formatVersion()).isZero(); + assertThat(result.recordCount()) + .isEqualTo(ADDED_FILES_COUNT + EXISTING_FILES_COUNT + DELETED_FILES_COUNT); Review Comment: The success test calls `newManifestFile(content)`. The count constants live in the helper. I removed the multi-arg overload because the argument list was hard to read. ########## core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java: ########## @@ -671,6 +733,312 @@ void unknownSpecIdThrows() { .hasMessageContaining("Cannot find partition spec for spec ID"); } + @Test + void dataTrackedFileAdapterFromDataFile() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE); + + assertThat(result.tracking()).isNull(); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + assertThatThrownBy(result::formatVersion) + .isInstanceOf(IllegalStateException.class) + .hasMessage("Format version is assigned at write time"); + } + + @Test + void dataTrackedFileAdapterFromExistingManifestEntry() { + TrackedFile result = + TrackedFileAdapters.forDataFile(TABLE_SCHEMA) + .wrap( + newEntry() + .wrapExisting( + SNAPSHOT_ID, DATA_SEQUENCE_NUMBER, FILE_SEQUENCE_NUMBER, DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().dataSequenceNumber()).isEqualTo(DATA_SEQUENCE_NUMBER); + assertThat(result.tracking().fileSequenceNumber()).isEqualTo(FILE_SEQUENCE_NUMBER); + assertThat(result.tracking().firstRowId()).isEqualTo(FIRST_ROW_ID); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterReuse() { + TrackedFileAdapters.DataTrackedFile adapter = TrackedFileAdapters.forDataFile(TABLE_SCHEMA); + + adapter.wrap(DATA_FILE); + assertWrappedDataFileMatchesFileFields(adapter, DATA_FILE); + assertThat(adapter.tracking()).isNull(); + + DataFile file2 = + new GenericDataFile( + UNPARTITIONED_SPEC.specId(), + "s3://bucket/data/file2.parquet", + FileFormat.PARQUET, + PartitionData.EMPTY, + 2048L, + new Metrics(200L, null, null, null, null), + null, + ImmutableList.of(0L), + null, + null); + assignManifestPosition(file2, "s3://bucket/table/manifest-2.parquet", 8L); + adapter.wrap( + newEntry().wrapExisting(SNAPSHOT_ID, DATA_SEQUENCE_NUMBER, FILE_SEQUENCE_NUMBER, file2)); + assertWrappedDataFileMatchesFileFields(adapter, file2); + assertThat(adapter.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertManifestPosition(adapter.tracking(), file2); + } + + @Test + void dataTrackedFileAdapterRejectsNullFile() { + TrackedFileAdapters.DataTrackedFile adapter = TrackedFileAdapters.forDataFile(TABLE_SCHEMA); + assertThatThrownBy(() -> adapter.wrap((DataFile) null)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("Invalid file: null"); + } + + @Test + void dataTrackedFileAdapterContentStats() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE_WITH_METRICS); + + ContentStats stats = result.contentStats(); + assertThat(stats).isNotNull(); + assertThat(stats.fieldStats()).extracting(FieldStats::fieldId).containsExactlyInAnyOrder(1, 2); + + FieldStats<?> idStats = stats.statsFor(1); + assertThat(idStats.valueCount()).isEqualTo(100L); + assertThat(idStats.lowerBound()).isEqualTo(1); + assertThat(idStats.upperBound()).isEqualTo(1000); + } + + @Test + void dataTrackedFileAdapterWithoutMetricsHasNoContentStats() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE); + + assertThat(result.contentStats()).isNull(); + } + + @Test + void dataTrackedFileAdapterKeepsPartitionTuple() { + DataFile partitioned = + new GenericDataFile( + PARTITIONED_SPEC.specId(), + DATA_FILE_LOCATION, + FileFormat.PARQUET, + PARTITION, + 1024L, + new Metrics(100L, null, null, null, null), + null, + ImmutableList.of(0L), + null, + null); + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(partitioned); + + assertThat(result.partition()) + .usingComparator(Comparators.forType(PARTITIONED_SPEC.partitionType())) + .isEqualTo(PARTITION); + assertThat(TrackedFileAdapters.asDataFile(result, specsById(PARTITIONED_SPEC)).partition()) + .usingComparator(Comparators.forType(PARTITIONED_SPEC.partitionType())) + .isEqualTo(PARTITION); + } + + @Test + void dataTrackedFileAdapterUnwrapsToOriginalTrackedFile() { + TrackedFile source = trackedFile(FileContent.DATA); + DataFile dataFile = TrackedFileAdapters.asDataFile(source, UNPARTITIONED); + TrackedFile roundTripped = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(dataFile); + assertThat(roundTripped).isSameAs(source); + } + + @Test + void manifestTrackedFileAdapterUnwrapsToOriginalTrackedFile() { + TrackedFile original = trackedFile(FileContent.DATA_MANIFEST, 0); + ManifestFile adapted = TrackedFileAdapters.asManifestFile(original); + TrackedFile result = TrackedFileAdapters.forManifestFile().wrap(adapted); + assertThat(result).isSameAs(original); + assertThat(result.formatVersion()).isZero(); + } + + @Test + void dataManifestTrackedFileAdapter() { + ManifestFile manifest = writeManifestFile(ManifestContent.DATA); + TrackedFile result = TrackedFileAdapters.forManifestFile().wrap(manifest); + + assertThat(result.contentType()).isEqualTo(FileContent.DATA_MANIFEST); + assertThat(result.formatVersion()).isZero(); + assertThat(result.location()).isEqualTo(MANIFEST_LOCATION); + assertThat(result.fileFormat()).isEqualTo(FileFormat.AVRO); + assertThat(result.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().firstRowId()).isNull(); + assertThat(result.recordCount()) + .isEqualTo(ADDED_FILES_COUNT + EXISTING_FILES_COUNT + DELETED_FILES_COUNT); + assertThat(result.manifestInfo()).isNotNull(); + assertThat(result.manifestInfo().addedFilesCount()).isEqualTo(ADDED_FILES_COUNT); + assertThat(result.manifestInfo().existingFilesCount()).isEqualTo(EXISTING_FILES_COUNT); + assertThat(result.manifestInfo().deletedFilesCount()).isEqualTo(DELETED_FILES_COUNT); + assertThat(result.manifestInfo().addedRowsCount()).isEqualTo(ADDED_ROWS_COUNT); + assertThat(result.manifestInfo().replacedFilesCount()).isEqualTo(0); + assertThat(result.manifestInfo().replacedRowsCount()).isEqualTo(0L); + assertThat(result.manifestInfo().modifiedFilesCount()).isEqualTo(0); + assertThat(result.manifestInfo().modifiedRowsCount()).isEqualTo(0L); + } + + @Test + void deleteManifestTrackedFileAdapter() { + ManifestFile manifest = writeManifestFile(ManifestContent.DELETES); + TrackedFile result = TrackedFileAdapters.forManifestFile().wrap(manifest); + + assertThat(result.contentType()).isEqualTo(FileContent.DELETE_MANIFEST); + assertThat(result.formatVersion()).isZero(); + assertThat(result.recordCount()) + .isEqualTo(ADDED_FILES_COUNT + EXISTING_FILES_COUNT + DELETED_FILES_COUNT); + assertThat(result.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertThat(result.tracking().firstRowId()).isNull(); + } + + @Test + void manifestTrackedFileAdapterFailsWhenAddedFilesCountMissing() { Review Comment: Existing and deleted file counts are in `InvalidCount`, along with null and non-zero replaced and modified file counts. Row counts are not in that enum. ########## core/src/test/java/org/apache/iceberg/TestMapBackedContentStats.java: ########## @@ -0,0 +1,382 @@ +/* + * 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.math.BigDecimal; +import java.nio.ByteBuffer; +import java.util.Comparator; +import java.util.List; +import java.util.Map; +import java.util.UUID; +import java.util.stream.Stream; +import org.apache.iceberg.geospatial.GeospatialBound; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; +import org.apache.iceberg.types.Comparators; +import org.apache.iceberg.types.Conversions; +import org.apache.iceberg.types.Type; +import org.apache.iceberg.types.Types; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.MethodSource; + +class TestMapBackedContentStats { + + private static final Schema SCHEMA = + new Schema( + Types.NestedField.required(1, "id", Types.IntegerType.get()), + Types.NestedField.optional(2, "score", Types.FloatType.get()), + Types.NestedField.optional(3, "ts", Types.LongType.get()), + Types.NestedField.optional(4, "name", Types.StringType.get()), + Types.NestedField.optional(5, "flag", Types.BooleanType.get())); + + private static final PartitionData EMPTY_PARTITION = + new PartitionData(PartitionSpec.unpartitioned().partitionType()); + + /** + * Stats on fields 1-4; field 5 is absent. Field 3 has no value-count entry, so valueCount() + * throws on unboxing. + */ + private static final DataFile FILE_WITH_STATS = + dataFile( + 100L, + ImmutableMap.of(1, 100L, 2, 100L, 4, 100L), + ImmutableMap.of(2, 5L, 3, 1L, 4, 2L), + ImmutableMap.of(2, 3L), + ImmutableMap.of( + 1, buf(Types.IntegerType.get(), 1), + 2, buf(Types.FloatType.get(), 1.5f), + 3, buf(Types.LongType.get(), 100L), + 4, buf(Types.StringType.get(), "aaa")), + ImmutableMap.of( + 1, buf(Types.IntegerType.get(), 1000), + 2, buf(Types.FloatType.get(), 9.5f), + 3, buf(Types.LongType.get(), 999L), + 4, buf(Types.StringType.get(), "zzz"))); + + @Test + void typeMatchesIdsPresentInMaps() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA); + + assertThat(stats.type().fields()).isEmpty(); + + stats.wrap(FILE_WITH_STATS); + + assertThat(stats.type().fields()) + .containsExactlyInAnyOrderElementsOf( + StatsUtil.statsReadSchema(SCHEMA, List.of(1, 2, 3, 4)).fields()); + } + + @Test + void wrapInvalidatesType() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + Types.StructType firstType = stats.type(); + assertThat(firstType.fields()) + .containsExactlyInAnyOrderElementsOf( + StatsUtil.statsReadSchema(SCHEMA, List.of(1, 2, 3, 4)).fields()); + + DataFile file2 = + dataFile( + 100L, + ImmutableMap.of(1, 50L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(1, buf(Types.IntegerType.get(), 500)), + ImmutableMap.of(1, buf(Types.IntegerType.get(), 5000))); + stats.wrap(file2); + + assertThat(stats.type().fields()) + .containsExactlyInAnyOrderElementsOf( + StatsUtil.statsReadSchema(SCHEMA, List.of(1)).fields()); + assertThat(stats.type()).isNotEqualTo(firstType); + } + + @ParameterizedTest + @MethodSource("typesAndBounds") + void boundDecodingPerType(Type type, Object lower, Object upper) { + int fieldId = 1; + Schema schema = new Schema(Types.NestedField.optional(fieldId, "col", type)); + DataFile file = + dataFile( + 100L, + ImmutableMap.of(fieldId, 1L), + ImmutableMap.of(), + ImmutableMap.of(), + bound(fieldId, Conversions.toByteBuffer(type, lower)), + bound(fieldId, Conversions.toByteBuffer(type, upper))); + FieldStats<?> stats = new MapBackedContentStats(schema).wrap(file).statsFor(fieldId); + + Comparator<Object> comparator = Comparators.forType(type.asPrimitiveType()); + assertThat(comparator.compare(stats.lowerBound(), lower)).isZero(); + assertThat(comparator.compare(stats.upperBound(), upper)).isZero(); + } + + private static Stream<Arguments> typesAndBounds() { + return Stream.of( + Arguments.of(Types.BooleanType.get(), false, true), + Arguments.of(Types.IntegerType.get(), -5, 100), + Arguments.of(Types.LongType.get(), 0L, 1_000L), + Arguments.of(Types.FloatType.get(), 1.5f, 9.5f), + Arguments.of(Types.DoubleType.get(), 0.0d, 25.0d), + Arguments.of(Types.DateType.get(), 100, 200), + Arguments.of(Types.TimeType.get(), 1_000L, 2_000L), + Arguments.of(Types.TimestampType.withZone(), 111L, 222L), + Arguments.of(Types.TimestampNanoType.withZone(), 111L, 222L), + Arguments.of(Types.StringType.get(), "a", "z"), + Arguments.of( + Types.UUIDType.get(), + UUID.fromString("07ceab48-62b2-4219-9172-856e687c92ad"), + UUID.fromString("a9a7c24d-2869-4c66-8803-b1f6a36256ed")), + Arguments.of( + Types.FixedType.ofLength(4), + ByteBuffer.wrap(new byte[] {0, 1, 2, 3}), + ByteBuffer.wrap(new byte[] {4, 5, 6, 7})), + Arguments.of( + Types.BinaryType.get(), + ByteBuffer.wrap(new byte[] {1, 2}), + ByteBuffer.wrap(new byte[] {3, 4, 5})), + Arguments.of(Types.DecimalType.of(9, 2), new BigDecimal("1.23"), new BigDecimal("9.99")), + Arguments.of(Types.UnknownType.get(), null, null)); + } + + @Test + void geoBoundsDecode() { + GeospatialBound lower = GeospatialBound.createXY(1.0, 2.0); + GeospatialBound upper = GeospatialBound.createXYZM(3.0, 4.0, 5.0, 6.0); + Schema schema = new Schema(Types.NestedField.optional(10, "geom", Types.GeometryType.crs84())); + DataFile file = + dataFile( + 100L, + ImmutableMap.of(10, 26L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(10, lower.toByteBuffer()), + ImmutableMap.of(10, upper.toByteBuffer())); + MapBackedContentStats stats = new MapBackedContentStats(schema).wrap(file); + + FieldStats<?> geom = stats.statsFor(10); + assertThat(geom.lowerBound()).isInstanceOf(GeospatialBound.class).isEqualTo(lower); + assertThat(geom.upperBound()).isInstanceOf(GeospatialBound.class).isEqualTo(upper); + } + + @Test + void missingBoundsDecodeToNull() { + DataFile file = + dataFile( + 100L, + ImmutableMap.of(1, 100L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of()); + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(file); + + FieldStats<?> id = stats.statsFor(1); + assertThat(id.lowerBound()).isNull(); + assertThat(id.upperBound()).isNull(); + assertThat(id.valueCount()).isEqualTo(100L); + } + + @Test + void countsAndTightBounds() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + + FieldStats<?> id = stats.statsFor(1); + assertThat(id.hasValueCount()).isTrue(); + assertThat(id.valueCount()).isEqualTo(100L); + assertThat(id.hasNullValueCount()).isFalse(); + assertThat(id.hasNanValueCount()).isFalse(); + assertThatThrownBy(id::nullValueCount) + .isInstanceOf(NullPointerException.class) + .hasMessageContaining("Long.longValue()"); + assertThatThrownBy(id::nanValueCount) + .isInstanceOf(NullPointerException.class) + .hasMessageContaining("Long.longValue()"); + // ContentFile maps have no tight-bounds flag so the view always reports false. + assertThat(id.tightBounds()).isFalse(); + // FILE_WITH_STATS has no avg-size map. + assertThat(id.avgValueSizeInBytes()).isNull(); + + FieldStats<?> score = stats.statsFor(2); + assertThat(score.hasNullValueCount()).isTrue(); + assertThat(score.nullValueCount()).isEqualTo(5L); + assertThat(score.hasNanValueCount()).isTrue(); + assertThat(score.nanValueCount()).isEqualTo(3L); + + FieldStats<?> ts = stats.statsFor(3); + assertThat(ts.hasValueCount()).isFalse(); + assertThatThrownBy(ts::valueCount) + .isInstanceOf(NullPointerException.class) + .hasMessageContaining("Long.longValue()"); + assertThat(ts.hasNullValueCount()).isTrue(); + assertThat(ts.nullValueCount()).isEqualTo(1L); + + FieldStats<?> name = stats.statsFor(4); + assertThat(name.hasValueCount()).isTrue(); + assertThat(name.valueCount()).isEqualTo(100L); + assertThat(name.hasNullValueCount()).isTrue(); + assertThat(name.nullValueCount()).isEqualTo(2L); + assertThat(name.hasNanValueCount()).isFalse(); + assertThatThrownBy(name::nanValueCount) + .isInstanceOf(NullPointerException.class) + .hasMessageContaining("Long.longValue()"); + } + + @Test + void fieldWithoutStatsIsExcluded() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + assertThat(stats.type().field(StatsUtil.toBaseId(5))).isNull(); + assertThat(stats.statsFor(5)).isNull(); + } + + @Test + void unknownFieldIdInMaps() { + DataFile file = + dataFile( + 100L, + ImmutableMap.of(99, 1L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of()); + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(file); + + assertThat(stats.statsFor(99)).isNull(); + assertThat(stats.type().fields()).isEmpty(); + assertThat(stats.fieldStats()).isEmpty(); + } + + @Test + void reuseRebindsFields() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA); + + stats.wrap(FILE_WITH_STATS); + assertThat(stats.statsFor(1).lowerBound()).isEqualTo(1); + assertThat(stats.statsFor(1).upperBound()).isEqualTo(1000); + assertThat(stats.statsFor(1).valueCount()).isEqualTo(100L); + + DataFile file2 = + dataFile( + 100L, + ImmutableMap.of(1, 50L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(1, buf(Types.IntegerType.get(), 500)), + ImmutableMap.of(1, buf(Types.IntegerType.get(), 5000))); + stats.wrap(file2); + + assertThat(stats.statsFor(1).lowerBound()).isEqualTo(500); + assertThat(stats.statsFor(1).upperBound()).isEqualTo(5000); + assertThat(stats.statsFor(1).valueCount()).isEqualTo(50L); + assertThat(stats.statsFor(2)).isNull(); Review Comment: Done. `reuseRebindsFields` asserts `statsFor(2)` is not null after `FILE_WITH_STATS`. ########## core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java: ########## @@ -671,6 +733,312 @@ void unknownSpecIdThrows() { .hasMessageContaining("Cannot find partition spec for spec ID"); } + @Test + void dataTrackedFileAdapterFromDataFile() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE); + + assertThat(result.tracking()).isNull(); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + assertThatThrownBy(result::formatVersion) + .isInstanceOf(IllegalStateException.class) + .hasMessage("Format version is assigned at write time"); + } + + @Test + void dataTrackedFileAdapterFromExistingManifestEntry() { + TrackedFile result = + TrackedFileAdapters.forDataFile(TABLE_SCHEMA) + .wrap( + newEntry() + .wrapExisting( + SNAPSHOT_ID, DATA_SEQUENCE_NUMBER, FILE_SEQUENCE_NUMBER, DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().dataSequenceNumber()).isEqualTo(DATA_SEQUENCE_NUMBER); + assertThat(result.tracking().fileSequenceNumber()).isEqualTo(FILE_SEQUENCE_NUMBER); + assertThat(result.tracking().firstRowId()).isEqualTo(FIRST_ROW_ID); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterReuse() { + TrackedFileAdapters.DataTrackedFile adapter = TrackedFileAdapters.forDataFile(TABLE_SCHEMA); + + adapter.wrap(DATA_FILE); + assertWrappedDataFileMatchesFileFields(adapter, DATA_FILE); + assertThat(adapter.tracking()).isNull(); + + DataFile file2 = + new GenericDataFile( + UNPARTITIONED_SPEC.specId(), + "s3://bucket/data/file2.parquet", + FileFormat.PARQUET, + PartitionData.EMPTY, + 2048L, + new Metrics(200L, null, null, null, null), + null, + ImmutableList.of(0L), + null, + null); + assignManifestPosition(file2, "s3://bucket/table/manifest-2.parquet", 8L); + adapter.wrap( + newEntry().wrapExisting(SNAPSHOT_ID, DATA_SEQUENCE_NUMBER, FILE_SEQUENCE_NUMBER, file2)); + assertWrappedDataFileMatchesFileFields(adapter, file2); + assertThat(adapter.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertManifestPosition(adapter.tracking(), file2); + } + + @Test + void dataTrackedFileAdapterRejectsNullFile() { + TrackedFileAdapters.DataTrackedFile adapter = TrackedFileAdapters.forDataFile(TABLE_SCHEMA); + assertThatThrownBy(() -> adapter.wrap((DataFile) null)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("Invalid file: null"); + } + + @Test + void dataTrackedFileAdapterContentStats() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE_WITH_METRICS); + + ContentStats stats = result.contentStats(); + assertThat(stats).isNotNull(); + assertThat(stats.fieldStats()).extracting(FieldStats::fieldId).containsExactlyInAnyOrder(1, 2); + + FieldStats<?> idStats = stats.statsFor(1); + assertThat(idStats.valueCount()).isEqualTo(100L); + assertThat(idStats.lowerBound()).isEqualTo(1); + assertThat(idStats.upperBound()).isEqualTo(1000); + } + + @Test + void dataTrackedFileAdapterWithoutMetricsHasNoContentStats() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE); + + assertThat(result.contentStats()).isNull(); + } + + @Test + void dataTrackedFileAdapterKeepsPartitionTuple() { + DataFile partitioned = + new GenericDataFile( + PARTITIONED_SPEC.specId(), + DATA_FILE_LOCATION, + FileFormat.PARQUET, + PARTITION, + 1024L, + new Metrics(100L, null, null, null, null), + null, + ImmutableList.of(0L), + null, + null); + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(partitioned); + + assertThat(result.partition()) + .usingComparator(Comparators.forType(PARTITIONED_SPEC.partitionType())) + .isEqualTo(PARTITION); + assertThat(TrackedFileAdapters.asDataFile(result, specsById(PARTITIONED_SPEC)).partition()) + .usingComparator(Comparators.forType(PARTITIONED_SPEC.partitionType())) + .isEqualTo(PARTITION); + } + + @Test + void dataTrackedFileAdapterUnwrapsToOriginalTrackedFile() { + TrackedFile source = trackedFile(FileContent.DATA); + DataFile dataFile = TrackedFileAdapters.asDataFile(source, UNPARTITIONED); + TrackedFile roundTripped = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(dataFile); + assertThat(roundTripped).isSameAs(source); + } + + @Test + void manifestTrackedFileAdapterUnwrapsToOriginalTrackedFile() { + TrackedFile original = trackedFile(FileContent.DATA_MANIFEST, 0); + ManifestFile adapted = TrackedFileAdapters.asManifestFile(original); + TrackedFile result = TrackedFileAdapters.forManifestFile().wrap(adapted); + assertThat(result).isSameAs(original); + assertThat(result.formatVersion()).isZero(); + } + + @Test + void dataManifestTrackedFileAdapter() { + ManifestFile manifest = writeManifestFile(ManifestContent.DATA); + TrackedFile result = TrackedFileAdapters.forManifestFile().wrap(manifest); + + assertThat(result.contentType()).isEqualTo(FileContent.DATA_MANIFEST); + assertThat(result.formatVersion()).isZero(); + assertThat(result.location()).isEqualTo(MANIFEST_LOCATION); + assertThat(result.fileFormat()).isEqualTo(FileFormat.AVRO); + assertThat(result.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().firstRowId()).isNull(); + assertThat(result.recordCount()) + .isEqualTo(ADDED_FILES_COUNT + EXISTING_FILES_COUNT + DELETED_FILES_COUNT); + assertThat(result.manifestInfo()).isNotNull(); + assertThat(result.manifestInfo().addedFilesCount()).isEqualTo(ADDED_FILES_COUNT); + assertThat(result.manifestInfo().existingFilesCount()).isEqualTo(EXISTING_FILES_COUNT); + assertThat(result.manifestInfo().deletedFilesCount()).isEqualTo(DELETED_FILES_COUNT); + assertThat(result.manifestInfo().addedRowsCount()).isEqualTo(ADDED_ROWS_COUNT); + assertThat(result.manifestInfo().replacedFilesCount()).isEqualTo(0); + assertThat(result.manifestInfo().replacedRowsCount()).isEqualTo(0L); + assertThat(result.manifestInfo().modifiedFilesCount()).isEqualTo(0); + assertThat(result.manifestInfo().modifiedRowsCount()).isEqualTo(0L); + } + + @Test + void deleteManifestTrackedFileAdapter() { + ManifestFile manifest = writeManifestFile(ManifestContent.DELETES); + TrackedFile result = TrackedFileAdapters.forManifestFile().wrap(manifest); + + assertThat(result.contentType()).isEqualTo(FileContent.DELETE_MANIFEST); + assertThat(result.formatVersion()).isZero(); + assertThat(result.recordCount()) + .isEqualTo(ADDED_FILES_COUNT + EXISTING_FILES_COUNT + DELETED_FILES_COUNT); + assertThat(result.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertThat(result.tracking().firstRowId()).isNull(); + } + + @Test + void manifestTrackedFileAdapterFailsWhenAddedFilesCountMissing() { + ManifestFile manifest = mock(ManifestFile.class); + when(manifest.path()).thenReturn(MANIFEST_LOCATION); + when(manifest.content()).thenReturn(ManifestContent.DATA); + when(manifest.addedFilesCount()).thenReturn(null); + when(manifest.existingFilesCount()).thenReturn(1); + when(manifest.deletedFilesCount()).thenReturn(0); + when(manifest.replacedFilesCount()).thenReturn(0); + + assertThatThrownBy(() -> TrackedFileAdapters.forManifestFile().wrap(manifest)) + .isInstanceOf(NullPointerException.class) + .hasMessageContaining("missing added files count"); + } + + @Test + void manifestTrackedFileAdapterRejectsNullReplacedFilesCount() { + ManifestFile manifest = manifestWithCounts(1, 1, 0, null, 0); + + assertThatThrownBy(() -> TrackedFileAdapters.forManifestFile().wrap(manifest)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage( + "Cannot convert manifest %s: Invalid replaced file count: null", MANIFEST_LOCATION); + } + + @Test + void manifestTrackedFileAdapterRejectsNonZeroReplacedFilesCount() { + ManifestFile manifest = manifestWithCounts(1, 1, 0, 1, 0); + + assertThatThrownBy(() -> TrackedFileAdapters.forManifestFile().wrap(manifest)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage( + "Cannot convert manifest %s: Invalid replaced file count: 1", MANIFEST_LOCATION); + } + + @Test + void manifestTrackedFileAdapterRejectsNullModifiedFilesCount() { + ManifestFile manifest = manifestWithCounts(1, 1, 0, 0, null); + + assertThatThrownBy(() -> TrackedFileAdapters.forManifestFile().wrap(manifest)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage( + "Cannot convert manifest %s: Invalid modified file count: null", MANIFEST_LOCATION); + } + + @Test + void manifestTrackedFileAdapterRejectsNonZeroModifiedFilesCount() { + ManifestFile manifest = manifestWithCounts(1, 1, 0, 0, 1); + + assertThatThrownBy(() -> TrackedFileAdapters.forManifestFile().wrap(manifest)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage( + "Cannot convert manifest %s: Invalid modified file count: 1", MANIFEST_LOCATION); + } + + @Test + void manifestTrackedFileAdapterFailsWhenAddedRowsCountMissing() { Review Comment: `writeManifestFile` is now `newManifestFile`. `wrap()` calls `manifestRecordCount` that reads file counts. Invalid file counts would cause wrap to fail. File-count failures are now covered by `manifestTrackedFileAdapterRejectsInvalidFilesCount`. The file-count cases now share one mock, `newManifestFileWithInvalidCount`. It stubs every count to the normal constant, then overrides only the injected file-count failure. I dropped the row-count cases, which were testing NPE when accessing null row counts after wrap. that is not an interesting test scenario. ########## core/src/test/java/org/apache/iceberg/TestMapBackedContentStats.java: ########## @@ -0,0 +1,382 @@ +/* + * 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.math.BigDecimal; +import java.nio.ByteBuffer; +import java.util.Comparator; +import java.util.List; +import java.util.Map; +import java.util.UUID; +import java.util.stream.Stream; +import org.apache.iceberg.geospatial.GeospatialBound; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList; +import org.apache.iceberg.relocated.com.google.common.collect.ImmutableMap; +import org.apache.iceberg.types.Comparators; +import org.apache.iceberg.types.Conversions; +import org.apache.iceberg.types.Type; +import org.apache.iceberg.types.Types; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.MethodSource; + +class TestMapBackedContentStats { + + private static final Schema SCHEMA = + new Schema( + Types.NestedField.required(1, "id", Types.IntegerType.get()), + Types.NestedField.optional(2, "score", Types.FloatType.get()), + Types.NestedField.optional(3, "ts", Types.LongType.get()), + Types.NestedField.optional(4, "name", Types.StringType.get()), + Types.NestedField.optional(5, "flag", Types.BooleanType.get())); + + private static final PartitionData EMPTY_PARTITION = + new PartitionData(PartitionSpec.unpartitioned().partitionType()); + + /** + * Stats on fields 1-4; field 5 is absent. Field 3 has no value-count entry, so valueCount() + * throws on unboxing. + */ + private static final DataFile FILE_WITH_STATS = + dataFile( + 100L, + ImmutableMap.of(1, 100L, 2, 100L, 4, 100L), + ImmutableMap.of(2, 5L, 3, 1L, 4, 2L), + ImmutableMap.of(2, 3L), + ImmutableMap.of( + 1, buf(Types.IntegerType.get(), 1), + 2, buf(Types.FloatType.get(), 1.5f), + 3, buf(Types.LongType.get(), 100L), + 4, buf(Types.StringType.get(), "aaa")), + ImmutableMap.of( + 1, buf(Types.IntegerType.get(), 1000), + 2, buf(Types.FloatType.get(), 9.5f), + 3, buf(Types.LongType.get(), 999L), + 4, buf(Types.StringType.get(), "zzz"))); + + @Test + void typeMatchesIdsPresentInMaps() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA); + + assertThat(stats.type().fields()).isEmpty(); + + stats.wrap(FILE_WITH_STATS); + + assertThat(stats.type().fields()) + .containsExactlyInAnyOrderElementsOf( + StatsUtil.statsReadSchema(SCHEMA, List.of(1, 2, 3, 4)).fields()); + } + + @Test + void wrapInvalidatesType() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + Types.StructType firstType = stats.type(); + assertThat(firstType.fields()) + .containsExactlyInAnyOrderElementsOf( + StatsUtil.statsReadSchema(SCHEMA, List.of(1, 2, 3, 4)).fields()); + + DataFile file2 = + dataFile( + 100L, + ImmutableMap.of(1, 50L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(1, buf(Types.IntegerType.get(), 500)), + ImmutableMap.of(1, buf(Types.IntegerType.get(), 5000))); + stats.wrap(file2); + + assertThat(stats.type().fields()) + .containsExactlyInAnyOrderElementsOf( + StatsUtil.statsReadSchema(SCHEMA, List.of(1)).fields()); + assertThat(stats.type()).isNotEqualTo(firstType); + } + + @ParameterizedTest + @MethodSource("typesAndBounds") + void boundDecodingPerType(Type type, Object lower, Object upper) { + int fieldId = 1; + Schema schema = new Schema(Types.NestedField.optional(fieldId, "col", type)); + DataFile file = + dataFile( + 100L, + ImmutableMap.of(fieldId, 1L), + ImmutableMap.of(), + ImmutableMap.of(), + bound(fieldId, Conversions.toByteBuffer(type, lower)), + bound(fieldId, Conversions.toByteBuffer(type, upper))); + FieldStats<?> stats = new MapBackedContentStats(schema).wrap(file).statsFor(fieldId); + + Comparator<Object> comparator = Comparators.forType(type.asPrimitiveType()); + assertThat(comparator.compare(stats.lowerBound(), lower)).isZero(); + assertThat(comparator.compare(stats.upperBound(), upper)).isZero(); + } + + private static Stream<Arguments> typesAndBounds() { + return Stream.of( + Arguments.of(Types.BooleanType.get(), false, true), + Arguments.of(Types.IntegerType.get(), -5, 100), + Arguments.of(Types.LongType.get(), 0L, 1_000L), + Arguments.of(Types.FloatType.get(), 1.5f, 9.5f), + Arguments.of(Types.DoubleType.get(), 0.0d, 25.0d), + Arguments.of(Types.DateType.get(), 100, 200), + Arguments.of(Types.TimeType.get(), 1_000L, 2_000L), + Arguments.of(Types.TimestampType.withZone(), 111L, 222L), + Arguments.of(Types.TimestampNanoType.withZone(), 111L, 222L), + Arguments.of(Types.StringType.get(), "a", "z"), + Arguments.of( + Types.UUIDType.get(), + UUID.fromString("07ceab48-62b2-4219-9172-856e687c92ad"), + UUID.fromString("a9a7c24d-2869-4c66-8803-b1f6a36256ed")), + Arguments.of( + Types.FixedType.ofLength(4), + ByteBuffer.wrap(new byte[] {0, 1, 2, 3}), + ByteBuffer.wrap(new byte[] {4, 5, 6, 7})), + Arguments.of( + Types.BinaryType.get(), + ByteBuffer.wrap(new byte[] {1, 2}), + ByteBuffer.wrap(new byte[] {3, 4, 5})), + Arguments.of(Types.DecimalType.of(9, 2), new BigDecimal("1.23"), new BigDecimal("9.99")), + Arguments.of(Types.UnknownType.get(), null, null)); + } + + @Test + void geoBoundsDecode() { + GeospatialBound lower = GeospatialBound.createXY(1.0, 2.0); + GeospatialBound upper = GeospatialBound.createXYZM(3.0, 4.0, 5.0, 6.0); + Schema schema = new Schema(Types.NestedField.optional(10, "geom", Types.GeometryType.crs84())); + DataFile file = + dataFile( + 100L, + ImmutableMap.of(10, 26L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(10, lower.toByteBuffer()), + ImmutableMap.of(10, upper.toByteBuffer())); + MapBackedContentStats stats = new MapBackedContentStats(schema).wrap(file); + + FieldStats<?> geom = stats.statsFor(10); + assertThat(geom.lowerBound()).isInstanceOf(GeospatialBound.class).isEqualTo(lower); + assertThat(geom.upperBound()).isInstanceOf(GeospatialBound.class).isEqualTo(upper); + } + + @Test + void missingBoundsDecodeToNull() { + DataFile file = + dataFile( + 100L, + ImmutableMap.of(1, 100L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of()); + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(file); + + FieldStats<?> id = stats.statsFor(1); + assertThat(id.lowerBound()).isNull(); + assertThat(id.upperBound()).isNull(); + assertThat(id.valueCount()).isEqualTo(100L); + } + + @Test + void countsAndTightBounds() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + + FieldStats<?> id = stats.statsFor(1); + assertThat(id.hasValueCount()).isTrue(); + assertThat(id.valueCount()).isEqualTo(100L); + assertThat(id.hasNullValueCount()).isFalse(); + assertThat(id.hasNanValueCount()).isFalse(); + assertThatThrownBy(id::nullValueCount) + .isInstanceOf(NullPointerException.class) + .hasMessageContaining("Long.longValue()"); + assertThatThrownBy(id::nanValueCount) + .isInstanceOf(NullPointerException.class) + .hasMessageContaining("Long.longValue()"); + // ContentFile maps have no tight-bounds flag so the view always reports false. + assertThat(id.tightBounds()).isFalse(); + // FILE_WITH_STATS has no avg-size map. + assertThat(id.avgValueSizeInBytes()).isNull(); + + FieldStats<?> score = stats.statsFor(2); + assertThat(score.hasNullValueCount()).isTrue(); + assertThat(score.nullValueCount()).isEqualTo(5L); + assertThat(score.hasNanValueCount()).isTrue(); + assertThat(score.nanValueCount()).isEqualTo(3L); + + FieldStats<?> ts = stats.statsFor(3); + assertThat(ts.hasValueCount()).isFalse(); + assertThatThrownBy(ts::valueCount) + .isInstanceOf(NullPointerException.class) + .hasMessageContaining("Long.longValue()"); + assertThat(ts.hasNullValueCount()).isTrue(); + assertThat(ts.nullValueCount()).isEqualTo(1L); + + FieldStats<?> name = stats.statsFor(4); + assertThat(name.hasValueCount()).isTrue(); + assertThat(name.valueCount()).isEqualTo(100L); + assertThat(name.hasNullValueCount()).isTrue(); + assertThat(name.nullValueCount()).isEqualTo(2L); + assertThat(name.hasNanValueCount()).isFalse(); + assertThatThrownBy(name::nanValueCount) + .isInstanceOf(NullPointerException.class) + .hasMessageContaining("Long.longValue()"); + } + + @Test + void fieldWithoutStatsIsExcluded() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + assertThat(stats.type().field(StatsUtil.toBaseId(5))).isNull(); + assertThat(stats.statsFor(5)).isNull(); + } + + @Test + void unknownFieldIdInMaps() { + DataFile file = + dataFile( + 100L, + ImmutableMap.of(99, 1L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of()); + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(file); + + assertThat(stats.statsFor(99)).isNull(); + assertThat(stats.type().fields()).isEmpty(); + assertThat(stats.fieldStats()).isEmpty(); + } + + @Test + void reuseRebindsFields() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA); + + stats.wrap(FILE_WITH_STATS); + assertThat(stats.statsFor(1).lowerBound()).isEqualTo(1); + assertThat(stats.statsFor(1).upperBound()).isEqualTo(1000); + assertThat(stats.statsFor(1).valueCount()).isEqualTo(100L); + + DataFile file2 = + dataFile( + 100L, + ImmutableMap.of(1, 50L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(1, buf(Types.IntegerType.get(), 500)), + ImmutableMap.of(1, buf(Types.IntegerType.get(), 5000))); + stats.wrap(file2); + + assertThat(stats.statsFor(1).lowerBound()).isEqualTo(500); + assertThat(stats.statsFor(1).upperBound()).isEqualTo(5000); + assertThat(stats.statsFor(1).valueCount()).isEqualTo(50L); + assertThat(stats.statsFor(2)).isNull(); + } + + @Test + void absentIdIsRereadWhenPresentAgain() { Review Comment: Done. Folded this into `reuseRebindsFields` and removed `absentIdIsRereadWhenPresentAgain`. Field 5 is absent on `FILE_WITH_STATS` and present on file2. Field 2 stays off file2. ########## core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java: ########## @@ -671,6 +733,312 @@ void unknownSpecIdThrows() { .hasMessageContaining("Cannot find partition spec for spec ID"); } + @Test + void dataTrackedFileAdapterFromDataFile() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE); + + assertThat(result.tracking()).isNull(); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + assertThatThrownBy(result::formatVersion) + .isInstanceOf(IllegalStateException.class) + .hasMessage("Format version is assigned at write time"); + } + + @Test + void dataTrackedFileAdapterFromExistingManifestEntry() { + TrackedFile result = + TrackedFileAdapters.forDataFile(TABLE_SCHEMA) + .wrap( + newEntry() + .wrapExisting( + SNAPSHOT_ID, DATA_SEQUENCE_NUMBER, FILE_SEQUENCE_NUMBER, DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().dataSequenceNumber()).isEqualTo(DATA_SEQUENCE_NUMBER); + assertThat(result.tracking().fileSequenceNumber()).isEqualTo(FILE_SEQUENCE_NUMBER); + assertThat(result.tracking().firstRowId()).isEqualTo(FIRST_ROW_ID); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterReuse() { + TrackedFileAdapters.DataTrackedFile adapter = TrackedFileAdapters.forDataFile(TABLE_SCHEMA); + + adapter.wrap(DATA_FILE); + assertWrappedDataFileMatchesFileFields(adapter, DATA_FILE); + assertThat(adapter.tracking()).isNull(); + + DataFile file2 = + new GenericDataFile( + UNPARTITIONED_SPEC.specId(), + "s3://bucket/data/file2.parquet", + FileFormat.PARQUET, + PartitionData.EMPTY, + 2048L, + new Metrics(200L, null, null, null, null), + null, + ImmutableList.of(0L), + null, + null); + assignManifestPosition(file2, "s3://bucket/table/manifest-2.parquet", 8L); + adapter.wrap( + newEntry().wrapExisting(SNAPSHOT_ID, DATA_SEQUENCE_NUMBER, FILE_SEQUENCE_NUMBER, file2)); + assertWrappedDataFileMatchesFileFields(adapter, file2); + assertThat(adapter.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertManifestPosition(adapter.tracking(), file2); + } + + @Test + void dataTrackedFileAdapterRejectsNullFile() { + TrackedFileAdapters.DataTrackedFile adapter = TrackedFileAdapters.forDataFile(TABLE_SCHEMA); + assertThatThrownBy(() -> adapter.wrap((DataFile) null)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("Invalid file: null"); + } + + @Test + void dataTrackedFileAdapterContentStats() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE_WITH_METRICS); + + ContentStats stats = result.contentStats(); + assertThat(stats).isNotNull(); + assertThat(stats.fieldStats()).extracting(FieldStats::fieldId).containsExactlyInAnyOrder(1, 2); + + FieldStats<?> idStats = stats.statsFor(1); + assertThat(idStats.valueCount()).isEqualTo(100L); + assertThat(idStats.lowerBound()).isEqualTo(1); + assertThat(idStats.upperBound()).isEqualTo(1000); + } + + @Test + void dataTrackedFileAdapterWithoutMetricsHasNoContentStats() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE); + + assertThat(result.contentStats()).isNull(); + } + + @Test + void dataTrackedFileAdapterKeepsPartitionTuple() { + DataFile partitioned = + new GenericDataFile( + PARTITIONED_SPEC.specId(), + DATA_FILE_LOCATION, + FileFormat.PARQUET, + PARTITION, + 1024L, + new Metrics(100L, null, null, null, null), + null, + ImmutableList.of(0L), + null, + null); + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(partitioned); + + assertThat(result.partition()) + .usingComparator(Comparators.forType(PARTITIONED_SPEC.partitionType())) + .isEqualTo(PARTITION); + assertThat(TrackedFileAdapters.asDataFile(result, specsById(PARTITIONED_SPEC)).partition()) + .usingComparator(Comparators.forType(PARTITIONED_SPEC.partitionType())) + .isEqualTo(PARTITION); + } + + @Test + void dataTrackedFileAdapterUnwrapsToOriginalTrackedFile() { + TrackedFile source = trackedFile(FileContent.DATA); + DataFile dataFile = TrackedFileAdapters.asDataFile(source, UNPARTITIONED); + TrackedFile roundTripped = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(dataFile); + assertThat(roundTripped).isSameAs(source); + } + + @Test + void manifestTrackedFileAdapterUnwrapsToOriginalTrackedFile() { + TrackedFile original = trackedFile(FileContent.DATA_MANIFEST, 0); + ManifestFile adapted = TrackedFileAdapters.asManifestFile(original); + TrackedFile result = TrackedFileAdapters.forManifestFile().wrap(adapted); + assertThat(result).isSameAs(original); + assertThat(result.formatVersion()).isZero(); + } + + @Test + void dataManifestTrackedFileAdapter() { + ManifestFile manifest = writeManifestFile(ManifestContent.DATA); + TrackedFile result = TrackedFileAdapters.forManifestFile().wrap(manifest); + + assertThat(result.contentType()).isEqualTo(FileContent.DATA_MANIFEST); + assertThat(result.formatVersion()).isZero(); + assertThat(result.location()).isEqualTo(MANIFEST_LOCATION); + assertThat(result.fileFormat()).isEqualTo(FileFormat.AVRO); + assertThat(result.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().firstRowId()).isNull(); + assertThat(result.recordCount()) + .isEqualTo(ADDED_FILES_COUNT + EXISTING_FILES_COUNT + DELETED_FILES_COUNT); + assertThat(result.manifestInfo()).isNotNull(); + assertThat(result.manifestInfo().addedFilesCount()).isEqualTo(ADDED_FILES_COUNT); + assertThat(result.manifestInfo().existingFilesCount()).isEqualTo(EXISTING_FILES_COUNT); + assertThat(result.manifestInfo().deletedFilesCount()).isEqualTo(DELETED_FILES_COUNT); + assertThat(result.manifestInfo().addedRowsCount()).isEqualTo(ADDED_ROWS_COUNT); + assertThat(result.manifestInfo().replacedFilesCount()).isEqualTo(0); + assertThat(result.manifestInfo().replacedRowsCount()).isEqualTo(0L); + assertThat(result.manifestInfo().modifiedFilesCount()).isEqualTo(0); + assertThat(result.manifestInfo().modifiedRowsCount()).isEqualTo(0L); + } + + @Test + void deleteManifestTrackedFileAdapter() { + ManifestFile manifest = writeManifestFile(ManifestContent.DELETES); + TrackedFile result = TrackedFileAdapters.forManifestFile().wrap(manifest); + + assertThat(result.contentType()).isEqualTo(FileContent.DELETE_MANIFEST); + assertThat(result.formatVersion()).isZero(); + assertThat(result.recordCount()) + .isEqualTo(ADDED_FILES_COUNT + EXISTING_FILES_COUNT + DELETED_FILES_COUNT); + assertThat(result.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertThat(result.tracking().firstRowId()).isNull(); + } + + @Test + void manifestTrackedFileAdapterFailsWhenAddedFilesCountMissing() { + ManifestFile manifest = mock(ManifestFile.class); + when(manifest.path()).thenReturn(MANIFEST_LOCATION); + when(manifest.content()).thenReturn(ManifestContent.DATA); + when(manifest.addedFilesCount()).thenReturn(null); + when(manifest.existingFilesCount()).thenReturn(1); + when(manifest.deletedFilesCount()).thenReturn(0); + when(manifest.replacedFilesCount()).thenReturn(0); + + assertThatThrownBy(() -> TrackedFileAdapters.forManifestFile().wrap(manifest)) + .isInstanceOf(NullPointerException.class) + .hasMessageContaining("missing added files count"); + } + + @Test + void manifestTrackedFileAdapterRejectsNullReplacedFilesCount() { + ManifestFile manifest = manifestWithCounts(1, 1, 0, null, 0); + + assertThatThrownBy(() -> TrackedFileAdapters.forManifestFile().wrap(manifest)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage( + "Cannot convert manifest %s: Invalid replaced file count: null", MANIFEST_LOCATION); + } + + @Test + void manifestTrackedFileAdapterRejectsNonZeroReplacedFilesCount() { + ManifestFile manifest = manifestWithCounts(1, 1, 0, 1, 0); + + assertThatThrownBy(() -> TrackedFileAdapters.forManifestFile().wrap(manifest)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage( + "Cannot convert manifest %s: Invalid replaced file count: 1", MANIFEST_LOCATION); + } + + @Test + void manifestTrackedFileAdapterRejectsNullModifiedFilesCount() { + ManifestFile manifest = manifestWithCounts(1, 1, 0, 0, null); + + assertThatThrownBy(() -> TrackedFileAdapters.forManifestFile().wrap(manifest)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage( + "Cannot convert manifest %s: Invalid modified file count: null", MANIFEST_LOCATION); + } + + @Test + void manifestTrackedFileAdapterRejectsNonZeroModifiedFilesCount() { + ManifestFile manifest = manifestWithCounts(1, 1, 0, 0, 1); + + assertThatThrownBy(() -> TrackedFileAdapters.forManifestFile().wrap(manifest)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage( + "Cannot convert manifest %s: Invalid modified file count: 1", MANIFEST_LOCATION); + } + + @Test + void manifestTrackedFileAdapterFailsWhenAddedRowsCountMissing() { + ManifestFile manifest = + writeManifestFile( + ManifestContent.DATA, MANIFEST_SEQUENCE_NUMBER, MANIFEST_MIN_SEQUENCE_NUMBER, null); + TrackedFile tracked = TrackedFileAdapters.forManifestFile().wrap(manifest); + + assertThatThrownBy(() -> tracked.manifestInfo().addedRowsCount()) + .isInstanceOf(NullPointerException.class) + .hasMessageContaining("null"); + } + + private static void assertWrappedDataFileMatchesFileFields(TrackedFile result, DataFile file) { + assertThat(result.contentType()).isEqualTo(FileContent.DATA); + assertThat(result.location()).isEqualTo(file.location()); + assertThat(result.fileFormat()).isEqualTo(file.format()); + assertThat(result.recordCount()).isEqualTo(file.recordCount()); + assertThat(result.fileSizeInBytes()).isEqualTo(file.fileSizeInBytes()); + assertThat(result.specId()).isEqualTo(file.specId()); + assertThat(result.sortOrderId()).isEqualTo(file.sortOrderId()); + assertThat(result.keyMetadata()).isEqualTo(file.keyMetadata()); + assertThat(result.splitOffsets()).isEqualTo(file.splitOffsets()); + assertThat(result.manifestInfo()).isNull(); + assertThat(result.deletionVector()).isNull(); + assertThat(result.equalityIds()).isNull(); + } + + private static ManifestFile manifestWithCounts( + Integer addedFilesCount, + Integer existingFilesCount, + Integer deletedFilesCount, + Integer replacedFilesCount, + Integer modifiedFilesCount) { + ManifestFile manifest = mock(ManifestFile.class); + when(manifest.path()).thenReturn(MANIFEST_LOCATION); + when(manifest.content()).thenReturn(ManifestContent.DATA); + when(manifest.addedFilesCount()).thenReturn(addedFilesCount); + when(manifest.existingFilesCount()).thenReturn(existingFilesCount); + when(manifest.deletedFilesCount()).thenReturn(deletedFilesCount); + when(manifest.replacedFilesCount()).thenReturn(replacedFilesCount); + when(manifest.modifiedFilesCount()).thenReturn(modifiedFilesCount); + return manifest; + } + + private static ManifestFile writeManifestFile(ManifestContent content) { + return writeManifestFile( + content, MANIFEST_SEQUENCE_NUMBER, MANIFEST_MIN_SEQUENCE_NUMBER, ADDED_ROWS_COUNT); + } + + private static ManifestFile writeManifestFile( Review Comment: Renamed to `newManifestFile`. Only `newManifestFile(ManifestContent)` remains, and the counts are constants inside the helper. ########## core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java: ########## @@ -671,6 +733,312 @@ void unknownSpecIdThrows() { .hasMessageContaining("Cannot find partition spec for spec ID"); } + @Test + void dataTrackedFileAdapterFromDataFile() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE); + + assertThat(result.tracking()).isNull(); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + assertThatThrownBy(result::formatVersion) + .isInstanceOf(IllegalStateException.class) + .hasMessage("Format version is assigned at write time"); + } + + @Test + void dataTrackedFileAdapterFromExistingManifestEntry() { + TrackedFile result = + TrackedFileAdapters.forDataFile(TABLE_SCHEMA) + .wrap( + newEntry() + .wrapExisting( + SNAPSHOT_ID, DATA_SEQUENCE_NUMBER, FILE_SEQUENCE_NUMBER, DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().dataSequenceNumber()).isEqualTo(DATA_SEQUENCE_NUMBER); + assertThat(result.tracking().fileSequenceNumber()).isEqualTo(FILE_SEQUENCE_NUMBER); + assertThat(result.tracking().firstRowId()).isEqualTo(FIRST_ROW_ID); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterReuse() { + TrackedFileAdapters.DataTrackedFile adapter = TrackedFileAdapters.forDataFile(TABLE_SCHEMA); + + adapter.wrap(DATA_FILE); + assertWrappedDataFileMatchesFileFields(adapter, DATA_FILE); + assertThat(adapter.tracking()).isNull(); + + DataFile file2 = + new GenericDataFile( + UNPARTITIONED_SPEC.specId(), + "s3://bucket/data/file2.parquet", + FileFormat.PARQUET, + PartitionData.EMPTY, + 2048L, + new Metrics(200L, null, null, null, null), + null, + ImmutableList.of(0L), + null, + null); + assignManifestPosition(file2, "s3://bucket/table/manifest-2.parquet", 8L); + adapter.wrap( + newEntry().wrapExisting(SNAPSHOT_ID, DATA_SEQUENCE_NUMBER, FILE_SEQUENCE_NUMBER, file2)); + assertWrappedDataFileMatchesFileFields(adapter, file2); + assertThat(adapter.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertManifestPosition(adapter.tracking(), file2); + } + + @Test + void dataTrackedFileAdapterRejectsNullFile() { + TrackedFileAdapters.DataTrackedFile adapter = TrackedFileAdapters.forDataFile(TABLE_SCHEMA); + assertThatThrownBy(() -> adapter.wrap((DataFile) null)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("Invalid file: null"); + } + + @Test + void dataTrackedFileAdapterContentStats() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE_WITH_METRICS); + + ContentStats stats = result.contentStats(); + assertThat(stats).isNotNull(); + assertThat(stats.fieldStats()).extracting(FieldStats::fieldId).containsExactlyInAnyOrder(1, 2); + + FieldStats<?> idStats = stats.statsFor(1); + assertThat(idStats.valueCount()).isEqualTo(100L); + assertThat(idStats.lowerBound()).isEqualTo(1); + assertThat(idStats.upperBound()).isEqualTo(1000); + } + + @Test + void dataTrackedFileAdapterWithoutMetricsHasNoContentStats() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE); + + assertThat(result.contentStats()).isNull(); + } + + @Test + void dataTrackedFileAdapterKeepsPartitionTuple() { + DataFile partitioned = + new GenericDataFile( + PARTITIONED_SPEC.specId(), + DATA_FILE_LOCATION, + FileFormat.PARQUET, + PARTITION, + 1024L, + new Metrics(100L, null, null, null, null), + null, + ImmutableList.of(0L), + null, + null); + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(partitioned); + + assertThat(result.partition()) + .usingComparator(Comparators.forType(PARTITIONED_SPEC.partitionType())) + .isEqualTo(PARTITION); + assertThat(TrackedFileAdapters.asDataFile(result, specsById(PARTITIONED_SPEC)).partition()) Review Comment: Removed the `asDataFile` partition assert. The partition check remains on `dataTrackedFileAdapterKeepsPartitionTuple`, using `Comparators.forType`. ########## core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java: ########## @@ -671,6 +733,312 @@ void unknownSpecIdThrows() { .hasMessageContaining("Cannot find partition spec for spec ID"); } + @Test + void dataTrackedFileAdapterFromDataFile() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE); + + assertThat(result.tracking()).isNull(); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + assertThatThrownBy(result::formatVersion) + .isInstanceOf(IllegalStateException.class) + .hasMessage("Format version is assigned at write time"); + } + + @Test + void dataTrackedFileAdapterFromExistingManifestEntry() { + TrackedFile result = + TrackedFileAdapters.forDataFile(TABLE_SCHEMA) + .wrap( + newEntry() + .wrapExisting( + SNAPSHOT_ID, DATA_SEQUENCE_NUMBER, FILE_SEQUENCE_NUMBER, DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().dataSequenceNumber()).isEqualTo(DATA_SEQUENCE_NUMBER); + assertThat(result.tracking().fileSequenceNumber()).isEqualTo(FILE_SEQUENCE_NUMBER); + assertThat(result.tracking().firstRowId()).isEqualTo(FIRST_ROW_ID); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterReuse() { + TrackedFileAdapters.DataTrackedFile adapter = TrackedFileAdapters.forDataFile(TABLE_SCHEMA); + + adapter.wrap(DATA_FILE); + assertWrappedDataFileMatchesFileFields(adapter, DATA_FILE); + assertThat(adapter.tracking()).isNull(); + + DataFile file2 = + new GenericDataFile( + UNPARTITIONED_SPEC.specId(), + "s3://bucket/data/file2.parquet", + FileFormat.PARQUET, + PartitionData.EMPTY, + 2048L, + new Metrics(200L, null, null, null, null), + null, + ImmutableList.of(0L), + null, + null); + assignManifestPosition(file2, "s3://bucket/table/manifest-2.parquet", 8L); + adapter.wrap( + newEntry().wrapExisting(SNAPSHOT_ID, DATA_SEQUENCE_NUMBER, FILE_SEQUENCE_NUMBER, file2)); + assertWrappedDataFileMatchesFileFields(adapter, file2); + assertThat(adapter.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertManifestPosition(adapter.tracking(), file2); + } + + @Test + void dataTrackedFileAdapterRejectsNullFile() { + TrackedFileAdapters.DataTrackedFile adapter = TrackedFileAdapters.forDataFile(TABLE_SCHEMA); + assertThatThrownBy(() -> adapter.wrap((DataFile) null)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("Invalid file: null"); + } + + @Test + void dataTrackedFileAdapterContentStats() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE_WITH_METRICS); + + ContentStats stats = result.contentStats(); + assertThat(stats).isNotNull(); + assertThat(stats.fieldStats()).extracting(FieldStats::fieldId).containsExactlyInAnyOrder(1, 2); + + FieldStats<?> idStats = stats.statsFor(1); + assertThat(idStats.valueCount()).isEqualTo(100L); + assertThat(idStats.lowerBound()).isEqualTo(1); + assertThat(idStats.upperBound()).isEqualTo(1000); + } + + @Test + void dataTrackedFileAdapterWithoutMetricsHasNoContentStats() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE); + + assertThat(result.contentStats()).isNull(); + } + + @Test + void dataTrackedFileAdapterKeepsPartitionTuple() { + DataFile partitioned = + new GenericDataFile( + PARTITIONED_SPEC.specId(), + DATA_FILE_LOCATION, + FileFormat.PARQUET, + PARTITION, + 1024L, + new Metrics(100L, null, null, null, null), + null, + ImmutableList.of(0L), + null, + null); + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(partitioned); + + assertThat(result.partition()) + .usingComparator(Comparators.forType(PARTITIONED_SPEC.partitionType())) + .isEqualTo(PARTITION); + assertThat(TrackedFileAdapters.asDataFile(result, specsById(PARTITIONED_SPEC)).partition()) + .usingComparator(Comparators.forType(PARTITIONED_SPEC.partitionType())) + .isEqualTo(PARTITION); + } + + @Test + void dataTrackedFileAdapterUnwrapsToOriginalTrackedFile() { + TrackedFile source = trackedFile(FileContent.DATA); + DataFile dataFile = TrackedFileAdapters.asDataFile(source, UNPARTITIONED); + TrackedFile roundTripped = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(dataFile); + assertThat(roundTripped).isSameAs(source); + } + + @Test + void manifestTrackedFileAdapterUnwrapsToOriginalTrackedFile() { + TrackedFile original = trackedFile(FileContent.DATA_MANIFEST, 0); + ManifestFile adapted = TrackedFileAdapters.asManifestFile(original); + TrackedFile result = TrackedFileAdapters.forManifestFile().wrap(adapted); + assertThat(result).isSameAs(original); + assertThat(result.formatVersion()).isZero(); + } + + @Test + void dataManifestTrackedFileAdapter() { + ManifestFile manifest = writeManifestFile(ManifestContent.DATA); + TrackedFile result = TrackedFileAdapters.forManifestFile().wrap(manifest); + + assertThat(result.contentType()).isEqualTo(FileContent.DATA_MANIFEST); + assertThat(result.formatVersion()).isZero(); + assertThat(result.location()).isEqualTo(MANIFEST_LOCATION); + assertThat(result.fileFormat()).isEqualTo(FileFormat.AVRO); + assertThat(result.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().firstRowId()).isNull(); + assertThat(result.recordCount()) + .isEqualTo(ADDED_FILES_COUNT + EXISTING_FILES_COUNT + DELETED_FILES_COUNT); + assertThat(result.manifestInfo()).isNotNull(); + assertThat(result.manifestInfo().addedFilesCount()).isEqualTo(ADDED_FILES_COUNT); + assertThat(result.manifestInfo().existingFilesCount()).isEqualTo(EXISTING_FILES_COUNT); + assertThat(result.manifestInfo().deletedFilesCount()).isEqualTo(DELETED_FILES_COUNT); + assertThat(result.manifestInfo().addedRowsCount()).isEqualTo(ADDED_ROWS_COUNT); + assertThat(result.manifestInfo().replacedFilesCount()).isEqualTo(0); + assertThat(result.manifestInfo().replacedRowsCount()).isEqualTo(0L); + assertThat(result.manifestInfo().modifiedFilesCount()).isEqualTo(0); + assertThat(result.manifestInfo().modifiedRowsCount()).isEqualTo(0L); + } + + @Test + void deleteManifestTrackedFileAdapter() { + ManifestFile manifest = writeManifestFile(ManifestContent.DELETES); + TrackedFile result = TrackedFileAdapters.forManifestFile().wrap(manifest); + + assertThat(result.contentType()).isEqualTo(FileContent.DELETE_MANIFEST); + assertThat(result.formatVersion()).isZero(); + assertThat(result.recordCount()) + .isEqualTo(ADDED_FILES_COUNT + EXISTING_FILES_COUNT + DELETED_FILES_COUNT); + assertThat(result.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertThat(result.tracking().firstRowId()).isNull(); + } + + @Test + void manifestTrackedFileAdapterFailsWhenAddedFilesCountMissing() { + ManifestFile manifest = mock(ManifestFile.class); + when(manifest.path()).thenReturn(MANIFEST_LOCATION); + when(manifest.content()).thenReturn(ManifestContent.DATA); + when(manifest.addedFilesCount()).thenReturn(null); + when(manifest.existingFilesCount()).thenReturn(1); + when(manifest.deletedFilesCount()).thenReturn(0); + when(manifest.replacedFilesCount()).thenReturn(0); + + assertThatThrownBy(() -> TrackedFileAdapters.forManifestFile().wrap(manifest)) + .isInstanceOf(NullPointerException.class) + .hasMessageContaining("missing added files count"); + } + + @Test + void manifestTrackedFileAdapterRejectsNullReplacedFilesCount() { + ManifestFile manifest = manifestWithCounts(1, 1, 0, null, 0); + + assertThatThrownBy(() -> TrackedFileAdapters.forManifestFile().wrap(manifest)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage( + "Cannot convert manifest %s: Invalid replaced file count: null", MANIFEST_LOCATION); + } + + @Test + void manifestTrackedFileAdapterRejectsNonZeroReplacedFilesCount() { + ManifestFile manifest = manifestWithCounts(1, 1, 0, 1, 0); + + assertThatThrownBy(() -> TrackedFileAdapters.forManifestFile().wrap(manifest)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage( + "Cannot convert manifest %s: Invalid replaced file count: 1", MANIFEST_LOCATION); + } + + @Test + void manifestTrackedFileAdapterRejectsNullModifiedFilesCount() { + ManifestFile manifest = manifestWithCounts(1, 1, 0, 0, null); + + assertThatThrownBy(() -> TrackedFileAdapters.forManifestFile().wrap(manifest)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage( + "Cannot convert manifest %s: Invalid modified file count: null", MANIFEST_LOCATION); + } + + @Test + void manifestTrackedFileAdapterRejectsNonZeroModifiedFilesCount() { + ManifestFile manifest = manifestWithCounts(1, 1, 0, 0, 1); + + assertThatThrownBy(() -> TrackedFileAdapters.forManifestFile().wrap(manifest)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage( + "Cannot convert manifest %s: Invalid modified file count: 1", MANIFEST_LOCATION); + } + + @Test + void manifestTrackedFileAdapterFailsWhenAddedRowsCountMissing() { + ManifestFile manifest = + writeManifestFile( + ManifestContent.DATA, MANIFEST_SEQUENCE_NUMBER, MANIFEST_MIN_SEQUENCE_NUMBER, null); + TrackedFile tracked = TrackedFileAdapters.forManifestFile().wrap(manifest); + + assertThatThrownBy(() -> tracked.manifestInfo().addedRowsCount()) + .isInstanceOf(NullPointerException.class) + .hasMessageContaining("null"); + } + + private static void assertWrappedDataFileMatchesFileFields(TrackedFile result, DataFile file) { + assertThat(result.contentType()).isEqualTo(FileContent.DATA); + assertThat(result.location()).isEqualTo(file.location()); + assertThat(result.fileFormat()).isEqualTo(file.format()); + assertThat(result.recordCount()).isEqualTo(file.recordCount()); + assertThat(result.fileSizeInBytes()).isEqualTo(file.fileSizeInBytes()); + assertThat(result.specId()).isEqualTo(file.specId()); + assertThat(result.sortOrderId()).isEqualTo(file.sortOrderId()); + assertThat(result.keyMetadata()).isEqualTo(file.keyMetadata()); + assertThat(result.splitOffsets()).isEqualTo(file.splitOffsets()); + assertThat(result.manifestInfo()).isNull(); + assertThat(result.deletionVector()).isNull(); + assertThat(result.equalityIds()).isNull(); + } + + private static ManifestFile manifestWithCounts( + Integer addedFilesCount, + Integer existingFilesCount, + Integer deletedFilesCount, + Integer replacedFilesCount, + Integer modifiedFilesCount) { + ManifestFile manifest = mock(ManifestFile.class); + when(manifest.path()).thenReturn(MANIFEST_LOCATION); + when(manifest.content()).thenReturn(ManifestContent.DATA); + when(manifest.addedFilesCount()).thenReturn(addedFilesCount); + when(manifest.existingFilesCount()).thenReturn(existingFilesCount); + when(manifest.deletedFilesCount()).thenReturn(deletedFilesCount); + when(manifest.replacedFilesCount()).thenReturn(replacedFilesCount); + when(manifest.modifiedFilesCount()).thenReturn(modifiedFilesCount); + return manifest; + } + + private static ManifestFile writeManifestFile(ManifestContent content) { + return writeManifestFile( + content, MANIFEST_SEQUENCE_NUMBER, MANIFEST_MIN_SEQUENCE_NUMBER, ADDED_ROWS_COUNT); + } + + private static ManifestFile writeManifestFile( + ManifestContent content, long sequenceNumber, long minSequenceNumber, Long addedRowsCount) { + List<ManifestFile.PartitionFieldSummary> partitions = ImmutableList.of(); + return new GenericManifestFile( + MANIFEST_LOCATION, + MANIFEST_FILE_SIZE, + UNPARTITIONED_SPEC.specId(), + content, + sequenceNumber, + minSequenceNumber, + SNAPSHOT_ID, + partitions, + null, Review Comment: Labeled both nulls. * The first is key metadata, not sort order id. * The last is first row id. This helper is shared by DATA and DELETES, and a delete manifest's first row id must be null. Set it to null for both manifest types. -- 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]
