diff --git a/parquet-column/src/main/java/org/apache/parquet/io/api/FileConverters.java b/parquet-column/src/main/java/org/apache/parquet/io/api/FileConverters.java new file mode 100644 index 0000000000..a0cc77e71b --- /dev/null +++ b/parquet-column/src/main/java/org/apache/parquet/io/api/FileConverters.java @@ -0,0 +1,128 @@ +/* + * 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.io.api; + +import java.util.Objects; +import java.util.function.Consumer; +import org.apache.parquet.schema.GroupType; +import org.apache.parquet.schema.LogicalTypeAnnotation; + +/** Converters for assembling values represented by a Parquet {@code FILE}-annotated group. */ +public final class FileConverters { + private FileConverters() {} + + /** + * Creates a converter for a FILE group. + * + *

The consumer receives {@code null} when an invalid FILE reference is encountered, as + * permitted by the FILE logical-type specification. + */ + public static GroupConverter newFileConverter(GroupType schema, Consumer consumer) { + Objects.requireNonNull(schema, "schema cannot be null"); + Objects.requireNonNull(consumer, "consumer cannot be null"); + FileValueValidator.validateSchema(schema); + return new FileGroupConverter(schema, consumer); + } + + private static final class FileGroupConverter extends GroupConverter { + private final Converter[] converters; + private final Consumer consumer; + private FileValue.Builder builder; + + private FileGroupConverter(GroupType schema, Consumer consumer) { + this.converters = new Converter[schema.getFieldCount()]; + this.consumer = consumer; + for (int index = 0; index < schema.getFieldCount(); index++) { + String fieldName = schema.getFieldName(index); + switch (fieldName) { + case LogicalTypeAnnotation.FileLogicalTypeAnnotation.URI_FIELD: + converters[index] = stringConverter(value -> builder.withUri(value)); + break; + case LogicalTypeAnnotation.FileLogicalTypeAnnotation.OFFSET_FIELD: + converters[index] = longConverter(value -> builder.withOffset(value)); + break; + case LogicalTypeAnnotation.FileLogicalTypeAnnotation.SIZE_FIELD: + converters[index] = longConverter(value -> builder.withSize(value)); + break; + case LogicalTypeAnnotation.FileLogicalTypeAnnotation.CONTENT_TYPE_FIELD: + converters[index] = stringConverter(value -> builder.withContentType(value)); + break; + case LogicalTypeAnnotation.FileLogicalTypeAnnotation.CHECKSUM_FIELD: + converters[index] = stringConverter(value -> builder.withChecksum(value)); + break; + case LogicalTypeAnnotation.FileLogicalTypeAnnotation.INLINE_FIELD: + converters[index] = binaryConverter(value -> builder.withInline(value)); + break; + default: + throw new IllegalArgumentException("Unrecognized FILE field: " + fieldName); + } + } + } + + @Override + public Converter getConverter(int fieldIndex) { + return converters[fieldIndex]; + } + + @Override + public void start() { + builder = FileValue.builder(); + } + + @Override + public void end() { + FileValue value = builder.build(); + boolean valid = true; + try { + FileValueValidator.validateValue(value); + } catch (IllegalArgumentException e) { + valid = false; + } + builder = null; + consumer.accept(valid ? value : null); + } + } + + private static PrimitiveConverter stringConverter(Consumer consumer) { + return new PrimitiveConverter() { + @Override + public void addBinary(Binary value) { + consumer.accept(value.toStringUsingUTF8()); + } + }; + } + + private static PrimitiveConverter binaryConverter(Consumer consumer) { + return new PrimitiveConverter() { + @Override + public void addBinary(Binary value) { + consumer.accept(Binary.fromConstantByteArray(value.getBytes())); + } + }; + } + + private static PrimitiveConverter longConverter(Consumer consumer) { + return new PrimitiveConverter() { + @Override + public void addLong(long value) { + consumer.accept(value); + } + }; + } +} diff --git a/parquet-column/src/main/java/org/apache/parquet/io/api/FileValue.java b/parquet-column/src/main/java/org/apache/parquet/io/api/FileValue.java new file mode 100644 index 0000000000..96a7d39458 --- /dev/null +++ b/parquet-column/src/main/java/org/apache/parquet/io/api/FileValue.java @@ -0,0 +1,112 @@ +/* + * 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.io.api; + +/** A value represented by a Parquet {@code FILE}-annotated group. */ +public final class FileValue { + private final String uri; + private final Long offset; + private final Long size; + private final String contentType; + private final String checksum; + private final Binary inline; + + private FileValue(Builder builder) { + this.uri = builder.uri; + this.offset = builder.offset; + this.size = builder.size; + this.contentType = builder.contentType; + this.checksum = builder.checksum; + this.inline = builder.inline; + } + + public static Builder builder() { + return new Builder(); + } + + public String getUri() { + return uri; + } + + public Long getOffset() { + return offset; + } + + public Long getSize() { + return size; + } + + public String getContentType() { + return contentType; + } + + public String getChecksum() { + return checksum; + } + + public Binary getInline() { + return inline; + } + + /** Builder for {@link FileValue}. */ + public static final class Builder { + private String uri; + private Long offset; + private Long size; + private String contentType; + private String checksum; + private Binary inline; + + private Builder() {} + + public Builder withUri(String uri) { + this.uri = uri; + return this; + } + + public Builder withOffset(long offset) { + this.offset = offset; + return this; + } + + public Builder withSize(long size) { + this.size = size; + return this; + } + + public Builder withContentType(String contentType) { + this.contentType = contentType; + return this; + } + + public Builder withChecksum(String checksum) { + this.checksum = checksum; + return this; + } + + public Builder withInline(Binary inline) { + this.inline = inline; + return this; + } + + public FileValue build() { + return new FileValue(this); + } + } +} diff --git a/parquet-column/src/main/java/org/apache/parquet/io/api/FileValueValidator.java b/parquet-column/src/main/java/org/apache/parquet/io/api/FileValueValidator.java new file mode 100644 index 0000000000..2ff9b0982f --- /dev/null +++ b/parquet-column/src/main/java/org/apache/parquet/io/api/FileValueValidator.java @@ -0,0 +1,100 @@ +/* + * 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.io.api; + +import org.apache.parquet.schema.GroupType; +import org.apache.parquet.schema.LogicalTypeAnnotation; +import org.apache.parquet.schema.PrimitiveType; +import org.apache.parquet.schema.Type; + +final class FileValueValidator { + private FileValueValidator() {} + + static void validateSchema(GroupType schema) { + if (!(schema.getLogicalTypeAnnotation() instanceof LogicalTypeAnnotation.FileLogicalTypeAnnotation)) { + throw new IllegalArgumentException( + "Cannot use a FILE value with a group without the FILE logical type: " + schema.getName()); + } + for (Type field : schema.getFields()) { + String fieldName = field.getName(); + if (!LogicalTypeAnnotation.FileLogicalTypeAnnotation.FIELD_NAMES.contains(fieldName)) { + throw new IllegalArgumentException( + "Unrecognized field '" + fieldName + "' in FILE group '" + schema.getName() + "'"); + } + if (!field.isPrimitive() || field.getRepetition() != Type.Repetition.OPTIONAL) { + throw new IllegalArgumentException("FILE type field '" + fieldName + + "' must be an optional primitive in group '" + schema.getName() + "'"); + } + validatePhysicalType(schema.getName(), field.asPrimitiveType()); + } + } + + static void validateValue(FileValue value) { + boolean hasUri = value.getUri() != null && !value.getUri().isEmpty(); + boolean hasInline = value.getInline() != null; + + if (!hasUri && !hasInline) { + throw new IllegalArgumentException("FILE value must set at least one of 'inline' or non-empty 'uri'"); + } + if (value.getOffset() != null && !hasUri) { + throw new IllegalArgumentException("FILE value field 'offset' may only be set together with 'uri'"); + } + if (value.getOffset() != null && value.getSize() == null) { + throw new IllegalArgumentException("FILE value field 'size' must be set whenever 'offset' is set"); + } + if (value.getOffset() != null && value.getOffset() < 0) { + throw new IllegalArgumentException("FILE value field 'offset' must not be negative: " + value.getOffset()); + } + if (value.getSize() != null && value.getSize() < 0) { + throw new IllegalArgumentException("FILE value field 'size' must not be negative: " + value.getSize()); + } + } + + private static void validatePhysicalType(String groupName, PrimitiveType field) { + String fieldName = field.getName(); + PrimitiveType.PrimitiveTypeName physicalType = field.getPrimitiveTypeName(); + switch (fieldName) { + case LogicalTypeAnnotation.FileLogicalTypeAnnotation.URI_FIELD: + case LogicalTypeAnnotation.FileLogicalTypeAnnotation.CONTENT_TYPE_FIELD: + case LogicalTypeAnnotation.FileLogicalTypeAnnotation.CHECKSUM_FIELD: + if (physicalType != PrimitiveType.PrimitiveTypeName.BINARY + || !(field.getLogicalTypeAnnotation() instanceof LogicalTypeAnnotation.StringLogicalTypeAnnotation)) { + throw new IllegalArgumentException( + "FILE type field '" + fieldName + + "' must be a STRING (BINARY annotated as STRING) in group '" + groupName + "'"); + } + break; + case LogicalTypeAnnotation.FileLogicalTypeAnnotation.OFFSET_FIELD: + case LogicalTypeAnnotation.FileLogicalTypeAnnotation.SIZE_FIELD: + if (physicalType != PrimitiveType.PrimitiveTypeName.INT64) { + throw new IllegalArgumentException( + "FILE type field '" + fieldName + "' must be an INT64 in group '" + groupName + "'"); + } + break; + case LogicalTypeAnnotation.FileLogicalTypeAnnotation.INLINE_FIELD: + if (physicalType != PrimitiveType.PrimitiveTypeName.BINARY) { + throw new IllegalArgumentException("FILE type field '" + fieldName + + "' must be a BYTE_ARRAY (BINARY) in group '" + groupName + "'"); + } + break; + default: + break; + } + } +} diff --git a/parquet-column/src/main/java/org/apache/parquet/io/api/FileValueWriter.java b/parquet-column/src/main/java/org/apache/parquet/io/api/FileValueWriter.java new file mode 100644 index 0000000000..433b628e5f --- /dev/null +++ b/parquet-column/src/main/java/org/apache/parquet/io/api/FileValueWriter.java @@ -0,0 +1,112 @@ +/* + * 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.io.api; + +import java.util.Objects; +import org.apache.parquet.schema.GroupType; +import org.apache.parquet.schema.LogicalTypeAnnotation; + +/** Writes and validates values represented by a Parquet {@code FILE}-annotated group. */ +public final class FileValueWriter { + private FileValueWriter() {} + + /** + * Writes one FILE value. The caller is responsible for starting and ending the field containing + * the FILE group. + */ + public static void write(RecordConsumer recordConsumer, GroupType schema, FileValue value) { + Objects.requireNonNull(recordConsumer, "recordConsumer cannot be null"); + Objects.requireNonNull(schema, "schema cannot be null"); + Objects.requireNonNull(value, "value cannot be null"); + + FileValueValidator.validateSchema(schema); + FileValueValidator.validateValue(value); + validateSchemaContainsSetFields(schema, value); + + recordConsumer.startGroup(); + for (int index = 0; index < schema.getFieldCount(); index++) { + String fieldName = schema.getFieldName(index); + switch (fieldName) { + case LogicalTypeAnnotation.FileLogicalTypeAnnotation.URI_FIELD: + writeString(recordConsumer, fieldName, index, value.getUri()); + break; + case LogicalTypeAnnotation.FileLogicalTypeAnnotation.OFFSET_FIELD: + writeLong(recordConsumer, fieldName, index, value.getOffset()); + break; + case LogicalTypeAnnotation.FileLogicalTypeAnnotation.SIZE_FIELD: + writeLong(recordConsumer, fieldName, index, value.getSize()); + break; + case LogicalTypeAnnotation.FileLogicalTypeAnnotation.CONTENT_TYPE_FIELD: + writeString(recordConsumer, fieldName, index, value.getContentType()); + break; + case LogicalTypeAnnotation.FileLogicalTypeAnnotation.CHECKSUM_FIELD: + writeString(recordConsumer, fieldName, index, value.getChecksum()); + break; + case LogicalTypeAnnotation.FileLogicalTypeAnnotation.INLINE_FIELD: + writeBinary(recordConsumer, fieldName, index, value.getInline()); + break; + default: + throw new IllegalArgumentException( + "Unrecognized field '" + fieldName + "' in FILE group '" + schema.getName() + "'"); + } + } + recordConsumer.endGroup(); + } + + private static void validateSchemaContainsSetFields(GroupType schema, FileValue value) { + requireField(schema, LogicalTypeAnnotation.FileLogicalTypeAnnotation.URI_FIELD, value.getUri()); + requireField(schema, LogicalTypeAnnotation.FileLogicalTypeAnnotation.OFFSET_FIELD, value.getOffset()); + requireField(schema, LogicalTypeAnnotation.FileLogicalTypeAnnotation.SIZE_FIELD, value.getSize()); + requireField( + schema, LogicalTypeAnnotation.FileLogicalTypeAnnotation.CONTENT_TYPE_FIELD, value.getContentType()); + requireField(schema, LogicalTypeAnnotation.FileLogicalTypeAnnotation.CHECKSUM_FIELD, value.getChecksum()); + requireField(schema, LogicalTypeAnnotation.FileLogicalTypeAnnotation.INLINE_FIELD, value.getInline()); + } + + private static void requireField(GroupType schema, String fieldName, Object value) { + if (value != null && !schema.containsField(fieldName)) { + throw new IllegalArgumentException( + "FILE value sets field '" + fieldName + "' which is absent from group '" + schema.getName() + "'"); + } + } + + private static void writeString(RecordConsumer consumer, String name, int index, String value) { + if (value != null) { + consumer.startField(name, index); + consumer.addBinary(Binary.fromString(value)); + consumer.endField(name, index); + } + } + + private static void writeLong(RecordConsumer consumer, String name, int index, Long value) { + if (value != null) { + consumer.startField(name, index); + consumer.addLong(value); + consumer.endField(name, index); + } + } + + private static void writeBinary(RecordConsumer consumer, String name, int index, Binary value) { + if (value != null) { + consumer.startField(name, index); + consumer.addBinary(value); + consumer.endField(name, index); + } + } +} diff --git a/parquet-column/src/main/java/org/apache/parquet/schema/LogicalTypeAnnotation.java b/parquet-column/src/main/java/org/apache/parquet/schema/LogicalTypeAnnotation.java index 625e9fd9d3..bf7d582d43 100644 --- a/parquet-column/src/main/java/org/apache/parquet/schema/LogicalTypeAnnotation.java +++ b/parquet-column/src/main/java/org/apache/parquet/schema/LogicalTypeAnnotation.java @@ -188,6 +188,12 @@ protected LogicalTypeAnnotation fromString(List params) { protected LogicalTypeAnnotation fromString(List params) { return unknownType(); } + }, + FILE { + @Override + protected LogicalTypeAnnotation fromString(List params) { + return fileType(); + } }; protected abstract LogicalTypeAnnotation fromString(List params); @@ -378,6 +384,10 @@ public static UnknownLogicalTypeAnnotation unknownType() { return UnknownLogicalTypeAnnotation.INSTANCE; } + public static FileLogicalTypeAnnotation fileType() { + return FileLogicalTypeAnnotation.INSTANCE; + } + public static class StringLogicalTypeAnnotation extends LogicalTypeAnnotation { private static final StringLogicalTypeAnnotation INSTANCE = new StringLogicalTypeAnnotation(); @@ -1229,6 +1239,85 @@ public boolean equals(Object obj) { } } + /** + * File logical type annotation. Annotates a group (struct) that represents a reference to a + * range of bytes, which may be stored inline in the value or in an external file. Every field is + * optional, both in the schema (a writer may omit any field from the group definition) and in the + * data (any field that is present has a field repetition type of {@code OPTIONAL}). Fields are + * identified by name (case sensitively), not by field order. A group need only define the fields + * it uses. The group may contain the following fields: + *

+ * No fields with names other than the above are permitted. Each declared field must match its + * required physical type. Rules involving which fields are set and their values are enforced per + * value rather than by the schema. + */ + public static class FileLogicalTypeAnnotation extends LogicalTypeAnnotation { + private static final FileLogicalTypeAnnotation INSTANCE = new FileLogicalTypeAnnotation(); + + /** Field name holding the URI-reference of an external file. */ + public static final String URI_FIELD = "uri"; + + /** Field name holding the start of the byte range. */ + public static final String OFFSET_FIELD = "offset"; + + /** Field name holding the byte length of the referenced data. */ + public static final String SIZE_FIELD = "size"; + + /** Field name holding the media (MIME) type of the resolved bytes. */ + public static final String CONTENT_TYPE_FIELD = "content_type"; + + /** Field name holding the integrity token for the resolved bytes. */ + public static final String CHECKSUM_FIELD = "checksum"; + + /** Field name holding the referenced bytes stored inline. */ + public static final String INLINE_FIELD = "inline"; + + /** All recognized field names in a FILE-annotated group. All fields are optional. */ + public static final Set FIELD_NAMES = + Set.of(URI_FIELD, OFFSET_FIELD, SIZE_FIELD, CONTENT_TYPE_FIELD, CHECKSUM_FIELD, INLINE_FIELD); + + private FileLogicalTypeAnnotation() {} + + @Override + public OriginalType toOriginalType() { + return null; + } + + @Override + public Optional accept(LogicalTypeAnnotationVisitor logicalTypeAnnotationVisitor) { + return logicalTypeAnnotationVisitor.visit(this); + } + + @Override + LogicalTypeToken getType() { + return LogicalTypeToken.FILE; + } + + @Override + public boolean equals(Object obj) { + return obj instanceof FileLogicalTypeAnnotation; + } + + @Override + public int hashCode() { + return getClass().hashCode(); + } + } + public static class GeometryLogicalTypeAnnotation extends LogicalTypeAnnotation { private final String crs; @@ -1434,5 +1523,9 @@ default Optional visit(GeographyLogicalTypeAnnotation geographyLogicalType) { default Optional visit(UnknownLogicalTypeAnnotation unknownLogicalTypeAnnotation) { return empty(); } + + default Optional visit(FileLogicalTypeAnnotation fileLogicalType) { + return empty(); + } } } 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 2d6f0cbf8f..27cf96bd5e 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 @@ -847,12 +847,73 @@ public THIS addFields(Type... types) { @Override protected GroupType build(String name) { if (newLogicalTypeSet) { + if (logicalTypeAnnotation instanceof LogicalTypeAnnotation.FileLogicalTypeAnnotation) { + validateFileTypeFields(name, fields); + } return new GroupType(repetition, name, logicalTypeAnnotation, fields, id); } else { return new GroupType(repetition, name, getOriginalType(), fields, id); } } + private static void validateFileTypeFields(String name, List fields) { + for (Type field : fields) { + String fieldName = field.getName(); + if (!LogicalTypeAnnotation.FileLogicalTypeAnnotation.FIELD_NAMES.contains(fieldName)) { + throw new IllegalArgumentException("FILE type group '" + name + "' contains unrecognized field '" + + fieldName + "'. Valid fields are: " + + String.join(", ", LogicalTypeAnnotation.FileLogicalTypeAnnotation.FIELD_NAMES)); + } + Preconditions.checkArgument( + field.isPrimitive() && field.getRepetition() == Type.Repetition.OPTIONAL, + "FILE type field '%s' must be an optional primitive in group '%s'", + fieldName, + name); + validateFileTypeFieldPhysicalType(name, field.asPrimitiveType()); + } + } + + /** + * Validates that a declared FILE field uses the physical type required by the spec: + * {@code uri}, {@code content_type}, and {@code checksum} are STRING (BINARY), {@code offset} + * and {@code size} are INT64, and {@code inline} is BYTE_ARRAY (BINARY). + */ + private static void validateFileTypeFieldPhysicalType(String name, PrimitiveType field) { + String fieldName = field.getName(); + PrimitiveType.PrimitiveTypeName physicalType = field.getPrimitiveTypeName(); + switch (fieldName) { + case LogicalTypeAnnotation.FileLogicalTypeAnnotation.URI_FIELD: + case LogicalTypeAnnotation.FileLogicalTypeAnnotation.CONTENT_TYPE_FIELD: + case LogicalTypeAnnotation.FileLogicalTypeAnnotation.CHECKSUM_FIELD: + Preconditions.checkArgument( + physicalType == PrimitiveType.PrimitiveTypeName.BINARY + && field.getLogicalTypeAnnotation() + instanceof LogicalTypeAnnotation.StringLogicalTypeAnnotation, + "FILE type field '%s' must be a STRING (BINARY annotated as STRING) in group '%s'", + fieldName, + name); + break; + case LogicalTypeAnnotation.FileLogicalTypeAnnotation.OFFSET_FIELD: + case LogicalTypeAnnotation.FileLogicalTypeAnnotation.SIZE_FIELD: + Preconditions.checkArgument( + physicalType == PrimitiveType.PrimitiveTypeName.INT64, + "FILE type field '%s' must be an INT64 in group '%s'", + fieldName, + name); + break; + case LogicalTypeAnnotation.FileLogicalTypeAnnotation.INLINE_FIELD: + Preconditions.checkArgument( + physicalType == PrimitiveType.PrimitiveTypeName.BINARY, + "FILE type field '%s' must be a BYTE_ARRAY (BINARY) in group '%s'", + fieldName, + name); + break; + default: + // Unreachable: field names are validated against FIELD_NAMES before this call. + break; + } + } + public MapBuilder map(Type.Repetition repetition) { return new MapBuilder<>(self()).repetition(repetition); } diff --git a/parquet-column/src/test/java/org/apache/parquet/io/api/TestFileValueWriter.java b/parquet-column/src/test/java/org/apache/parquet/io/api/TestFileValueWriter.java new file mode 100644 index 0000000000..e3601bfc6f --- /dev/null +++ b/parquet-column/src/test/java/org/apache/parquet/io/api/TestFileValueWriter.java @@ -0,0 +1,210 @@ +/* + * 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.io.api; + +import static org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName.BINARY; +import static org.apache.parquet.schema.PrimitiveType.PrimitiveTypeName.INT64; +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +import java.util.ArrayDeque; +import java.util.Arrays; +import java.util.Deque; +import org.apache.parquet.io.ExpectationValidatingRecordConsumer; +import org.apache.parquet.schema.GroupType; +import org.apache.parquet.schema.LogicalTypeAnnotation; +import org.apache.parquet.schema.Types; +import org.junit.jupiter.api.Test; + +public class TestFileValueWriter { + private static final GroupType ALL_FIELDS_SCHEMA = Types.optionalGroup() + .as(LogicalTypeAnnotation.fileType()) + .optional(BINARY) + .as(LogicalTypeAnnotation.stringType()) + .named("uri") + .optional(INT64) + .named("offset") + .optional(INT64) + .named("size") + .optional(BINARY) + .as(LogicalTypeAnnotation.stringType()) + .named("content_type") + .optional(BINARY) + .as(LogicalTypeAnnotation.stringType()) + .named("checksum") + .optional(BINARY) + .named("inline") + .named("file"); + + @Test + public void testWritesSetFieldsInSchemaOrder() { + FileValue value = FileValue.builder() + .withUri("s3://bucket/file") + .withOffset(10) + .withSize(20) + .withContentType("image/png") + .withChecksum("ETAG:abc") + .withInline(Binary.fromString("bytes")) + .build(); + Deque expectations = new ArrayDeque<>(Arrays.asList( + "startGroup()", + "startField(uri, 0)", + "addBinary(s3://bucket/file)", + "endField(uri, 0)", + "startField(offset, 1)", + "addLong(10)", + "endField(offset, 1)", + "startField(size, 2)", + "addLong(20)", + "endField(size, 2)", + "startField(content_type, 3)", + "addBinary(image/png)", + "endField(content_type, 3)", + "startField(checksum, 4)", + "addBinary(ETAG:abc)", + "endField(checksum, 4)", + "startField(inline, 5)", + "addBinary(bytes)", + "endField(inline, 5)", + "endGroup()")); + + FileValueWriter.write(new ExpectationValidatingRecordConsumer(expectations), ALL_FIELDS_SCHEMA, value); + + assertThat(expectations).isEmpty(); + } + + @Test + public void testWritesInlineOnlyValue() { + GroupType schema = Types.optionalGroup() + .as(LogicalTypeAnnotation.fileType()) + .optional(BINARY) + .named("inline") + .named("file"); + FileValue value = FileValue.builder() + .withInline(Binary.fromConstantByteArray(new byte[0])) + .build(); + Deque expectations = new ArrayDeque<>(Arrays.asList( + "startGroup()", "startField(inline, 0)", "addBinary()", "endField(inline, 0)", "endGroup()")); + + FileValueWriter.write(new ExpectationValidatingRecordConsumer(expectations), schema, value); + + assertThat(expectations).isEmpty(); + } + + @Test + public void testRejectsUnresolvableValue() { + FileValue value = FileValue.builder().withContentType("image/png").build(); + + assertThatThrownBy(() -> FileValueWriter.write(noopConsumer(), ALL_FIELDS_SCHEMA, value)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("FILE value must set at least one of 'inline' or non-empty 'uri'"); + } + + @Test + public void testRejectsOffsetWithoutUri() { + FileValue value = FileValue.builder() + .withOffset(1) + .withSize(2) + .withInline(Binary.fromString("bytes")) + .build(); + + assertThatThrownBy(() -> FileValueWriter.write(noopConsumer(), ALL_FIELDS_SCHEMA, value)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("FILE value field 'offset' may only be set together with 'uri'"); + } + + @Test + public void testRejectsOffsetWithoutSize() { + FileValue value = FileValue.builder().withUri("file").withOffset(1).build(); + + assertThatThrownBy(() -> FileValueWriter.write(noopConsumer(), ALL_FIELDS_SCHEMA, value)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("FILE value field 'size' must be set whenever 'offset' is set"); + } + + @Test + public void testRejectsNegativeOffsetAndSize() { + FileValue negativeOffset = + FileValue.builder().withUri("file").withOffset(-1).withSize(1).build(); + FileValue negativeSize = + FileValue.builder().withUri("file").withSize(-1).build(); + + assertThatThrownBy(() -> FileValueWriter.write(noopConsumer(), ALL_FIELDS_SCHEMA, negativeOffset)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("FILE value field 'offset' must not be negative: -1"); + assertThatThrownBy(() -> FileValueWriter.write(noopConsumer(), ALL_FIELDS_SCHEMA, negativeSize)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("FILE value field 'size' must not be negative: -1"); + } + + @Test + public void testRejectsSetFieldMissingFromSchema() { + GroupType schema = Types.optionalGroup() + .as(LogicalTypeAnnotation.fileType()) + .optional(BINARY) + .as(LogicalTypeAnnotation.stringType()) + .named("uri") + .named("file"); + FileValue value = FileValue.builder().withUri("file").withSize(10).build(); + + assertThatThrownBy(() -> FileValueWriter.write(noopConsumer(), schema, value)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessage("FILE value sets field 'size' which is absent from group 'file'"); + } + + private static RecordConsumer noopConsumer() { + return new RecordConsumer() { + @Override + public void startMessage() {} + + @Override + public void endMessage() {} + + @Override + public void startField(String field, int index) {} + + @Override + public void endField(String field, int index) {} + + @Override + public void startGroup() {} + + @Override + public void endGroup() {} + + @Override + public void addInteger(int value) {} + + @Override + public void addLong(long value) {} + + @Override + public void addBoolean(boolean value) {} + + @Override + public void addBinary(Binary value) {} + + @Override + public void addFloat(float value) {} + + @Override + public void addDouble(double value) {} + }; + } +} 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 fbc79cc81e..ef469c0416 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 @@ -587,4 +587,263 @@ public void testVariantLogicalTypeWithShredded() { assertThat(((LogicalTypeAnnotation.VariantLogicalTypeAnnotation) annotation).getSpecVersion()) .isEqualTo(specVersion); } + + @Test + public void testFileLogicalTypeUriOnly() { + String name = "file_field"; + GroupType file = new GroupType( + REQUIRED, + name, + LogicalTypeAnnotation.fileType(), + Types.optional(BINARY).as(LogicalTypeAnnotation.stringType()).named("uri")); + + assertThat(file.toString()) + .isEqualTo("required group file_field (FILE) {\n" + " optional binary uri (STRING);\n" + "}"); + + LogicalTypeAnnotation annotation = file.getLogicalTypeAnnotation(); + assertThat(annotation.getType()).isEqualTo(LogicalTypeAnnotation.LogicalTypeToken.FILE); + assertThat(annotation.toOriginalType()).isNull(); + assertThat(annotation).isInstanceOf(LogicalTypeAnnotation.FileLogicalTypeAnnotation.class); + } + + @Test + public void testFileLogicalTypeAllFields() { + String name = "file_field"; + GroupType file = Types.requiredGroup() + .as(LogicalTypeAnnotation.fileType()) + .optional(BINARY) + .as(LogicalTypeAnnotation.stringType()) + .named("uri") + .optional(INT64) + .named("offset") + .optional(INT64) + .named("size") + .optional(BINARY) + .as(LogicalTypeAnnotation.stringType()) + .named("content_type") + .optional(BINARY) + .as(LogicalTypeAnnotation.stringType()) + .named("checksum") + .optional(BINARY) + .named("inline") + .named(name); + + LogicalTypeAnnotation annotation = file.getLogicalTypeAnnotation(); + assertThat(annotation).isInstanceOf(LogicalTypeAnnotation.FileLogicalTypeAnnotation.class); + assertThat(file.getFieldCount()).isEqualTo(6); + assertThat(file.getType("uri").getName()).isEqualTo("uri"); + assertThat(file.getType("offset").getName()).isEqualTo("offset"); + assertThat(file.getType("size").getName()).isEqualTo("size"); + assertThat(file.getType("content_type").getName()).isEqualTo("content_type"); + assertThat(file.getType("checksum").getName()).isEqualTo("checksum"); + assertThat(file.getType("inline").getName()).isEqualTo("inline"); + } + + @Test + public void testFileLogicalTypeInlineOnly() { + // Every field is optional, so an inline-only group is valid (spec inline case). + GroupType file = Types.requiredGroup() + .as(LogicalTypeAnnotation.fileType()) + .optional(BINARY) + .named("inline") + .named("inline_file"); + + assertThat(file.getLogicalTypeAnnotation()).isInstanceOf(LogicalTypeAnnotation.FileLogicalTypeAnnotation.class); + assertThat(file.getFieldCount()).isEqualTo(1); + assertThat(file.getType("inline").getName()).isEqualTo("inline"); + } + + @Test + public void testFileLogicalTypeAllowsOffsetWithoutUri() { + GroupType file = Types.requiredGroup() + .as(LogicalTypeAnnotation.fileType()) + .optional(INT64) + .named("offset") + .optional(INT64) + .named("size") + .optional(BINARY) + .named("inline") + .named("file_offset_without_uri"); + + assertThat(file.getFieldCount()).isEqualTo(3); + } + + @Test + public void testFileLogicalTypeExternalRangedReferenceWithoutInline() { + // An external ranged reference declares 'uri' + 'offset' + 'size' to point at a byte range of + // an external file. Because 'uri' is declared, the schema is treated as an external-reference + // schema and is not required to declare 'inline', even though it declares 'offset'. + GroupType file = Types.requiredGroup() + .as(LogicalTypeAnnotation.fileType()) + .optional(BINARY) + .as(LogicalTypeAnnotation.stringType()) + .named("uri") + .optional(INT64) + .named("offset") + .optional(INT64) + .named("size") + .named("external_ranged_file"); + + assertThat(file.getLogicalTypeAnnotation()).isInstanceOf(LogicalTypeAnnotation.FileLogicalTypeAnnotation.class); + assertThat(file.getFieldCount()).isEqualTo(3); + } + + @Test + public void testFileLogicalTypeAllowsMetadataOnlySchema() { + GroupType file = Types.requiredGroup() + .as(LogicalTypeAnnotation.fileType()) + .optional(BINARY) + .as(LogicalTypeAnnotation.stringType()) + .named("content_type") + .optional(BINARY) + .as(LogicalTypeAnnotation.stringType()) + .named("checksum") + .named("file_metadata_only"); + + assertThat(file.getFieldCount()).isEqualTo(2); + } + + @Test + public void testFileLogicalTypeAllowsSizeOnlySchema() { + GroupType file = Types.requiredGroup() + .as(LogicalTypeAnnotation.fileType()) + .optional(INT64) + .named("size") + .named("file_size_only"); + + assertThat(file.getFieldCount()).isEqualTo(1); + } + + @Test + public void testFileLogicalTypeAllowsOffsetWithoutSize() { + GroupType file = Types.requiredGroup() + .as(LogicalTypeAnnotation.fileType()) + .optional(BINARY) + .as(LogicalTypeAnnotation.stringType()) + .named("uri") + .optional(INT64) + .named("offset") + .named("file_offset_without_size"); + + assertThat(file.getFieldCount()).isEqualTo(2); + } + + @Test + public void testFileLogicalTypeOffsetWithSize() { + // 'offset' accompanied by 'size' is valid. + GroupType file = Types.requiredGroup() + .as(LogicalTypeAnnotation.fileType()) + .optional(BINARY) + .as(LogicalTypeAnnotation.stringType()) + .named("uri") + .optional(INT64) + .named("offset") + .optional(INT64) + .named("size") + .named("file_offset_with_size"); + + assertThat(file.getLogicalTypeAnnotation()).isInstanceOf(LogicalTypeAnnotation.FileLogicalTypeAnnotation.class); + assertThat(file.getFieldCount()).isEqualTo(3); + } + + @Test + public void testFileLogicalTypeSizeWithoutOffset() { + // 'uri' + 'size' (without 'offset') is valid: an external reference to '[0, size)'. + GroupType file = Types.requiredGroup() + .as(LogicalTypeAnnotation.fileType()) + .optional(BINARY) + .as(LogicalTypeAnnotation.stringType()) + .named("uri") + .optional(INT64) + .named("size") + .named("file_size_without_offset"); + + assertThat(file.getLogicalTypeAnnotation()).isInstanceOf(LogicalTypeAnnotation.FileLogicalTypeAnnotation.class); + assertThat(file.getFieldCount()).isEqualTo(2); + } + + @Test + public void testFileLogicalTypeRejectsUnrecognizedField() { + assertThatThrownBy(() -> Types.requiredGroup() + .as(LogicalTypeAnnotation.fileType()) + .optional(BINARY) + .as(LogicalTypeAnnotation.stringType()) + .named("uri") + .optional(BINARY) + .named("unknown_field") + .named("file_with_bad_field")) + .isInstanceOf(IllegalArgumentException.class); + } + + @Test + public void testFileLogicalTypeRejectsRequiredField() { + // All FILE fields must have OPTIONAL repetition under the current spec. + assertThatThrownBy(() -> Types.requiredGroup() + .as(LogicalTypeAnnotation.fileType()) + .required(BINARY) + .as(LogicalTypeAnnotation.stringType()) + .named("uri") + .named("file_with_required_uri")) + .isInstanceOf(IllegalArgumentException.class); + } + + @Test + public void testFileLogicalTypeRejectsGroupField() { + // FILE fields must be primitives, not nested groups. + assertThatThrownBy(() -> Types.requiredGroup() + .as(LogicalTypeAnnotation.fileType()) + .optionalGroup() + .optional(BINARY) + .named("nested") + .named("uri") + .named("file_with_group_field")) + .isInstanceOf(IllegalArgumentException.class); + } + + @Test + public void testFileLogicalTypeRejectsWrongStringPhysicalType() { + // 'uri' must be a STRING (BINARY annotated as STRING); an INT64 is rejected. + assertThatThrownBy(() -> Types.requiredGroup() + .as(LogicalTypeAnnotation.fileType()) + .optional(INT64) + .named("uri") + .named("file_uri_wrong_type")) + .isInstanceOf(IllegalArgumentException.class); + } + + @Test + public void testFileLogicalTypeRejectsUnannotatedStringField() { + // A STRING field must carry the STRING logical annotation; plain BINARY is rejected. + assertThatThrownBy(() -> Types.requiredGroup() + .as(LogicalTypeAnnotation.fileType()) + .optional(BINARY) + .named("uri") + .named("file_uri_unannotated")) + .isInstanceOf(IllegalArgumentException.class); + } + + @Test + public void testFileLogicalTypeRejectsWrongInt64PhysicalType() { + // 'offset' and 'size' must be INT64; an INT32 is rejected. + assertThatThrownBy(() -> Types.requiredGroup() + .as(LogicalTypeAnnotation.fileType()) + .optional(BINARY) + .as(LogicalTypeAnnotation.stringType()) + .named("uri") + .optional(INT32) + .named("size") + .named("file_size_wrong_type")) + .isInstanceOf(IllegalArgumentException.class); + } + + @Test + public void testFileLogicalTypeRejectsWrongInlinePhysicalType() { + // 'inline' must be a BYTE_ARRAY (BINARY); an INT64 is rejected. + assertThatThrownBy(() -> Types.requiredGroup() + .as(LogicalTypeAnnotation.fileType()) + .optional(INT64) + .named("inline") + .named("file_inline_wrong_type")) + .isInstanceOf(IllegalArgumentException.class); + } } diff --git a/parquet-format-structures/src/main/java/org/apache/parquet/format/LogicalTypes.java b/parquet-format-structures/src/main/java/org/apache/parquet/format/LogicalTypes.java index 8956d3944e..8aa21e0ae3 100644 --- a/parquet-format-structures/src/main/java/org/apache/parquet/format/LogicalTypes.java +++ b/parquet-format-structures/src/main/java/org/apache/parquet/format/LogicalTypes.java @@ -60,4 +60,5 @@ public static LogicalType VARIANT(byte specificationVersion) { public static final LogicalType BSON = LogicalType.BSON(new BsonType()); public static final LogicalType FLOAT16 = LogicalType.FLOAT16(new Float16Type()); public static final LogicalType UUID = LogicalType.UUID(new UUIDType()); + public static final LogicalType FILE = LogicalType.FILE(new FileType()); } 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 4252523fa0..0d5b0a7aa1 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 @@ -600,6 +600,11 @@ public Optional visit(LogicalTypeAnnotation.GeographyLogicalTypeAnn geographyType.setAlgorithm(fromParquetEdgeInterpolationAlgorithm(geographyLogicalType.getAlgorithm())); return of(LogicalType.GEOGRAPHY(geographyType)); } + + @Override + public Optional visit(LogicalTypeAnnotation.FileLogicalTypeAnnotation fileLogicalType) { + return of(LogicalTypes.FILE); + } } private void addRowGroup( @@ -1417,9 +1422,7 @@ LogicalTypeAnnotation getLogicalTypeAnnotation(LogicalType type) { VariantType variant = type.getVARIANT(); return LogicalTypeAnnotation.variantType(variant.getSpecification_version()); case FILE: - // Present in the format but not mapped to a LogicalTypeAnnotation yet. Ignore it to - // preserve the physical type, as an unrecognised logical type would be. - return null; + return LogicalTypeAnnotation.fileType(); default: throw new RuntimeException("Unknown logical type " + 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 46ffccf7fa..3137ff3848 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 @@ -573,12 +573,10 @@ public void testLogicalToConvertedTypeConversion() { } @Test - public void testFileLogicalTypeIsIgnoredRatherThanFailing() { + public void testFileLogicalTypeIsConverted() { ParquetMetadataConverter converter = new ParquetMetadataConverter(); - // FILE has no LogicalTypeAnnotation yet, so it must degrade to the physical type the way an - // unrecognised logical type does, rather than throwing. assertThat(converter.getLogicalTypeAnnotation(LogicalType.FILE(new FileType()))) - .isNull(); + .isEqualTo(LogicalTypeAnnotation.fileType()); } @Test @@ -2392,6 +2390,59 @@ public void testColumnIndexNanCountsRoundTrip() { assertThat(roundTrip.getNanCounts()).containsExactly(1L, 0L, 0L); } + @Test + public void testFileLogicalType() { + ParquetMetadataConverter parquetMetadataConverter = new ParquetMetadataConverter(); + + MessageType expected = Types.buildMessage() + .requiredGroup() + .as(LogicalTypeAnnotation.fileType()) + .optional(PrimitiveTypeName.BINARY) + .as(LogicalTypeAnnotation.stringType()) + .named("uri") + .optional(PrimitiveTypeName.INT64) + .named("offset") + .optional(PrimitiveTypeName.INT64) + .named("size") + .optional(PrimitiveTypeName.BINARY) + .as(LogicalTypeAnnotation.stringType()) + .named("content_type") + .optional(PrimitiveTypeName.BINARY) + .as(LogicalTypeAnnotation.stringType()) + .named("checksum") + .optional(PrimitiveTypeName.BINARY) + .named("inline") + .named("f") + .named("example"); + + List parquetSchema = parquetMetadataConverter.toParquetSchema(expected); + MessageType schema = parquetMetadataConverter.fromParquetSchema(parquetSchema, null); + assertThat(schema).isEqualTo(expected); + LogicalTypeAnnotation logicalType = schema.getType("f").getLogicalTypeAnnotation(); + assertThat(logicalType).isInstanceOf(LogicalTypeAnnotation.FileLogicalTypeAnnotation.class); + assertThat(logicalType).isEqualTo(LogicalTypeAnnotation.fileType()); + } + + @Test + public void testFileLogicalTypeRoundTripUriOnly() { + ParquetMetadataConverter parquetMetadataConverter = new ParquetMetadataConverter(); + + MessageType expected = Types.buildMessage() + .requiredGroup() + .as(LogicalTypeAnnotation.fileType()) + .optional(PrimitiveTypeName.BINARY) + .as(LogicalTypeAnnotation.stringType()) + .named("uri") + .named("f") + .named("example"); + + List parquetSchema = parquetMetadataConverter.toParquetSchema(expected); + MessageType schema = parquetMetadataConverter.fromParquetSchema(parquetSchema, null); + assertThat(schema).isEqualTo(expected); + LogicalTypeAnnotation logicalType = schema.getType("f").getLogicalTypeAnnotation(); + assertThat(logicalType).isInstanceOf(LogicalTypeAnnotation.FileLogicalTypeAnnotation.class); + } + @Test public void testV2StatsDoNotTriggerCorruptStatisticsCheck() { // Regression test: when V2 stats (min_value/max_value) are present,