メッセージの種類

gRPC経由でZerobus Ingest SDKを導入すると、特定のストリーム上でレコードのエンコード方法を選択します。 Zerobus Ingestは3つのメッセージ形式(JSON、プロトコルバッファ(protobuf)、Apache Arrowをサポートしているため、作業負荷に応じてシンプルさ、型の安全性、スループットをトレードオフできます。 すべての形式は、データが永続化される前に、Deltaテーブルのスキーマに対して検証されます。 スキーマ管理を参照してください。

同じ1,000レコードを3つのメッセージ形式で使用します。protobufはコンパクトな型付き行エンコーディング、Apache Arrowは一つの列単位バッチでメタデータとバッファがバッチ全体に分散され、JSONは繰り返しフィールド名が重複する読みやすいテキストです

どのフォーマットを使うべきでしょうか?

フォーマット 最適な用途 メモ
JSON 概要およびシンプルなプロデューサー。 最もシンプルな選択肢で、コンパイルするスキーマ定義が不要です。 便利ですが、大量作業用のバイナリフォーマットより遅いです。
プロトコルバッファ 本番、行指向の、大容量のストリーム。 型安全でコンパクトなバイナリ符号化。 ほとんどの本社ワークロードに推奨されています。 コンパイルされたスキーマが必要です。
アパッチ・アロー 列指向またはバッチ指向のワークロード。 Apache Arrowレコードバッチを直接送信し、行ごとのシリアライズを回避します。 データがすでにカラム形式で行われているか、バッチで取り込んでいる場合が最適です。 ベータ版です。 Zerobus Ingest で Arrow Flight を使用するを参照してください。

以下のスニペットは、Python SDK内の各フォーマットの形状を示しています。 彼らは、すでにSDKクライアントを作成し、ターゲットテーブルを知っていると仮定します。 すべての言語における完全なセットアップ(エンドポイント、テーブル、サービスプリンシパル)および例については、「 Use Zerobus Ingest」を参照してください。

JSON

JSONは最もシンプルな始め方です。スキーマ定義なしでコンパイルするレコードをJSONオブジェクトとして送信します。 迅速な立ち上げやプロトタイピング、そして純粋な処理性能よりも利便性を重視するユーザーに最適です。 大量生産ワークロードには、バイナリフォーマット(protobufまたはArrow)の方が効率的です。

ディスクリプタなしでテーブル名を TableProperties に渡してJSONレコード用のストリームを作成します:

table_properties = TableProperties(TABLE_NAME)
stream = sdk.create_stream(CLIENT_ID, CLIENT_SECRET, table_properties)

stream.ingest_record_offset({"device_name": "sensor-1", "temp": 22, "humidity": 55})

JSONの全攻略については 「クライアントを書く」をご覧ください。

プロトコルバッファ

Protobufは型安全でコンパクトなバイナリエンコーディングを提供し、ほとんどの本番の行指向ワークロードで推奨されるフォーマットです。 ターゲットのデルタテーブルに合うprotobufスキーマを定義し( Protobufスキーマを参照)、それをコンパイルするとSDKがgRPC上でレコードごとにレコードを取り込みます。

protobufを使うには3つのステップがあります:テーブルに合った .proto スキーマを生成し、言語モジュールにコンパイルし、ディスクリプタを TablePropertiesに渡してレコードを取り込みます。 以下の例はPython SDKを使用しています。

1. テーブルから .proto スキーマを生成する。 Python SDKには、デルタテーブルを読み取り、マッチングスキーマを書き込むgenerate_protoツールが含まれています。

python -m zerobus.tools.generate_proto \
    --uc-endpoint "https://<workspace-id>.cloud.databricks.com" \
    --client-id "<client-id>" \
    --client-secret "<client-secret>" \
    --table "main.default.air_quality" \
    --output "record.proto" \
    --proto-msg "AirQuality"

生成されたスキーマは proto2 構文を使用し、各Delta列ごとにオプションフィールドがあります:

syntax = "proto2";
message AirQuality {
    optional string device_name = 1;
    optional int32 temp = 2;
    optional int64 humidity = 3;
}

2. スキーマをprotobufコンパイラでPythonモジュールにコンパイルします:

pip install "grpcio-tools>=1.60.0,<2.0"
python -m grpc_tools.protoc --python_out=. --proto_path=. record.proto

これにより、 record_pb2.pyが生成されます。

3. コンパイルされたディスクリプタをTableProperties(protobufのデフォルト)に渡してレコードを取り込む。 SDKは各レコードをシリアル化するためにディスクリプタを使用します:

import record_pb2

descriptor_bytes = record_pb2.AirQuality.DESCRIPTOR.file.serialized_pb
table_properties = TableProperties(TABLE_NAME, descriptor_bytes)
stream = sdk.create_stream(CLIENT_ID, CLIENT_SECRET, table_properties)

record = record_pb2.AirQuality(device_name="sensor-1", temp=22, humidity=55)
stream.ingest_record_offset(record)

上記の例はPython SDKを使用しています。 ツールは言語によって異なります。あるSDKはテーブルから.protoを生成するgenerate_protoユーティリティを付属させ、他のSDKは既存の.protoをコンパイルします。 言語ごとの手順については、「 クライアントを書く」の各SDKタブにあるprotobufノートをご覧ください。 ツールのソースや完全な例については 、Zerobus SDKリポジトリをご覧ください。

アパッチ・アロー

Important

Apache Arrow取り込み機能はベータ版です。

Apache Arrowの取り込みは、各行をJSONやprotobufに変換するのではなく、同じgRPC接続経由で ArrowRecordBatch データを直接送信します。 アプリケーションがすでにArrowデータを生成している場合や、行をバッチで取り込む場合、特に広く数値重視のスキーマや分析志向のスキーマで行ごとのシリアライゼーションがオーバーヘッドになる場合に最適です。

Arrowは、非常に大きなバッチにも向いています。 JSONやprotobufのバッチ方式は全か無かで、メッセージあたりのサイズ制限で制限されますが、Arrow Flightパスは大きなバッチをより小さなトランスポートメッセージに分割し、それぞれ個別に送信・確認します。 Arrow Flight のバッチは例外およびZerobus Ingest で Arrow Flight を使用するを参照してください。

pyarrow.SchemaでArrowストリームを開き、RecordBatchデータを取り込みます:

stream = sdk.create_arrow_stream(TABLE_NAME, schema, CLIENT_ID, CLIENT_SECRET)

stream.ingest_batch(batch)

スキーマ定義、バッチ処理、圧縮を含むArrow Flightの全攻略については、「 Use Arrow Flight with Zerobus Ingest」をご覧ください。