ClickHouse / ClickHouse/clickhouse-java

Support Protobuf insert

Open
#2,187 7 comments 0 reactions 0 assignees View on GitHub
area:format client-api-v2 enhancement
Dominant language
Java
Stars
1.6k
Forks
636
Avg merge
2d 16h
Merged PRs (30d)
28

Description

### Use case

We are used to work with `Protocol Buffer` data in our `Beam Dataflow` ingestion pipeline. We tried many ways to ingest data and the easiest was with this new `client-v2`. While using `JSONEachRow` we can insert our protobuf messages through `com.google.protobuf.util.JsonFormat.Printer#print()` but it is not efficient as mentioned in [FastFomats](https://clickhouse.com/blog/clickhouse-input-format-matchup-which-is-fastest-most-efficient).

So I tried serializing data into `RowBinary` format. It works with some adjustments to existing code. I'm still wondering how to serialize into `Native` format to serialize as columns and be closer to `MergeTree format`. If you have any tip on it I would be very happy to make some tries.

### Describe the solution you'd like

1. Integrate a `ProtobufSerializer` in the source code
2. Detail how to serialize into `Native` format. `RowBinary` is good but still row oriented. Compared to what we had to ingest into `BigQuery` we are already way more efficient with `ClickHouse` so maybe we can keep improving this part

### Describe the alternatives you've considered

From my experience below `serializeProto` function is enough to insert protobuf messages. This is at very experimental phase.

```
public static void serializeProto(GeneratedMessageV3 message, ClickHouseColumn column, OutputStream out) throws IOException {

Object value = null;

// Get the descriptor of the field to serialize
Descriptors.FieldDescriptor descriptor = getFieldDescriptor(message, column.getColumnName());

if (descriptor != null) {
if (descriptor.isRepeated() || message.hasField(descriptor)) {
// Ensure the field has a value otherwise Protobuf returns the default value for type which is not necessarily
// the default value ClickHouse uses in its table schema definition.
value = message.getField(descriptor);

LOGGER.debug("{}({}) - {} - nullable:{} - default:{}",
descriptor.getName(), descriptor.getType(), value, column.isNullable(), column.hasDefault());

if (Set.of(ClickHouseDataType.DateTime, ClickHouseDataType.DateTime64).contains(column.getDataType())) {
// Some events DateTime fields are long. We need to convert them.
// We use `DateTime` and `DateTime64` only, but we may want to support more in the future.
// TODO: add more checking when needed as some Protobuf fields could be `com.google.protobuf.Timestamp`
value = convertFromInt((Long) value);
}
} else if (!message.hasField(descriptor) && !column.isNullable() && !column.hasDefault()) {
// TODO: integrate this condition with the rest to also handle default dates
// The field exists in the proto but it is not set. If it is not nullable or has no default value in
// ClickHouse then we should take default proto value
value = message.getField(descriptor);
}
} else {
// In this case the protobuf does not have this column. In our case it's `insertTime` and nothing needs
// to be done. This is an extra column not part of original data but set at run time by CH for further monitoring
LOGGER.debug("{}({}) - {} - nulable:{} - default:{}",
column.getColumnName(), column.getDataType(), value, column.isNullable(), column.hasDefault());
}

// Same code as in `RowBinaryFormatWriter.commitRow`
if (RowBinaryFormatSerializer.writeValuePreamble(out, true, column, value)) {
ABSerializerUtils.serializeData(out, value, column);
}
}
```

This method is paired with a single modification in `com.clickhouse.client.api.data_formats.internal.SerializerUtils#serializeTupleData`:
```
else if (value instanceof GeneratedMessageV3) {
GeneratedMessageV3 message = (GeneratedMessageV3) value;
// Start: added section
// From Protobuf messages Tuple are actually Protobuf messages themselves. So we need more capabilities to
// translate them.
for (ClickHouseColumn nestedColumn : column.getNestedColumns()) {
// TODO: support cases when nested value is also a `DateTime` or `DateTime64` as above.
Descriptors.FieldDescriptor descriptor = getFieldDescriptor(message, nestedColumn.getColumnName());
Object nestedValue = message.getField(descriptor);
serializeData(stream, nestedValue, nestedColumn);
}
// End: added section
}
```

Then inserting using the client is very straightforward

```
// Get ClickHouse table schema
TableSchema schema = clickhouseClient.getTableSchema(tableName);

// Serialise messages into an output stream
ByteArrayOutputStream out = new ByteArrayOutputStream();

for (GeneratedMessageV3 event: c.element().getValue()) {
List columnList = schema.getColumns();
for (ClickHouseColumn column : columnList) {
serializeProto(event, column, out);
}
}

// Prepare insert
InputStream inputStream = new ByteArrayInputStream(out.toByteArray());
InsertSettings settings = new InsertSettings();
ClickHouseFormat format = ClickHouseFormat.RowBinaryWithDefaults;
```

### Additional context

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.