Zerobus Ingestをご利用ください

このページでは、Lakeflow ConnectのZerobus Ingestを使ってデータを取り込む方法を説明しています。

Zerobus Ingestを始めましょう

始める前に、Zerobus Ingestがあなたのワークスペースの地域で利用可能かを確認してください。 インジェストの可用性を参照してください。

  1. Zerobus インジェスト URL を取得します。
  2. データを取り込むテーブルを作成または識別します。
  3. サービス プリンシパルを作成し、テーブルに権限を付与します。
  4. クライアントまたはエクスポーターを接続して、データの送信を開始します。

ユース ケースのガイドを選択します。

  • 独自のデータを取り込む: 定義したスキーマで Zerobus Ingest SDK または REST API を使用します。 このページで説明されている手順を実行します。

  • OpenTelemetry データを取り込む: 標準の OpenTelemetry SDK またはコレクターを使用して、トレース、ログ、メトリックを定義済みのテーブル スキーマに送信します。 完全な手順については、「 Zerobus Ingest を使用して OpenTelemetry データを取り込む」を参照してください。

インターフェイスを選択する

Zerobus Ingestは複数のインターフェースをサポートしており、すべてUnity CatalogのDeltaテーブルに直接書き込みます。 要約すると、次のようになります。

  • gRPC上のSDKは、持続的なスループットが最も高く、大量ストリーミング制作者に最適です。
  • REST: ステートレスで、軽量なエッジデバイスや「頻繁に通信する」エッジデバイスを多数展開する環境に最適です。
  • OpenTelemetry(OTLP):すでにOpenTelemetryのトレース、ログ、メトリクスを発信しているシステム向けです。 Zerobus IngestによるOpenTelemetryデータの取り込みを参照してください。

