Skip to content
Original file line number Diff line number Diff line change
@@ -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.
*
* <p>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<FileValue> 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<FileValue> consumer;
private FileValue.Builder builder;

private FileGroupConverter(GroupType schema, Consumer<FileValue> 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<String> consumer) {
return new PrimitiveConverter() {
@Override
public void addBinary(Binary value) {
consumer.accept(value.toStringUsingUTF8());
}
};
}

private static PrimitiveConverter binaryConverter(Consumer<Binary> consumer) {
return new PrimitiveConverter() {
@Override
public void addBinary(Binary value) {
consumer.accept(Binary.fromConstantByteArray(value.getBytes()));
}
};
}

private static PrimitiveConverter longConverter(Consumer<Long> consumer) {
return new PrimitiveConverter() {
@Override
public void addLong(long value) {
consumer.accept(value);
}
};
}
}
112 changes: 112 additions & 0 deletions parquet-column/src/main/java/org/apache/parquet/io/api/FileValue.java
Original file line number Diff line number Diff line change
@@ -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);
}
}
}
Original file line number Diff line number Diff line change
@@ -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;
}
}
}
Loading
Loading