Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
45 changes: 45 additions & 0 deletions parquet-protobuf/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -21,3 +21,48 @@ parquet-protobuf
================

Protocol Buffer support for Parquet columnar format.

## Message fields stored as proto bytes

Two kinds of message fields cannot be mapped to a Parquet group, so `ProtoSchemaConverter`
terminates them as the **serialized protobuf message** instead:

* **Fields of an empty message type** (`message Stub {}`) — Parquet forbids empty groups.
An empty message serializes to zero bytes, so the column is cheap, and field presence still
round-trips: `null` means the field was unset, an empty value means it was set.
* **Recursive fields beyond `parquet.proto.maxRecursion`** (default 5) — the remaining
sub-tree is stored as the serialized message instead of expanding the schema forever.

The Parquet type is an unannotated `BINARY` column that keeps the field's own repetition (or,
for repeated fields and map values, sits inside the standard `LIST` / `MAP` wrappers when
`parquet.proto.writeSpecsCompliant` is set):

```
message Trees.StubBox {
optional binary stub = 1; // Stub stub = 1;
optional group stubs (LIST) = 2 { // repeated Stub stubs = 2;
repeated group list {
required binary element;
}
}
optional group stub_map (MAP) = 3 { // map<string, Stub> stub_map = 3;
repeated group key_value {
required binary key (STRING);
optional binary value;
}
}
}
```

Readers that do not know about protobuf simply see opaque bytes (all of them empty for an
empty message type). Readers that have the generated message class can parse the bytes back into
the message; `ProtoParquetReader` does this automatically, resolving the class from the
`parquet.proto.class` footer key (or from the class configured for reading). The writer also stores
the message descriptor in the footer under `parquet.proto.descriptor`, which tools that do not
have the generated class can use to interpret the bytes; `ProtoParquetReader` itself does not read
it.

Note that the column type follows the proto schema at write time: if an empty message type later
gains fields, or `parquet.proto.maxRecursion` is changed, new files store the field as a group
where old files store `BINARY`. Tools that merge schemas across such files will report a type
conflict, the same way they do for any other field whose type changed.
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@
import com.google.protobuf.FloatValue;
import com.google.protobuf.Int32Value;
import com.google.protobuf.Int64Value;
import com.google.protobuf.InvalidProtocolBufferException;
import com.google.protobuf.Message;
import com.google.protobuf.StringValue;
import com.google.protobuf.UInt32Value;
Expand Down Expand Up @@ -347,6 +348,9 @@ protected Converter newScalarConverter(
if (messageType.equals(BytesValue.getDescriptor())) {
return new ProtoBytesValueConverter(pvc);
}
// Otherwise the column holds the serialized message itself: ProtoSchemaConverter stores
// fields of empty message types and recursion beyond maxRecursion as proto bytes.
return new ProtoBinaryMessageConverter(pvc, parentBuilder, fieldDescriptor);
}
Message.Builder subBuilder = parentBuilder.newBuilderForField(fieldDescriptor);
return new ProtoMessageConverter(conf, pvc, subBuilder, parquetType.asGroupType(), extraMetadata);
Expand Down Expand Up @@ -495,6 +499,39 @@ public void setDictionary(Dictionary dictionary) {
}
}

/**
* Reads a message field that {@link ProtoSchemaConverter} stored as the serialized proto bytes
* (a field of an empty message type, or recursion truncated at maxRecursion) back into the
* message.
*/
static final class ProtoBinaryMessageConverter extends PrimitiveConverter {

private final ParentValueContainer parent;
private final Message.Builder parentBuilder;
private final Descriptors.FieldDescriptor fieldDescriptor;

ProtoBinaryMessageConverter(
ParentValueContainer parent,
Message.Builder parentBuilder,
Descriptors.FieldDescriptor fieldDescriptor) {
this.parent = parent;
this.parentBuilder = parentBuilder;
this.fieldDescriptor = fieldDescriptor;
}

@Override
public void addBinary(Binary binary) {
Message.Builder builder = parentBuilder.newBuilderForField(fieldDescriptor);
try {
builder.mergeFrom(ByteString.copyFrom(binary.toByteBuffer()));
} catch (InvalidProtocolBufferException e) {
throw new ParquetDecodingException(
"Cannot parse field " + fieldDescriptor.getFullName() + " from its serialized proto bytes", e);
}
parent.add(builder.build());
}
}