詳細な比較方法や選択方法については 、APIプロトコルをご覧ください。 SDKの上で、レコード形式(JSON、プロトコルバッファ(protobuf)、またはApache Arrowなども選択できます。 メッセージ の種類を参照してください。 このページの残りの部分はSDKとREST APIを使っています。

ワークスペースの URL と Zerobus 取り込みエンドポイントを取得する

ログインすると、ワークスペースの URL がブラウザーに表示されます。 完全な URL は https://<databricks-instance>.net/o=XXXXXの形式に従います。ワークスペースの URL は、 /o=XXXXXの前のすべての URL で構成されます。 たとえば、次の完全な URL を指定すると、ワークスペースの URL とワークスペース ID を決定できます。

  • 完全な URL: https://abcd-teste2-test-spcse2.azuredatabricks.net/?o=2281745829657864#
  • ワークスペース URL: https://abcd-teste2-test-spcse2.azuredatabricks.net
  • ワークスペース ID: 2281745829657864

サーバー エンドポイントは、ワークスペースとリージョンによって異なります。

  • サーバー エンドポイント: <workspace-id>.zerobus.<region>.azuredatabricks.net

ワークスペースのリージョンを見つけるには、Databricks UIの上部ナビゲーションバーにあるワークスペーススイッチャーを開きます。 リージョンは各ワークスペース名の下に表示されます(例: eastus)。 アカウント コンソール[ワークスペース] で見つけることもできます。

利用可能なリージョンについては、Zerobus Ingest quotasを参照してください。

ターゲット テーブルを作成または識別する

データを取り込むターゲット テーブルを特定します。 新しいターゲット テーブルを作成するには、 CREATE TABLE SQL コマンドを実行します。 たとえば、 unity.default.air_qualityという名前の新しいテーブルを作成します。

    CREATE TABLE unity.default.air_quality (
    device_name STRING, temp INT, humidity LONG);

OpenTelemetry インジェストの場合、テーブルでは、シグナルの種類 (トレース、ログ、メトリック) ごとに定義済みのスキーマを使用する必要があります。 Unity カタログでのターゲット テーブルの作成を参照してください。

テーブルスキーマはZerobus Ingestが受け入れる契約であり、Zerobus Ingestはそれを自動進化させません。 スキーマ変更は事前に計画しましょう。まずテーブルを変更し、その後でプロデューサーを更新します。 Zerobus Ingest は、テーブルに互換性のない変更が加えられた後、適合しなくなったレコードを破棄せず、永続的なフォールバック先に書き込みます。 スキーマ管理耐久フォールバック場所からのデータの復旧を参照してください。

デフォルトでは、Zerobus Ingestはターゲットテーブルのスキーマと一致しないフィールドのレコードを拒否します。 これらのフィールドが失われないように取り込むには、レスキュー列を設定します。 Zerobus レスキュー列を参照してください。

ストリーミングテーブルに取り込む

Important

Zerobus Ingestを使った ストリーミングテーブル へのインジェストは ベータ版です。 ワークスペース管理者は、[ プレビュー] ページからこの機能へのアクセスを制御できます。 Manage Azure Databricks プレビューを参照してください。

新しいストリーミング テーブルを作成するには、 CREATE STREAMING TABLE SQL コマンドを実行します。 例えば次が挙げられます。

CREATE STREAMING TABLE unity.default.air_quality (
device_name STRING, temp INT, humidity LONG);

ストリーミング テーブルが作成されたら、「 クライアントの書き込み」のいずれかのインターフェイスを使用して、標準の Delta テーブルの場合とまったく同じように取り込みます。 ストリーミングテーブルへの書き込みは、管理されたDeltaテーブルへの書き込みと同じ方法で動作し、制限とクォータも同じです。

サービス プリンシパルを作成し、アクセス許可を付与する

サービス プリンシパルは、パーソナライズされたアカウントよりもセキュリティを提供する特殊な ID です。 サービスプリンシパルや認証への利用方法の詳細については、「OAuthでサービスプリンシパルアクセスをAzure Databricksに許可する」をご覧ください。

サービスプリンシパルは、Azure Databricks REST APIやSDKを使ってプログラム的に作成・管理、または下記のワークスペースUIを通じて行えます。 このセクションの最後に付与される権限は、どのクライアントからでも実行可能なSQLです。

  1. サービスプリンシパルを作成するには、「 設定>アイデンティティとアクセス」へ行ってください。

  2. [ サービス プリンシパル] で 、[ 管理] を選択します。

  3. サービスプリンシパルの追加をクリックします。

  4. [ サービス プリンシパルの追加 ] ウィンドウで、[新規追加] をクリックして しいサービス プリンシパルを作成します。

  5. サービス プリンシパルのクライアント ID とクライアント シークレットを生成して保存します。

  6. カタログ、スキーマ、およびテーブルに必要なアクセス許可をサービス プリンシパルに付与します。

    1. Service principal ページで「Configurations」タブに移動します。
    2. アプリケーション ID (UUID) をコピーします。
    3. 次の SQL を使用してアクセス許可を付与します。必要に応じて、UUID とカタログ、スキーマ名、テーブル名の例を置き換えます。
    GRANT USE CATALOG ON CATALOG <catalog> TO `<UUID>`;
    GRANT USE SCHEMA ON SCHEMA <catalog.schema> TO `<UUID>`;
    GRANT MODIFY, SELECT ON TABLE <catalog.schema.table_name> TO `<UUID>`;
    

クライアントを書き込む

任意のプログラミング言語で Zerobus SDK を使用するか、REST API を使用してターゲット テーブルにデータを取り込みます。 SDKはオープンソースです。 完全なライブラリ、言語固有のドキュメント、追加の例については 、Zerobus SDKリポジトリをご覧ください。

以下の例は ingest_record_offsetを使っており、これはレコードを送信する順序を保持します。

Python SDK

Python 3.9 以上が必要です。 SDKは非同期ランタイムを通じて高スループットかつ効率的なネットワークI/Oを提供します。 JSON (最も単純) とプロトコル バッファー (運用環境に推奨) がサポートされています。 SDKは同期および非同期の両方の実装、オフセットベースおよび未来ベースのインジェスション方式もサポートしています。

pip install databricks-zerobus-ingest-sdk

JSON の例:

import logging
from zerobus.sdk.sync import ZerobusSdk
from zerobus.sdk.shared import TableProperties

# See "Get your workspace URL and Zerobus Ingest endpoint" for information on obtaining these values.
SERVER_ENDPOINT="https://1234567890123456.zerobus.eastus.azuredatabricks.net"
DATABRICKS_WORKSPACE_URL="https://adb-1234567890123456.12.azuredatabricks.net"
TABLE_NAME="main.default.air_quality"
CLIENT_ID="your-client-id"
CLIENT_SECRET="your-client-secret"

sdk = ZerobusSdk(
    SERVER_ENDPOINT,
    DATABRICKS_WORKSPACE_URL
)

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

try:
    for i in range(1000):
        record_dict = {
            "device_name": f"sensor-{i}",
            "temp": 20 + i % 15,
            "humidity": 50 + i % 40
        }
        stream.ingest_record_offset(record_dict)
finally:
    stream.close()

上記の例は、返還オフセットを待たずにオフセットベースの ingest_record_offset 法を使用しています。 利用可能な取り込み方法、オフセットの耐久性確認を待つタイミング、そして確認応答による進捗追跡方法については、「 メッセージブロッキングと確認応答」をご覧ください。

プロトコルバッファ: 型安全なインジェスティションのために、protobufディスクリプタを TableProperties に渡します(フォーマットは自動的に選択されます)。 テーブルから generate_proto ツールを使ってスキーマを作成し、 protocでコンパイルし、コンパイルしたディスクリプタを渡してストリームを作成します。

矢印フライト (ベータ): 同じ gRPC 接続を介した Apache Arrow RecordBatch データの列指向またはバッチ指向の取り込みについては、「 Zerobus 取り込みでの矢印フライトの使用」を参照してください。 [arrow]拡張機能が必要です: pip install "databricks-zerobus-ingest-sdk[arrow]" pyarrow

完全なドキュメント、構成オプション、バッチ インジェスト、プロトコル バッファーの例については、Python SDK リポジトリを参照してください。

Rust SDK

Rust 1.70 以上が必要です。 SDKは高スループットの取り込みのために非同期I/OとgRPCを使用しています。 JSON (最も単純) とプロトコル バッファー (運用環境に推奨) がサポートされています。

まず、パッケージをインポートします。

cargo add databricks-zerobus-ingest-sdk

または、 Cargo.tomlに追加します。

[dependencies]
databricks-zerobus-ingest-sdk = "2.0.0" # Latest version at time of publication

JSON の例:

  use databricks_zerobus_ingest_sdk::{JsonString, ZerobusSdk};
  use std::error::Error;

  // See "Get your workspace URL and Zerobus Ingest endpoint" for information on obtaining these values.
  const DATABRICKS_WORKSPACE_URL: &str = "https://adb-1234567890123456.12.azuredatabricks.net";
  const SERVER_ENDPOINT: &str = "1234567890123456.zerobus.eastus.azuredatabricks.net";
  const TABLE_NAME: &str = "main.default.air_quality";
  const CLIENT_ID: &str = "your-client-id";
  const CLIENT_SECRET: &str = "your-client-secret";


  #[tokio::main]
  async fn main() -> Result<(), Box<dyn Error>> {
      let sdk_handle = ZerobusSdk::builder()
          .endpoint(SERVER_ENDPOINT)
          .unity_catalog_url(DATABRICKS_WORKSPACE_URL)
          .build()?;

      let mut stream = sdk_handle
          .stream_builder()
          .table(TABLE_NAME)
          .oauth(CLIENT_ID, CLIENT_SECRET)
          .json()
          .max_inflight_requests(100)
          .build()
          .await?;

      stream.ingest_record_offset(
        JsonString("{
          \"device_name\": \"sensor\",
          \"temp\": 22,
          \"humidity\": 55}".to_string())).await?;

      println!("Record ingested successfully");
      stream.close().await?;
      println!("Stream closed successfully");

      Ok(())
  }

Protocol Buffers: 型安全な取り込みを行うには、.json() ではなく、ストリームビルダー上の .compiled_proto(descriptor) を介して Protocol Buffers を使用してください。ここで、descriptorprost_types::DescriptorProto です。 generate_proto ツールを使用して必要なファイルを生成し、プロジェクトにインポートします。 矢印フライト (ベータ): 同じ gRPC 接続を介した Apache Arrow RecordBatch データの列指向またはバッチ指向の取り込みについては、「 Zerobus 取り込みでの矢印フライトの使用」を参照してください。 Cargo 機能を使用して有効にする: cargo add databricks-zerobus-ingest-sdk --features arrow-flight

完全なドキュメント、構成オプション、バッチ インジェスト、 generate_proto ツール、プロトコル バッファーの例については、 Rust SDK リポジトリを参照してください。

Java SDK

Java 8 以上が必要です。 SDKは低遅延と効率的なネットワークI/Oを提供し、高スループットの取り込みを実現します。 JSON (最も単純) とプロトコル バッファー (運用環境に推奨) がサポートされています。

Maven:

<dependency>
    <groupId>com.databricks</groupId>
    <artifactId>zerobus-ingest-sdk</artifactId>
    <version>0.2.0</version>
</dependency>

JSON の例:

import com.databricks.zerobus.*;

public class ZerobusClient {

// See "Get your workspace URL and Zerobus Ingest endpoint" for information on obtaining these values.
    private static final String SERVER_ENDPOINT =
        "https://1234567890123456.zerobus.eastus.azuredatabricks.net";
    private static final String DATABRICKS_WORKSPACE_URL =
        "https://adb-1234567890123456.12.azuredatabricks.net";
    private static final String TABLE_NAME = "main.default.air_quality";
    private static final String CLIENT_ID = "your-client-id";
    private static final String CLIENT_SECRET = "your-client-secret";

    public static void main(String[] args) throws Exception {
        ZerobusSdk sdk = new ZerobusSdk(
            SERVER_ENDPOINT,
            DATABRICKS_WORKSPACE_URL
        );

        ZerobusJsonStream stream = sdk.streamBuilder()
            .table(TABLE_NAME)
            .oauth(CLIENT_ID, CLIENT_SECRET)
            .json()
            .build()
            .join();

        try {
            for (int i = 0; i < 100; i++) {
                String record = String.format(
                    "{\"device_name\": \"sensor-%d\", \"temp\": 22, \"humidity\": 55}", i
                );
                stream.ingestRecordOffset(record);
            }
        } finally {
            stream.close();
        }
    }
}

プロトコル バッファ: 型安全な取り込みを行うには、.compiledProto(...)ZerobusProtoStream を使用して streamBuilder() を作成します。 バンドルされた JAR ツールを使用してテーブルからスキーマを生成し、 protocでコンパイルします。

矢印フライト (ベータ): 同じ gRPC 接続を介した Apache Arrow RecordBatch データの列指向またはバッチ指向の取り込みについては、「 Zerobus 取り込みでの矢印フライトの使用」を参照してください。

完全なドキュメント、構成オプション、バッチ インジェスト、プロトコル バッファーの例については、Java SDK リポジトリを参照してください。

Go SDK (ソフトウェア開発キット)

Go のバージョン 1.21 以降が必要です。 SDKはストリーミングの取り込みに高いスループットとパフォーマンスを提供します。 JSON (最も単純) とプロトコル バッファー (運用環境に推奨) がサポートされています。

go get github.com/databricks/zerobus-sdk/go@latest

JSON の例:

わかりやすくするために、ここではエラーは無視されます。 運用コードでは、常にエラーを確認します。

package main

import (
	"fmt"

	zerobus "github.com/databricks/zerobus-sdk/go"
)

// See "Get your workspace URL and Zerobus Ingest endpoint" for information on obtaining these values.
const (
	ServerEndpoint         = "https://1234567890123456.zerobus.eastus.azuredatabricks.net"
	DatabricksWorkspaceURL = "https://adb-1234567890123456.12.azuredatabricks.net"
	TableName              = "main.default.air_quality"
	ClientID               = "your-client-id"
	ClientSecret           = "your-client-secret"
)

func main() {
	sdk, _ := zerobus.NewZerobusSdk(
		ServerEndpoint,
		DatabricksWorkspaceURL,
	)
	defer sdk.Free()

	options := zerobus.DefaultStreamConfigurationOptions()
	options.RecordType = zerobus.RecordTypeJson

	stream, _ := sdk.CreateStream(
		zerobus.TableProperties{
			TableName: TableName,
		},
		ClientID,
		ClientSecret,
		options,
	)
	defer stream.Close()

	_, _ = stream.IngestRecordOffset(`{
		"device_name": "sensor-001",
		"temp": 20,
		"humidity": 60
	}`)

  fmt.Println("Record ingested successfully")

  _ = stream.Close()
  fmt.Println("Stream closed successfully")
}

プロトコル バッファー: タイプ セーフなインジェストの場合は、 RecordTypeProto (既定) でプロトコル バッファーを使用し、テーブルのプロパティに descriptorProto を指定します。 テーブル スキーマに一致する .proto ファイルを作成し、 generate_proto スクリプトを実行して、ファイルをプロジェクトにインポートできるようにします。

矢印フライト (ベータ): 同じ gRPC 接続を介した Apache Arrow RecordBatch データの列指向またはバッチ指向の取り込みについては、「 Zerobus 取り込みでの矢印フライトの使用」を参照してください。

完全なドキュメント、構成オプション、バッチ インジェスト、generate_proto ツール、プロトコル バッファーの例については、 Go SDK リポジトリを参照してください。

C++ SDK

Important

C++ SDK は ベータ版です

C++17 以降が必要です。 SDKはネイティブのgRPCストリーミング、OAuth、RAII C++インターフェースを通じた自動復旧を提供します。 これは、運用ワークロード用の単純なセットアップとプロトコル バッファーの JSON をサポートします。

SDK は事前構築済みのプラットフォームごとのリリース バンドルとして付属しているため、それを使用するために Rust ツールチェーンは必要ありません。 リリース ページからプラットフォームのバンドル (macOS、Linux (musl を含む)、またはWindows) をダウンロードし、それを抽出して、バンドルされた FFI アーカイブで CMake をポイントします。 アーカイブの名前は、macOS と Linux では libzerobus_ffi.a、Windowsではzerobus_ffi.lib

# macOS and Linux
cmake -S cpp -B build \
  -DZEROBUS_FFI_LIBRARY="$PWD/lib/libzerobus_ffi.a" \
  -DZEROBUS_FFI_HEADER_DIR="$PWD/lib"
cmake --build build -j

Windows (PowerShell) で、代わりに .lib アーカイブをポイントします。

cmake -S cpp -B build `
  -DZEROBUS_FFI_LIBRARY="$PWD/lib/zerobus_ffi.lib" `
  -DZEROBUS_FFI_HEADER_DIR="$PWD/lib"
cmake --build build -j

独自の CMake プロジェクトのソース チェックアウトから SDK をビルドするには、それをサブディレクトリとして追加し、ターゲットをリンクします。 FetchContent を使用して、構成時点でそれを取得することもできます。 これにより Rust ソースから FFI がビルドされるため、Rust ツールチェーンが必要です。

add_subdirectory(path/to/zerobus-sdk/cpp)
target_link_libraries(your_app PRIVATE zerobus::zerobus)

代わりに add_subdirectory 経由で事前ビルド済みバンドルを利用するには、まず FFI パスを設定して、CMake が存在しない Rust ソースからビルドしようとするのではなく、バンドルされたアーカイブをリンクするようにします。 Windowsでzerobus_ffi.libを使用します。

set(ZEROBUS_FFI_LIBRARY "/path/to/bundle/lib/libzerobus_ffi.a")
set(ZEROBUS_FFI_HEADER_DIR "/path/to/bundle/lib")
add_subdirectory(path/to/zerobus-sdk/cpp zerobus-cpp)
target_link_libraries(your_app PRIVATE zerobus::zerobus)

JSON の例:

インジェストは非同期でパイプライン化されます。 ingest_* メソッドはレコードをキューに入れ、すぐに戻ります。 各レコードの後で待機するのではなく、バッチをキューに入れ、 flush() を 1 回呼び出します。

#include "zerobus/zerobus.hpp"
#include <string>
#include <vector>

int main() {
  // See "Get your workspace URL and Zerobus Ingest endpoint" for information on obtaining these values.
  const std::string SERVER_ENDPOINT = "https://1234567890123456.zerobus.eastus.azuredatabricks.net";
  const std::string DATABRICKS_WORKSPACE_URL = "https://adb-1234567890123456.12.azuredatabricks.net";
  const std::string TABLE_NAME = "main.default.air_quality";
  const std::string CLIENT_ID = "your-client-id";
  const std::string CLIENT_SECRET = "your-client-secret";

  zerobus::Sdk sdk = zerobus::Sdk::builder()
                         .endpoint(SERVER_ENDPOINT)
                         .unity_catalog_url(DATABRICKS_WORKSPACE_URL)
                         .application_name("my-app")
                         .build();

  zerobus::TableProperties table;
  table.table_name = TABLE_NAME;   // empty descriptor => JSON stream

  zerobus::StreamOptions options;
  options.record_type = zerobus::RecordType::Json;

  zerobus::Stream stream =
      sdk.create_stream(table, CLIENT_ID, CLIENT_SECRET, options);

  std::vector<std::string> batch = {
      R"({"device_name": "sensor-001", "temp": 20, "humidity": 60})",
      R"({"device_name": "sensor-002", "temp": 22, "humidity": 55})",
  };
  stream.ingest_json_records(batch);   // queue the batch — no per-record wait
  stream.flush();                      // wait once for all acks
  stream.close();

  return 0;
}

失敗ごとに zerobus::ZerobusException をスローします。そこにはメッセージとis_retryable() フラグが伴います。 ブロックせずに連続ストリームの持続性を追跡するには、AckCallbackを使用してStreamOptions::ack_callbackを登録します。 コールバックはバックグラウンド スレッドでシリアル化されて実行され、 noexceptする必要があります。 完全なスレッド処理、ドレイン ポリシー、および有効期間コントラクトについては、 C++ SDK のドキュメント を参照してください。

タイプ セーフなインジェストの場合は、次の 2 つの方法のいずれかでプロトコル バッファーを使用できます。

  • ProtoSchema::from_uc_json()を使用して Unity カタログからスキーマを生成します。 これにより、テーブルのメタデータから直接記述子と JSON 間エンコーダーがビルドされるため、 .proto ファイルや protocは必要ありません。

    1. Get a table API () からテーブルのメタデータ JSON をGET /api/2.1/unity-catalog/tables/{full_name}します。 サービス プリンシパルは、テーブルに SELECT する必要があります。
    2. メタデータを ProtoSchema::from_uc_json() に渡して、記述子とエンコーダーを構築します。
    3. TableProperties::descriptor_proto設定し、ingest_proto_records()を使用して取り込む。
  • チェックイン済み .proto は、コンパイル時の型指定のための protoc とともにコンパイルします。

実行可能なチュートリアルについては、 プロトコル バッファーの例を参照してください。

同じ gRPC 接続を介して Apache Arrow レコードバッチを列指向またはバッチ指向で取り込む方法については、Use Arrow Flight with Zerobus Ingestを参照してください。

完全なドキュメント、構成オプション、バッチ インジェスト、プロトコル バッファーの例については、 C++ SDK リポジトリを参照してください。

C# SDK

Important

C# / .NET SDKはベータ版です。 Databricks.Zerobusパッケージはプレリリースです。

.NET 8.0以上の条件が必要です。 SDKはネイティブのgRPCストリーミング、OAuth、自動リカバリーを提供します。 これは、運用ワークロード用の単純なセットアップとプロトコル バッファーの JSON をサポートします。 Arrow FlightはC# SDKには含まれていません。

Databricks.Zerobus パッケージをプロジェクトに追加します。

dotnet add package Databricks.Zerobus

JSON の例:

using Databricks.Zerobus;

// See "Get your workspace URL and Zerobus Ingest endpoint" for information on obtaining these values.
const string SERVER_ENDPOINT = "https://1234567890123456.zerobus.eastus.azuredatabricks.net";
const string DATABRICKS_WORKSPACE_URL = "https://adb-1234567890123456.12.azuredatabricks.net";
const string TABLE_NAME = "main.default.air_quality";
const string CLIENT_ID = "your-client-id";
const string CLIENT_SECRET = "your-client-secret";

using var sdk = ZerobusSdk.CreateBuilder()
    .Endpoint(SERVER_ENDPOINT)
    .UnityCatalogUrl(DATABRICKS_WORKSPACE_URL)
    .Build();

using var stream = sdk.CreateJsonStream(TABLE_NAME, CLIENT_ID, CLIENT_SECRET);

long offset = stream.IngestRecord(
    """{"device_name": "sensor-1", "temp": 22, "humidity": 55}""");
stream.WaitForOffset(offset);
stream.Close();

IngestRecord はレコードのオフセットを返し、WaitForOffset はそのレコードが永続化されるまでブロックします。 バッチを取り込むには IngestRecordsを使い、これはレコード配列を受け取り、最後のオフセットを返します。 オフセットでブロックすることは省略可能です。 メッセージのブロックと確認応答を参照してください。

Protocol Buffers: 型安全に取り込むには、sdk.CreateProtoStream(TABLE_NAME, descriptorProto, CLIENT_ID, CLIENT_SECRET) でストリームを作成します。ここで、descriptorProto はコンパイル済みメッセージをシリアル化した DescriptorProto バイトであり、stream.IngestRecord(protoBytes) を使用して取り込みます。

完全なドキュメント、設定オプション、プロトコルバッファの例については 、C# SDKリポジトリをご覧ください。

TypeScript SDK

Node.js 16 以上が必要です。 SDKはJavaScript Promisesを通じて非同期サポートと高性能を提供します。 JSON (最も単純) とプロトコル バッファー (運用環境に推奨) がサポートされています。

npm install @databricks/zerobus-ingest-sdk

JSON の例:

import { ZerobusSdk, RecordType } from '@databricks/zerobus-ingest-sdk';

// See "Get your workspace URL and Zerobus Ingest endpoint" for information on obtaining these values.
const SERVER_ENDPOINT = 'https://1234567890123456.zerobus.eastus.azuredatabricks.net';
const DATABRICKS_WORKSPACE_URL = 'https://adb-1234567890123456.12.azuredatabricks.net';
const TABLE_NAME = 'main.default.air_quality';
const CLIENT_ID = 'your-client-id';
const CLIENT_SECRET = 'your-client-secret';

const sdk = new ZerobusSdk(SERVER_ENDPOINT, DATABRICKS_WORKSPACE_URL);

const stream = await sdk.createStream({ tableName: TABLE_NAME }, CLIENT_ID, CLIENT_SECRET, {
  recordType: RecordType.Json,
});

try {
  for (let i = 0; i < 100; i++) {
    const record = { device_name: `sensor-${i}`, temp: 22, humidity: 55 };
    await stream.ingestRecordOffset(record);
  }
} finally {
  await stream.close();
}

プロトコル バッファー: タイプ セーフなインジェストの場合は、 RecordType.Proto (既定) でプロトコル バッファーを使用し、テーブルのプロパティに descriptorProto を指定します。

矢印フライト (ベータ): 同じ gRPC 接続を介した Apache Arrow RecordBatch データの列指向またはバッチ指向の取り込みについては、「 Zerobus 取り込みでの矢印フライトの使用」を参照してください。

完全なドキュメント、構成オプション、バッチ インジェスト、プロトコル バッファーの例については、 TypeScript SDK リポジトリを参照してください。

REST API

REST API を使用すると、 エンドポイントに HTTP POST 要求を送信することで、1 つのレコードを取り込めます。 レコード自体は要求本文に含まれており、JSON 形式である必要があります。

この例では、CURL を使用して REST API を使用して Zerobus Ingest にデータをプッシュする方法について説明します。

ヘッダー

要求を正しく認証して書式設定するには、2 つの特定の HTTP ヘッダーが必要です。

  • Content-Type: application/json
    • コンテンツ タイプを指定するための必須フィールド。 現在、サポートされている唯一のメッセージ形式は JSON です。
  • 認証: ベアラー <トークン>
    • <token>を、後で提供される curl コマンドを使用してフェッチした OAuth トークンに置き換えます。

OAuth トークンのフェッチ: これらのトークンは 1 時間ごとに期限切れになり、更新する必要があります。 それらを更新するには、OAuth トークンを再取得します。

次のパラメーターを入力します。

  • $CATALOG$SCHEMA$TABLE$WORKSPACE_ID$WORKSPACE_URL
  • $DATABRICKS_CLIENT_ID および $DATABRICKS_CLIENT_SECRET
    • これら 2 つのパラメーターは、作成したサービス プリンシパルに対応します。
authorization_details=$(cat <<EOF
[{
  "type": "unity_catalog_privileges",
  "privileges": ["USE CATALOG"],
  "object_type": "CATALOG",
  "object_full_path": "$CATALOG"
},
{
  "type": "unity_catalog_privileges",
  "privileges": ["USE SCHEMA"],
  "object_type": "SCHEMA",
  "object_full_path": "$CATALOG.$SCHEMA"
},
{
  "type": "unity_catalog_privileges",
  "privileges": ["SELECT", "MODIFY"],
  "object_type": "TABLE",
  "object_full_path": "$CATALOG.$SCHEMA.$TABLE"
}]
EOF
)

export OAUTH_TOKEN=$(curl -X POST \
  -u "$DATABRICKS_CLIENT_ID:$DATABRICKS_CLIENT_SECRET" \
  -d "grant_type=client_credentials" \
  -d "scope=all-apis" \
  -d "resource=api://databricks/workspaces/$WORKSPACE_ID/zerobusDirectWriteApi" \
  --data-urlencode "authorization_details=$authorization_details" \
  "$WORKSPACE_URL/oidc/v1/token" | jq -r '.access_token')

記録の取り込み:

次のパラメーターを入力します。

要求本文は、JSON オブジェクトの一覧である必要があります。

curl -X POST \
  "$ZEROBUS_ENDPOINT/zerobus/v1/tables/$CATALOG.$SCHEMA.$TABLE/insert" \
  -H "Content-Type: application/json" \
  -H "Authorization: Bearer $OAUTH_TOKEN" \
  -d '[{ "device_name": "device_num_1", "temp": 28, "humidity": 60 },
       { "device_name": "device_num_1", "temp": 28, "humidity": 60 }]'

すべての情報が正しく入力されている場合は、HTTP 状態コードが 200 の空の JSON 応答を受け取る必要があります。

エラーを処理する

上記の例は幸せな道を示しています。 本番環境では、取り込み処理をエラーハンドリングで囲んでください。 SDKはネットワーク問題などの一時的なエラーを内蔵のリカバリーを通じて自動で再試行します。 無効な認証情報やテーブルの欠落など、回復できない障害は以下の ZerobusExceptionとして現れます。

from zerobus.sdk.shared import ZerobusException

try:
    stream.ingest_record_offset(record)
except ZerobusException as e:
    # Handle the failure: log it, fix the cause, recover on a new stream, or stop.
    ...

SDKは一時的な障害からも自動で回復し、ストリームが永久に故障した際に未確認のレコードを救出することも可能です。 レジリエントクライアントパターンおよび完全なエラー参照については、 Recovery and retry patterns および Zerobus Ingestエラー処理を参照してください。

次のステップ