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
35 changes: 35 additions & 0 deletions parquet-column/src/main/java/org/apache/parquet/schema/Types.java
Original file line number Diff line number Diff line change
Expand Up @@ -344,6 +344,9 @@ public abstract static class BasePrimitiveBuilder<P, THIS extends BasePrimitiveB
private int precision = NOT_SET;
private int scale = NOT_SET;
private ColumnOrder columnOrder;
// When true and an unsupported logical/physical type combination is encountered, the
// annotation is dropped and the column order is forced to "undefined" so stats are ignored.
private boolean dropUnsupportedLogicalTypeCombinations = false;

private BasePrimitiveBuilder(P parent, PrimitiveTypeName type) {
super(parent);
Expand Down Expand Up @@ -426,8 +429,40 @@ public THIS columnOrder(ColumnOrder columnOrder) {
return self();
}

/**
* When set, an unsupported combination results in the logical type annotation being dropped
* rather than throwing. The associated statistics are also forcefully ignored by setting the
* column order to {@link ColumnOrderName#UNDEFINED}.
*
* @return this builder for method chaining
*/
public THIS dropUnsupportedLogicalTypeCombinations() {
this.dropUnsupportedLogicalTypeCombinations = true;
return self();
}

@Override
protected PrimitiveType build(String name) {
try {
return validateAndBuild(name);
} catch (IllegalStateException e) {
boolean canDrop = dropUnsupportedLogicalTypeCombinations
&& !(logicalTypeAnnotation instanceof LogicalTypeAnnotation.DecimalLogicalTypeAnnotation);
if (!canDrop) {
throw e;
}

LOGGER.warn(
"Dropping unsupported logical type annotation {} on physical type {}: {}",
logicalTypeAnnotation,
primitiveType,
e.getMessage());
return new PrimitiveType(
repetition, primitiveType, length, name, null, null, id, ColumnOrder.undefined());
}
}

private PrimitiveType validateAndBuild(String name) {
if (length == 0 && logicalTypeAnnotation instanceof LogicalTypeAnnotation.UUIDLogicalTypeAnnotation) {
length = LogicalTypeAnnotation.UUIDLogicalTypeAnnotation.BYTES;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1605,4 +1605,17 @@ public void testGeographyLogicalTypeWithoutEdgeInterpolationAlgorithm() {
Types.optional(BINARY).as(LogicalTypeAnnotation.geographyType()).named("aGeography");
assertThat(optionalGeographyActual).isEqualTo(optionalGeographyExpected);
}

@Test
public void testDropUnsupportedLogicalTypeCombinations() {
// Other tests already validate that unsupported type combinations throw by default, so this
// test only validates that the dropUnsupportedLogicalTypeCombinations flag works.
PrimitiveType pt = Types.required(BOOLEAN)
.dropUnsupportedLogicalTypeCombinations()
.as(LogicalTypeAnnotation.timestampType(true, MILLIS))
.named("bool_ts");
assertThat(pt.getPrimitiveTypeName()).isEqualTo(BOOLEAN);
assertThat(pt.getLogicalTypeAnnotation()).isNull(); // Dropped
assertThat(pt.columnOrder().getColumnOrderName()).isEqualTo(ColumnOrder.ColumnOrderName.UNDEFINED);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -2062,6 +2062,8 @@ private void buildChildren(
}
primitiveBuilder.columnOrder(columnOrder);
}
// Gracefully handle unsupported logical type combinations on the read path.
primitiveBuilder.dropUnsupportedLogicalTypeCombinations();
childBuilder = primitiveBuilder;
} else {
childBuilder = builder.group(fromParquetRepetition(schemaElement.repetition_type));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -106,12 +106,16 @@
import org.apache.parquet.format.GeospatialStatistics;
import org.apache.parquet.format.LogicalType;
import org.apache.parquet.format.MapType;
import org.apache.parquet.format.MilliSeconds;
import org.apache.parquet.format.PageHeader;
import org.apache.parquet.format.PageType;
import org.apache.parquet.format.RowGroup;
import org.apache.parquet.format.SchemaElement;
import org.apache.parquet.format.StringType;
import org.apache.parquet.format.TimeUnit;
import org.apache.parquet.format.TimestampType;
import org.apache.parquet.format.Type;
import org.apache.parquet.format.TypeDefinedOrder;
import org.apache.parquet.format.Util;
import org.apache.parquet.hadoop.ParquetReader;
import org.apache.parquet.hadoop.ParquetWriter;
Expand Down Expand Up @@ -2279,4 +2283,62 @@ public void testColumnIndexNanCountsRoundTrip() {
assertThat(roundTrip).isNotNull();
assertThat(roundTrip.getNanCounts()).containsExactly(1L, 0L, 0L);
}

@Test
public void testUnsupportedTypeCombinationDropsAnnotationAndStats() {
ParquetMetadataConverter converter = new ParquetMetadataConverter();
TimeUnit unit = new TimeUnit();
unit.setMILLIS(new MilliSeconds());
SchemaElement leaf = new SchemaElement("bool_ts")
.setRepetition_type(FieldRepetitionType.OPTIONAL)
.setType(Type.BOOLEAN)
.setLogicalType(LogicalType.TIMESTAMP(new TimestampType(true, unit)));
List<SchemaElement> parquetSchema = Lists.newArrayList(new SchemaElement("Message").setNum_children(1), leaf);
List<org.apache.parquet.format.ColumnOrder> columnOrders =
Lists.newArrayList(new org.apache.parquet.format.ColumnOrder());
columnOrders.get(0).setTYPE_ORDER(new TypeDefinedOrder());

MessageType schema = converter.fromParquetSchema(parquetSchema, columnOrders);

PrimitiveType result = schema.getType("bool_ts").asPrimitiveType();
assertThat(result.getPrimitiveTypeName()).isEqualTo(PrimitiveTypeName.BOOLEAN);
assertThat(result.getLogicalTypeAnnotation()).isNull();
assertThat(result.columnOrder().getColumnOrderName()).isEqualTo(ColumnOrder.ColumnOrderName.UNDEFINED);
}

private static PrimitiveType droppedAnnotationInt32() {
return Types.optional(PrimitiveTypeName.INT32)
.columnOrder(ColumnOrder.undefined())
.named("ts_int32");
}

@Test
public void testDroppedAnnotationIgnoresStats() {
ParquetMetadataConverter converter = new ParquetMetadataConverter();
org.apache.parquet.format.Statistics stats = new org.apache.parquet.format.Statistics();
stats.setMin_value(new byte[] {1, 2, 3, 4});
stats.setMax_value(new byte[] {0, 1, 2, 3});
stats.setNull_count(3L);

Statistics<?> result = converter.fromParquetStatistics(Version.FULL_VERSION, stats, droppedAnnotationInt32());

assertThat(result.hasNonNullValue()).isFalse();
assertThat(result.isNumNullsSet()).isTrue();
assertThat(result.getNumNulls()).isEqualTo(3L);
}

@Test
public void testDroppedAnnotationColumnIndexIsNull() {
PrimitiveType int32Type = Types.required(PrimitiveTypeName.INT32).named("i32");
ColumnIndexBuilder cb = ColumnIndexBuilder.getBuilder(int32Type, Integer.MAX_VALUE);
Statistics<?> stats = Statistics.createStats(int32Type);
stats.updateStats(-100);
stats.updateStats(100);
cb.add(stats, null);
org.apache.parquet.format.ColumnIndex parquetColumnIndex =
ParquetMetadataConverter.toParquetColumnIndex(int32Type, cb.build());

assertThat(ParquetMetadataConverter.fromParquetColumnIndex(droppedAnnotationInt32(), parquetColumnIndex))
.isNull();
}
}