static final class ProtoBinaryConverter extends PrimitiveConverter {

final ParentValueContainer parent;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -70,6 +70,22 @@

/**
* Converts a Protocol Buffer Descriptor into a Parquet schema.
* <p>
* Message fields normally become Parquet groups. Two kinds of message fields cannot, and are
* instead terminated as an unannotated {@code BINARY} column holding the serialized proto message
* (keeping the field's repetition, or sitting inside the usual LIST/MAP wrappers):
* <ul>
* <li>fields of an <em>empty</em> message type, because Parquet forbids empty groups; the value is
* zero bytes when the field is set and {@code null} when it is not, so presence still
* round-trips;</li>
* <li>recursive fields nested deeper than {@code maxRecursion}.</li>
* </ul>
* Readers unaware of protobuf see opaque bytes. {@code ProtoParquetReader} parses them back into the
* message using the generated class it resolves from the {@code parquet.proto.class} footer key (or
* the class configured for reading). Since the column type follows
* the proto schema at write time, an empty message type that later gains fields (or a changed
* {@code maxRecursion}) produces a group where older files hold {@code BINARY}, like any other
* field whose type changed. See the parquet-protobuf README for details.
*/
public class ProtoSchemaConverter {

Expand Down Expand Up @@ -312,21 +328,30 @@ private <T> Builder<? extends Builder<?, GroupBuilder<T>>, GroupBuilder<T>> addM
final GroupBuilder<T> builder,
ImmutableSetMultimap<String, Integer> seen,
int depth) {
// Prevent recursion by terminating with optional proto bytes.
// Terminate with proto bytes anything a static parquet schema cannot represent - recursion
// beyond maxRecursion and empty message types (parquet forbids empty groups) - preserving the
// field's repetition so the write path (Array/Repeated/MapWriter) still matches the schema.
depth += 1;
String typeName = getInnerTypeName(descriptor);
LOG.trace("addMessageField: {} type: {} depth: {}", descriptor.getFullName(), typeName, depth);
if (typeName != null) {
if (seen.get(typeName).size() > maxRecursion) {
return builder.primitive(BINARY, Type.Repetition.OPTIONAL).as((LogicalTypeAnnotation) null);
}
}

if (descriptor.isMapField() && parquetSpecsCompliant) {
// the old schema style did not include the MAP wrapper around map groups
// the old schema style did not include the MAP wrapper around map groups.
// The MAP structure is always preserved; a recursive or empty value type is truncated to
// proto bytes by the check below when addMapField recurses into the value field.
return addMapField(descriptor, builder, seen, depth);
}

boolean emptyMessage = descriptor.getMessageType().getFields().isEmpty();
if (emptyMessage || (typeName != null && seen.get(typeName).size() > maxRecursion)) {
if (descriptor.isRepeated() && parquetSpecsCompliant) {
// LIST-wrap the truncated bytes the same way any repeated primitive is wrapped
return addRepeatedPrimitive(BINARY, null, builder);
}
// optional, required, or repeated in the old schema style
return builder.primitive(BINARY, getRepetition(descriptor)).as((LogicalTypeAnnotation) null);
}

seen = ImmutableSetMultimap.<String, Integer>builder()
.putAll(seen)
.put(typeName, depth)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -356,42 +356,45 @@ private FieldWriter createMessageWriter(FieldDescriptor fieldDescriptor, Type ty
}

// This can happen now that recursive schemas get truncated to bytes. Write the bytes.
if (type.isPrimitive()
&& type.asPrimitiveType().getPrimitiveTypeName() == PrimitiveType.PrimitiveTypeName.BINARY) {
// The truncated type keeps the field's shape, so it may sit behind a LIST wrapper
// (repeated field) or be the value inside a MAP's key_value group.
Type contentType = getContentType(type);
if (contentType.isPrimitive()
&& contentType.asPrimitiveType().getPrimitiveTypeName() == PrimitiveType.PrimitiveTypeName.BINARY) {
return new BinaryWriter();
}

return new MessageWriter(fieldDescriptor.getMessageType(), getGroupType(type));
return new MessageWriter(fieldDescriptor.getMessageType(), contentType.asGroupType());
}

private GroupType getGroupType(Type type) {
/** Unwraps the LIST/MAP wrapper groups to the type holding the message content itself. */
private Type getContentType(Type type) {
if (type.isPrimitive()) {
return type;
}
LogicalTypeAnnotation logicalTypeAnnotation = type.getLogicalTypeAnnotation();
if (logicalTypeAnnotation == null) {
return type.asGroupType();
return type;
}
return logicalTypeAnnotation
.accept(new LogicalTypeAnnotation.LogicalTypeAnnotationVisitor<GroupType>() {
.accept(new LogicalTypeAnnotation.LogicalTypeAnnotationVisitor<Type>() {
@Override
public Optional<GroupType> visit(
LogicalTypeAnnotation.ListLogicalTypeAnnotation listLogicalType) {
public Optional<Type> visit(LogicalTypeAnnotation.ListLogicalTypeAnnotation listLogicalType) {
return ofNullable(type.asGroupType()
.getType("list")
.asGroupType()
.getType("element")
.asGroupType());
.getType("element"));
}

@Override
public Optional<GroupType> visit(
LogicalTypeAnnotation.MapLogicalTypeAnnotation mapLogicalType) {
public Optional<Type> visit(LogicalTypeAnnotation.MapLogicalTypeAnnotation mapLogicalType) {
return ofNullable(type.asGroupType()
.getType("key_value")
.asGroupType()
.getType("value")
.asGroupType());
.getType("value"));
}
})
.orElse(type.asGroupType());
.orElse(type);
}

private MapWriter createMapWriter(FieldDescriptor fieldDescriptor, Type type) {
Expand Down
Loading
Loading