メインコンテンツまでスキップ

Zerobus Ingest のコンセプト

このページでは、Lakeflow Connect における Zerobus Ingest の主要な概念(サービスの仕組み、ストリーム、サーバー、クライアント、およびサポートされているデータ型)について説明します。

コンセプトにジャンプ:

Zerobus Ingest の仕組み

データプロデューサーは、まず Zerobus Ingest API へのストリームを開いてターゲットの Delta テーブルを指定し、そのスキーマに一致するメッセージを構築してから、開いたストリームを通じてメッセージをプッシュします。サービスはデータを永続化し、クライアントのメッセージを承諾します。その後、別のステップとして、最適化された方法でデータを Delta テーブルにマテリアライズします。承諾は永続性を確認するものであり、クエリ可能性を確認するものではありません。これがどのように機能し、クライアントにとって何を意味するかについては、非同期通信を参照してください。

Zerobus Ingest は、ワークロードに合わせて弾力的にスケーリングする Serverless サービスです。スケーリングの詳細については、以下の How Zerobus Ingest scales を参照してください。

Zerobus Ingestの仕組み

このセクションでは、Zerobus Ingest への接続方法と、データの形式についても説明します。

  • APIプロトコル: APIプロトコル(SDK を使用した gRPC、REST、OpenTelemetry)と、それぞれの使用場面について。
  • メッセージ型: レコード形式である JSON、プロトコル バッファー (protobuf)、および Apache Arrow と、それぞれの使用場面について。

サーバー

Zerobus Ingest サービスは、テーブルを自動的に作成または操作しません。ユーザー自身がテーブルを作成する必要があります。テーブルとそのスキーマは、受信データの期待値に関する信頼できるソースです。

Zerobus Ingest サーバーは、クライアントから送信されたデータを受け取り、それがターゲットテーブルのスキーマに適合していることを検証します。レコードが適合する場合、サーバーはそれを耐久性のあるものにし、クライアントに応答を返します。レコードを Delta テーブルにマテリアライズしてクエリー可能にする処理は、その直後の別のステップとして実行されます。

サービス責任には以下が含まれます:

  • テーブルに対するメッセージのスキーマ検証。
  • レコードを耐久性のあるものにし、クライアントに対して承諾を通知します。この承諾は耐久性を確認するものであり、レコードがまだクエリー照会可能であることを確認するものではありません。
  • ターゲットテーブルへのデータのマテリアライズをタイムリーに行うことで、クエリが可能になります。レイテンシーの数値については、レイテンシーを参照してください。

クライアント

クライアントは Zerobus Ingest に接続してレコードを送信し、それらが永続化されたことを確認します。Zerobus Ingest SDK を使用する場合、SDK がそのほとんどを処理するため、構成するものと SDK が自動的に行うものを区別しておくと役立ちます。

構成または実装するもの:

  • ターゲットテーブルの選択。
  • Zerobus Ingestサービスへのストリームを開いています。
  • スキーマ互換性のあるメッセージを構築して送信します。

SDK は自動的に以下を処理します:

  • メッセージの承諾。 SDKが代わりに承諾ループを実行し、オフセットまたは 承諾コールバック を通じて耐久性の確認を提示します。アプリケーションが必要とする場合にのみ、特定のレコードでブロックします。非同期通信を参照してください。
  • リカバリ。 defaultでは、一時的なエラーが発生した場合、SDK は再接続を行い、未確認のレコードをリプレイします。
    • 組み込みのリカバリをオフにして、代わりに独自のリカバリメカニズムを実装することもできます。リカバリを Trigger するもの、構成オプション、およびカスタムリカバリパターンについては、「Recovery and retry patterns」を参照してください。

SDK を使用する場合、確認ロジックやリカバリロジックを手動で記述する必要はありません。SDK を使用しないカスタム統合の場合、 Zerobus SDK リポジトリが統合構造とリカバリ処理の参照先となります。

Stream

ストリームは、クライアントと Zerobus Ingest サーバー間の直接接続であり、永続的な双方向 gRPC 接続を介して確立されます。SDK はストリームを使用して、長期間持続する高 throughput の接続を促進します。

  • ストリームは、SDK を使用した gRPC API でのみ使用されます。
  • ストリームは、単一のターゲットテーブルにデータを取り込みます。
  • 追加のストリームを開いて異なるテーブルに書き込むか、ワークロードの要件に合わせて単一クライアントのthroughputを拡張します。

ストリームは順序付けの単位でもあり(順序の保証を参照)、Zerobus Ingest がスケールする単位でもあります(Zerobus Ingest のスケーリング方法を参照)。

順序の保証

順序はストリームごとに保証されます。レコードは、単一のストリームでエンキューされた順序でターゲットテーブルにコミットされます。ストリーム間でのグローバルな順序付けはありません。これに基づき、いくつかの設計上のポイントがあります。

  • レコードを複数のストリームに分散させた場合(ラウンドロビンなど)、それらのストリーム間での順序は保証されません。
  • ユースケースで多数のプロデューサーやストリーム全体にわたる単一の全順序が必要な場合は、取り込み順序に依存するのではなく、アプリケーションでその順序を強制してください(例えば、クエリー対象のTimestampやシーケンス番号を使用します)。

gRPC ストリーミングの理由

ストリームの gRPC 接続は開いたままになるため、クライアントはステートレスプロトコルのリクエストごとのセットアップコストを回避でき、単一のチャンネルを通じて連続的かつ大容量のレコードフローをプッシュできます。これが、SDK が最も高い throughput で取り込みを行える理由です。その他のインターフェイス (REST および OpenTelemetry) とそれぞれの選択時期については、「API プロトコル」を参照してください。

