diff --git a/parquet-column/src/main/java/org/apache/parquet/schema/Types.java b/parquet-column/src/main/java/org/apache/parquet/schema/Types.java index ad550cbc77..7b02a6f021 100644 --- a/parquet-column/src/main/java/org/apache/parquet/schema/Types.java +++ b/parquet-column/src/main/java/org/apache/parquet/schema/Types.java @@ -22,9 +22,11 @@ import java.util.ArrayList; import java.util.Collections; +import java.util.EnumSet; import java.util.List; import java.util.Objects; import java.util.Optional; +import java.util.stream.Collectors; import org.apache.parquet.Preconditions; import org.apache.parquet.schema.ColumnOrder.ColumnOrderName; import org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName; @@ -339,11 +341,156 @@ public abstract static class BasePrimitiveBuilder + ALLOWED_PHYSICAL_TYPES = + new LogicalTypeAnnotation.LogicalTypeAnnotationVisitor() { + @Override + public Optional visit( + LogicalTypeAnnotation.StringLogicalTypeAnnotation stringLogicalType) { + return AllowedPhysicalTypes.of(PrimitiveTypeName.BINARY); + } + + @Override + public Optional visit( + LogicalTypeAnnotation.JsonLogicalTypeAnnotation jsonLogicalType) { + return AllowedPhysicalTypes.of(PrimitiveTypeName.BINARY); + } + + @Override + public Optional visit( + LogicalTypeAnnotation.BsonLogicalTypeAnnotation bsonLogicalType) { + return AllowedPhysicalTypes.of(PrimitiveTypeName.BINARY); + } + + @Override + public Optional visit( + LogicalTypeAnnotation.UUIDLogicalTypeAnnotation uuidLogicalType) { + return AllowedPhysicalTypes.fixed( + LogicalTypeAnnotation.UUIDLogicalTypeAnnotation.BYTES); + } + + @Override + public Optional visit( + LogicalTypeAnnotation.Float16LogicalTypeAnnotation float16LogicalType) { + return AllowedPhysicalTypes.fixed( + LogicalTypeAnnotation.Float16LogicalTypeAnnotation.BYTES); + } + + @Override + public Optional visit( + LogicalTypeAnnotation.UnknownLogicalTypeAnnotation unknownLogicalType) { + return Optional.of(AllowedPhysicalTypes.ANY); + } + + @Override + public Optional visit( + LogicalTypeAnnotation.DecimalLogicalTypeAnnotation decimalLogicalType) { + return AllowedPhysicalTypes.of( + PrimitiveTypeName.INT32, + PrimitiveTypeName.INT64, + PrimitiveTypeName.BINARY, + PrimitiveTypeName.FIXED_LEN_BYTE_ARRAY); + } + + @Override + public Optional visit( + LogicalTypeAnnotation.DateLogicalTypeAnnotation dateLogicalType) { + return AllowedPhysicalTypes.of(PrimitiveTypeName.INT32); + } + + @Override + public Optional visit( + LogicalTypeAnnotation.TimeLogicalTypeAnnotation timeLogicalType) { + return timeLogicalType.getUnit() == LogicalTypeAnnotation.TimeUnit.MILLIS + ? AllowedPhysicalTypes.of(PrimitiveTypeName.INT32) + : AllowedPhysicalTypes.of(PrimitiveTypeName.INT64); + } + + @Override + public Optional visit( + LogicalTypeAnnotation.IntLogicalTypeAnnotation intLogicalType) { + return intLogicalType.getBitWidth() == 64 + ? AllowedPhysicalTypes.of(PrimitiveTypeName.INT64) + : AllowedPhysicalTypes.of(PrimitiveTypeName.INT32); + } + + @Override + public Optional visit( + LogicalTypeAnnotation.TimestampLogicalTypeAnnotation timestampLogicalType) { + return AllowedPhysicalTypes.of(PrimitiveTypeName.INT64); + } + + @Override + public Optional visit( + LogicalTypeAnnotation.IntervalLogicalTypeAnnotation intervalLogicalType) { + return AllowedPhysicalTypes.fixed(12); + } + + @Override + public Optional visit( + LogicalTypeAnnotation.EnumLogicalTypeAnnotation enumLogicalType) { + return AllowedPhysicalTypes.of(PrimitiveTypeName.BINARY); + } + + @Override + public Optional visit( + LogicalTypeAnnotation.GeometryLogicalTypeAnnotation geometryLogicalType) { + return AllowedPhysicalTypes.of(PrimitiveTypeName.BINARY); + } + + @Override + public Optional visit( + LogicalTypeAnnotation.GeographyLogicalTypeAnnotation geographyLogicalType) { + return AllowedPhysicalTypes.of(PrimitiveTypeName.BINARY); + } + }; + + private static final class AllowedPhysicalTypes { + private static final AllowedPhysicalTypes NONE = + new AllowedPhysicalTypes(EnumSet.noneOf(PrimitiveTypeName.class), NOT_SET); + private static final AllowedPhysicalTypes ANY = + new AllowedPhysicalTypes(EnumSet.allOf(PrimitiveTypeName.class), NOT_SET); + + private final EnumSet types; + private final int requiredLength; + + private AllowedPhysicalTypes(EnumSet types, int requiredLength) { + this.types = types; + this.requiredLength = requiredLength; + } + + private static Optional of(PrimitiveTypeName first, PrimitiveTypeName... rest) { + return Optional.of(new AllowedPhysicalTypes(EnumSet.of(first, rest), NOT_SET)); + } + + private static Optional fixed(int requiredLength) { + return Optional.of( + new AllowedPhysicalTypes(EnumSet.of(PrimitiveTypeName.FIXED_LEN_BYTE_ARRAY), requiredLength)); + } + + private boolean accepts(PrimitiveTypeName type, int length) { + return types.contains(type) && (requiredLength == NOT_SET || length == requiredLength); + } + + private boolean isEmpty() { + return types.isEmpty(); + } + + @Override + public String toString() { + if (requiredLength != NOT_SET) { + return PrimitiveTypeName.FIXED_LEN_BYTE_ARRAY + "(" + requiredLength + ")"; + } + return types.stream().map(Enum::name).collect(Collectors.joining(", ")); + } + } + private final PrimitiveTypeName primitiveType; private int length = NOT_SET; private int precision = NOT_SET; private int scale = NOT_SET; private ColumnOrder columnOrder; + private boolean ignoreUnsupportedLogicalAnnotations = false; private BasePrimitiveBuilder(P parent, PrimitiveTypeName type) { super(parent); @@ -428,6 +575,18 @@ public THIS columnOrder(ColumnOrder columnOrder) { return self(); } + /** + * When set, an unsupported logical type annotation or logical/physical type combination + * results in the 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 ignoreUnsupportedLogicalAnnotations() { + this.ignoreUnsupportedLogicalAnnotations = true; + return self(); + } + @Override protected PrimitiveType build(String name) { if (length == 0 && logicalTypeAnnotation instanceof LogicalTypeAnnotation.UUIDLogicalTypeAnnotation) { @@ -439,198 +598,27 @@ protected PrimitiveType build(String name) { DecimalMetadata meta = decimalMetadata(); - // validate type annotations and required metadata if (logicalTypeAnnotation != null) { - logicalTypeAnnotation - .accept(new LogicalTypeAnnotation.LogicalTypeAnnotationVisitor() { - @Override - public Optional visit( - LogicalTypeAnnotation.StringLogicalTypeAnnotation stringLogicalType) { - return checkBinaryPrimitiveType(stringLogicalType); - } - - @Override - public Optional visit( - LogicalTypeAnnotation.JsonLogicalTypeAnnotation jsonLogicalType) { - return checkBinaryPrimitiveType(jsonLogicalType); - } - - @Override - public Optional visit( - LogicalTypeAnnotation.BsonLogicalTypeAnnotation bsonLogicalType) { - return checkBinaryPrimitiveType(bsonLogicalType); - } - - @Override - public Optional visit( - LogicalTypeAnnotation.UUIDLogicalTypeAnnotation uuidLogicalType) { - return checkFixedPrimitiveType( - LogicalTypeAnnotation.UUIDLogicalTypeAnnotation.BYTES, uuidLogicalType); - } - - @Override - public Optional visit( - LogicalTypeAnnotation.Float16LogicalTypeAnnotation float16LogicalType) { - return checkFixedPrimitiveType( - LogicalTypeAnnotation.Float16LogicalTypeAnnotation.BYTES, float16LogicalType); - } - - @Override - public Optional visit( - LogicalTypeAnnotation.UnknownLogicalTypeAnnotation unknownLogicalType) { - return Optional.of(true); - } - - @Override - public Optional visit( - LogicalTypeAnnotation.DecimalLogicalTypeAnnotation decimalLogicalType) { - Preconditions.checkState( - (primitiveType == PrimitiveTypeName.INT32) - || (primitiveType == PrimitiveTypeName.INT64) - || (primitiveType == PrimitiveTypeName.BINARY) - || (primitiveType == PrimitiveTypeName.FIXED_LEN_BYTE_ARRAY), - "DECIMAL can only annotate INT32, INT64, BINARY, and FIXED"); - if (primitiveType == PrimitiveTypeName.INT32) { - Preconditions.checkState( - meta.getPrecision() <= MAX_PRECISION_INT32, - "INT32 cannot store %s digits (max %s)", - meta.getPrecision(), - MAX_PRECISION_INT32); - } else if (primitiveType == PrimitiveTypeName.INT64) { - Preconditions.checkState( - meta.getPrecision() <= MAX_PRECISION_INT64, - "INT64 cannot store %s digits (max %s)", - meta.getPrecision(), - MAX_PRECISION_INT64); - if (meta.getPrecision() <= MAX_PRECISION_INT32) { - LOGGER.warn( - "Decimal with {} digits is stored in an INT64, but fits in an INT32. See {}.", - precision, - LOGICAL_TYPES_DOC_URL); - } - } else if (primitiveType == PrimitiveTypeName.FIXED_LEN_BYTE_ARRAY) { - Preconditions.checkState( - meta.getPrecision() <= maxPrecision(length), - "FIXED(%s) cannot store %s digits (max %s)", - length, - meta.getPrecision(), - maxPrecision(length)); - } - return Optional.of(true); - } - - @Override - public Optional visit( - LogicalTypeAnnotation.DateLogicalTypeAnnotation dateLogicalType) { - return checkInt32PrimitiveType(dateLogicalType); - } - - @Override - public Optional visit( - LogicalTypeAnnotation.TimeLogicalTypeAnnotation timeLogicalType) { - LogicalTypeAnnotation.TimeUnit unit = timeLogicalType.getUnit(); - switch (unit) { - case MILLIS: - checkInt32PrimitiveType(timeLogicalType); - break; - case MICROS: - case NANOS: - checkInt64PrimitiveType(timeLogicalType); - break; - default: - throw new RuntimeException("Invalid time unit: " + unit); - } - return Optional.of(true); - } - - @Override - public Optional visit( - LogicalTypeAnnotation.IntLogicalTypeAnnotation intLogicalType) { - int bitWidth = intLogicalType.getBitWidth(); - switch (bitWidth) { - case 8: - case 16: - case 32: - checkInt32PrimitiveType(intLogicalType); - break; - case 64: - checkInt64PrimitiveType(intLogicalType); - break; - default: - throw new RuntimeException("Invalid bit width: " + bitWidth); - } - return Optional.of(true); - } - - @Override - public Optional visit( - LogicalTypeAnnotation.TimestampLogicalTypeAnnotation timestampLogicalType) { - return checkInt64PrimitiveType(timestampLogicalType); - } - - @Override - public Optional visit( - LogicalTypeAnnotation.IntervalLogicalTypeAnnotation intervalLogicalType) { - return checkFixedPrimitiveType(12, intervalLogicalType); - } - - @Override - public Optional visit( - LogicalTypeAnnotation.EnumLogicalTypeAnnotation enumLogicalType) { - return checkBinaryPrimitiveType(enumLogicalType); - } - - @Override - public Optional visit( - LogicalTypeAnnotation.GeometryLogicalTypeAnnotation geometryLogicalType) { - return checkBinaryPrimitiveType(geometryLogicalType); - } - - @Override - public Optional visit( - LogicalTypeAnnotation.GeographyLogicalTypeAnnotation geographyLogicalType) { - return checkBinaryPrimitiveType(geographyLogicalType); - } - - private Optional checkFixedPrimitiveType( - int l, LogicalTypeAnnotation logicalTypeAnnotation) { - Preconditions.checkState( - primitiveType == PrimitiveTypeName.FIXED_LEN_BYTE_ARRAY && length == l, - "%s can only annotate FIXED_LEN_BYTE_ARRAY(%s)", - logicalTypeAnnotation, - l); - return Optional.of(true); - } - - private Optional checkBinaryPrimitiveType( - LogicalTypeAnnotation logicalTypeAnnotation) { - Preconditions.checkState( - primitiveType == PrimitiveTypeName.BINARY, - "%s can only annotate BINARY", - logicalTypeAnnotation); - return Optional.of(true); - } - - private Optional checkInt32PrimitiveType( - LogicalTypeAnnotation logicalTypeAnnotation) { - Preconditions.checkState( - primitiveType == PrimitiveTypeName.INT32, - "%s can only annotate INT32", - logicalTypeAnnotation); - return Optional.of(true); - } - - private Optional checkInt64PrimitiveType( - LogicalTypeAnnotation logicalTypeAnnotation) { - Preconditions.checkState( - primitiveType == PrimitiveTypeName.INT64, - "%s can only annotate INT64", - logicalTypeAnnotation); - return Optional.of(true); - } - }) - .orElseThrow(() -> new IllegalStateException( - logicalTypeAnnotation + " can not be applied to a primitive type")); + String annotation = newLogicalTypeSet + ? logicalTypeAnnotation.toString() + : getOriginalType().toString(); + AllowedPhysicalTypes allowed = + logicalTypeAnnotation.accept(ALLOWED_PHYSICAL_TYPES).orElse(AllowedPhysicalTypes.NONE); + if (!allowed.accepts(primitiveType, length)) { + if (!ignoreUnsupportedLogicalAnnotations) { + throw new IllegalStateException( + allowed.isEmpty() + ? annotation + " can not be applied to a primitive type" + : String.format("%s can only annotate [%s]", annotation, allowed)); + } + LOGGER.warn( + "Dropping unsupported logical type annotation {} on physical type {}", + logicalTypeAnnotation, + primitiveType); + return new PrimitiveType( + repetition, primitiveType, length, name, null, null, id, ColumnOrder.undefined()); + } + validateDecimalPrecision(meta); } if (newLogicalTypeSet) { @@ -642,6 +630,38 @@ private Optional checkInt64PrimitiveType( } } + private void validateDecimalPrecision(DecimalMetadata meta) { + if (meta == null) { + return; + } + if (primitiveType == PrimitiveTypeName.INT32) { + Preconditions.checkState( + meta.getPrecision() <= MAX_PRECISION_INT32, + "INT32 cannot store %s digits (max %s)", + meta.getPrecision(), + MAX_PRECISION_INT32); + } else if (primitiveType == PrimitiveTypeName.INT64) { + Preconditions.checkState( + meta.getPrecision() <= MAX_PRECISION_INT64, + "INT64 cannot store %s digits (max %s)", + meta.getPrecision(), + MAX_PRECISION_INT64); + if (meta.getPrecision() <= MAX_PRECISION_INT32) { + LOGGER.warn( + "Decimal with {} digits is stored in an INT64, but fits in an INT32. See {}.", + precision, + LOGICAL_TYPES_DOC_URL); + } + } else if (primitiveType == PrimitiveTypeName.FIXED_LEN_BYTE_ARRAY) { + Preconditions.checkState( + meta.getPrecision() <= maxPrecision(length), + "FIXED(%s) cannot store %s digits (max %s)", + length, + meta.getPrecision(), + maxPrecision(length)); + } + } + private static long maxPrecision(int numBytes) { return Math.round( // convert double to long Math.floor( diff --git a/parquet-column/src/test/java/org/apache/parquet/schema/TestTypeBuilders.java b/parquet-column/src/test/java/org/apache/parquet/schema/TestTypeBuilders.java index d0d00898c0..dfc4fd6feb 100644 --- a/parquet-column/src/test/java/org/apache/parquet/schema/TestTypeBuilders.java +++ b/parquet-column/src/test/java/org/apache/parquet/schema/TestTypeBuilders.java @@ -467,7 +467,7 @@ public void testDECIMALAnnotationRejectsUnsupportedTypes() { .scale(2) .named("d")) .isInstanceOf(IllegalStateException.class) - .hasMessage("DECIMAL can only annotate INT32, INT64, BINARY, and FIXED"); + .hasMessage("DECIMAL can only annotate [INT64, INT32, BINARY, FIXED_LEN_BYTE_ARRAY]"); } } @@ -489,16 +489,14 @@ public void testBinaryAnnotationsRejectsNonBinary() { for (final PrimitiveTypeName type : nonBinary) { assertThatThrownBy(() -> Types.required(type).as(logicalType).named("col")) .isInstanceOf(IllegalStateException.class) - .hasMessage(LogicalTypeAnnotation.fromOriginalType(logicalType, null) - + " can only annotate BINARY"); + .hasMessage(logicalType + " can only annotate [BINARY]"); } assertThatThrownBy(() -> Types.required(FIXED_LEN_BYTE_ARRAY) .length(1) .as(logicalType) .named("col")) .isInstanceOf(IllegalStateException.class) - .hasMessage( - LogicalTypeAnnotation.fromOriginalType(logicalType, null) + " can only annotate BINARY"); + .hasMessage(logicalType + " can only annotate [BINARY]"); } } @@ -520,15 +518,14 @@ public void testInt32AnnotationsRejectNonInt32() { for (final PrimitiveTypeName type : nonInt32) { assertThatThrownBy(() -> Types.required(type).as(logicalType).named("col")) .isInstanceOf(IllegalStateException.class) - .hasMessage( - LogicalTypeAnnotation.fromOriginalType(logicalType, null) + " can only annotate INT32"); + .hasMessage(logicalType + " can only annotate [INT32]"); } assertThatThrownBy(() -> Types.required(FIXED_LEN_BYTE_ARRAY) .length(1) .as(logicalType) .named("col")) .isInstanceOf(IllegalStateException.class) - .hasMessage(LogicalTypeAnnotation.fromOriginalType(logicalType, null) + " can only annotate INT32"); + .hasMessage(logicalType + " can only annotate [INT32]"); } } @@ -550,15 +547,14 @@ public void testInt64AnnotationsRejectNonInt64() { for (final PrimitiveTypeName type : nonInt64) { assertThatThrownBy(() -> Types.required(type).as(logicalType).named("col")) .isInstanceOf(IllegalStateException.class) - .hasMessage( - LogicalTypeAnnotation.fromOriginalType(logicalType, null) + " can only annotate INT64"); + .hasMessage(logicalType + " can only annotate [INT64]"); } assertThatThrownBy(() -> Types.required(FIXED_LEN_BYTE_ARRAY) .length(1) .as(logicalType) .named("col")) .isInstanceOf(IllegalStateException.class) - .hasMessage(LogicalTypeAnnotation.fromOriginalType(logicalType, null) + " can only annotate INT64"); + .hasMessage(logicalType + " can only annotate [INT64]"); } } @@ -576,7 +572,7 @@ public void testIntervalAnnotationRejectsNonFixed() { for (final PrimitiveTypeName type : nonFixed) { assertThatThrownBy(() -> Types.required(type).as(INTERVAL).named("interval")) .isInstanceOf(IllegalStateException.class) - .hasMessage("INTERVAL can only annotate FIXED_LEN_BYTE_ARRAY(12)"); + .hasMessage("INTERVAL can only annotate [FIXED_LEN_BYTE_ARRAY(12)]"); } } @@ -587,7 +583,7 @@ public void testIntervalAnnotationRejectsNonFixed12() { .as(INTERVAL) .named("interval")) .isInstanceOf(IllegalStateException.class) - .hasMessage("INTERVAL can only annotate FIXED_LEN_BYTE_ARRAY(12)"); + .hasMessage("INTERVAL can only annotate [FIXED_LEN_BYTE_ARRAY(12)]"); } @Test @@ -1605,4 +1601,20 @@ public void testGeographyLogicalTypeWithoutEdgeInterpolationAlgorithm() { Types.optional(BINARY).as(LogicalTypeAnnotation.geographyType()).named("aGeography"); assertThat(optionalGeographyActual).isEqualTo(optionalGeographyExpected); } + + @Test + public void testIgnoreUnsupportedLogicalAnnotations() { + LogicalTypeAnnotation[] annotations = { + LogicalTypeAnnotation.timestampType(true, MILLIS), LogicalTypeAnnotation.decimalType(2, 9) + }; + for (LogicalTypeAnnotation annotation : annotations) { + PrimitiveType pt = Types.required(BOOLEAN) + .ignoreUnsupportedLogicalAnnotations() + .as(annotation) + .named("unsupported"); + assertThat(pt.getPrimitiveTypeName()).isEqualTo(BOOLEAN); + assertThat(pt.getLogicalTypeAnnotation()).isNull(); + assertThat(pt.columnOrder().getColumnOrderName()).isEqualTo(ColumnOrder.ColumnOrderName.UNDEFINED); + } + } } diff --git a/parquet-column/src/test/java/org/apache/parquet/schema/TestTypeBuildersWithLogicalTypes.java b/parquet-column/src/test/java/org/apache/parquet/schema/TestTypeBuildersWithLogicalTypes.java index 0d7791a19b..1fc46d7fd8 100644 --- a/parquet-column/src/test/java/org/apache/parquet/schema/TestTypeBuildersWithLogicalTypes.java +++ b/parquet-column/src/test/java/org/apache/parquet/schema/TestTypeBuildersWithLogicalTypes.java @@ -203,7 +203,7 @@ public void testDECIMALAnnotationRejectsUnsupportedTypes() { for (final PrimitiveTypeName type : unsupported) { assertThatThrownBy(() -> Types.required(type).as(decimalType(2, 9)).named("d")) .isInstanceOf(IllegalStateException.class) - .hasMessage("DECIMAL can only annotate INT32, INT64, BINARY, and FIXED"); + .hasMessage("DECIMAL(9,2) can only annotate [INT64, INT32, BINARY, FIXED_LEN_BYTE_ARRAY]"); } } @@ -234,15 +234,15 @@ public void testBinaryAnnotationsRejectsNonBinary() { PrimitiveTypeName[] nonBinary = new PrimitiveTypeName[] {BOOLEAN, INT32, INT64, INT96, DOUBLE, FLOAT}; for (final PrimitiveTypeName type : nonBinary) { String expectedMessage = logicalType.equals(float16Type()) - ? "FLOAT16 can only annotate FIXED_LEN_BYTE_ARRAY(2)" - : logicalType + " can only annotate BINARY"; + ? "FLOAT16 can only annotate [FIXED_LEN_BYTE_ARRAY(2)]" + : logicalType + " can only annotate [BINARY]"; assertThatThrownBy(() -> Types.required(type).as(logicalType).named("col")) .isInstanceOf(IllegalStateException.class) .hasMessage(expectedMessage); } String fixedMessage = logicalType.equals(float16Type()) - ? "FLOAT16 can only annotate FIXED_LEN_BYTE_ARRAY(2)" - : logicalType + " can only annotate BINARY"; + ? "FLOAT16 can only annotate [FIXED_LEN_BYTE_ARRAY(2)]" + : logicalType + " can only annotate [BINARY]"; assertThatThrownBy(() -> Types.required(FIXED_LEN_BYTE_ARRAY) .length(1) .as(logicalType) @@ -278,14 +278,14 @@ public void testInt32AnnotationsRejectNonInt32() { for (final PrimitiveTypeName type : nonInt32) { assertThatThrownBy(() -> Types.required(type).as(logicalType).named("col")) .isInstanceOf(IllegalStateException.class) - .hasMessage(logicalType + " can only annotate INT32"); + .hasMessage(logicalType + " can only annotate [INT32]"); } assertThatThrownBy(() -> Types.required(FIXED_LEN_BYTE_ARRAY) .length(1) .as(logicalType) .named("col")) .isInstanceOf(IllegalStateException.class) - .hasMessage(logicalType + " can only annotate INT32"); + .hasMessage(logicalType + " can only annotate [INT32]"); } } @@ -321,14 +321,14 @@ public void testInt64AnnotationsRejectNonInt64() { for (final PrimitiveTypeName type : nonInt64) { assertThatThrownBy(() -> Types.required(type).as(logicalType).named("col")) .isInstanceOf(IllegalStateException.class) - .hasMessage(logicalType + " can only annotate INT64"); + .hasMessage(logicalType + " can only annotate [INT64]"); } assertThatThrownBy(() -> Types.required(FIXED_LEN_BYTE_ARRAY) .length(1) .as(logicalType) .named("col")) .isInstanceOf(IllegalStateException.class) - .hasMessage(logicalType + " can only annotate INT64"); + .hasMessage(logicalType + " can only annotate [INT64]"); } } @@ -340,7 +340,7 @@ public void testIntervalAnnotationRejectsNonFixed() { .as(LogicalTypeAnnotation.IntervalLogicalTypeAnnotation.getInstance()) .named("interval")) .isInstanceOf(IllegalStateException.class) - .hasMessage("INTERVAL can only annotate FIXED_LEN_BYTE_ARRAY(12)"); + .hasMessage("INTERVAL can only annotate [FIXED_LEN_BYTE_ARRAY(12)]"); } } @@ -351,7 +351,7 @@ public void testIntervalAnnotationRejectsNonFixed12() { .as(LogicalTypeAnnotation.IntervalLogicalTypeAnnotation.getInstance()) .named("interval")) .isInstanceOf(IllegalStateException.class) - .hasMessage("INTERVAL can only annotate FIXED_LEN_BYTE_ARRAY(12)"); + .hasMessage("INTERVAL can only annotate [FIXED_LEN_BYTE_ARRAY(12)]"); } @Test @@ -462,13 +462,13 @@ public void testUUIDLogicalType() { .named("uuid_field") .toString()) .isInstanceOf(IllegalStateException.class) - .hasMessage("UUID can only annotate FIXED_LEN_BYTE_ARRAY(16)"); + .hasMessage("UUID can only annotate [FIXED_LEN_BYTE_ARRAY(16)]"); assertThatThrownBy(() -> Types.required(BINARY) .as(uuidType()) .named("uuid_field") .toString()) .isInstanceOf(IllegalStateException.class) - .hasMessage("UUID can only annotate FIXED_LEN_BYTE_ARRAY(16)"); + .hasMessage("UUID can only annotate [FIXED_LEN_BYTE_ARRAY(16)]"); } @Test @@ -486,13 +486,13 @@ public void testFloat16LogicalType() { .named("float16_field") .toString()) .isInstanceOf(IllegalStateException.class) - .hasMessage("FLOAT16 can only annotate FIXED_LEN_BYTE_ARRAY(2)"); + .hasMessage("FLOAT16 can only annotate [FIXED_LEN_BYTE_ARRAY(2)]"); assertThatThrownBy(() -> Types.required(BINARY) .as(float16Type()) .named("float16_field") .toString()) .isInstanceOf(IllegalStateException.class) - .hasMessage("FLOAT16 can only annotate FIXED_LEN_BYTE_ARRAY(2)"); + .hasMessage("FLOAT16 can only annotate [FIXED_LEN_BYTE_ARRAY(2)]"); } @Test diff --git a/parquet-hadoop/src/main/java/org/apache/parquet/format/converter/ParquetMetadataConverter.java b/parquet-hadoop/src/main/java/org/apache/parquet/format/converter/ParquetMetadataConverter.java index f6ee73bbc0..4cd9ea1138 100644 --- a/parquet-hadoop/src/main/java/org/apache/parquet/format/converter/ParquetMetadataConverter.java +++ b/parquet-hadoop/src/main/java/org/apache/parquet/format/converter/ParquetMetadataConverter.java @@ -2112,6 +2112,8 @@ private void buildChildren( // any, were written under the legacy type-defined order and must be read under it. primitiveBuilder.columnOrder(org.apache.parquet.schema.ColumnOrder.typeDefined()); } + // Gracefully handle unsupported logical type combinations on the read path. + primitiveBuilder.ignoreUnsupportedLogicalAnnotations(); childBuilder = primitiveBuilder; } else { childBuilder = builder.group(fromParquetRepetition(schemaElement.repetition_type)); diff --git a/parquet-hadoop/src/test/java/org/apache/parquet/format/converter/TestParquetMetadataConverter.java b/parquet-hadoop/src/test/java/org/apache/parquet/format/converter/TestParquetMetadataConverter.java index f5222de828..c2e52837a8 100644 --- a/parquet-hadoop/src/test/java/org/apache/parquet/format/converter/TestParquetMetadataConverter.java +++ b/parquet-hadoop/src/test/java/org/apache/parquet/format/converter/TestParquetMetadataConverter.java @@ -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; @@ -2378,4 +2382,62 @@ public void testV2StatsDoNotTriggerCorruptStatisticsCheck() { assertThat(result.getMinBytes()).isEqualTo(new byte[] {0}); assertThat(result.getMaxBytes()).isEqualTo(new byte[] {1}); } + + @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 parquetSchema = Lists.newArrayList(new SchemaElement("Message").setNum_children(1), leaf); + List 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(); + } } diff --git a/parquet-hadoop/src/test/java/org/apache/parquet/hadoop/TestReadInvalidTypeCombination.java b/parquet-hadoop/src/test/java/org/apache/parquet/hadoop/TestReadInvalidTypeCombination.java new file mode 100644 index 0000000000..6b3fe6d06a --- /dev/null +++ b/parquet-hadoop/src/test/java/org/apache/parquet/hadoop/TestReadInvalidTypeCombination.java @@ -0,0 +1,71 @@ +/* + * 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 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 { + + // parquet-testing file with an invalid logical/physical type combination. + private static final String REFERENCE_FILE = "int32_with_uuid_logical_type.parquet"; + private static final String REFERENCE_CHANGESET = "4b1ce4502afff8d20c9b4bb08d07e04e21cdeff3"; + + private final InterOpTester interop = new InterOpTester(); + + @Test + public void testReadInvalidTypeCombinationSucceeds() throws Exception { + Configuration conf = new Configuration(); + Path file = interop.GetInterOpFile(REFERENCE_FILE, REFERENCE_CHANGESET); + + // 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 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); + } +}