From bfc36e027fe420b5ca24077ea5dc0fbc8be0c769 Mon Sep 17 00:00:00 2001 From: Ryan Blue Date: Fri, 24 Jan 2025 16:22:49 -0800 Subject: [PATCH 1/7] Parquet: Remove unnecessary types passed to RecordReader. --- .../iceberg/flink/data/FlinkParquetReaders.java | 13 +++---------- .../data/parquet/BaseParquetReaders.java | 12 +++++++++++- .../data/parquet/GenericParquetReaders.java | 5 ++--- .../iceberg/data/parquet/InternalReader.java | 5 ++--- .../parquet/ParquetAvroValueReaders.java | 9 +++------ .../iceberg/parquet/ParquetValueReaders.java | 16 ++++++++++++---- .../iceberg/spark/data/SparkParquetReaders.java | 17 ++++------------- 7 files changed, 37 insertions(+), 40 deletions(-) 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..f0101d0b5b9d 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 @@ -67,8 +67,18 @@ protected ParquetValueReader createReader( } } + /** + * @deprecated will be removed in 1.9.0; use {@link #createStructReader(List, Types.StructType)} + * instead. + */ + @Deprecated + protected ParquetValueReader createStructReader( + List types, List> fieldReaders, Types.StructType structType) { + return createStructReader(fieldReaders, structType); + } + protected abstract ParquetValueReader createStructReader( - List types, List> fieldReaders, Types.StructType structType); + List> fieldReaders, Types.StructType structType); protected ParquetValueReader fixedReader(ColumnDescriptor desc) { return new GenericParquetReaders.FixedReader(desc); 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..039189b0f634 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 @@ -38,7 +38,6 @@ import org.apache.iceberg.types.Types.StructType; import org.apache.parquet.column.ColumnDescriptor; import org.apache.parquet.schema.MessageType; -import org.apache.parquet.schema.Type; public class GenericParquetReaders extends BaseParquetReaders { @@ -58,8 +57,8 @@ public static ParquetValueReader buildReader( @Override protected ParquetValueReader createStructReader( - List types, List> fieldReaders, StructType structType) { - return ParquetValueReaders.recordReader(types, fieldReaders, structType); + List> fieldReaders, StructType 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..5599676ead00 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 @@ -29,7 +29,6 @@ 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 { @@ -52,9 +51,9 @@ public static ParquetValueReader create( @Override @SuppressWarnings("unchecked") protected ParquetValueReader createStructReader( - List types, List> fieldReaders, StructType structType) { + List> fieldReaders, StructType structType) { return (ParquetValueReader) - ParquetValueReaders.recordReader(types, fieldReaders, structType); + ParquetValueReaders.recordReader(fieldReaders, structType); } @Override 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/ParquetValueReaders.java b/parquet/src/main/java/org/apache/iceberg/parquet/ParquetValueReaders.java index 31f73b3bce74..da8ad6fb0c2d 100644 --- a/parquet/src/main/java/org/apache/iceberg/parquet/ParquetValueReaders.java +++ b/parquet/src/main/java/org/apache/iceberg/parquet/ParquetValueReaders.java @@ -86,8 +86,8 @@ public static ParquetValueReader millisAsTimestamps(ColumnDescriptor 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 { @@ -826,7 +826,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 = @@ -943,8 +951,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..8b3b02e41351 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 @@ -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(); } From 4a9425dea38e6374223808fd6cf51a26e198d584 Mon Sep 17 00:00:00 2001 From: Ryan Blue Date: Fri, 24 Jan 2025 16:46:05 -0800 Subject: [PATCH 2/7] Parquet: Simplify reader construction. --- .../data/parquet/BaseParquetReaders.java | 131 +++++++++--------- .../iceberg/data/parquet/InternalReader.java | 27 +--- .../iceberg/parquet/ParquetValueReaders.java | 94 +++++++++++-- .../spark/data/SparkParquetReaders.java | 4 +- 4 files changed, 159 insertions(+), 97 deletions(-) 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 f0101d0b5b9d..e200611eb0d9 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; @@ -88,8 +98,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); @@ -100,12 +114,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 @@ -158,96 +177,82 @@ 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( + LogicalTypeReadBuilder( ColumnDescriptor desc, - org.apache.iceberg.types.Type.PrimitiveType expected, - PrimitiveType primitive) { + 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) { + if (intLogicalType.isSigned()) { + // this will throw an UnsupportedOperationException + return Optional.empty(); + } + 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 @@ -398,7 +403,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( @@ -409,31 +414,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/InternalReader.java b/parquet/src/main/java/org/apache/iceberg/data/parquet/InternalReader.java index 5599676ead00..03585c55c9b6 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; public class InternalReader extends BaseParquetReaders { @@ -52,8 +50,7 @@ public static ParquetValueReader create( @SuppressWarnings("unchecked") protected ParquetValueReader createStructReader( List> fieldReaders, StructType structType) { - return (ParquetValueReader) - ParquetValueReaders.recordReader(fieldReaders, structType); + return (ParquetValueReader) ParquetValueReaders.recordReader(fieldReaders, structType); } @Override @@ -67,26 +64,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/ParquetValueReaders.java b/parquet/src/main/java/org/apache/iceberg/parquet/ParquetValueReaders.java index da8ad6fb0c2d..9621326d2819 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,81 @@ 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; @@ -77,14 +159,6 @@ 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); - } - public static ParquetValueReader recordReader( List> readers, Types.StructType struct) { return new RecordReader(readers, struct); @@ -145,7 +219,7 @@ public void setPageSource(PageReadStore pageStore, long rowPosition) {} 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; @@ -211,7 +285,7 @@ public void setPageSource(PageReadStore pageStore, long rowPosition) {} public void setPageSource(PageReadStore pageStore) {} } - static class PositionReader implements ParquetValueReader { + private static class PositionReader implements ParquetValueReader { private long rowOffset = -1; private long rowGroupStart; 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 8b3b02e41351..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 @@ -251,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(); From 13ce07e23ae977fcbaf20e779ccc8932ad8a84f2 Mon Sep 17 00:00:00 2001 From: Ryan Blue Date: Fri, 24 Jan 2025 16:47:44 -0800 Subject: [PATCH 3/7] Parquet: Remove useless overrides of setPageSource. --- .../iceberg/parquet/ParquetValueReader.java | 4 ++- .../iceberg/parquet/ParquetValueReaders.java | 36 ------------------- 2 files changed, 3 insertions(+), 37 deletions(-) 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 9621326d2819..509835457e64 100644 --- a/parquet/src/main/java/org/apache/iceberg/parquet/ParquetValueReaders.java +++ b/parquet/src/main/java/org/apache/iceberg/parquet/ParquetValueReaders.java @@ -212,9 +212,6 @@ public List> columns() { return COLUMNS; } - @Override - public void setPageSource(PageReadStore pageStore, long rowPosition) {} - @Override public void setPageSource(PageReadStore pageStore) {} } @@ -278,9 +275,6 @@ public List> columns() { return children; } - @Override - public void setPageSource(PageReadStore pageStore, long rowPosition) {} - @Override public void setPageSource(PageReadStore pageStore) {} } @@ -305,11 +299,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 = @@ -337,11 +326,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)); @@ -588,11 +572,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); @@ -638,11 +617,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); @@ -762,11 +736,6 @@ protected RepeatedKeyValueReader( .build(); } - @Override - public void setPageSource(PageReadStore pageStore, long rowPosition) { - setPageSource(pageStore); - } - @Override public void setPageSource(PageReadStore pageStore) { keyReader.setPageSource(pageStore); @@ -926,11 +895,6 @@ protected StructReader(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) { From 07138774d7b42bb7f9a9be9793e9dac6a6f13bf7 Mon Sep 17 00:00:00 2001 From: Ryan Blue Date: Fri, 24 Jan 2025 16:50:52 -0800 Subject: [PATCH 4/7] Parquet: Remove unnecessary ParquetValueReaders prefix. --- .../org/apache/iceberg/parquet/ParquetValueReaders.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) 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 509835457e64..34f6432c7456 100644 --- a/parquet/src/main/java/org/apache/iceberg/parquet/ParquetValueReaders.java +++ b/parquet/src/main/java/org/apache/iceberg/parquet/ParquetValueReaders.java @@ -152,11 +152,11 @@ 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); + return new TimestampInt96Reader(desc); } public static ParquetValueReader recordReader( @@ -497,7 +497,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); } From 0a8204ce79fb48e6aa08d710043fb0c038bab276 Mon Sep 17 00:00:00 2001 From: Ryan Blue Date: Sat, 25 Jan 2025 13:13:13 -0800 Subject: [PATCH 5/7] Apply spotless. --- .../org/apache/iceberg/data/parquet/BaseParquetReaders.java | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) 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 e200611eb0d9..bc76dd82e222 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 @@ -184,8 +184,7 @@ private class LogicalTypeReadBuilder private final org.apache.iceberg.types.Type.PrimitiveType expected; LogicalTypeReadBuilder( - ColumnDescriptor desc, - org.apache.iceberg.types.Type.PrimitiveType expected) { + ColumnDescriptor desc, org.apache.iceberg.types.Type.PrimitiveType expected) { this.desc = desc; this.expected = expected; } From f52fc15f8cea8d6d80cd50d98195f8686588c043 Mon Sep 17 00:00:00 2001 From: Ryan Blue Date: Mon, 27 Jan 2025 19:29:54 -0800 Subject: [PATCH 6/7] Remove new method to fix revapi. --- .../iceberg/data/parquet/BaseParquetReaders.java | 12 +----------- .../iceberg/data/parquet/GenericParquetReaders.java | 3 ++- .../apache/iceberg/data/parquet/InternalReader.java | 3 ++- 3 files changed, 5 insertions(+), 13 deletions(-) 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 bc76dd82e222..61b7f21d9513 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 @@ -77,18 +77,8 @@ protected ParquetValueReader createReader( } } - /** - * @deprecated will be removed in 1.9.0; use {@link #createStructReader(List, Types.StructType)} - * instead. - */ - @Deprecated - protected ParquetValueReader createStructReader( - List types, List> fieldReaders, Types.StructType structType) { - return createStructReader(fieldReaders, structType); - } - protected abstract ParquetValueReader createStructReader( - List> fieldReaders, Types.StructType structType); + List types, List> fieldReaders, Types.StructType structType); protected ParquetValueReader fixedReader(ColumnDescriptor desc) { return new GenericParquetReaders.FixedReader(desc); 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 039189b0f634..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 @@ -38,6 +38,7 @@ import org.apache.iceberg.types.Types.StructType; import org.apache.parquet.column.ColumnDescriptor; import org.apache.parquet.schema.MessageType; +import org.apache.parquet.schema.Type; public class GenericParquetReaders extends BaseParquetReaders { @@ -57,7 +58,7 @@ public static ParquetValueReader buildReader( @Override protected ParquetValueReader createStructReader( - List> fieldReaders, StructType structType) { + List types, List> fieldReaders, StructType structType) { return ParquetValueReaders.recordReader(fieldReaders, structType); } 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 03585c55c9b6..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 @@ -27,6 +27,7 @@ import org.apache.iceberg.types.Types.StructType; import org.apache.parquet.column.ColumnDescriptor; import org.apache.parquet.schema.MessageType; +import org.apache.parquet.schema.Type; public class InternalReader extends BaseParquetReaders { @@ -49,7 +50,7 @@ public static ParquetValueReader create( @Override @SuppressWarnings("unchecked") protected ParquetValueReader createStructReader( - List> fieldReaders, StructType structType) { + List types, List> fieldReaders, StructType structType) { return (ParquetValueReader) ParquetValueReaders.recordReader(fieldReaders, structType); } From 3a814f60319ed01c36a0bfbc8081192f95fb2ffd Mon Sep 17 00:00:00 2001 From: Ryan Blue Date: Tue, 28 Jan 2025 10:07:21 -0800 Subject: [PATCH 7/7] Fix nits. --- .../org/apache/iceberg/data/parquet/BaseParquetReaders.java | 6 ++---- .../org/apache/iceberg/parquet/ParquetValueReaders.java | 1 + 2 files changed, 3 insertions(+), 4 deletions(-) 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 61b7f21d9513..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 @@ -214,10 +214,8 @@ public Optional> visit( @Override public Optional> visit(IntLogicalTypeAnnotation intLogicalType) { if (intLogicalType.getBitWidth() == 64) { - if (intLogicalType.isSigned()) { - // this will throw an UnsupportedOperationException - return Optional.empty(); - } + Preconditions.checkArgument( + intLogicalType.isSigned(), "Cannot read UINT64 as a long value"); return Optional.of(new ParquetValueReaders.UnboxedReader<>(desc)); } 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 34f6432c7456..73ce83b9bfdd 100644 --- a/parquet/src/main/java/org/apache/iceberg/parquet/ParquetValueReaders.java +++ b/parquet/src/main/java/org/apache/iceberg/parquet/ParquetValueReaders.java @@ -96,6 +96,7 @@ public static ParquetValueReader bigDecimals(ColumnDescriptor desc) case INT32: return new IntegerAsDecimalReader(desc, scale); } + throw new IllegalArgumentException( "Invalid primitive type for decimal: " + desc.getPrimitiveType()); }