ゼロバス・インジェストの概念

このページでは、Lakeflow ConnectにおけるZerobus Ingestのコアコンセプト、すなわちサービスの動作、ストリーム、サーバー、クライアント、そしてサポートするデータ型について説明します。

コンセプトへ移動:

ゼロバス・インジェストの仕組み

データプロデューサーはまずストリームをZerobus Ingest APIに開き、ターゲットとなるデルタテーブルを指定し、そのスキーマに合致するメッセージを構築し、そのメッセージを開かれたストリームにプッシュします。 このサービスはデータの耐久性を確保し、クライアントのメッセージに確認応答します。 その後、データを最適化された形でデルタテーブルに別のステップとして物質化します。 確認応答は永続性を示すものであり、クエリ可能性を示すものではありません。 この機能やクライアントにとっての意味については、 非同期コミュニケーション を参照してください。

Zerobus Ingestは、ワークロードに応じて弾力的にスケールするサーバーレスサービスです。 そのスケーリング方法については、以下の Zerobus Ingest のスケーリング方法 を参照してください。

Zerobus インジェストのしくみ

このセクションでは、Zerobus Ingestへの接続方法やデータがどのような形をとるかについても説明します。

  • APIプロトコル:APIプロトコル、gRPCとSDK、REST、OpenTelemetry、そしてそれぞれの使用タイミングについてです。
  • メッセージの種類:レコード形式、JSON、プロトコルバッファ(protobuf)、Apache Arrow、そしてそれぞれの使用タイミング。

Server

Zerobus 取り込みサービスでは、テーブルが自動的に作成または操作されることはありません。 ユーザーはテーブル自体を作成する必要があります。 テーブルとそのスキーマは、受信データの期待に対する権限のあるソースです。

Zerobus Ingestサーバーはクライアントから送られたデータを受け入れ、ターゲットテーブルスキーマに適合していることを検証します。 レコードが合致すれば、サーバーはそれを耐久性のあるものにし、クライアントにそれを認識します。 レコードをDeltaテーブルに物質化し、クエリ可能にするのは、その後すぐに別のステップとして行われます。

サービスの責任は次のとおりです。

  • メッセージのスキーマ検証をテーブルに対して行います。
  • 記録を耐久性のあるものにし、クライアントにそれを認めることです。 承認は耐久性を証明していますが、記録がまだ問い合わせられるものではありません。
  • データをターゲットテーブルにタイムリーに具現化し、その時点でクエリ可能になります。 レイテンシの数値については レイテンシを参照してください。

Client

クライアントはZerobus Ingestに接続し、レコードを送信し、それが耐久性があることを確認します。 Zerobus Ingest SDKを使うと、SDKがほとんどの処理をしてくれます。つまり、設定内容とSDKが自動的に行うことを分けて管理しやすくなります。

あなたは以下を構成または実装します:

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

SDKは自動的に以下の処理を行います:

  • メッセージの確認応答。 SDKは確認ループを実行し、オフセットや 確認コールバックを通じて耐久性確認を表示します。 特定のレコードは、アプリケーションが必要な時だけブロックします。 非 同期通信を参照してください。
  • リカバリ。 デフォルトでは、SDKは一時的な障害時に未確認のレコードを再接続し再生します。
    • 内蔵のリカバリー機能をオフにして、自分だけのリカバリーメカニズムを実装することもできます。 リカバリーのトリガー、設定オプション、カスタムリカバリーパターンについては 「リカバリーおよびリトライパターン」を参照してください。

SDKを使う際に確認応答やリカバリーロジックを手書きする必要はありません。 SDKを使わないカスタム統合の場合、 Zerobus SDKリポジトリ は統合構造やリカバリー処理の参考資料となります。

ストリーム

