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
33 changes: 33 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,38 @@ 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() {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

how is this method expected to be called in practice?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

This is called directly in ParquetMetadataConverter.java here

this.dropUnsupportedLogicalTypeCombinations = true;
return self();
}

@Override
protected PrimitiveType build(String name) {
try {
return validateAndBuild(name);
} catch (IllegalStateException e) {
if (!dropUnsupportedLogicalTypeCombinations) {
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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

not sure where the right file for it is but it would be good to have a more end-to-end test.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Added a new TestReadInvalidTypeCombination

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();
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,77 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance
* with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing,
* software distributed under the License is distributed on an
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
* KIND, either express or implied. See the License for the
* specific language governing permissions and limitations
* under the License.
*/
package org.apache.parquet.hadoop;

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

import java.net.URISyntaxException;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.Path;
import org.apache.parquet.example.data.Group;
import org.apache.parquet.hadoop.example.GroupReadSupport;
import org.apache.parquet.hadoop.metadata.ParquetMetadata;
import org.apache.parquet.hadoop.util.HadoopInputFile;
import org.apache.parquet.schema.ColumnOrder.ColumnOrderName;
import org.apache.parquet.schema.PrimitiveType;
import org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName;
import org.junit.jupiter.api.Test;

public class TestReadInvalidTypeCombination {

// Path to a Parquet file that contains an invalid logical/physical type combination.
private static final String FILE_PATH = "/invalid_type_combination.parquet";

private static Path getFilePath() {
try {
return new Path(TestReadInvalidTypeCombination.class.getResource(FILE_PATH).toURI());
} catch (URISyntaxException e) {
throw new RuntimeException(e);
}
}

@Test
public void testReadInvalidTypeCombinationSucceeds() throws Exception {
Configuration conf = new Configuration();
Path file = getFilePath();

// The footer parse should succeed and drop the annotation and stats for the column.
try (ParquetFileReader reader = ParquetFileReader.open(HadoopInputFile.fromPath(file, conf))) {
ParquetMetadata footer = reader.getFooter();
PrimitiveType column =
footer.getFileMetaData().getSchema().getType("int32_uuid").asPrimitiveType();

assertThat(column.getPrimitiveTypeName()).isEqualTo(PrimitiveTypeName.INT32);
assertThat(column.getLogicalTypeAnnotation()).isNull();
assertThat(column.columnOrder().getColumnOrderName()).isEqualTo(ColumnOrderName.UNDEFINED);
}

// The physical values are still fully readable.
int rows = 0;
try (ParquetReader<Group> reader = ParquetReader.builder(new GroupReadSupport(), file)
.withConf(conf)
.build()) {
Group g;
while ((g = reader.read()) != null) {
assertThat(g.getInteger("int32_uuid", 0)).isEqualTo(rows);
rows++;
}
}
assertThat(rows).isEqualTo(10);
}
}
Binary file not shown.