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:
+ *
+ * - {@code uri} (STRING): a URI-reference (RFC 3986) that identifies an external file, for
+ * example {@code s3://bucket/file.jpg}.
+ * - {@code offset} (INT64): start of the byte range within the external file identified by
+ * {@code uri}; if not set, treated as 0. Must not be negative and may only be set together
+ * with {@code uri}.
+ * - {@code size} (INT64): byte length of the referenced data within the external file. May be
+ * omitted for a whole-file external reference, in which case the range runs to the end of
+ * the referenced file. Must be set whenever {@code offset} is set. Must not be negative.
+ * - {@code content_type} (STRING): the media (MIME) type (RFC 2046) of the resolved bytes;
+ * when not set, {@code application/octet-stream} is assumed.
+ * - {@code checksum} (STRING): a self-describing integrity token for the resolved bytes, of
+ * the form {@code :}.
+ * - {@code inline} (BYTE_ARRAY): the referenced bytes stored inline in the value.
+ *
+ * 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,