Zerobus Ingest のスケーリング方法

Zerobus Ingest は高いスケーラビリティを実現するように設計されており、容量計画を立てることなくその規模に到達できます。これを可能にする 2 つの設計上の選択肢があります。

  • Serverlessです。 このサービスは負荷の変化に応じて自動的に容量を追加および削除するため、ブローカーのサイズ設定やパーティションのプロビジョニングを行う必要はありません。ワークロードが必要とする数だけ、並列ストリームを開き、テーブルに書き込むことができます。
  • ストリームは動的なパーティション単位です。 スケールアウトのために再パーティション化や再調整が必要となる固定されたパーティションセットとは異なり、ストリームは開いたり、閉じたり、ローテーションしたりすることができます。ストリームをローテーションすることで、需要の変化に応じてサービスが容量とリソースを再調整できるようになります。そのため、ストリームを増やしてプロデューサーを多く実行することでスケールし、残りの処理はサービス側で吸収されます。

実質的な結果として、「hello world」クライアントとペタバイト規模のワークロードは、本質的に同じコードを実行します。違いは、実行するプロデューサーとストリームの数です。この設計により、単一の Delta テーブルへの1兆件を超えるレコードの取り込みが維持されています。技術的な背景については、Ingesting the Milky Way: Petabyte-Scale with Zerobus Ingest ブログ記事を参照してください。

テーブル要件

Zerobus Ingest は、ユーザーが作成および所有する Delta テーブルに書き込みを行います。ターゲットテーブルとワークスペースは、以下の要件を満たしている必要があります。

  • Zerobus Ingest は、管理された Delta テーブルにのみ書き込みを行います。defaultストレージへの書き込みはサポートされていません。
  • Zerobus Ingest は、プライベートEndpoint経由で保護されたストレージには書き込みを行いません。
  • Zerobus Ingest は、ターゲットテーブルの再作成をサポートしていません。
  • テーブル名に使用できるのはASCII文字、数字、アンダースコアのみです。
  • ワークスペースとターゲットテーブルは、両方ともサポートされているリージョンのいずれかにある必要があります。

テーブルスキーマに対してレコードがどのように検証されるかについては、「スキーマ管理」を参照してください。パーティション分割やリキッドクラスタリングなどのテーブル機能については、「Deltaテーブルの機能」を参照してください。

サポートされているデータ型

次の表は、取り込み用にサポートされている Delta タイプと、それに対応する Protobuf タイプを示しています。

Delta 型

Protobuf 型

INTEGER

int32

STRING

string

FLOAT

float

LONG

int64

SHORT

int32

DOUBLE

double

DECIMAL(p, s)

Decimalテキスト(例:)「123.45」、「1e2」など。

string

BOOLEAN

bool

BINARY

bytes

BYTE (TINYINT)

int32

DATE

int32 (エポックからの日数)に変換する必要があります。

int32

TIMESTAMP

int64 (エポック時間、マイクロ秒単位) に変換する必要があります。

int64

TIMESTAMPNTZ

int64 (エポック時間、マイクロ秒単位) に変換する必要があります。

int64

ARRAY<TYPE>

repeated TYPE

MAP<K,V>

map<K,V>

map Protobuf の糖衣構文(syntactic sugar)は、Protobuf コンパイラのバージョン 3 以降でのみ利用可能です。

STRUCT<FIELDS>

message Nested { FIELDS }

VARIANT

gRPC SDK および REST を介して、STRING 型のキーを持つ JSON エンコードされた文字列としてバリアント値を取り込みます。Zerobus Ingest は、そのデータをシュレッディングせずに列に書き込みます。Apache Arrow Flight の場合、クライアントは代わりにバリアント列の基盤となる metadata フィールドと value フィールドを構築します。VARIANT 列の取り込みを参照してください。

サポートされている形式は以下の通りです:

  • オブジェクト: "{\"id\":0,\"example\":\"this is variant example\"}"
  • プリミティブ: "5""3.14""\"string\""
  • 配列: "[1,2,3]"

string

Delta 型

Protobuf 型

INTEGER

int32

STRING

string

FLOAT

float

LONG

int64

SHORT

int32

DOUBLE

double

DECIMAL(p, s)

Decimalテキスト(例:)「123.45」、「1e2」など。

string

BOOLEAN

bool

BINARY

bytes

BYTE (TINYINT)

int32

DATE

int32 (エポックからの日数)に変換する必要があります。

int32

TIMESTAMP

int64 (エポック時間、マイクロ秒単位) に変換する必要があります。

int64

TIMESTAMPNTZ

int64 (エポック時間、マイクロ秒単位) に変換する必要があります。

int64

ARRAY<TYPE>

repeated TYPE

MAP<K,V>

map<K,V>

map Protobuf の糖衣構文(syntactic sugar)は、Protobuf コンパイラのバージョン 3 以降でのみ利用可能です。

STRUCT<FIELDS>

message Nested { FIELDS }

VARIANT

gRPC SDK および REST を介して、STRING 型のキーを持つ JSON エンコードされた文字列としてバリアント値を取り込みます。Zerobus Ingest は、そのデータをシュレッディングせずに列に書き込みます。Apache Arrow Flight の場合、クライアントは代わりにバリアント列の基盤となる metadata フィールドと value フィールドを構築します。VARIANT 列の取り込みを参照してください。

サポートされている形式は以下の通りです:

  • オブジェクト: "{\"id\":0,\"example\":\"this is variant example\"}"
  • プリミティブ: "5""3.14""\"string\""
  • 配列: "[1,2,3]"

string