stevenzwu commented on code in PR #16936: URL: https://github.com/apache/iceberg/pull/16936#discussion_r4202321282
########## core/src/main/java/org/apache/iceberg/MapBackedContentStats.java: ########## @@ -0,0 +1,244 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.iceberg; + +import java.nio.ByteBuffer; +import java.util.Map; +import java.util.Objects; +import java.util.Set; +import org.apache.iceberg.relocated.com.google.common.base.Preconditions; +import org.apache.iceberg.relocated.com.google.common.collect.Iterables; +import org.apache.iceberg.relocated.com.google.common.collect.Maps; +import org.apache.iceberg.relocated.com.google.common.collect.Sets; +import org.apache.iceberg.types.Conversions; +import org.apache.iceberg.types.Type; +import org.apache.iceberg.types.Types; + +/** Reusable {@link ContentStats} view over a {@link ContentFile}'s stat maps. */ +class MapBackedContentStats implements ContentStats { + private final Schema tableSchema; + private final Map<Integer, FieldStats<?>> statsById = Maps.newHashMap(); + + private Types.StructType type; + private Map<Integer, Long> valueCounts; + private Map<Integer, Long> nullValueCounts; + private Map<Integer, Long> nanValueCounts; + private Map<Integer, Integer> avgValueSizes; + private Map<Integer, ByteBuffer> lowerBounds; + private Map<Integer, ByteBuffer> upperBounds; + + MapBackedContentStats(Schema tableSchema) { + Preconditions.checkArgument(tableSchema != null, "Invalid table schema: null"); + this.tableSchema = tableSchema; + } + + MapBackedContentStats wrap(ContentFile<?> file) { + this.valueCounts = file.valueCounts(); + this.nullValueCounts = file.nullValueCounts(); + this.nanValueCounts = file.nanValueCounts(); + this.avgValueSizes = file.avgValueSizes(); + this.lowerBounds = file.lowerBounds(); + this.upperBounds = file.upperBounds(); + this.type = null; + return this; + } + + @Override + public Iterable<FieldStats<?>> fieldStats() { + return Iterables.filter(Iterables.transform(statsFieldIds(), this::statsFor), Objects::nonNull); + } + + @Override + @SuppressWarnings("unchecked") + public <T> FieldStats<T> statsFor(int fieldId) { + // Schema is fixed for this instance; wrap() rebinds the metric maps. An id absent + // from the current maps is left uncached so a later file can still surface it. + if (!containsFieldInMaps(fieldId)) { + return null; + } + + if (!statsById.containsKey(fieldId)) { + statsById.put(fieldId, createFieldStats(fieldId)); + } + + return (FieldStats<T>) statsById.get(fieldId); + } + + FieldStats<?> createFieldStats(int fieldId) { Review Comment: fixed ########## core/src/test/java/org/apache/iceberg/TestMapBackedContentStats.java: ########## @@ -0,0 +1,445 @@ +/* + * 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.nio.ByteBuffer; +import java.util.Map; +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.relocated.com.google.common.collect.ImmutableSet; +import org.apache.iceberg.types.Conversions; +import org.apache.iceberg.types.Type; +import org.apache.iceberg.types.Types; +import org.junit.jupiter.api.Test; + +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 the absent-count + * getter throws. + */ + private static final DataFile FILE_WITH_STATS = + dataFile( + 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 contentStatsTypeBuiltLazilyFromMapIds() { Review Comment: Renamed this to typeMatchesIdsPresentInMaps. It checks the struct before wrap and the fields after wrap. It does not assert laziness. ########## core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java: ########## @@ -671,6 +729,399 @@ void unknownSpecIdThrows() { .hasMessageContaining("Cannot find partition spec for spec ID"); } + @Test + void dataTrackedFileAdapterHasNoTracking() { + 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 dataTrackedFileAdapterFromAddedManifestEntry() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(addedEntry(DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.ADDED); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().dataSequenceNumber()).isNull(); + assertThat(result.tracking().fileSequenceNumber()).isNull(); + assertThat(result.tracking().firstRowId()).isNull(); Review Comment: The existing-entry test builds the file with FIRST_ROW_ID and asserts tracking().firstRowId(). ########## core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java: ########## @@ -671,6 +729,399 @@ void unknownSpecIdThrows() { .hasMessageContaining("Cannot find partition spec for spec ID"); } + @Test + void dataTrackedFileAdapterHasNoTracking() { + 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 dataTrackedFileAdapterFromAddedManifestEntry() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(addedEntry(DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.ADDED); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().dataSequenceNumber()).isNull(); + assertThat(result.tracking().fileSequenceNumber()).isNull(); + assertThat(result.tracking().firstRowId()).isNull(); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterFromAddedManifestEntryWithNullSnapshotId() { + TrackedFile result = + TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(newEntry().wrapAppend(null, DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.ADDED); + assertThat(result.tracking().snapshotId()).isNull(); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterFromExistingManifestEntry() { Review Comment: Kept dataTrackedFileAdapterFromExistingManifestEntry and removed the added-entry test. The existing-entry test inlines wrapExisting and checks status, snapshot id, both sequence numbers, first row id, manifest position, and the file fields. ########## core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java: ########## @@ -671,6 +729,399 @@ void unknownSpecIdThrows() { .hasMessageContaining("Cannot find partition spec for spec ID"); } + @Test + void dataTrackedFileAdapterHasNoTracking() { + 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 dataTrackedFileAdapterFromAddedManifestEntry() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(addedEntry(DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.ADDED); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().dataSequenceNumber()).isNull(); + assertThat(result.tracking().fileSequenceNumber()).isNull(); + assertThat(result.tracking().firstRowId()).isNull(); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterFromAddedManifestEntryWithNullSnapshotId() { + TrackedFile result = + TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(newEntry().wrapAppend(null, DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.ADDED); + assertThat(result.tracking().snapshotId()).isNull(); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterFromExistingManifestEntry() { + TrackedFile result = + TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(existingEntry(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); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterFromDeletedManifestEntry() { + TrackedFile result = + TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(deletedEntry(DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.DELETED); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().dataSequenceNumber()).isEqualTo(DATA_SEQUENCE_NUMBER); + assertThat(result.tracking().fileSequenceNumber()).isEqualTo(FILE_SEQUENCE_NUMBER); + assertManifestPosition(result.tracking(), DATA_FILE); + } + + @Test + void dataTrackedFileAdapterReuse() { + TrackedFileAdapters.DataTrackedFile adapter = TrackedFileAdapters.forDataFile(TABLE_SCHEMA); + + adapter.wrap(DATA_FILE); + assertThat(adapter.location()).isEqualTo(DATA_FILE_LOCATION); + assertThat(adapter.recordCount()).isEqualTo(100L); + 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(existingEntry(file2)); + assertThat(adapter.location()).isEqualTo("s3://bucket/data/file2.parquet"); + assertThat(adapter.recordCount()).isEqualTo(200L); + assertThat(adapter.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertManifestPosition(adapter.tracking(), file2); + assertThat(adapter.fileSizeInBytes()).isEqualTo(2048L); + assertThat(adapter.location()).isNotEqualTo(DATA_FILE.location()); + } + + @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().get(0, String.class)).isEqualTo("books"); + assertThat( + TrackedFileAdapters.asDataFile(result, specsById(PARTITIONED_SPEC)) + .partition() + .get(0, String.class)) + .isEqualTo("books"); + } + + @Test + void dataFileDoubleWrapRoundTrip() { + DataFile source = DATA_FILE_WITH_METRICS; + Map<Integer, PartitionSpec> specs = + ImmutableMap.of(UNPARTITIONED_SPEC.specId(), UNPARTITIONED_SPEC); + + TrackedFile tracked = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(existingEntry(source)); + + DataFile roundTripped = TrackedFileAdapters.asDataFile(tracked, specs); + + assertThat(roundTripped.content()).isEqualTo(FileContent.DATA); + assertThat(roundTripped.location()).isEqualTo(source.location()); + assertThat(roundTripped.format()).isEqualTo(source.format()); + assertThat(roundTripped.recordCount()).isEqualTo(source.recordCount()); + assertThat(roundTripped.fileSizeInBytes()).isEqualTo(source.fileSizeInBytes()); + assertThat(roundTripped.specId()).isEqualTo(source.specId()); + assertThat(roundTripped.partition()).isEqualTo(source.partition()); + assertThat(roundTripped.sortOrderId()).isEqualTo(source.sortOrderId()); + assertThat(roundTripped.splitOffsets()).isEqualTo(source.splitOffsets()); + assertThat(roundTripped.keyMetadata()).isEqualTo(source.keyMetadata()); + assertThat(roundTripped.dataSequenceNumber()).isEqualTo(DATA_SEQUENCE_NUMBER); + assertThat(roundTripped.fileSequenceNumber()).isEqualTo(FILE_SEQUENCE_NUMBER); + assertThat(roundTripped.firstRowId()).isNull(); + assertThat(roundTripped.manifestLocation()).isEqualTo(source.manifestLocation()); + assertThat(roundTripped.pos()).isEqualTo(source.pos()); + assertThat(roundTripped.valueCounts()).containsAllEntriesOf(source.valueCounts()); + assertThat(roundTripped.nullValueCounts()).containsAllEntriesOf(source.nullValueCounts()); + assertThat(roundTripped.nanValueCounts()).isNull(); + assertThat(roundTripped.lowerBounds()).containsAllEntriesOf(source.lowerBounds()); + assertThat(roundTripped.upperBounds()).containsAllEntriesOf(source.upperBounds()); + assertThat(roundTripped.columnSizes()).isNull(); + } + + @Test + void dataTrackedFileAdapterUnwrapsToOriginalTrackedFile() { + TrackedFile original = dummyTrackedFile(FileContent.DATA); + DataFile adapted = TrackedFileAdapters.asDataFile(original, UNPARTITIONED); + + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(adapted); + + assertThat(result).isSameAs(original); + } + + @Test + void manifestTrackedFileAdapterUnwrapsToOriginalTrackedFile() { + TrackedFile original = dummyTrackedFile(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(6L); + assertThat(result.manifestInfo()).isNotNull(); + assertThat(result.manifestInfo().addedFilesCount()).isEqualTo(2); + assertThat(result.manifestInfo().existingFilesCount()).isEqualTo(3); + assertThat(result.manifestInfo().deletedFilesCount()).isEqualTo(1); + assertThat(result.manifestInfo().addedRowsCount()).isEqualTo(200L); + 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(6L); + 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: replaced files count must be 0 for v3 or earlier manifests, was 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: replaced files count must be 0 for v3 or earlier manifests, was 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: modified files count must be 0 for v3 or earlier manifests, was 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: modified files count must be 0 for v3 or earlier manifests, was 1", + MANIFEST_LOCATION); + } + + @Test + void manifestTrackedFileAdapterFailsWhenAddedRowsCountMissing() { + ManifestFile manifest = writeManifestFile(ManifestContent.DATA, 5L, 4L, 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, 5L, 4L, 200L); + } + + private static ManifestFile writeManifestFile( + ManifestContent content, long sequenceNumber, long minSequenceNumber) { + return writeManifestFile(content, sequenceNumber, minSequenceNumber, 200L); Review Comment: Removed the unused 3-arg writeManifestFile and hoisted the manifest counts and sequence numbers the tests assert. The 4-arg overload stays because one rejection test passes a null added-rows count. manifestWithCounts stays for those rejection cases. ########## core/src/test/java/org/apache/iceberg/TestMapBackedContentStats.java: ########## @@ -0,0 +1,445 @@ +/* + * 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.nio.ByteBuffer; +import java.util.Map; +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.relocated.com.google.common.collect.ImmutableSet; +import org.apache.iceberg.types.Conversions; +import org.apache.iceberg.types.Type; +import org.apache.iceberg.types.Types; +import org.junit.jupiter.api.Test; + +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 the absent-count Review Comment: Reworded the comment. Field 3 has no value-count entry, so valueCount() throws when it unboxes a null Long. ########## core/src/main/java/org/apache/iceberg/MapBackedContentStats.java: ########## @@ -0,0 +1,244 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.iceberg; + +import java.nio.ByteBuffer; +import java.util.Map; +import java.util.Objects; +import java.util.Set; +import org.apache.iceberg.relocated.com.google.common.base.Preconditions; +import org.apache.iceberg.relocated.com.google.common.collect.Iterables; +import org.apache.iceberg.relocated.com.google.common.collect.Maps; +import org.apache.iceberg.relocated.com.google.common.collect.Sets; +import org.apache.iceberg.types.Conversions; +import org.apache.iceberg.types.Type; +import org.apache.iceberg.types.Types; + +/** Reusable {@link ContentStats} view over a {@link ContentFile}'s stat maps. */ +class MapBackedContentStats implements ContentStats { + private final Schema tableSchema; + private final Map<Integer, FieldStats<?>> statsById = Maps.newHashMap(); + + private Types.StructType type; + private Map<Integer, Long> valueCounts; + private Map<Integer, Long> nullValueCounts; + private Map<Integer, Long> nanValueCounts; + private Map<Integer, Integer> avgValueSizes; + private Map<Integer, ByteBuffer> lowerBounds; + private Map<Integer, ByteBuffer> upperBounds; + + MapBackedContentStats(Schema tableSchema) { + Preconditions.checkArgument(tableSchema != null, "Invalid table schema: null"); + this.tableSchema = tableSchema; + } + + MapBackedContentStats wrap(ContentFile<?> file) { + this.valueCounts = file.valueCounts(); + this.nullValueCounts = file.nullValueCounts(); + this.nanValueCounts = file.nanValueCounts(); + this.avgValueSizes = file.avgValueSizes(); + this.lowerBounds = file.lowerBounds(); + this.upperBounds = file.upperBounds(); + this.type = null; + return this; + } + + @Override + public Iterable<FieldStats<?>> fieldStats() { + return Iterables.filter(Iterables.transform(statsFieldIds(), this::statsFor), Objects::nonNull); + } + + @Override + @SuppressWarnings("unchecked") + public <T> FieldStats<T> statsFor(int fieldId) { + // Schema is fixed for this instance; wrap() rebinds the metric maps. An id absent + // from the current maps is left uncached so a later file can still surface it. + if (!containsFieldInMaps(fieldId)) { + return null; + } + + if (!statsById.containsKey(fieldId)) { + statsById.put(fieldId, createFieldStats(fieldId)); + } + + return (FieldStats<T>) statsById.get(fieldId); + } + + FieldStats<?> createFieldStats(int fieldId) { + Types.NestedField field = tableSchema.findField(fieldId); + // A file can carry metrics for an id this schema does not have. Skip it. + if (field == null) { + return null; + } + + Type fieldType = field.type(); + Types.StructType struct = + StatsUtil.fieldStatsStruct(fieldType, StatsUtil.toBaseId(fieldId), MetricsModes.Full.get()); + // null means a struct, list, or map, or an id outside the stats window. The schema is bound + // once during wrapper creation, so the cached null stays valid across wrap(). + if (struct == null) { + return null; + } + + return new MapBackedFieldStats<>(fieldId, fieldType, struct); + } + + @Override + public Types.StructType type() { + if (type == null) { + this.type = StatsUtil.statsReadSchema(tableSchema, statsFieldIds()); + } + + return type; + } + + @Override + public ContentStats copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + + @Override + public ContentStats copy(Set<Integer> fieldIds) { + throw new UnsupportedOperationException("copy is not implemented"); + } + + boolean containsFieldInMaps(int fieldId) { Review Comment: fixed ########## core/src/test/java/org/apache/iceberg/TestMapBackedContentStats.java: ########## @@ -0,0 +1,445 @@ +/* + * 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.nio.ByteBuffer; +import java.util.Map; +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.relocated.com.google.common.collect.ImmutableSet; +import org.apache.iceberg.types.Conversions; +import org.apache.iceberg.types.Type; +import org.apache.iceberg.types.Types; +import org.junit.jupiter.api.Test; + +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 the absent-count + * getter throws. + */ + private static final DataFile FILE_WITH_STATS = + dataFile( + 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 contentStatsTypeBuiltLazilyFromMapIds() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA); + + assertThat(stats.type().fields()).isEmpty(); + + stats.wrap(FILE_WITH_STATS); + + assertThat(stats.type().fields()) + .extracting(field -> StatsUtil.toFieldId(field.fieldId())) + .containsExactlyInAnyOrder(1, 2, 3, 4); Review Comment: The type is now compared to statsReadSchema for the ids in the file. The comparison is order-independent because the id set is a HashSet. wrapInvalidatesType does the same for the single remaining id. ########## core/src/test/java/org/apache/iceberg/TestMapBackedContentStats.java: ########## @@ -0,0 +1,445 @@ +/* + * 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.nio.ByteBuffer; +import java.util.Map; +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.relocated.com.google.common.collect.ImmutableSet; +import org.apache.iceberg.types.Conversions; +import org.apache.iceberg.types.Type; +import org.apache.iceberg.types.Types; +import org.junit.jupiter.api.Test; + +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 the absent-count + * getter throws. + */ + private static final DataFile FILE_WITH_STATS = + dataFile( + 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 contentStatsTypeBuiltLazilyFromMapIds() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA); + + assertThat(stats.type().fields()).isEmpty(); + + stats.wrap(FILE_WITH_STATS); + + assertThat(stats.type().fields()) + .extracting(field -> StatsUtil.toFieldId(field.fieldId())) + .containsExactlyInAnyOrder(1, 2, 3, 4); + } + + @Test + void wrapInvalidatesType() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + Types.StructType firstType = stats.type(); + assertThat(firstType.fields()) + .extracting(field -> StatsUtil.toFieldId(field.fieldId())) + .containsExactlyInAnyOrder(1, 2, 3, 4); + + DataFile file2 = + dataFile( + 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()) + .extracting(field -> StatsUtil.toFieldId(field.fieldId())) + .containsExactlyInAnyOrder(1); + assertThat(stats.type()).isNotEqualTo(firstType); + } + + @Test + void boundDecodingPerType() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + + FieldStats<?> id = stats.statsFor(1); + assertThat(id.lowerBound()).isInstanceOf(Integer.class).isEqualTo(1); + assertThat(id.upperBound()).isInstanceOf(Integer.class).isEqualTo(1000); + + FieldStats<?> score = stats.statsFor(2); + assertThat(score.lowerBound()).isInstanceOf(Float.class).isEqualTo(1.5f); + assertThat(score.upperBound()).isInstanceOf(Float.class).isEqualTo(9.5f); + + FieldStats<?> ts = stats.statsFor(3); + assertThat(ts.lowerBound()).isInstanceOf(Long.class).isEqualTo(100L); + assertThat(ts.upperBound()).isInstanceOf(Long.class).isEqualTo(999L); + + FieldStats<?> name = stats.statsFor(4); + assertThat(name.lowerBound()).isInstanceOf(CharSequence.class); + assertThat(name.lowerBound().toString()).isEqualTo("aaa"); + assertThat(name.upperBound()).isInstanceOf(CharSequence.class); + assertThat(name.upperBound().toString()).isEqualTo("zzz"); + } + + @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( + 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( + 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); + } + + @Test + void fieldWithoutStatsIsExcluded() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + + assertThat(stats.containsFieldInMaps(5)).isFalse(); + assertThat(stats.type().field(StatsUtil.toBaseId(5))).isNull(); + assertThat(stats.statsFor(5)).isNull(); + assertThat(stats.fieldStats()) + .extracting(FieldStats::fieldId) + .containsExactlyInAnyOrder(1, 2, 3, 4); + } + + @Test + void unknownFieldIdInMaps() { + DataFile file = + dataFile( + ImmutableMap.of(99, 1L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of()); + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(file); + + assertThat(stats.containsFieldInMaps(99)).isTrue(); + assertThat(stats.statsFor(99)).isNull(); + assertThat(stats.type().fields()).isEmpty(); + assertThat(stats.fieldStats()).isEmpty(); + } + + @Test + void copyNotSupported() { Review Comment: Removed copyNotSupported and fieldStatsCopyNotSupported. ########## core/src/main/java/org/apache/iceberg/TrackedFileAdapters.java: ########## @@ -529,6 +566,427 @@ public ManifestFile copy() { } } + /** Adapts a {@link DataFile} to {@link TrackedFile}. */ + static class DataTrackedFile implements TrackedFile { + private final MapBackedContentStats statsWrapper; + private Tracking tracking; + private DataFile file; + private ContentStats stats; + + DataTrackedFile(Schema tableSchema, MetricsConfig metricsConfig) { + this.statsWrapper = new MapBackedContentStats(tableSchema, metricsConfig); + } + + /** Re-points this adapter at a {@link DataFile} from the public API. Tracking is unset. */ + public TrackedFile wrap(DataFile newFile) { + return wrapFile(newFile, null); + } + + /** + * Re-points this adapter at a {@link ManifestEntry}. Converts the contained data file and the + * entry's tracking fields. + */ + public TrackedFile wrap(ManifestEntry<DataFile> entry) { + Preconditions.checkArgument(entry != null, "Invalid entry: null"); + return wrapFile(entry.file(), trackingFrom(entry, entry.file())); + } + + private TrackedFile wrapFile(DataFile newFile, Tracking newTracking) { + if (newFile instanceof TrackedDataFile tracked) { + return tracked.file(); + } + + Preconditions.checkArgument(newFile != null, "Invalid file: null"); + Preconditions.checkArgument( + newFile.content() == FileContent.DATA, + "Invalid content for data file: %s", + newFile.content()); + + this.file = newFile; + this.stats = hasContentStats(newFile) ? statsWrapper.wrap(newFile) : null; + this.tracking = newTracking; + return this; + } + + @Override + public Tracking tracking() { + return tracking; + } + + @Override + public FileContent contentType() { + return FileContent.DATA; + } + + @Override + public int formatVersion() { + throw new IllegalStateException("Format version is assigned at write time"); + } + + @Override + public String location() { + return file.location(); + } + + @Override + public FileFormat fileFormat() { + return file.format(); + } + + @Override + public long recordCount() { + return file.recordCount(); + } + + @Override + public long fileSizeInBytes() { + return file.fileSizeInBytes(); + } + + @Override + public Integer specId() { + // Files in one manifest may use different specs; this is the spec for this data file only. + return file.specId(); + } + + @Override + public StructLike partition() { + return file.partition(); + } + + @Override + public ContentStats contentStats() { + return stats; + } + + @Override + public Integer sortOrderId() { + return file.sortOrderId(); + } + + @Override + public DeletionVector deletionVector() { + return file.deletionVector(); + } + + @Override + public ManifestInfo manifestInfo() { + return null; + } + + @Override + public ByteBuffer keyMetadata() { + return file.keyMetadata(); + } + + @Override + public List<Long> splitOffsets() { + return file.splitOffsets(); + } + + @Override + public List<Integer> equalityIds() { + return null; + } + + @Override + public TrackedFile copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + + @Override + public TrackedFile copyWithStats(Set<Integer> requestedColumnIds) { + throw new UnsupportedOperationException("copy is not implemented"); + } + } + + /** Adapts a {@link ManifestFile} to {@link TrackedFile}. */ + static class ManifestTrackedFile implements TrackedFile { + private final WrappedManifestInfo manifestInfo = new WrappedManifestInfo(); + private Tracking tracking; + private ManifestFile manifest; + private long recordCount; + private FileContent contentType; + + ManifestTrackedFile() {} + + /** + * Re-points this adapter at {@code newManifest}. Converts the manifest's own fields only; + * write-time tracking updates are applied by the versioned writer. + */ + public TrackedFile wrap(ManifestFile newManifest) { + if (newManifest instanceof TrackedManifestFile tracked) { + return tracked.file(); + } + + Preconditions.checkArgument(newManifest != null, "Invalid manifest file: null"); + + this.manifest = newManifest; + this.contentType = + newManifest.content() == ManifestContent.DATA + ? FileContent.DATA_MANIFEST + : FileContent.DELETE_MANIFEST; + this.recordCount = manifestRecordCount(newManifest); + this.tracking = trackingFrom(newManifest); + this.manifestInfo.wrap(newManifest); + return this; + } + + @Override + public Tracking tracking() { + return tracking; + } + + @Override + public FileContent contentType() { + return contentType; + } + + @Override + public int formatVersion() { + return manifest.formatVersion(); + } + + @Override + public String location() { + return manifest.path(); + } + + @Override + public FileFormat fileFormat() { + return FileFormat.fromFileName(manifest.path()); + } + + @Override + public long recordCount() { + // Number of TrackedFile rows stored in the manifest. + return recordCount; + } + + @Override + public long fileSizeInBytes() { + return manifest.length(); + } + + @Override + public Integer specId() { + // Spec the wrapped manifest was written with. Data file entries in a v4 manifest file may + // use different specs. + return manifest.partitionSpecId(); + } + + @Override + public StructLike partition() { + return null; + } + + @Override + public ContentStats contentStats() { + return null; + } + + @Override + public Integer sortOrderId() { + return null; + } + + @Override + public DeletionVector deletionVector() { + return null; + } + + @Override + public ManifestInfo manifestInfo() { + return manifestInfo; + } + + @Override + public ByteBuffer keyMetadata() { + return manifest.keyMetadata(); + } + + @Override + public List<Long> splitOffsets() { + return null; + } + + @Override + public List<Integer> equalityIds() { + return null; + } + + @Override + public TrackedFile copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + + @Override + public TrackedFile copyWithStats(Set<Integer> requestedColumnIds) { + throw new UnsupportedOperationException("copy is not implemented"); + } + } + + /** Reusable {@link ManifestInfo} view over a {@link ManifestFile}'s counts. */ + private static class WrappedManifestInfo implements ManifestInfo { + private ManifestFile manifest; + + void wrap(ManifestFile newManifest) { + this.manifest = newManifest; + } + + @Override + public int addedFilesCount() { + return manifest.addedFilesCount(); + } + + @Override + public int existingFilesCount() { + return manifest.existingFilesCount(); + } + + @Override + public int deletedFilesCount() { + return manifest.deletedFilesCount(); + } + + @Override + public int replacedFilesCount() { + return manifest.replacedFilesCount(); + } + + @Override + public int modifiedFilesCount() { + return manifest.modifiedFilesCount(); + } + + @Override + public long addedRowsCount() { + return manifest.addedRowsCount(); + } + + @Override + public long existingRowsCount() { + return manifest.existingRowsCount(); + } + + @Override + public long deletedRowsCount() { + return manifest.deletedRowsCount(); + } + + @Override + public long replacedRowsCount() { + return manifest.replacedRowsCount(); + } + + @Override + public long modifiedRowsCount() { + return manifest.modifiedRowsCount(); + } + + @Override + public long minSequenceNumber() { + return manifest.minSequenceNumber(); + } + + @Override + public ManifestBitmap manifestDeletionVector() { + return manifest.manifestDeletionVector(); + } + + @Override + public ManifestInfo copy() { + throw new UnsupportedOperationException("copy is not implemented"); + } + } + + private static Tracking trackingFrom(ManifestEntry<?> entry, ContentFile<?> file) { + return new TrackingStruct( + entryStatus(entry.status()), + entry.snapshotId(), + entry.dataSequenceNumber(), + entry.fileSequenceNumber(), + null, + file.firstRowId(), + null, + null); + } + + private static Tracking trackingFrom(ManifestFile manifest) { + return new TrackingStruct( + null, + manifest.snapshotId(), + manifest.sequenceNumber(), + manifest.sequenceNumber(), + null, + manifest.firstRowId(), + null, + null); + } + + private static EntryStatus entryStatus(ManifestEntry.Status status) { + return switch (status) { + case EXISTING -> EntryStatus.EXISTING; + case ADDED -> EntryStatus.ADDED; + case DELETED -> EntryStatus.DELETED; + }; + } + + private static boolean hasContentStats(ContentFile<?> file) { + return isPresent(file.valueCounts()) + || isPresent(file.nullValueCounts()) + || isPresent(file.nanValueCounts()) + || isPresent(file.avgValueSizes()) + || isPresent(file.lowerBounds()) + || isPresent(file.upperBounds()); + } + + private static boolean isPresent(Map<?, ?> map) { + return map != null && !map.isEmpty(); + } + + /** + * Record count of a manifest is the number of TrackedFile rows it stores: the sum of per-status + * file counts. Missing counts fail rather than producing an incorrect total. + */ + private static long manifestRecordCount(ManifestFile manifest) { + Preconditions.checkNotNull( + manifest.addedFilesCount(), + "Cannot convert manifest %s: missing added files count", + manifest.path()); + Preconditions.checkNotNull( + manifest.existingFilesCount(), + "Cannot convert manifest %s: missing existing files count", + manifest.path()); + Preconditions.checkNotNull( + manifest.deletedFilesCount(), + "Cannot convert manifest %s: missing deleted files count", + manifest.path()); + Preconditions.checkNotNull( + manifest.replacedFilesCount(), + "Cannot convert manifest %s: missing replaced files count", + manifest.path()); + Preconditions.checkNotNull( + manifest.modifiedFilesCount(), + "Cannot convert manifest %s: missing modified files count", + manifest.path()); + Preconditions.checkArgument( + manifest.replacedFilesCount() == 0, + "Cannot convert manifest %s: replaced files count must be 0 for v3 or earlier manifests, was %s", Review Comment: Shortened both messages. I used "Invalid" rather than "Non-zero" because a null count fails the same Objects.equals check: ``` Cannot convert manifest %s: Invalid replaced file count: %s Cannot convert manifest %s: Invalid modified file count: %s ``` ########## core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java: ########## @@ -532,6 +583,13 @@ void manifestFileAdapterDelegation(FileContent contentType) { .hasMessage("v4 manifests are not bound to a single partition spec"); } + @Test + void manifestFileAdapterKeepsZeroFormatVersion() { + TrackedFile original = dummyTrackedFile(FileContent.DATA_MANIFEST, 0); Review Comment: Removed manifestFileAdapterKeepsZeroFormatVersion. dataManifestTrackedFileAdapter already wraps a real manifest whose format version is 0. ########## core/src/test/java/org/apache/iceberg/TestMapBackedContentStats.java: ########## @@ -0,0 +1,445 @@ +/* + * 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.nio.ByteBuffer; +import java.util.Map; +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.relocated.com.google.common.collect.ImmutableSet; +import org.apache.iceberg.types.Conversions; +import org.apache.iceberg.types.Type; +import org.apache.iceberg.types.Types; +import org.junit.jupiter.api.Test; + +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 the absent-count + * getter throws. + */ + private static final DataFile FILE_WITH_STATS = + dataFile( + 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 contentStatsTypeBuiltLazilyFromMapIds() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA); + + assertThat(stats.type().fields()).isEmpty(); + + stats.wrap(FILE_WITH_STATS); + + assertThat(stats.type().fields()) + .extracting(field -> StatsUtil.toFieldId(field.fieldId())) + .containsExactlyInAnyOrder(1, 2, 3, 4); + } + + @Test + void wrapInvalidatesType() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + Types.StructType firstType = stats.type(); + assertThat(firstType.fields()) + .extracting(field -> StatsUtil.toFieldId(field.fieldId())) + .containsExactlyInAnyOrder(1, 2, 3, 4); + + DataFile file2 = + dataFile( + 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()) + .extracting(field -> StatsUtil.toFieldId(field.fieldId())) + .containsExactlyInAnyOrder(1); + assertThat(stats.type()).isNotEqualTo(firstType); + } + + @Test + void boundDecodingPerType() { Review Comment: `boundDecodingPerType` is now a parameterized round trip over the same type and bound pairs as `TestFieldStatsStruct.TYPES_AND_BOUNDS` is private, so the tuples are copied here. Unknown still expects null bounds. Geometry stays in geoBoundsDecode. ########## core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java: ########## @@ -671,6 +729,399 @@ void unknownSpecIdThrows() { .hasMessageContaining("Cannot find partition spec for spec ID"); } + @Test + void dataTrackedFileAdapterHasNoTracking() { Review Comment: Kept this as the single wrap(DataFile) test and renamed it to dataTrackedFileAdapterFromDataFile. It still expects tracking to be null. ########## core/src/test/java/org/apache/iceberg/TestMapBackedContentStats.java: ########## @@ -0,0 +1,445 @@ +/* + * 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.nio.ByteBuffer; +import java.util.Map; +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.relocated.com.google.common.collect.ImmutableSet; +import org.apache.iceberg.types.Conversions; +import org.apache.iceberg.types.Type; +import org.apache.iceberg.types.Types; +import org.junit.jupiter.api.Test; + +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 the absent-count + * getter throws. + */ + private static final DataFile FILE_WITH_STATS = + dataFile( + 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 contentStatsTypeBuiltLazilyFromMapIds() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA); + + assertThat(stats.type().fields()).isEmpty(); + + stats.wrap(FILE_WITH_STATS); + + assertThat(stats.type().fields()) + .extracting(field -> StatsUtil.toFieldId(field.fieldId())) + .containsExactlyInAnyOrder(1, 2, 3, 4); + } + + @Test + void wrapInvalidatesType() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + Types.StructType firstType = stats.type(); + assertThat(firstType.fields()) + .extracting(field -> StatsUtil.toFieldId(field.fieldId())) + .containsExactlyInAnyOrder(1, 2, 3, 4); + + DataFile file2 = + dataFile( + 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()) + .extracting(field -> StatsUtil.toFieldId(field.fieldId())) + .containsExactlyInAnyOrder(1); + assertThat(stats.type()).isNotEqualTo(firstType); + } + + @Test + void boundDecodingPerType() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + + FieldStats<?> id = stats.statsFor(1); + assertThat(id.lowerBound()).isInstanceOf(Integer.class).isEqualTo(1); + assertThat(id.upperBound()).isInstanceOf(Integer.class).isEqualTo(1000); + + FieldStats<?> score = stats.statsFor(2); + assertThat(score.lowerBound()).isInstanceOf(Float.class).isEqualTo(1.5f); + assertThat(score.upperBound()).isInstanceOf(Float.class).isEqualTo(9.5f); + + FieldStats<?> ts = stats.statsFor(3); + assertThat(ts.lowerBound()).isInstanceOf(Long.class).isEqualTo(100L); + assertThat(ts.upperBound()).isInstanceOf(Long.class).isEqualTo(999L); + + FieldStats<?> name = stats.statsFor(4); + assertThat(name.lowerBound()).isInstanceOf(CharSequence.class); + assertThat(name.lowerBound().toString()).isEqualTo("aaa"); + assertThat(name.upperBound()).isInstanceOf(CharSequence.class); + assertThat(name.upperBound().toString()).isEqualTo("zzz"); + } + + @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( + 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( + 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); + } + + @Test + void fieldWithoutStatsIsExcluded() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + + assertThat(stats.containsFieldInMaps(5)).isFalse(); + assertThat(stats.type().field(StatsUtil.toBaseId(5))).isNull(); + assertThat(stats.statsFor(5)).isNull(); + assertThat(stats.fieldStats()) + .extracting(FieldStats::fieldId) + .containsExactlyInAnyOrder(1, 2, 3, 4); + } + + @Test + void unknownFieldIdInMaps() { + DataFile file = + dataFile( + ImmutableMap.of(99, 1L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of()); + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(file); + + assertThat(stats.containsFieldInMaps(99)).isTrue(); + assertThat(stats.statsFor(99)).isNull(); + assertThat(stats.type().fields()).isEmpty(); + assertThat(stats.fieldStats()).isEmpty(); + } + + @Test + void copyNotSupported() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + + assertThatThrownBy(stats::copy) + .isInstanceOf(UnsupportedOperationException.class) + .hasMessage("copy is not implemented"); + + assertThatThrownBy(() -> stats.copy(ImmutableSet.of(1))) + .isInstanceOf(UnsupportedOperationException.class) + .hasMessage("copy is not implemented"); + } + + @Test + void fieldStatsCopyNotSupported() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + + assertThatThrownBy(stats.statsFor(1)::copy) + .isInstanceOf(UnsupportedOperationException.class) + .hasMessage("copy is not implemented"); + } + + @Test + void reuseRebindsBounds() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA); + + stats.wrap(FILE_WITH_STATS); + assertThat(stats.statsFor(1).lowerBound()).isEqualTo(1); + + DataFile file2 = + dataFile( + 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 wrapReusesFieldStats() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA); + stats.wrap(FILE_WITH_STATS); + FieldStats<?> id = stats.statsFor(1); + FieldStats<?> score = stats.statsFor(2); + FieldStats<?> ts = stats.statsFor(3); + FieldStats<?> name = stats.statsFor(4); + + DataFile file2 = + dataFile( + ImmutableMap.of(1, 50L, 2, 40L, 3, 20L, 4, 30L), + ImmutableMap.of(1, 9L, 2, 8L, 3, 7L, 4, 6L), + ImmutableMap.of(2, 11L), + ImmutableMap.of( + 1, buf(Types.IntegerType.get(), 500), + 2, buf(Types.FloatType.get(), 2.5f), + 3, buf(Types.LongType.get(), 200L), + 4, buf(Types.StringType.get(), "bbb")), + ImmutableMap.of( + 1, buf(Types.IntegerType.get(), 5000), + 2, buf(Types.FloatType.get(), 8.5f), + 3, buf(Types.LongType.get(), 800L), + 4, buf(Types.StringType.get(), "yyy"))); + stats.wrap(file2); + + assertThat(stats.statsFor(1)).isSameAs(id); + assertThat(id.valueCount()).isEqualTo(50L); + assertThat(id.nullValueCount()).isEqualTo(9L); + assertThat(id.lowerBound()).isEqualTo(500); + assertThat(id.upperBound()).isEqualTo(5000); + + assertThat(stats.statsFor(2)).isSameAs(score); + assertThat(score.valueCount()).isEqualTo(40L); + assertThat(score.nullValueCount()).isEqualTo(8L); + assertThat(score.nanValueCount()).isEqualTo(11L); + assertThat(score.lowerBound()).isEqualTo(2.5f); + assertThat(score.upperBound()).isEqualTo(8.5f); + + assertThat(stats.statsFor(3)).isSameAs(ts); + assertThat(ts.valueCount()).isEqualTo(20L); + assertThat(ts.nullValueCount()).isEqualTo(7L); + assertThat(ts.lowerBound()).isEqualTo(200L); + assertThat(ts.upperBound()).isEqualTo(800L); + + assertThat(stats.statsFor(4)).isSameAs(name); + assertThat(name.valueCount()).isEqualTo(30L); + assertThat(name.nullValueCount()).isEqualTo(6L); + assertThat(name.lowerBound().toString()).isEqualTo("bbb"); + assertThat(name.upperBound().toString()).isEqualTo("yyy"); + } + + @Test + void absentIdIsRereadWhenPresentAgain() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA); + stats.wrap(FILE_WITH_STATS); + FieldStats<?> id = stats.statsFor(1); + assertThat(id).isNotNull(); + + stats.wrap( + dataFile( + ImmutableMap.of(2, 10L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(2, buf(Types.FloatType.get(), 0.0f)), + ImmutableMap.of(2, buf(Types.FloatType.get(), 1.0f)))); + assertThat(stats.statsFor(1)).isNull(); + + stats.wrap( + dataFile( + ImmutableMap.of(1, 7L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(1, buf(Types.IntegerType.get(), 42)), + ImmutableMap.of(1, buf(Types.IntegerType.get(), 43)))); + assertThat(stats.statsFor(1)).isSameAs(id); + assertThat(id.lowerBound()).isEqualTo(42); + assertThat(id.upperBound()).isEqualTo(43); + assertThat(id.valueCount()).isEqualTo(7L); + } + + @Test + void outOfRangeFieldStatsAreNullAndCached() { Review Comment: Removed outOfRangeFieldStatsAreNullAndCached and CountingContentStats. ########## core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java: ########## @@ -114,17 +121,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, null, null, null, null), + null, Review Comment: DATA_FILE now uses new Metrics(100L), KEY_METADATA, a sort order id, and FIRST_ROW_ID, and the field helper asserts those values. firstRowId is checked on the existing-entry path. wrapAppend still suppresses it. ########## core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java: ########## @@ -671,6 +729,399 @@ void unknownSpecIdThrows() { .hasMessageContaining("Cannot find partition spec for spec ID"); } + @Test + void dataTrackedFileAdapterHasNoTracking() { + 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 dataTrackedFileAdapterFromAddedManifestEntry() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(addedEntry(DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.ADDED); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().dataSequenceNumber()).isNull(); Review Comment: Those nulls came from wrapAppend, which stores a null data sequence number. The existing-entry test now asserts DATA_SEQUENCE_NUMBER and FILE_SEQUENCE_NUMBER. ########## core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java: ########## @@ -671,6 +729,399 @@ void unknownSpecIdThrows() { .hasMessageContaining("Cannot find partition spec for spec ID"); } + @Test + void dataTrackedFileAdapterHasNoTracking() { + 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"); Review Comment: I think this should keep throwing. A v3 DataFile has no format version to copy. On the manifest path, 0 is forwarded from `ManifestFile.formatVersion` when the manifest predates format tracking. We also discussed on Monday's sync that format version should be moved to manifest_info struct that should be set at write time. ########## core/src/test/java/org/apache/iceberg/TestMapBackedContentStats.java: ########## @@ -0,0 +1,445 @@ +/* + * 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.nio.ByteBuffer; +import java.util.Map; +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.relocated.com.google.common.collect.ImmutableSet; +import org.apache.iceberg.types.Conversions; +import org.apache.iceberg.types.Type; +import org.apache.iceberg.types.Types; +import org.junit.jupiter.api.Test; + +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 the absent-count + * getter throws. + */ + private static final DataFile FILE_WITH_STATS = + dataFile( + 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 contentStatsTypeBuiltLazilyFromMapIds() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA); + + assertThat(stats.type().fields()).isEmpty(); + + stats.wrap(FILE_WITH_STATS); + + assertThat(stats.type().fields()) + .extracting(field -> StatsUtil.toFieldId(field.fieldId())) + .containsExactlyInAnyOrder(1, 2, 3, 4); + } + + @Test + void wrapInvalidatesType() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + Types.StructType firstType = stats.type(); + assertThat(firstType.fields()) + .extracting(field -> StatsUtil.toFieldId(field.fieldId())) + .containsExactlyInAnyOrder(1, 2, 3, 4); + + DataFile file2 = + dataFile( + 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()) + .extracting(field -> StatsUtil.toFieldId(field.fieldId())) + .containsExactlyInAnyOrder(1); + assertThat(stats.type()).isNotEqualTo(firstType); + } + + @Test + void boundDecodingPerType() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + + FieldStats<?> id = stats.statsFor(1); + assertThat(id.lowerBound()).isInstanceOf(Integer.class).isEqualTo(1); + assertThat(id.upperBound()).isInstanceOf(Integer.class).isEqualTo(1000); + + FieldStats<?> score = stats.statsFor(2); + assertThat(score.lowerBound()).isInstanceOf(Float.class).isEqualTo(1.5f); + assertThat(score.upperBound()).isInstanceOf(Float.class).isEqualTo(9.5f); + + FieldStats<?> ts = stats.statsFor(3); + assertThat(ts.lowerBound()).isInstanceOf(Long.class).isEqualTo(100L); + assertThat(ts.upperBound()).isInstanceOf(Long.class).isEqualTo(999L); + + FieldStats<?> name = stats.statsFor(4); + assertThat(name.lowerBound()).isInstanceOf(CharSequence.class); + assertThat(name.lowerBound().toString()).isEqualTo("aaa"); + assertThat(name.upperBound()).isInstanceOf(CharSequence.class); + assertThat(name.upperBound().toString()).isEqualTo("zzz"); + } + + @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( + 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( + 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); + } + + @Test + void fieldWithoutStatsIsExcluded() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + + assertThat(stats.containsFieldInMaps(5)).isFalse(); + assertThat(stats.type().field(StatsUtil.toBaseId(5))).isNull(); + assertThat(stats.statsFor(5)).isNull(); + assertThat(stats.fieldStats()) + .extracting(FieldStats::fieldId) + .containsExactlyInAnyOrder(1, 2, 3, 4); + } + + @Test + void unknownFieldIdInMaps() { + DataFile file = + dataFile( + ImmutableMap.of(99, 1L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of()); + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(file); + + assertThat(stats.containsFieldInMaps(99)).isTrue(); + assertThat(stats.statsFor(99)).isNull(); + assertThat(stats.type().fields()).isEmpty(); + assertThat(stats.fieldStats()).isEmpty(); + } + + @Test + void copyNotSupported() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + + assertThatThrownBy(stats::copy) + .isInstanceOf(UnsupportedOperationException.class) + .hasMessage("copy is not implemented"); + + assertThatThrownBy(() -> stats.copy(ImmutableSet.of(1))) + .isInstanceOf(UnsupportedOperationException.class) + .hasMessage("copy is not implemented"); + } + + @Test + void fieldStatsCopyNotSupported() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + + assertThatThrownBy(stats.statsFor(1)::copy) + .isInstanceOf(UnsupportedOperationException.class) + .hasMessage("copy is not implemented"); + } + + @Test + void reuseRebindsBounds() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA); + + stats.wrap(FILE_WITH_STATS); + assertThat(stats.statsFor(1).lowerBound()).isEqualTo(1); + + DataFile file2 = + dataFile( + 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 wrapReusesFieldStats() { Review Comment: Removed wrapReusesFieldStats. absentIdIsRereadWhenPresentAgain no longer checks instance identity. It still checks that the id is absent, then present again with the new file's values. ########## core/src/test/java/org/apache/iceberg/TestMapBackedContentStats.java: ########## @@ -0,0 +1,445 @@ +/* + * 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.nio.ByteBuffer; +import java.util.Map; +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.relocated.com.google.common.collect.ImmutableSet; +import org.apache.iceberg.types.Conversions; +import org.apache.iceberg.types.Type; +import org.apache.iceberg.types.Types; +import org.junit.jupiter.api.Test; + +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 the absent-count + * getter throws. + */ + private static final DataFile FILE_WITH_STATS = + dataFile( + 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 contentStatsTypeBuiltLazilyFromMapIds() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA); + + assertThat(stats.type().fields()).isEmpty(); + + stats.wrap(FILE_WITH_STATS); + + assertThat(stats.type().fields()) + .extracting(field -> StatsUtil.toFieldId(field.fieldId())) + .containsExactlyInAnyOrder(1, 2, 3, 4); + } + + @Test + void wrapInvalidatesType() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + Types.StructType firstType = stats.type(); + assertThat(firstType.fields()) + .extracting(field -> StatsUtil.toFieldId(field.fieldId())) + .containsExactlyInAnyOrder(1, 2, 3, 4); + + DataFile file2 = + dataFile( + 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()) + .extracting(field -> StatsUtil.toFieldId(field.fieldId())) + .containsExactlyInAnyOrder(1); + assertThat(stats.type()).isNotEqualTo(firstType); + } + + @Test + void boundDecodingPerType() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + + FieldStats<?> id = stats.statsFor(1); + assertThat(id.lowerBound()).isInstanceOf(Integer.class).isEqualTo(1); + assertThat(id.upperBound()).isInstanceOf(Integer.class).isEqualTo(1000); + + FieldStats<?> score = stats.statsFor(2); + assertThat(score.lowerBound()).isInstanceOf(Float.class).isEqualTo(1.5f); + assertThat(score.upperBound()).isInstanceOf(Float.class).isEqualTo(9.5f); + + FieldStats<?> ts = stats.statsFor(3); + assertThat(ts.lowerBound()).isInstanceOf(Long.class).isEqualTo(100L); + assertThat(ts.upperBound()).isInstanceOf(Long.class).isEqualTo(999L); + + FieldStats<?> name = stats.statsFor(4); + assertThat(name.lowerBound()).isInstanceOf(CharSequence.class); + assertThat(name.lowerBound().toString()).isEqualTo("aaa"); + assertThat(name.upperBound()).isInstanceOf(CharSequence.class); + assertThat(name.upperBound().toString()).isEqualTo("zzz"); + } + + @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( + 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( + 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); + } + + @Test + void fieldWithoutStatsIsExcluded() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + + assertThat(stats.containsFieldInMaps(5)).isFalse(); + assertThat(stats.type().field(StatsUtil.toBaseId(5))).isNull(); + assertThat(stats.statsFor(5)).isNull(); + assertThat(stats.fieldStats()) + .extracting(FieldStats::fieldId) + .containsExactlyInAnyOrder(1, 2, 3, 4); + } + + @Test + void unknownFieldIdInMaps() { + DataFile file = + dataFile( + ImmutableMap.of(99, 1L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of()); + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(file); + + assertThat(stats.containsFieldInMaps(99)).isTrue(); + assertThat(stats.statsFor(99)).isNull(); + assertThat(stats.type().fields()).isEmpty(); + assertThat(stats.fieldStats()).isEmpty(); + } + + @Test + void copyNotSupported() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + + assertThatThrownBy(stats::copy) + .isInstanceOf(UnsupportedOperationException.class) + .hasMessage("copy is not implemented"); + + assertThatThrownBy(() -> stats.copy(ImmutableSet.of(1))) + .isInstanceOf(UnsupportedOperationException.class) + .hasMessage("copy is not implemented"); + } + + @Test + void fieldStatsCopyNotSupported() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + + assertThatThrownBy(stats.statsFor(1)::copy) + .isInstanceOf(UnsupportedOperationException.class) + .hasMessage("copy is not implemented"); + } + + @Test + void reuseRebindsBounds() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA); + + stats.wrap(FILE_WITH_STATS); + assertThat(stats.statsFor(1).lowerBound()).isEqualTo(1); + + DataFile file2 = + dataFile( + 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 wrapReusesFieldStats() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA); + stats.wrap(FILE_WITH_STATS); + FieldStats<?> id = stats.statsFor(1); + FieldStats<?> score = stats.statsFor(2); + FieldStats<?> ts = stats.statsFor(3); + FieldStats<?> name = stats.statsFor(4); + + DataFile file2 = + dataFile( + ImmutableMap.of(1, 50L, 2, 40L, 3, 20L, 4, 30L), + ImmutableMap.of(1, 9L, 2, 8L, 3, 7L, 4, 6L), + ImmutableMap.of(2, 11L), + ImmutableMap.of( + 1, buf(Types.IntegerType.get(), 500), + 2, buf(Types.FloatType.get(), 2.5f), + 3, buf(Types.LongType.get(), 200L), + 4, buf(Types.StringType.get(), "bbb")), + ImmutableMap.of( + 1, buf(Types.IntegerType.get(), 5000), + 2, buf(Types.FloatType.get(), 8.5f), + 3, buf(Types.LongType.get(), 800L), + 4, buf(Types.StringType.get(), "yyy"))); + stats.wrap(file2); + + assertThat(stats.statsFor(1)).isSameAs(id); + assertThat(id.valueCount()).isEqualTo(50L); + assertThat(id.nullValueCount()).isEqualTo(9L); + assertThat(id.lowerBound()).isEqualTo(500); + assertThat(id.upperBound()).isEqualTo(5000); + + assertThat(stats.statsFor(2)).isSameAs(score); + assertThat(score.valueCount()).isEqualTo(40L); + assertThat(score.nullValueCount()).isEqualTo(8L); + assertThat(score.nanValueCount()).isEqualTo(11L); + assertThat(score.lowerBound()).isEqualTo(2.5f); + assertThat(score.upperBound()).isEqualTo(8.5f); + + assertThat(stats.statsFor(3)).isSameAs(ts); + assertThat(ts.valueCount()).isEqualTo(20L); + assertThat(ts.nullValueCount()).isEqualTo(7L); + assertThat(ts.lowerBound()).isEqualTo(200L); + assertThat(ts.upperBound()).isEqualTo(800L); + + assertThat(stats.statsFor(4)).isSameAs(name); + assertThat(name.valueCount()).isEqualTo(30L); + assertThat(name.nullValueCount()).isEqualTo(6L); + assertThat(name.lowerBound().toString()).isEqualTo("bbb"); + assertThat(name.upperBound().toString()).isEqualTo("yyy"); + } + + @Test + void absentIdIsRereadWhenPresentAgain() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA); + stats.wrap(FILE_WITH_STATS); + FieldStats<?> id = stats.statsFor(1); + assertThat(id).isNotNull(); + + stats.wrap( + dataFile( + ImmutableMap.of(2, 10L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(2, buf(Types.FloatType.get(), 0.0f)), + ImmutableMap.of(2, buf(Types.FloatType.get(), 1.0f)))); + assertThat(stats.statsFor(1)).isNull(); + + stats.wrap( + dataFile( + ImmutableMap.of(1, 7L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(1, buf(Types.IntegerType.get(), 42)), + ImmutableMap.of(1, buf(Types.IntegerType.get(), 43)))); + assertThat(stats.statsFor(1)).isSameAs(id); + assertThat(id.lowerBound()).isEqualTo(42); + assertThat(id.upperBound()).isEqualTo(43); + assertThat(id.valueCount()).isEqualTo(7L); + } + + @Test + void outOfRangeFieldStatsAreNullAndCached() { + int fieldId = 999_950; + Schema schema = + new Schema(Types.NestedField.optional(fieldId, "too_high", Types.IntegerType.get())); + DataFile file = + dataFile( + ImmutableMap.of(fieldId, 1L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(fieldId, buf(Types.IntegerType.get(), 1)), + ImmutableMap.of(fieldId, buf(Types.IntegerType.get(), 2))); + CountingContentStats stats = new CountingContentStats(schema); + stats.wrap(file); + + assertThat(stats.statsFor(fieldId)).isNull(); + assertThat(stats.creates).isEqualTo(1); + + DataFile file2 = + dataFile( + ImmutableMap.of(fieldId, 9L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(fieldId, buf(Types.IntegerType.get(), 30)), + ImmutableMap.of(fieldId, buf(Types.IntegerType.get(), 40))); + stats.wrap(file2); + + assertThat(stats.statsFor(fieldId)).isNull(); + assertThat(stats.creates).isEqualTo(1); + assertThat(stats.fieldStats()).isEmpty(); + assertThat(stats.creates).isEqualTo(1); + } + + @Test + void listElementFieldStats() { + Schema schema = + new Schema( + Types.NestedField.required( + 1, "nums", Types.ListType.ofRequired(2, Types.IntegerType.get()))); + DataFile file = + dataFile( + ImmutableMap.of(2, 4L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(2, buf(Types.IntegerType.get(), 1)), + ImmutableMap.of(2, buf(Types.IntegerType.get(), 9))); + MapBackedContentStats stats = new MapBackedContentStats(schema).wrap(file); + + FieldStats<?> element = stats.statsFor(2); + assertThat(element.lowerBound()).isEqualTo(1); + assertThat(element.upperBound()).isEqualTo(9); + assertThat(stats.fieldStats()).extracting(FieldStats::fieldId).containsExactly(2); + } + + private static final class CountingContentStats extends MapBackedContentStats { + private int creates; + + private CountingContentStats(Schema tableSchema) { + super(tableSchema); + } + + @Override + FieldStats<?> createFieldStats(int fieldId) { + creates += 1; + return super.createFieldStats(fieldId); + } + } + + private static ByteBuffer buf(Type type, Object value) { + return Conversions.toByteBuffer(type, value); + } + + private static DataFile dataFile( + Map<Integer, Long> valueCounts, + Map<Integer, Long> nullValueCounts, + Map<Integer, Long> nanValueCounts, + Map<Integer, ByteBuffer> lowerBounds, + Map<Integer, ByteBuffer> upperBounds) { + Metrics metrics = + new Metrics( + 100L, null, valueCounts, nullValueCounts, nanValueCounts, lowerBounds, upperBounds); Review Comment: dataFile helper now takes the record count and passes it into Metrics. Each call site passes 100L. ########## core/src/test/java/org/apache/iceberg/TestMapBackedContentStats.java: ########## @@ -0,0 +1,445 @@ +/* + * 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.nio.ByteBuffer; +import java.util.Map; +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.relocated.com.google.common.collect.ImmutableSet; +import org.apache.iceberg.types.Conversions; +import org.apache.iceberg.types.Type; +import org.apache.iceberg.types.Types; +import org.junit.jupiter.api.Test; + +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 the absent-count + * getter throws. + */ + private static final DataFile FILE_WITH_STATS = + dataFile( + 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 contentStatsTypeBuiltLazilyFromMapIds() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA); + + assertThat(stats.type().fields()).isEmpty(); + + stats.wrap(FILE_WITH_STATS); + + assertThat(stats.type().fields()) + .extracting(field -> StatsUtil.toFieldId(field.fieldId())) + .containsExactlyInAnyOrder(1, 2, 3, 4); + } + + @Test + void wrapInvalidatesType() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + Types.StructType firstType = stats.type(); + assertThat(firstType.fields()) + .extracting(field -> StatsUtil.toFieldId(field.fieldId())) + .containsExactlyInAnyOrder(1, 2, 3, 4); + + DataFile file2 = + dataFile( + 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()) + .extracting(field -> StatsUtil.toFieldId(field.fieldId())) + .containsExactlyInAnyOrder(1); + assertThat(stats.type()).isNotEqualTo(firstType); + } + + @Test + void boundDecodingPerType() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + + FieldStats<?> id = stats.statsFor(1); + assertThat(id.lowerBound()).isInstanceOf(Integer.class).isEqualTo(1); + assertThat(id.upperBound()).isInstanceOf(Integer.class).isEqualTo(1000); + + FieldStats<?> score = stats.statsFor(2); + assertThat(score.lowerBound()).isInstanceOf(Float.class).isEqualTo(1.5f); + assertThat(score.upperBound()).isInstanceOf(Float.class).isEqualTo(9.5f); + + FieldStats<?> ts = stats.statsFor(3); + assertThat(ts.lowerBound()).isInstanceOf(Long.class).isEqualTo(100L); + assertThat(ts.upperBound()).isInstanceOf(Long.class).isEqualTo(999L); + + FieldStats<?> name = stats.statsFor(4); + assertThat(name.lowerBound()).isInstanceOf(CharSequence.class); + assertThat(name.lowerBound().toString()).isEqualTo("aaa"); + assertThat(name.upperBound()).isInstanceOf(CharSequence.class); + assertThat(name.upperBound().toString()).isEqualTo("zzz"); + } + + @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( + 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( + 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); + } + + @Test + void fieldWithoutStatsIsExcluded() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + + assertThat(stats.containsFieldInMaps(5)).isFalse(); + assertThat(stats.type().field(StatsUtil.toBaseId(5))).isNull(); + assertThat(stats.statsFor(5)).isNull(); + assertThat(stats.fieldStats()) + .extracting(FieldStats::fieldId) + .containsExactlyInAnyOrder(1, 2, 3, 4); + } + + @Test + void unknownFieldIdInMaps() { + DataFile file = + dataFile( + ImmutableMap.of(99, 1L), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of(), + ImmutableMap.of()); + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(file); + + assertThat(stats.containsFieldInMaps(99)).isTrue(); + assertThat(stats.statsFor(99)).isNull(); + assertThat(stats.type().fields()).isEmpty(); + assertThat(stats.fieldStats()).isEmpty(); + } + + @Test + void copyNotSupported() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + + assertThatThrownBy(stats::copy) + .isInstanceOf(UnsupportedOperationException.class) + .hasMessage("copy is not implemented"); + + assertThatThrownBy(() -> stats.copy(ImmutableSet.of(1))) + .isInstanceOf(UnsupportedOperationException.class) + .hasMessage("copy is not implemented"); + } + + @Test + void fieldStatsCopyNotSupported() { + MapBackedContentStats stats = new MapBackedContentStats(SCHEMA).wrap(FILE_WITH_STATS); + + assertThatThrownBy(stats.statsFor(1)::copy) + .isInstanceOf(UnsupportedOperationException.class) + .hasMessage("copy is not implemented"); + } + + @Test + void reuseRebindsBounds() { Review Comment: Renamed to reuseRebindsFields. Both files now assert lower bound, upper bound, and value count for field 1. The second file still expects statsFor(2) to be null. ########## core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java: ########## @@ -671,6 +729,399 @@ void unknownSpecIdThrows() { .hasMessageContaining("Cannot find partition spec for spec ID"); } + @Test + void dataTrackedFileAdapterHasNoTracking() { + 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 dataTrackedFileAdapterFromAddedManifestEntry() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(addedEntry(DATA_FILE)); Review Comment: Removed addedEntry, existingEntry, and deletedEntry. The remaining tests call wrapExisting directly. ########## core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java: ########## @@ -671,6 +729,399 @@ void unknownSpecIdThrows() { .hasMessageContaining("Cannot find partition spec for spec ID"); } + @Test + void dataTrackedFileAdapterHasNoTracking() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(DATA_FILE); + + assertThat(result.tracking()).isNull(); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); Review Comment: Renamed to dataTrackedFileAdapterFromDataFile. ########## core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java: ########## @@ -671,6 +729,399 @@ void unknownSpecIdThrows() { .hasMessageContaining("Cannot find partition spec for spec ID"); } + @Test + void dataTrackedFileAdapterHasNoTracking() { + 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 dataTrackedFileAdapterFromAddedManifestEntry() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(addedEntry(DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.ADDED); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().dataSequenceNumber()).isNull(); + assertThat(result.tracking().fileSequenceNumber()).isNull(); + assertThat(result.tracking().firstRowId()).isNull(); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterFromAddedManifestEntryWithNullSnapshotId() { + TrackedFile result = + TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(newEntry().wrapAppend(null, DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.ADDED); + assertThat(result.tracking().snapshotId()).isNull(); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterFromExistingManifestEntry() { + TrackedFile result = + TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(existingEntry(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); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterFromDeletedManifestEntry() { + TrackedFile result = + TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(deletedEntry(DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.DELETED); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().dataSequenceNumber()).isEqualTo(DATA_SEQUENCE_NUMBER); + assertThat(result.tracking().fileSequenceNumber()).isEqualTo(FILE_SEQUENCE_NUMBER); + assertManifestPosition(result.tracking(), DATA_FILE); + } + + @Test + void dataTrackedFileAdapterReuse() { + TrackedFileAdapters.DataTrackedFile adapter = TrackedFileAdapters.forDataFile(TABLE_SCHEMA); + + adapter.wrap(DATA_FILE); + assertThat(adapter.location()).isEqualTo(DATA_FILE_LOCATION); Review Comment: dataTrackedFileAdapterReuse now calls assertWrappedDataFileMatchesFileFields after each wrap. The test only asserts tracking itself: null, then EXISTING and the manifest position. ########## core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java: ########## @@ -671,6 +729,399 @@ void unknownSpecIdThrows() { .hasMessageContaining("Cannot find partition spec for spec ID"); } + @Test + void dataTrackedFileAdapterHasNoTracking() { + 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 dataTrackedFileAdapterFromAddedManifestEntry() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(addedEntry(DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.ADDED); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().dataSequenceNumber()).isNull(); + assertThat(result.tracking().fileSequenceNumber()).isNull(); + assertThat(result.tracking().firstRowId()).isNull(); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterFromAddedManifestEntryWithNullSnapshotId() { Review Comment: Removed this. wrapAppend stores the snapshot id it is given, and the existing-entry test already checks a non-null id. ########## core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java: ########## @@ -692,10 +1143,14 @@ private static PartitionData partition(String category) { /** Minimal file with no tracking, used by the rejection and null-tracking tests. */ private static TrackedFileStruct dummyTrackedFile(FileContent contentType) { + return dummyTrackedFile(contentType, FORMAT_VERSION_V4); + } + + private static TrackedFileStruct dummyTrackedFile(FileContent contentType, int formatVersion) { Review Comment: Renamed dummyTrackedFile to trackedFile. ########## core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java: ########## @@ -671,6 +729,399 @@ void unknownSpecIdThrows() { .hasMessageContaining("Cannot find partition spec for spec ID"); } + @Test + void dataTrackedFileAdapterHasNoTracking() { + 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 dataTrackedFileAdapterFromAddedManifestEntry() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(addedEntry(DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.ADDED); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().dataSequenceNumber()).isNull(); + assertThat(result.tracking().fileSequenceNumber()).isNull(); + assertThat(result.tracking().firstRowId()).isNull(); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterFromAddedManifestEntryWithNullSnapshotId() { + TrackedFile result = + TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(newEntry().wrapAppend(null, DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.ADDED); + assertThat(result.tracking().snapshotId()).isNull(); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterFromExistingManifestEntry() { + TrackedFile result = + TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(existingEntry(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); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterFromDeletedManifestEntry() { + TrackedFile result = + TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(deletedEntry(DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.DELETED); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().dataSequenceNumber()).isEqualTo(DATA_SEQUENCE_NUMBER); + assertThat(result.tracking().fileSequenceNumber()).isEqualTo(FILE_SEQUENCE_NUMBER); + assertManifestPosition(result.tracking(), DATA_FILE); + } + + @Test + void dataTrackedFileAdapterReuse() { + TrackedFileAdapters.DataTrackedFile adapter = TrackedFileAdapters.forDataFile(TABLE_SCHEMA); + + adapter.wrap(DATA_FILE); + assertThat(adapter.location()).isEqualTo(DATA_FILE_LOCATION); + assertThat(adapter.recordCount()).isEqualTo(100L); + 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(existingEntry(file2)); + assertThat(adapter.location()).isEqualTo("s3://bucket/data/file2.parquet"); + assertThat(adapter.recordCount()).isEqualTo(200L); + assertThat(adapter.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertManifestPosition(adapter.tracking(), file2); + assertThat(adapter.fileSizeInBytes()).isEqualTo(2048L); + assertThat(adapter.location()).isNotEqualTo(DATA_FILE.location()); + } + + @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().get(0, String.class)).isEqualTo("books"); + assertThat( + TrackedFileAdapters.asDataFile(result, specsById(PARTITIONED_SPEC)) + .partition() + .get(0, String.class)) + .isEqualTo("books"); Review Comment: The partition is compared with Comparators.forType on the tracked file and again after asDataFile. ########## core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java: ########## @@ -671,6 +729,399 @@ void unknownSpecIdThrows() { .hasMessageContaining("Cannot find partition spec for spec ID"); } + @Test + void dataTrackedFileAdapterHasNoTracking() { + 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 dataTrackedFileAdapterFromAddedManifestEntry() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(addedEntry(DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.ADDED); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().dataSequenceNumber()).isNull(); + assertThat(result.tracking().fileSequenceNumber()).isNull(); + assertThat(result.tracking().firstRowId()).isNull(); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterFromAddedManifestEntryWithNullSnapshotId() { + TrackedFile result = + TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(newEntry().wrapAppend(null, DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.ADDED); + assertThat(result.tracking().snapshotId()).isNull(); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterFromExistingManifestEntry() { + TrackedFile result = + TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(existingEntry(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); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterFromDeletedManifestEntry() { Review Comment: Removed dataTrackedFileAdapterFromDeletedManifestEntry. ########## core/src/test/java/org/apache/iceberg/TestTrackedFileAdapters.java: ########## @@ -671,6 +729,399 @@ void unknownSpecIdThrows() { .hasMessageContaining("Cannot find partition spec for spec ID"); } + @Test + void dataTrackedFileAdapterHasNoTracking() { + 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 dataTrackedFileAdapterFromAddedManifestEntry() { + TrackedFile result = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(addedEntry(DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.ADDED); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().dataSequenceNumber()).isNull(); + assertThat(result.tracking().fileSequenceNumber()).isNull(); + assertThat(result.tracking().firstRowId()).isNull(); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterFromAddedManifestEntryWithNullSnapshotId() { + TrackedFile result = + TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(newEntry().wrapAppend(null, DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.ADDED); + assertThat(result.tracking().snapshotId()).isNull(); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterFromExistingManifestEntry() { + TrackedFile result = + TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(existingEntry(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); + assertManifestPosition(result.tracking(), DATA_FILE); + assertWrappedDataFileMatchesFileFields(result, DATA_FILE); + } + + @Test + void dataTrackedFileAdapterFromDeletedManifestEntry() { + TrackedFile result = + TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(deletedEntry(DATA_FILE)); + + assertThat(result.tracking().status()).isEqualTo(EntryStatus.DELETED); + assertThat(result.tracking().snapshotId()).isEqualTo(SNAPSHOT_ID); + assertThat(result.tracking().dataSequenceNumber()).isEqualTo(DATA_SEQUENCE_NUMBER); + assertThat(result.tracking().fileSequenceNumber()).isEqualTo(FILE_SEQUENCE_NUMBER); + assertManifestPosition(result.tracking(), DATA_FILE); + } + + @Test + void dataTrackedFileAdapterReuse() { + TrackedFileAdapters.DataTrackedFile adapter = TrackedFileAdapters.forDataFile(TABLE_SCHEMA); + + adapter.wrap(DATA_FILE); + assertThat(adapter.location()).isEqualTo(DATA_FILE_LOCATION); + assertThat(adapter.recordCount()).isEqualTo(100L); + 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(existingEntry(file2)); + assertThat(adapter.location()).isEqualTo("s3://bucket/data/file2.parquet"); + assertThat(adapter.recordCount()).isEqualTo(200L); + assertThat(adapter.tracking().status()).isEqualTo(EntryStatus.EXISTING); + assertManifestPosition(adapter.tracking(), file2); + assertThat(adapter.fileSizeInBytes()).isEqualTo(2048L); + assertThat(adapter.location()).isNotEqualTo(DATA_FILE.location()); + } + + @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().get(0, String.class)).isEqualTo("books"); + assertThat( + TrackedFileAdapters.asDataFile(result, specsById(PARTITIONED_SPEC)) + .partition() + .get(0, String.class)) + .isEqualTo("books"); + } + + @Test + void dataFileDoubleWrapRoundTrip() { + DataFile source = DATA_FILE_WITH_METRICS; + Map<Integer, PartitionSpec> specs = + ImmutableMap.of(UNPARTITIONED_SPEC.specId(), UNPARTITIONED_SPEC); + + TrackedFile tracked = TrackedFileAdapters.forDataFile(TABLE_SCHEMA).wrap(existingEntry(source)); + + DataFile roundTripped = TrackedFileAdapters.asDataFile(tracked, specs); + + assertThat(roundTripped.content()).isEqualTo(FileContent.DATA); Review Comment: Renamed this to trackedFileDoubleWrapRoundTrip. It builds a fresh TrackedFile, converts it with asDataFile, and asserts wrap returns that same TrackedFile. asDataFile always allocates a new TrackedDataFile, so the DataFile to TrackedFile to DataFile path is not identity. The original-file check is only in wrap, when the argument is already a TrackedDataFile. -- 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]
