Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 20 additions & 0 deletions api/src/main/java/org/apache/iceberg/ManifestFile.java
Original file line number Diff line number Diff line change
Expand Up @@ -186,6 +186,26 @@ default boolean hasDeletedFiles() {
/** Returns the total number of rows in all files with status DELETED in the manifest file. */
Long deletedRowsCount();

/** Returns the number of files with status REPLACED in the manifest file. */
default Integer replacedFilesCount() {
return 0;
}

/** Returns the total number of rows in all files with status REPLACED in the manifest file. */
default Long replacedRowsCount() {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

What about modified rows?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agreed. Added modifiedFilesCount() / modifiedRowsCount() on ManifestFile, defaulting to 0 like replaced.

TrackedManifestFile forwards both from ManifestInfo. Pre-v4 wrap requires modified files count to be 0 and includes it in the record-count sum.

return 0L;
}

/** Returns the number of files with status MODIFIED in the manifest file. */
default Integer modifiedFilesCount() {
return 0;
}

/** Returns the total number of rows in all files with status MODIFIED in the manifest file. */
default Long modifiedRowsCount() {
return 0L;
}

/**
* Returns a list of {@link PartitionFieldSummary partition field summaries}.
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@
import java.nio.ByteBuffer;
import java.nio.ByteOrder;
import java.util.Objects;
import org.apache.iceberg.StructLike;
import org.apache.iceberg.relocated.com.google.common.base.Preconditions;

/**
Expand All @@ -44,7 +45,7 @@
* <p>This class represents a lower or upper geospatial bound and handles serialization and
* deserialization of these bounds to/from byte arrays, conforming to the Iceberg specification.
*/
public class GeospatialBound {
public class GeospatialBound implements StructLike {
/**
* Parses a geospatial bound from a byte buffer according to Iceberg spec.
*
Expand Down Expand Up @@ -282,6 +283,27 @@ public boolean hasM() {
return !Double.isNaN(m);
}

@Override
public int size() {
return 4;
}

@Override
public <T> T get(int pos, Class<T> javaClass) {
return switch (pos) {
case 0 -> javaClass.cast(x);
case 1 -> javaClass.cast(y);
case 2 -> hasZ() ? javaClass.cast(z) : null;
case 3 -> hasM() ? javaClass.cast(m) : null;
default -> throw new IllegalArgumentException("Invalid position: " + pos);
};
}

@Override
public <T> void set(int pos, T value) {
throw new UnsupportedOperationException("GeospatialBound is read only");

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Not part of this PR, but should we use GeospatialBound when reading? I think right now we use a generic struct for it but we could use this.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I created a separate PR to use GeospatialBound for the read path. #18394

}

@Override
public String toString() {
return "GeospatialBound(" + simpleString() + ")";
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,9 +19,11 @@
package org.apache.iceberg.geospatial;

import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;

import java.nio.ByteBuffer;
import java.nio.ByteOrder;
import org.apache.iceberg.StructLike;
import org.apache.iceberg.util.ByteBuffers;
import org.junit.jupiter.api.Test;

Expand Down Expand Up @@ -71,6 +73,34 @@ public void testCreateXYZM() {
assertThat(bound.hasM()).isTrue();
}

@Test
public void testStructLikeXY() {
StructLike bound = GeospatialBound.createXY(1.0, 2.0);
assertThat(bound.size()).isEqualTo(4);
assertThat(bound.get(0, Double.class)).isEqualTo(1.0);
assertThat(bound.get(1, Double.class)).isEqualTo(2.0);
assertThat(bound.get(2, Double.class)).isNull();
assertThat(bound.get(3, Double.class)).isNull();
}

@Test
public void testStructLikeXYZM() {
StructLike bound = GeospatialBound.createXYZM(3.0, 4.0, 5.0, 6.0);
assertThat(bound.size()).isEqualTo(4);
assertThat(bound.get(0, Double.class)).isEqualTo(3.0);
assertThat(bound.get(1, Double.class)).isEqualTo(4.0);
assertThat(bound.get(2, Double.class)).isEqualTo(5.0);
assertThat(bound.get(3, Double.class)).isEqualTo(6.0);
}

@Test
public void testSet() {
GeospatialBound bound = GeospatialBound.createXY(1.0, 2.0);
assertThatThrownBy(() -> bound.set(0, 9.0))
.isInstanceOf(UnsupportedOperationException.class)
.hasMessage("GeospatialBound is read only");
}

@Test
public void testEqualsAndHashCode() {
GeospatialBound xy1 = GeospatialBound.createXY(1.0, 2.0);
Expand Down
244 changes: 244 additions & 0 deletions core/src/main/java/org/apache/iceberg/MapBackedContentStats.java
Original file line number Diff line number Diff line change
@@ -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);
}

private 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() {
Comment thread
rdblue marked this conversation as resolved.
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");
}

private boolean containsFieldInMaps(int fieldId) {
return containsId(valueCounts, fieldId)
|| containsId(nullValueCounts, fieldId)
|| containsId(nanValueCounts, fieldId)
|| containsId(avgValueSizes, fieldId)
|| containsId(lowerBounds, fieldId)
|| containsId(upperBounds, fieldId);
}

private Set<Integer> statsFieldIds() {
Set<Integer> ids =
Sets.newHashSetWithExpectedSize(valueCounts == null ? 0 : valueCounts.size());
addKeys(ids, valueCounts);
addKeys(ids, nullValueCounts);
addKeys(ids, nanValueCounts);
addKeys(ids, avgValueSizes);
addKeys(ids, lowerBounds);
addKeys(ids, upperBounds);
return ids;
}

private static void addKeys(Set<Integer> ids, Map<Integer, ?> map) {
if (map != null) {
ids.addAll(map.keySet());
}
}

private static boolean containsId(Map<Integer, ?> map, int id) {
return map != null && map.containsKey(id);
}

/** Reusable {@link FieldStats} view over one field's entries in a {@link ContentFile}'s maps. */
private class MapBackedFieldStats<T> implements FieldStats<T> {
private final int fieldId;
private final Types.StructType struct;
private final Type fieldType;
private final Type boundType;

private MapBackedFieldStats(int fieldId, Type fieldType, Types.StructType struct) {
this.fieldId = fieldId;
this.fieldType = fieldType;
this.struct = struct;
this.boundType = struct.fieldType(StatsUtil.LOWER_BOUND_NAME);
}

@Override
public int fieldId() {
return fieldId;
}

@Override
public Types.StructType type() {
return struct;
}

@Override
public T lowerBound() {
return decodeBound(lowerBounds);
Comment thread
rdblue marked this conversation as resolved.
}

@Override
public T upperBound() {
return decodeBound(upperBounds);
}

@SuppressWarnings("unchecked")
private T decodeBound(Map<Integer, ByteBuffer> bounds) {
if (boundType == null || bounds == null) {
return null;
}

return (T) Conversions.fromByteBuffer(fieldType, bounds.get(fieldId));
}

@Override
public boolean tightBounds() {
return false;
}

@Override
public boolean hasValueCount() {
return count(valueCounts) != null;
}

@Override
public long valueCount() {
return count(valueCounts);
}

@Override
public boolean hasNullValueCount() {
return count(nullValueCounts) != null;
}

@Override
public long nullValueCount() {
return count(nullValueCounts);
}

@Override
public boolean hasNanValueCount() {
return count(nanValueCounts) != null;
}

@Override
public long nanValueCount() {
return count(nanValueCounts);
}

@Override
public Integer avgValueSizeInBytes() {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

will adjust this once Eduard's avg -> total metric change landed.

return avgValueSizes == null ? null : avgValueSizes.get(fieldId);
}

@Override
public FieldStats<T> copy() {
throw new UnsupportedOperationException("copy is not implemented");
}

private Long count(Map<Integer, Long> counts) {
return counts == null ? null : counts.get(fieldId);
}
}
}
Loading
Loading