Skip to content
Merged
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
Original file line number Diff line number Diff line change
Expand Up @@ -116,7 +116,6 @@ public ParquetValueReader<RowData> struct(
expected != null ? expected.fields() : ImmutableList.of();
List<ParquetValueReader<?>> reorderedFields =
Lists.newArrayListWithExpectedSize(expectedFields.size());
List<Type> types = Lists.newArrayListWithExpectedSize(expectedFields.size());
// Defaulting to parent max definition level
int defaultMaxDefinitionLevel = type.getMaxDefinitionLevel(currentPath());
for (Types.NestedField field : expectedFields) {
Expand All @@ -128,32 +127,26 @@ public ParquetValueReader<RowData> 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
Expand Down Expand Up @@ -662,8 +655,8 @@ private static class RowDataReader
extends ParquetValueReaders.StructReader<RowData, GenericRowData> {
private final int numFields;

RowDataReader(List<Type> types, List<ParquetValueReader<?>> readers) {
super(types, readers);
RowDataReader(List<ParquetValueReader<?>> readers) {
super(readers);
this.numFields = readers.size();
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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);
Expand All @@ -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
Expand Down Expand Up @@ -148,96 +167,79 @@ public ParquetValueReader<?> struct(
}
}

private class LogicalTypeAnnotationParquetValueReaderVisitor
implements LogicalTypeAnnotation.LogicalTypeAnnotationVisitor<ParquetValueReader<?>> {
private class LogicalTypeReadBuilder
implements LogicalTypeAnnotationVisitor<ParquetValueReader<?>> {

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<ParquetValueReader<?>> visit(
LogicalTypeAnnotation.StringLogicalTypeAnnotation stringLogicalType) {
return Optional.of(new ParquetValueReaders.StringReader(desc));
public Optional<ParquetValueReader<?>> visit(StringLogicalTypeAnnotation stringLogicalType) {
return Optional.of(ParquetValueReaders.strings(desc));
}

@Override
public Optional<ParquetValueReader<?>> visit(
LogicalTypeAnnotation.EnumLogicalTypeAnnotation enumLogicalType) {
return Optional.of(new ParquetValueReaders.StringReader(desc));
public Optional<ParquetValueReader<?>> visit(EnumLogicalTypeAnnotation enumLogicalType) {
return Optional.of(ParquetValueReaders.strings(desc));
}

@Override
public Optional<ParquetValueReader<?>> 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<ParquetValueReader<?>> visit(
LogicalTypeAnnotation.DateLogicalTypeAnnotation dateLogicalType) {
public Optional<ParquetValueReader<?>> visit(DateLogicalTypeAnnotation dateLogicalType) {
return Optional.of(dateReader(desc));
}

@Override
public Optional<ParquetValueReader<?>> visit(
LogicalTypeAnnotation.TimeLogicalTypeAnnotation timeLogicalType) {
return Optional.of(timeReader(desc, timeLogicalType.getUnit()));
public Optional<ParquetValueReader<?>> visit(TimeLogicalTypeAnnotation timeLogicalType) {
return Optional.of(timeReader(desc));
}

@Override
public Optional<ParquetValueReader<?>> 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<ParquetValueReader<?>> visit(
LogicalTypeAnnotation.IntLogicalTypeAnnotation intLogicalType) {
public Optional<ParquetValueReader<?>> visit(IntLogicalTypeAnnotation intLogicalType) {
if (intLogicalType.getBitWidth() == 64) {
Preconditions.checkArgument(

@ajantha-bhat ajantha-bhat Jan 29, 2025

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

We wanted to throw the exception when it is signed right?

Now it throws when it is unsigned. check should be inverted?

@ajantha-bhat ajantha-bhat Jan 29, 2025

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

oh, got it. Based on the error message, we should not read unsigned. So, previous check was wrong!

While fixing nits, we actually fixed the check.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

I think we have lot of gaps in test coverage for parquet readers like milli-time and timestamp reading, int96 reading, all kinds of signed, unsigned reads. I will create a ticket to follow this up.

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<ParquetValueReader<?>> visit(
LogicalTypeAnnotation.JsonLogicalTypeAnnotation jsonLogicalType) {
return Optional.of(new ParquetValueReaders.StringReader(desc));
public Optional<ParquetValueReader<?>> visit(JsonLogicalTypeAnnotation jsonLogicalType) {
return Optional.of(ParquetValueReaders.strings(desc));
}

@Override
public Optional<ParquetValueReader<?>> visit(
LogicalTypeAnnotation.BsonLogicalTypeAnnotation bsonLogicalType) {
return Optional.of(new ParquetValueReaders.BytesReader(desc));
return Optional.of(ParquetValueReaders.byteBuffers(desc));
}

@Override
Expand Down Expand Up @@ -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(
Expand All @@ -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);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -59,7 +59,7 @@ public static ParquetValueReader<Record> buildReader(
@Override
protected ParquetValueReader<Record> createStructReader(
List<Type> types, List<ParquetValueReader<?>> fieldReaders, StructType structType) {
return ParquetValueReaders.recordReader(types, fieldReaders, structType);
return ParquetValueReaders.recordReader(fieldReaders, structType);
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<T extends StructLike> extends BaseParquetReaders<T> {
Expand All @@ -53,8 +51,7 @@ public static <T extends StructLike> ParquetValueReader<T> create(
@SuppressWarnings("unchecked")
protected ParquetValueReader<T> createStructReader(
List<Type> types, List<ParquetValueReader<?>> fieldReaders, StructType structType) {
return (ParquetValueReader<T>)
ParquetValueReaders.recordReader(types, fieldReaders, structType);
return (ParquetValueReader<T>) ParquetValueReaders.recordReader(fieldReaders, structType);
}

@Override
Expand All @@ -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);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -100,20 +100,17 @@ public ParquetValueReader<?> struct(
expected != null ? expected.fields() : ImmutableList.of();
List<ParquetValueReader<?>> reorderedFields =
Lists.newArrayListWithExpectedSize(expectedFields.size());
List<Type> 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
Expand Down Expand Up @@ -346,8 +343,8 @@ public long readLong() {
static class RecordReader extends StructReader<Record, Record> {
private final Schema schema;

RecordReader(List<Type> types, List<ParquetValueReader<?>> readers, Schema schema) {
super(types, readers);
RecordReader(List<ParquetValueReader<?>> readers, Schema schema) {
super(readers);
this.schema = schema;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,9 @@ public interface ParquetValueReader<T> {
* 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(
Expand Down
Loading