ストリームとは、クライアントとZerobus Ingestサーバーとの間に、永続的かつ双方向のgRPC接続によって確立される直接接続のことです。 SDK はストリームを使用して、有効期間が長くスループットの高い接続を容易にします。

  • ストリームは、SDK と共に gRPC API でのみ使用されます。
  • ストリームは、1 つのターゲット テーブルにデータを取り込みます。
  • 異なるテーブルに書き込みするための追加のストリームを開いたり、ワークロードの要求に応じて単一のクライアントのスループットを最大にスケールさせたりしましょう。

ストリームは順序付けの単位( 順序保証を参照)であり、ゼロバス・インジェストがスケールする単位でもあります(ゼロ バス・インジェストのスケール参照)。

順序保証

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

  • 複数のストリームにレコードを分散させる場合(例えばラウンドロビン)、それらのストリーム間での順序付けの保証はありません。
  • もしユースケースで複数の生産者やストリームにまたがる単一のトータルオーダーが必要な場合は、インジェスションの順序に頼るのではなく、アプリケーション内でその順序(例えばタイムスタンプやシーケンス番号で)を強制してください。

なぜgRPCストリーミングを送るのか

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

Zerobus Ingest のスケーリング方法

Zerobus Ingestは高いスケーラビリティを前提に設計されており、容量計画を求められずにそのスケールを実現できます。 これを可能にしているのは二つの設計選択です。

  • サーバーレスです。 サービスは負荷の変化に応じて自動的に容量を追加・削除するため、ブローカーのサイズ調整やパーティションのプロビジョニングは不要です。 ワークロードに必要な数だけ、同時ストリームを開いて、必要な数のテーブルに書き込むことができます。
  • ストリームは動的分割ユニットです。 固定されたパーティションのセットを再分割・再バランスさせてスケールアウトするのではなく、ストリームを開放・閉じ・回転させることが可能です。 ストリームをローテーションすることで、需要の変化に応じてサービスが能力と資源を再バランスできるため、より多くのストリームを開設し、より多くの生産者を稼働させながら、サービスが残りを吸収することでスケールを拡大できます。

実際の結果として、「Hello World」クライアントとペタバイト規模のワークロードは本質的に同じコードを実行します。 違いは、どれだけ多くのプロデューサーや配信を運営するかです。 この設計により、1兆件以上のレコードが単一のデルタテーブルに取り込まれることに耐えられています。 技術的背景については、「Ingesting the Milky Way: Petabyte-Scale with Zerobus Ingest」ブログ記事をご覧ください。

テーブル要件

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

  • Zerobus Ingestは管理されたDeltaテーブルにのみ書き込みます。 既定のストレージへの書き込みはサポートされていません。
  • Zerobus Ingestはプライベートエンドポイントを通じて保護されたストレージに書き込みを行いません。
  • Zerobus Ingestはターゲットテーブルの再作成をサポートしていません。
  • テーブル名はASCIIの文字、数字、アンダースコアのみをサポートします。
  • ワークスペースとターゲットテーブルの両方が サポートされているリージョンのいずれかに存在しなければなりません。

レコードをテーブルスキーマに対して検証する方法については、 スキーマ管理を参照してください。 パーティションやリキッドクラスタリングなどのテーブル機能については、 デルタテーブルの特徴を参照してください。

サポートされるデータ型

次の表は、インジェストでサポートされている Delta 型とそれに対応する Protobuf 型を示しています。

デルタ型 Protobuf 型
INTEGER int32
STRING string
FLOAT float
LONG int64
SHORT int32
DOUBLE double
DECIMAL(p, s)
10進テキスト、例えば「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上で、Variant値をJSONでエンコードされた文字列として取り込み、キーは STRING型で、Zerobus Ingestはデータをシュレッダせずにカラムに書き込みます。 Apache Arrow Flightでは、クライアントは代わりにVariant列のバック metadatavalue フィールドを構築します。 詳細は 「変異列の取り込み」を参照してください。
サポートされている形式は以下のとおりです。
  • オブジェクト: "{\"id\":0,\"example\":\"this is variant example\"}"
  • プリミティブ: "5""3.14""\"string\""
  • 配列: "[1,2,3]"
string