diff --git a/flink/v1.20/flink/src/main/java/org/apache/iceberg/flink/data/FlinkParquetReaders.java b/flink/v1.20/flink/src/main/java/org/apache/iceberg/flink/data/FlinkParquetReaders.java index a23fb2d6ee36..fc407fe2a1a8 100644 --- a/flink/v1.20/flink/src/main/java/org/apache/iceberg/flink/data/FlinkParquetReaders.java +++ b/flink/v1.20/flink/src/main/java/org/apache/iceberg/flink/data/FlinkParquetReaders.java @@ -116,7 +116,6 @@ public ParquetValueReader struct( expected != null ? expected.fields() : ImmutableList.of(); List> reorderedFields = Lists.newArrayListWithExpectedSize(expectedFields.size()); - List types = Lists.newArrayListWithExpectedSize(expectedFields.size()); // Defaulting to parent max definition level int defaultMaxDefinitionLevel = type.getMaxDefinitionLevel(currentPath()); for (Types.NestedField field : expectedFields) { @@ -128,32 +127,26 @@ public ParquetValueReader struct( maxDefinitionLevelsById.getOrDefault(id, defaultMaxDefinitionLevel); reorderedFields.add( ParquetValueReaders.constant(idToConstant.get(id), fieldMaxDefinitionLevel)); - types.add(null); } else if (id == MetadataColumns.ROW_POSITION.fieldId()) { reorderedFields.add(ParquetValueReaders.position()); - types.add(null); } else if (id == MetadataColumns.IS_DELETED.fieldId()) { reorderedFields.add(ParquetValueReaders.constant(false)); - types.add(null); } else if (reader != null) { reorderedFields.add(reader); - types.add(typesById.get(id)); } else if (field.initialDefault() != null) { reorderedFields.add( ParquetValueReaders.constant( RowDataUtil.convertConstant(field.type(), field.initialDefault()), maxDefinitionLevelsById.getOrDefault(id, defaultMaxDefinitionLevel))); - types.add(typesById.get(id)); } else if (field.isOptional()) { reorderedFields.add(ParquetValueReaders.nulls()); - types.add(null); } else { throw new IllegalArgumentException( String.format("Missing required field: %s", field.name())); } } - return new RowDataReader(types, reorderedFields); + return new RowDataReader(reorderedFields); } @Override @@ -662,8 +655,8 @@ private static class RowDataReader extends ParquetValueReaders.StructReader { private final int numFields; - RowDataReader(List types, List> readers) { - super(types, readers); + RowDataReader(List> readers) { + super(readers); this.numFields = readers.size(); } diff --git a/parquet/src/main/java/org/apache/iceberg/data/parquet/BaseParquetReaders.java b/parquet/src/main/java/org/apache/iceberg/data/parquet/BaseParquetReaders.java index 6fbcbb750bdf..70e6b3ff447e 100644 --- a/parquet/src/main/java/org/apache/iceberg/data/parquet/BaseParquetReaders.java +++ b/parquet/src/main/java/org/apache/iceberg/data/parquet/BaseParquetReaders.java @@ -27,15 +27,25 @@ import org.apache.iceberg.parquet.ParquetValueReader; import org.apache.iceberg.parquet.ParquetValueReaders; import org.apache.iceberg.parquet.TypeWithSchemaVisitor; +import org.apache.iceberg.relocated.com.google.common.base.Preconditions; 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.Lists; import org.apache.iceberg.relocated.com.google.common.collect.Maps; +import org.apache.iceberg.types.Type.TypeID; import org.apache.iceberg.types.Types; import org.apache.parquet.column.ColumnDescriptor; import org.apache.parquet.schema.GroupType; import org.apache.parquet.schema.LogicalTypeAnnotation; +import org.apache.parquet.schema.LogicalTypeAnnotation.DateLogicalTypeAnnotation; import org.apache.parquet.schema.LogicalTypeAnnotation.DecimalLogicalTypeAnnotation; +import org.apache.parquet.schema.LogicalTypeAnnotation.EnumLogicalTypeAnnotation; +import org.apache.parquet.schema.LogicalTypeAnnotation.IntLogicalTypeAnnotation; +import org.apache.parquet.schema.LogicalTypeAnnotation.JsonLogicalTypeAnnotation; +import org.apache.parquet.schema.LogicalTypeAnnotation.LogicalTypeAnnotationVisitor; +import org.apache.parquet.schema.LogicalTypeAnnotation.StringLogicalTypeAnnotation; +import org.apache.parquet.schema.LogicalTypeAnnotation.TimeLogicalTypeAnnotation; +import org.apache.parquet.schema.LogicalTypeAnnotation.TimestampLogicalTypeAnnotation; import org.apache.parquet.schema.MessageType; import org.apache.parquet.schema.PrimitiveType; import org.apache.parquet.schema.Type; @@ -78,8 +88,12 @@ protected ParquetValueReader dateReader(ColumnDescriptor desc) { return new GenericParquetReaders.DateReader(desc); } - protected ParquetValueReader timeReader( - ColumnDescriptor desc, LogicalTypeAnnotation.TimeUnit unit) { + protected ParquetValueReader timeReader(ColumnDescriptor desc) { + LogicalTypeAnnotation time = desc.getPrimitiveType().getLogicalTypeAnnotation(); + Preconditions.checkArgument( + time instanceof TimeLogicalTypeAnnotation, "Invalid time logical type: " + time); + + LogicalTypeAnnotation.TimeUnit unit = ((TimeLogicalTypeAnnotation) time).getUnit(); switch (unit) { case MICROS: return new GenericParquetReaders.TimeReader(desc); @@ -90,12 +104,17 @@ protected ParquetValueReader timeReader( } } - protected ParquetValueReader timestampReader( - ColumnDescriptor desc, LogicalTypeAnnotation.TimeUnit unit, boolean isAdjustedToUTC) { + protected ParquetValueReader timestampReader(ColumnDescriptor desc, boolean isAdjustedToUTC) { if (desc.getPrimitiveType().getPrimitiveTypeName() == PrimitiveType.PrimitiveTypeName.INT96) { return new GenericParquetReaders.TimestampInt96Reader(desc); } + LogicalTypeAnnotation timestamp = desc.getPrimitiveType().getLogicalTypeAnnotation(); + Preconditions.checkArgument( + timestamp instanceof TimestampLogicalTypeAnnotation, + "Invalid timestamp logical type: " + timestamp); + + LogicalTypeAnnotation.TimeUnit unit = ((TimestampLogicalTypeAnnotation) timestamp).getUnit(); switch (unit) { case MICROS: return isAdjustedToUTC @@ -148,96 +167,79 @@ public ParquetValueReader struct( } } - private class LogicalTypeAnnotationParquetValueReaderVisitor - implements LogicalTypeAnnotation.LogicalTypeAnnotationVisitor> { + private class LogicalTypeReadBuilder + implements LogicalTypeAnnotationVisitor> { private final ColumnDescriptor desc; private final org.apache.iceberg.types.Type.PrimitiveType expected; - private final PrimitiveType primitive; - LogicalTypeAnnotationParquetValueReaderVisitor( - ColumnDescriptor desc, - org.apache.iceberg.types.Type.PrimitiveType expected, - PrimitiveType primitive) { + LogicalTypeReadBuilder( + ColumnDescriptor desc, org.apache.iceberg.types.Type.PrimitiveType expected) { this.desc = desc; this.expected = expected; - this.primitive = primitive; } @Override - public Optional> visit( - LogicalTypeAnnotation.StringLogicalTypeAnnotation stringLogicalType) { - return Optional.of(new ParquetValueReaders.StringReader(desc)); + public Optional> visit(StringLogicalTypeAnnotation stringLogicalType) { + return Optional.of(ParquetValueReaders.strings(desc)); } @Override - public Optional> visit( - LogicalTypeAnnotation.EnumLogicalTypeAnnotation enumLogicalType) { - return Optional.of(new ParquetValueReaders.StringReader(desc)); + public Optional> visit(EnumLogicalTypeAnnotation enumLogicalType) { + return Optional.of(ParquetValueReaders.strings(desc)); } @Override public Optional> visit(DecimalLogicalTypeAnnotation decimalLogicalType) { - switch (primitive.getPrimitiveTypeName()) { - case BINARY: - case FIXED_LEN_BYTE_ARRAY: - return Optional.of( - new ParquetValueReaders.BinaryAsDecimalReader(desc, decimalLogicalType.getScale())); - case INT64: - return Optional.of( - new ParquetValueReaders.LongAsDecimalReader(desc, decimalLogicalType.getScale())); - case INT32: - return Optional.of( - new ParquetValueReaders.IntegerAsDecimalReader(desc, decimalLogicalType.getScale())); - default: - throw new UnsupportedOperationException( - "Unsupported base type for decimal: " + primitive.getPrimitiveTypeName()); - } + return Optional.of(ParquetValueReaders.bigDecimals(desc)); } @Override - public Optional> visit( - LogicalTypeAnnotation.DateLogicalTypeAnnotation dateLogicalType) { + public Optional> visit(DateLogicalTypeAnnotation dateLogicalType) { return Optional.of(dateReader(desc)); } @Override - public Optional> visit( - LogicalTypeAnnotation.TimeLogicalTypeAnnotation timeLogicalType) { - return Optional.of(timeReader(desc, timeLogicalType.getUnit())); + public Optional> visit(TimeLogicalTypeAnnotation timeLogicalType) { + return Optional.of(timeReader(desc)); } @Override public Optional> visit( - LogicalTypeAnnotation.TimestampLogicalTypeAnnotation timestampLogicalType) { + TimestampLogicalTypeAnnotation timestampLogicalType) { return Optional.of( - timestampReader( - desc, - timestampLogicalType.getUnit(), - ((Types.TimestampType) expected).shouldAdjustToUTC())); + timestampReader(desc, ((Types.TimestampType) expected).shouldAdjustToUTC())); } @Override - public Optional> visit( - LogicalTypeAnnotation.IntLogicalTypeAnnotation intLogicalType) { + public Optional> visit(IntLogicalTypeAnnotation intLogicalType) { if (intLogicalType.getBitWidth() == 64) { + Preconditions.checkArgument( + intLogicalType.isSigned(), "Cannot read UINT64 as a long value"); + return Optional.of(new ParquetValueReaders.UnboxedReader<>(desc)); } - return (expected.typeId() == org.apache.iceberg.types.Type.TypeID.LONG) - ? Optional.of(new ParquetValueReaders.IntAsLongReader(desc)) - : Optional.of(new ParquetValueReaders.UnboxedReader<>(desc)); + + if (expected.typeId() == TypeID.LONG) { + return Optional.of(new ParquetValueReaders.IntAsLongReader(desc)); + } + + Preconditions.checkArgument( + intLogicalType.isSigned() || intLogicalType.getBitWidth() < 32, + "Cannot read UINT32 as an int value"); + + return Optional.of(new ParquetValueReaders.UnboxedReader<>(desc)); } @Override - public Optional> visit( - LogicalTypeAnnotation.JsonLogicalTypeAnnotation jsonLogicalType) { - return Optional.of(new ParquetValueReaders.StringReader(desc)); + public Optional> visit(JsonLogicalTypeAnnotation jsonLogicalType) { + return Optional.of(ParquetValueReaders.strings(desc)); } @Override public Optional> visit( LogicalTypeAnnotation.BsonLogicalTypeAnnotation bsonLogicalType) { - return Optional.of(new ParquetValueReaders.BytesReader(desc)); + return Optional.of(ParquetValueReaders.byteBuffers(desc)); } @Override @@ -388,7 +390,7 @@ public ParquetValueReader primitive( if (primitive.getLogicalTypeAnnotation() != null) { return primitive .getLogicalTypeAnnotation() - .accept(new LogicalTypeAnnotationParquetValueReaderVisitor(desc, expected, primitive)) + .accept(new LogicalTypeReadBuilder(desc, expected)) .orElseThrow( () -> new UnsupportedOperationException( @@ -399,31 +401,31 @@ public ParquetValueReader primitive( case FIXED_LEN_BYTE_ARRAY: return fixedReader(desc); case BINARY: - if (expected.typeId() == org.apache.iceberg.types.Type.TypeID.STRING) { - return new ParquetValueReaders.StringReader(desc); + if (expected.typeId() == TypeID.STRING) { + return ParquetValueReaders.strings(desc); } else { - return new ParquetValueReaders.BytesReader(desc); + return ParquetValueReaders.byteBuffers(desc); } case INT32: - if (expected.typeId() == org.apache.iceberg.types.Type.TypeID.LONG) { - return new ParquetValueReaders.IntAsLongReader(desc); + if (expected.typeId() == TypeID.LONG) { + return ParquetValueReaders.intsAsLongs(desc); } else { - return new ParquetValueReaders.UnboxedReader<>(desc); + return ParquetValueReaders.unboxed(desc); } case FLOAT: - if (expected.typeId() == org.apache.iceberg.types.Type.TypeID.DOUBLE) { - return new ParquetValueReaders.FloatAsDoubleReader(desc); + if (expected.typeId() == TypeID.DOUBLE) { + return ParquetValueReaders.floatsAsDoubles(desc); } else { - return new ParquetValueReaders.UnboxedReader<>(desc); + return ParquetValueReaders.unboxed(desc); } case BOOLEAN: case INT64: case DOUBLE: - return new ParquetValueReaders.UnboxedReader<>(desc); + return ParquetValueReaders.unboxed(desc); case INT96: // Impala & Spark used to write timestamps as INT96 without a logical type. For backwards // compatibility we try to read INT96 as timestamps. - return timestampReader(desc, LogicalTypeAnnotation.TimeUnit.NANOS, true); + return timestampReader(desc, true); default: throw new UnsupportedOperationException("Unsupported type: " + primitive); } diff --git a/parquet/src/main/java/org/apache/iceberg/data/parquet/GenericParquetReaders.java b/parquet/src/main/java/org/apache/iceberg/data/parquet/GenericParquetReaders.java index 3bf924081ed7..8d22e2337cdf 100644 --- a/parquet/src/main/java/org/apache/iceberg/data/parquet/GenericParquetReaders.java +++ b/parquet/src/main/java/org/apache/iceberg/data/parquet/GenericParquetReaders.java @@ -59,7 +59,7 @@ public static ParquetValueReader buildReader( @Override protected ParquetValueReader createStructReader( List types, List> fieldReaders, StructType structType) { - return ParquetValueReaders.recordReader(types, fieldReaders, structType); + return ParquetValueReaders.recordReader(fieldReaders, structType); } @Override diff --git a/parquet/src/main/java/org/apache/iceberg/data/parquet/InternalReader.java b/parquet/src/main/java/org/apache/iceberg/data/parquet/InternalReader.java index 3bf0a4e80130..05613eb1de16 100644 --- a/parquet/src/main/java/org/apache/iceberg/data/parquet/InternalReader.java +++ b/parquet/src/main/java/org/apache/iceberg/data/parquet/InternalReader.java @@ -26,9 +26,7 @@ import org.apache.iceberg.parquet.ParquetValueReaders; import org.apache.iceberg.types.Types.StructType; import org.apache.parquet.column.ColumnDescriptor; -import org.apache.parquet.schema.LogicalTypeAnnotation; import org.apache.parquet.schema.MessageType; -import org.apache.parquet.schema.PrimitiveType; import org.apache.parquet.schema.Type; public class InternalReader extends BaseParquetReaders { @@ -53,8 +51,7 @@ public static ParquetValueReader create( @SuppressWarnings("unchecked") protected ParquetValueReader createStructReader( List types, List> fieldReaders, StructType structType) { - return (ParquetValueReader) - ParquetValueReaders.recordReader(types, fieldReaders, structType); + return (ParquetValueReader) ParquetValueReaders.recordReader(fieldReaders, structType); } @Override @@ -68,26 +65,12 @@ protected ParquetValueReader dateReader(ColumnDescriptor desc) { } @Override - protected ParquetValueReader timeReader( - ColumnDescriptor desc, LogicalTypeAnnotation.TimeUnit unit) { - if (unit == LogicalTypeAnnotation.TimeUnit.MILLIS) { - return ParquetValueReaders.millisAsTimes(desc); - } - - return new ParquetValueReaders.UnboxedReader<>(desc); + protected ParquetValueReader timeReader(ColumnDescriptor desc) { + return ParquetValueReaders.times(desc); } @Override - protected ParquetValueReader timestampReader( - ColumnDescriptor desc, LogicalTypeAnnotation.TimeUnit unit, boolean isAdjustedToUTC) { - if (desc.getPrimitiveType().getPrimitiveTypeName() == PrimitiveType.PrimitiveTypeName.INT96) { - return ParquetValueReaders.int96Timestamps(desc); - } - - if (unit == LogicalTypeAnnotation.TimeUnit.MILLIS) { - return ParquetValueReaders.millisAsTimestamps(desc); - } - - return new ParquetValueReaders.UnboxedReader<>(desc); + protected ParquetValueReader timestampReader(ColumnDescriptor desc, boolean isAdjustedToUTC) { + return ParquetValueReaders.timestamps(desc); } } diff --git a/parquet/src/main/java/org/apache/iceberg/parquet/ParquetAvroValueReaders.java b/parquet/src/main/java/org/apache/iceberg/parquet/ParquetAvroValueReaders.java index c30a57655ca5..98cc480d10c4 100644 --- a/parquet/src/main/java/org/apache/iceberg/parquet/ParquetAvroValueReaders.java +++ b/parquet/src/main/java/org/apache/iceberg/parquet/ParquetAvroValueReaders.java @@ -100,20 +100,17 @@ public ParquetValueReader struct( expected != null ? expected.fields() : ImmutableList.of(); List> reorderedFields = Lists.newArrayListWithExpectedSize(expectedFields.size()); - List types = Lists.newArrayListWithExpectedSize(expectedFields.size()); for (Types.NestedField field : expectedFields) { int id = field.fieldId(); ParquetValueReader reader = readersById.get(id); if (reader != null) { reorderedFields.add(reader); - types.add(typesById.get(id)); } else { reorderedFields.add(ParquetValueReaders.nulls()); - types.add(null); } } - return new RecordReader(types, reorderedFields, avroSchema); + return new RecordReader(reorderedFields, avroSchema); } @Override @@ -346,8 +343,8 @@ public long readLong() { static class RecordReader extends StructReader { private final Schema schema; - RecordReader(List types, List> readers, Schema schema) { - super(types, readers); + RecordReader(List> readers, Schema schema) { + super(readers); this.schema = schema; } diff --git a/parquet/src/main/java/org/apache/iceberg/parquet/ParquetValueReader.java b/parquet/src/main/java/org/apache/iceberg/parquet/ParquetValueReader.java index a2ade5336621..01d3e15bb43b 100644 --- a/parquet/src/main/java/org/apache/iceberg/parquet/ParquetValueReader.java +++ b/parquet/src/main/java/org/apache/iceberg/parquet/ParquetValueReader.java @@ -33,7 +33,9 @@ public interface ParquetValueReader { * instead. */ @Deprecated - void setPageSource(PageReadStore pageStore, long rowPosition); + default void setPageSource(PageReadStore pageStore, long rowPosition) { + setPageSource(pageStore); + } default void setPageSource(PageReadStore pageStore) { throw new UnsupportedOperationException( diff --git a/parquet/src/main/java/org/apache/iceberg/parquet/ParquetValueReaders.java b/parquet/src/main/java/org/apache/iceberg/parquet/ParquetValueReaders.java index 31f73b3bce74..73ce83b9bfdd 100644 --- a/parquet/src/main/java/org/apache/iceberg/parquet/ParquetValueReaders.java +++ b/parquet/src/main/java/org/apache/iceberg/parquet/ParquetValueReaders.java @@ -31,6 +31,7 @@ import java.util.UUID; import org.apache.iceberg.data.GenericRecord; import org.apache.iceberg.data.Record; +import org.apache.iceberg.relocated.com.google.common.base.Preconditions; import org.apache.iceberg.relocated.com.google.common.collect.ImmutableList; import org.apache.iceberg.relocated.com.google.common.collect.Lists; import org.apache.iceberg.relocated.com.google.common.collect.Maps; @@ -39,6 +40,12 @@ import org.apache.parquet.column.ColumnDescriptor; import org.apache.parquet.column.page.PageReadStore; import org.apache.parquet.io.api.Binary; +import org.apache.parquet.schema.LogicalTypeAnnotation; +import org.apache.parquet.schema.LogicalTypeAnnotation.DecimalLogicalTypeAnnotation; +import org.apache.parquet.schema.LogicalTypeAnnotation.TimeLogicalTypeAnnotation; +import org.apache.parquet.schema.LogicalTypeAnnotation.TimeUnit; +import org.apache.parquet.schema.LogicalTypeAnnotation.TimestampLogicalTypeAnnotation; +import org.apache.parquet.schema.PrimitiveType; import org.apache.parquet.schema.Type; public class ParquetValueReaders { @@ -52,6 +59,82 @@ public static ParquetValueReader option( return reader; } + public static ParquetValueReader unboxed(ColumnDescriptor desc) { + return new UnboxedReader<>(desc); + } + + public static ParquetValueReader strings(ColumnDescriptor desc) { + return new StringReader(desc); + } + + public static ParquetValueReader byteBuffers(ColumnDescriptor desc) { + return new BytesReader(desc); + } + + public static ParquetValueReader intsAsLongs(ColumnDescriptor desc) { + return new IntAsLongReader(desc); + } + + public static ParquetValueReader floatsAsDoubles(ColumnDescriptor desc) { + return new FloatAsDoubleReader(desc); + } + + public static ParquetValueReader bigDecimals(ColumnDescriptor desc) { + LogicalTypeAnnotation decimal = desc.getPrimitiveType().getLogicalTypeAnnotation(); + Preconditions.checkArgument( + decimal instanceof DecimalLogicalTypeAnnotation, + "Invalid timestamp logical type: " + decimal); + + int scale = ((DecimalLogicalTypeAnnotation) decimal).getScale(); + + switch (desc.getPrimitiveType().getPrimitiveTypeName()) { + case FIXED_LEN_BYTE_ARRAY: + case BINARY: + return new BinaryAsDecimalReader(desc, scale); + case INT64: + return new LongAsDecimalReader(desc, scale); + case INT32: + return new IntegerAsDecimalReader(desc, scale); + } + + throw new IllegalArgumentException( + "Invalid primitive type for decimal: " + desc.getPrimitiveType()); + } + + public static ParquetValueReader times(ColumnDescriptor desc) { + LogicalTypeAnnotation time = desc.getPrimitiveType().getLogicalTypeAnnotation(); + Preconditions.checkArgument( + time instanceof TimeLogicalTypeAnnotation, "Invalid time logical type: " + time); + + TimeUnit unit = ((TimeLogicalTypeAnnotation) time).getUnit(); + if (unit == LogicalTypeAnnotation.TimeUnit.MILLIS) { + return new TimeMillisReader(desc); + } + + return new UnboxedReader<>(desc); + } + + public static ParquetValueReader timestamps(ColumnDescriptor desc) { + if (desc.getPrimitiveType().getPrimitiveTypeName() == PrimitiveType.PrimitiveTypeName.INT96) { + return new TimestampInt96Reader(desc); + } + + LogicalTypeAnnotation timestamp = desc.getPrimitiveType().getLogicalTypeAnnotation(); + Preconditions.checkArgument( + timestamp instanceof TimestampLogicalTypeAnnotation, + "Invalid timestamp logical type: " + timestamp); + + TimeUnit unit = ((TimestampLogicalTypeAnnotation) timestamp).getUnit(); + switch (unit) { + case MILLIS: + return new TimestampMillisReader(desc); + case MICROS: + return new UnboxedReader<>(desc); + } + + throw new IllegalArgumentException("Unsupported timestamp unit: " + unit); + } + @SuppressWarnings("unchecked") public static ParquetValueReader nulls() { return (ParquetValueReader) NullReader.INSTANCE; @@ -70,24 +153,16 @@ public static ParquetValueReader position() { } public static ParquetValueReader uuids(ColumnDescriptor desc) { - return new ParquetValueReaders.UUIDReader(desc); + return new UUIDReader(desc); } public static ParquetValueReader int96Timestamps(ColumnDescriptor desc) { - return new ParquetValueReaders.TimestampInt96Reader(desc); - } - - public static ParquetValueReader millisAsTimes(ColumnDescriptor desc) { - return new ParquetValueReaders.TimeMillisReader(desc); - } - - public static ParquetValueReader millisAsTimestamps(ColumnDescriptor desc) { - return new ParquetValueReaders.TimestampMillisReader(desc); + return new TimestampInt96Reader(desc); } public static ParquetValueReader recordReader( - List types, List> readers, Types.StructType struct) { - return new RecordReader(types, readers, struct); + List> readers, Types.StructType struct) { + return new RecordReader(readers, struct); } private static class NullReader implements ParquetValueReader { @@ -138,14 +213,11 @@ public List> columns() { return COLUMNS; } - @Override - public void setPageSource(PageReadStore pageStore, long rowPosition) {} - @Override public void setPageSource(PageReadStore pageStore) {} } - static class ConstantReader implements ParquetValueReader { + private static class ConstantReader implements ParquetValueReader { private final C constantValue; private final TripleIterator column; private final List> children; @@ -204,14 +276,11 @@ public List> columns() { return children; } - @Override - public void setPageSource(PageReadStore pageStore, long rowPosition) {} - @Override public void setPageSource(PageReadStore pageStore) {} } - static class PositionReader implements ParquetValueReader { + private static class PositionReader implements ParquetValueReader { private long rowOffset = -1; private long rowGroupStart; @@ -231,11 +300,6 @@ public List> columns() { return NullReader.COLUMNS; } - @Override - public void setPageSource(PageReadStore pageStore, long rowPosition) { - setPageSource(pageStore); - } - @Override public void setPageSource(PageReadStore pageStore) { this.rowGroupStart = @@ -263,11 +327,6 @@ protected PrimitiveReader(ColumnDescriptor desc) { this.children = ImmutableList.of(column); } - @Override - public void setPageSource(PageReadStore pageStore, long rowPosition) { - setPageSource(pageStore); - } - @Override public void setPageSource(PageReadStore pageStore) { column.setPageSource(pageStore.getPageReader(desc)); @@ -439,7 +498,7 @@ public UUID read(UUID reuse) { } } - public static class ByteArrayReader extends ParquetValueReaders.PrimitiveReader { + public static class ByteArrayReader extends PrimitiveReader { public ByteArrayReader(ColumnDescriptor desc) { super(desc); } @@ -514,11 +573,6 @@ private static class OptionReader implements ParquetValueReader { this.children = reader.columns(); } - @Override - public void setPageSource(PageReadStore pageStore, long rowPosition) { - setPageSource(pageStore); - } - @Override public void setPageSource(PageReadStore pageStore) { reader.setPageSource(pageStore); @@ -564,11 +618,6 @@ protected RepeatedReader( this.children = reader.columns(); } - @Override - public void setPageSource(PageReadStore pageStore, long rowPosition) { - setPageSource(pageStore); - } - @Override public void setPageSource(PageReadStore pageStore) { reader.setPageSource(pageStore); @@ -688,11 +737,6 @@ protected RepeatedKeyValueReader( .build(); } - @Override - public void setPageSource(PageReadStore pageStore, long rowPosition) { - setPageSource(pageStore); - } - @Override public void setPageSource(PageReadStore pageStore) { keyReader.setPageSource(pageStore); @@ -826,7 +870,15 @@ public abstract static class StructReader implements ParquetValueReader private final TripleIterator column; private final List> children; + /** + * @deprecated will be removed in 1.9.0; use {@link #StructReader(List)} instead. + */ + @Deprecated protected StructReader(List types, List> readers) { + this(readers); + } + + protected StructReader(List> readers) { this.readers = (ParquetValueReader[]) Array.newInstance(ParquetValueReader.class, readers.size()); TripleIterator[] columns = @@ -844,11 +896,6 @@ protected StructReader(List types, List> readers) { this.column = firstNonNullColumn(children); } - @Override - public final void setPageSource(PageReadStore pageStore, long rowPosition) { - setPageSource(pageStore); - } - @Override public final void setPageSource(PageReadStore pageStore) { for (ParquetValueReader reader : readers) { @@ -943,8 +990,8 @@ private TripleIterator firstNonNullColumn(List> columns) { private static class RecordReader extends StructReader { private final GenericRecord template; - RecordReader(List types, List> readers, Types.StructType struct) { - super(types, readers); + RecordReader(List> readers, Types.StructType struct) { + super(readers); this.template = struct != null ? GenericRecord.create(struct) : null; } diff --git a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/data/SparkParquetReaders.java b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/data/SparkParquetReaders.java index 4fec047dc987..5a7fd50067b9 100644 --- a/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/data/SparkParquetReaders.java +++ b/spark/v3.5/spark/src/main/java/org/apache/iceberg/spark/data/SparkParquetReaders.java @@ -106,16 +106,14 @@ public ParquetValueReader struct( // the expected struct is ignored because nested fields are never found when the List> newFields = Lists.newArrayListWithExpectedSize(fieldReaders.size()); - List types = Lists.newArrayListWithExpectedSize(fieldReaders.size()); List fields = struct.getFields(); for (int i = 0; i < fields.size(); i += 1) { Type fieldType = fields.get(i); int fieldD = type().getMaxDefinitionLevel(path(fieldType.getName())) - 1; newFields.add(ParquetValueReaders.option(fieldType, fieldD, fieldReaders.get(i))); - types.add(fieldType); } - return new InternalRowReader(types, newFields); + return new InternalRowReader(newFields); } } @@ -159,7 +157,6 @@ public ParquetValueReader struct( expected != null ? expected.fields() : ImmutableList.of(); List> reorderedFields = Lists.newArrayListWithExpectedSize(expectedFields.size()); - List types = Lists.newArrayListWithExpectedSize(expectedFields.size()); // Defaulting to parent max definition level int defaultMaxDefinitionLevel = type.getMaxDefinitionLevel(currentPath()); for (Types.NestedField field : expectedFields) { @@ -171,32 +168,26 @@ public ParquetValueReader struct( maxDefinitionLevelsById.getOrDefault(id, defaultMaxDefinitionLevel); reorderedFields.add( ParquetValueReaders.constant(idToConstant.get(id), fieldMaxDefinitionLevel)); - types.add(null); } else if (id == MetadataColumns.ROW_POSITION.fieldId()) { reorderedFields.add(ParquetValueReaders.position()); - types.add(null); } else if (id == MetadataColumns.IS_DELETED.fieldId()) { reorderedFields.add(ParquetValueReaders.constant(false)); - types.add(null); } else if (reader != null) { reorderedFields.add(reader); - types.add(typesById.get(id)); } else if (field.initialDefault() != null) { reorderedFields.add( ParquetValueReaders.constant( SparkUtil.internalToSpark(field.type(), field.initialDefault()), maxDefinitionLevelsById.getOrDefault(id, defaultMaxDefinitionLevel))); - types.add(typesById.get(id)); } else if (field.isOptional()) { reorderedFields.add(ParquetValueReaders.nulls()); - types.add(null); } else { throw new IllegalArgumentException( String.format("Missing required field: %s", field.name())); } } - return new InternalRowReader(types, reorderedFields); + return new InternalRowReader(reorderedFields); } @Override @@ -260,10 +251,10 @@ public ParquetValueReader primitive( } case DATE: case INT_64: - case TIMESTAMP_MICROS: return new UnboxedReader<>(desc); + case TIMESTAMP_MICROS: case TIMESTAMP_MILLIS: - return ParquetValueReaders.millisAsTimestamps(desc); + return ParquetValueReaders.timestamps(desc); case DECIMAL: DecimalLogicalTypeAnnotation decimal = (DecimalLogicalTypeAnnotation) primitive.getLogicalTypeAnnotation(); @@ -518,8 +509,8 @@ protected MapData buildMap(ReusableMapData map) { private static class InternalRowReader extends StructReader { private final int numFields; - InternalRowReader(List types, List> readers) { - super(types, readers); + InternalRowReader(List> readers) { + super(readers); this.numFields = readers.size(); }