ゼロバスインジェストで矢印フライトを使用する

Important

Arrow Flight の取り込みはベータ段階です。

Arrow Flightの取り込みは、 Apache ArrowRecordBatch データをすべての行をJSONやプロトコルバッファ(protobuf)に変換する代わりに、直接 Zerobus Ingest に送ることができます。 これはZerobus SDKの3番目のレコードフォーマットオプションであり、JSONやprotobufと並んで、同じgRPC接続上で動作します。 同じ Zerobus エンドポイント、同じ OAuth フロー、および同じ x-databricks-zerobus-table-name ヘッダー規則を使用します。 ワイヤプロトコルはArrow FlightDoPutであり、gRPC 経由で Arrow IPC メッセージを伝送します。

矢印フライトを使用する場合

矢印フライトは、次のシナリオに最適です。

  • アプリケーションはすでに、pyarrow.Tablepyarrow.RecordBatch(Python)、arrow_array::RecordBatch クレート(Rust)の 、または VectorSchemaRoot(Java)などの Arrow データを生成しています。 PolarsDataFusionのようなArrow上で構築されたDataFrameライブラリは、この道筋に自然に適合します。
  • 一度に 1 つのレコードを送信するのではなく、行をバッチで取り込みます。
  • スキーマはワイド、数値負荷、または分析指向であり、行単位のシリアル化によって CPU オーバーヘッドが顕著になります。
  • 短い間隔でデータを集計し、1 列形式のバッチとして送信するコレクターまたはゲートウェイを構築しています。

通常、Arrow Flight は、疎な、1行ずつのトラフィックには最適な選択肢ではありません。 その場合、SDKのgRPCパスを経由するJSONやprotobufの方が通常より簡単です。 「 インターフェイスの選択」を参照してください。

インジェスト モデルのしくみ

Arrow Flight インジェストでは、1 つのストリームが 1 つのターゲット テーブルに書き込まれます。 データを取り込むには、次の順序に従います。

  1. 宛先の Delta テーブル スキーマと一致する Arrow スキーマを定義します。
  2. そのテーブルのゼロバス矢印ストリームを開きます。
  3. RecordBatch (またはTable) ペイロードを送信します。
  4. 最後のオフセットまたは呼び出し flush() を待って持続性を確認します。
  5. ストリームを閉じます。

Zerobus SDK を使用する場合、SDK は低レベルの Arrow Flight ワイヤの詳細を自動的に処理します。 ArrowのデータをIPC形式にシリアル化し、大きなバッチを小さなトランスポートメッセージに分割し、サーバーは個別にそれを認識します。

Arrow Flight は、論理的な全バッチに対して全か無かの耐久度を提供しているわけではありません。 Arrowのバッチは非常に大きくなり得ますが、SDKはそれを個別のトランスポートメッセージに分割し、受信時に確認応答されるため、途中で障害が発生した場合でも大きなバッチは部分的に耐久性を持つことがあります。 これは、原子的にコミットし、10MBのメッセージサイズで制限されるJSONやprotobufバッチとは異なります。 アローフライトのバッチは例外を参照してください。

論理オフセット抽象化はこのチャンク化の上にも依然として有効です。 ingest_batch() 提出したバッチに対して単一の論理オフセットを返し、そのオフセットの wait_for_offset() 処理はバッチを構成するすべてのトランスポートメッセージが確認された後にのみ完了します。 (メソッド名はPython SDKから取られており、他のSDKはJavaでingestBatchwaitForOffsetなど同等のメソッドを公開しています。)

Protobufスキーマルールと同様に、ストリームに渡すスキーマはターゲットのデルタテーブルに収まらなければなりません。最低でもすべての非ノールカラムを含む必要があります。 スキーマでは、Delta テーブルに存在する null 許容列を省略できますが (これは非破壊的スキーマの変更として扱われます)、その他の不一致は拒否されます。 各アローフィールドのタイプは、そのデルタ列と互換性がある必要があります。 サポートされるデルタ型については、「 サポートデータ型」を参照してください。

SDKは大規模なバッチをトランスポートメッセージに分割するため、ArrowのバッチはJSONやプロトコルバッファのバッチのように10MBのメッセージサイズ制限に縛られません。 Arrow Flightは同じgRPCトランスポート上で動作するため、スループット、レイテンシ、クォータ特性が適用され、これらはより高いワークロードに応じてスケールします。 「Zerobus Ingest クォータ」を参照してください。

クライアントを書き込む

以下の例では、air_quality の例で使用したものと同じ テーブルに対して、Arrow Flight ストリームを開始します。 簡潔にするためにPythonと Rust に表示されますが、すべての Zerobus SDK で同じビルダー、構成オプション、および呼び出しシーケンスを使用できます。 言語の構文を調整し、言語固有の矢印の種類については SDK リポジトリを参照してください。

Python SDK

Python SDK は、ストリームの作成時に pyarrow.Schema、取り込み呼び出しごとに pyarrow.RecordBatch または pyarrow.Table を受け入れます。

pip install "databricks-zerobus-ingest-sdk[arrow]" pyarrow
import pyarrow as pa

from zerobus.sdk.sync import ZerobusSdk

# See "Get your workspace URL and Zerobus Ingest endpoint" in zerobus-ingest.md.
SERVER_ENDPOINT = "https://1234567890123456.zerobus.us-west-2.cloud.databricks.com"
DATABRICKS_WORKSPACE_URL = "https://dbc-a1b2c3d4-e5f6.cloud.databricks.com"
TABLE_NAME = "main.default.air_quality"
CLIENT_ID = "your-client-id"
CLIENT_SECRET = "your-client-secret"

schema = pa.schema(
    [
        ("device_name", pa.large_utf8()),
        ("temp", pa.int32()),
        ("humidity", pa.int64()),
    ]
)

sdk = ZerobusSdk(SERVER_ENDPOINT, DATABRICKS_WORKSPACE_URL)

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

row_count = 1_000
batch = pa.record_batch(
    {
        "device_name": [f"sensor-{i}" for i in range(row_count)],
        "temp": [20 + (i % 5) for i in range(row_count)],
        "humidity": [55 + (i % 10) for i in range(row_count)],
    },
    schema=schema,
)

try:
    offset = stream.ingest_batch(batch)

    # Optional: block until the batch is durably written
    stream.wait_for_offset(offset)
finally:
    stream.close()

stream.ingest_batch() は、 pyarrow.Tableも受け入れます。 SDK は、送信する前に内部的に 1 つの RecordBatch に変換します。 各呼び出しは論理オフセットを返します。 オフセットでブロックすることは省略可能です。 待つタイミングや確認応答の仕組みについては、「 メッセージブロックと確認応答」を参照してください。

Rust SDK

Rust SDK は、stream_builder() Cargo 機能の背後にあるarrow-flight API を介して Arrow Flight を公開します。 コンパイル時に RecordBatch と配列の型が一致するように、SDK と同じ Arrow メジャー バージョンを使用します。

cargo add databricks-zerobus-ingest-sdk --features arrow-flight
cargo add arrow-array
cargo add arrow-schema
cargo add tokio --features macros,rt-multi-thread
use std::sync::Arc;

use arrow_array::{Int32Array, Int64Array, LargeStringArray, RecordBatch};
use arrow_schema::{DataType, Field, Schema as ArrowSchema};
use databricks_zerobus_ingest_sdk::ZerobusSdk;

const SERVER_ENDPOINT: &str = "https://1234567890123456.zerobus.us-west-2.cloud.databricks.com";
const DATABRICKS_WORKSPACE_URL: &str = "https://dbc-a1b2c3d4-e5f6.cloud.databricks.com";
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 std::error::Error>> {
    let schema = Arc::new(ArrowSchema::new(vec![
        Field::new("device_name", DataType::LargeUtf8, false),
        Field::new("temp", DataType::Int32, false),
        Field::new("humidity", DataType::Int64, false),
    ]));

    let sdk = ZerobusSdk::builder()
        .endpoint(SERVER_ENDPOINT)
        .unity_catalog_url(DATABRICKS_WORKSPACE_URL)
        .build()?;

    let mut stream = sdk
        .stream_builder()
        .table(TABLE_NAME)
        .oauth(CLIENT_ID, CLIENT_SECRET)
        .arrow(Arc::clone(&schema))
        .build_arrow()
        .await?;

    let row_count: i32 = 1_000;
    let batch = RecordBatch::try_new(
        Arc::clone(&schema),
        vec![
            Arc::new(LargeStringArray::from(
                (0..row_count)
                    .map(|i| format!("sensor-{i}"))
                    .collect::<Vec<_>>(),
            )),
            Arc::new(Int32Array::from(
                (0..row_count).map(|i| 20 + (i % 5)).collect::<Vec<_>>(),
            )),
            Arc::new(Int64Array::from(
                (0..row_count)
                    .map(|i| 55 + (i % 10) as i64)
                    .collect::<Vec<_>>(),
            )),
        ],
    )?;

    let offset = stream.ingest_batch(batch).await?;

    // Optional: block until the batch is durably written
    stream.wait_for_offset(offset).await?;
    stream.close().await?;

    Ok(())
}

ビルダーは、 .arrow(schema) で方向フライト形式を選択し、ストリームを .build_arrow() で最終処理し、 ZerobusArrowStreamを返します。 JSONとprotobufは引き続き .json() / .compiled_proto(...).build()を使用しています。

VARIANTカラムの取り込み

Apache Arrowにはネイティブの VARIANT タイプはありません。 Arrow Flight上の VARIANT カラムに取り込みたい場合は、そのカラムのバックイング metadatavalue フィールドを2つの LargeBinary カラムの構造体として作成し、その構造体を RecordBatchに含めてください。 gRPC SDKやRESTでは、代わりにVariant値をJSONでエンコードされた文字列として渡します。 「サポートされるデータ型」を参照してください。

以下のRustの例は、JSON行から VARIANT 構造体列を作成し、それを取り込みます。

fn variant_struct(json_rows: &[&str]) -> ArrayRef {
    let mut metas: Vec<Vec<u8>> = Vec::new();
    let mut vals: Vec<Vec<u8>> = Vec::new();
    for json in json_rows {
        let mut vb = VariantBuilder::new();
        vb.append_json(json).expect("invalid JSON for variant");
        let (metadata, value) = vb.finish();
        metas.push(metadata);
        vals.push(value);
    }
    let fields = Fields::from(vec![
        Field::new("metadata", DataType::LargeBinary, false),
        Field::new("value", DataType::LargeBinary, false),
    ]);
    let meta_arr = Arc::new(LargeBinaryArray::from_iter_values(metas)) as ArrayRef;
    let val_arr = Arc::new(LargeBinaryArray::from_iter_values(vals)) as ArrayRef;
    Arc::new(StructArray::try_new(fields, vec![meta_arr, val_arr], None).expect("variant struct"))
}

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    let client_id = std::env::var("DATABRICKS_CLIENT_ID")?;
    let client_secret = std::env::var("DATABRICKS_CLIENT_SECRET")?;

    let variant_type = DataType::Struct(Fields::from(vec![
        Field::new("metadata", DataType::LargeBinary, false),
        Field::new("value", DataType::LargeBinary, false),
    ]));
    let schema = Arc::new(ArrowSchema::new(vec![
        Field::new("id", DataType::Int32, true),
        Field::new("payload", variant_type, true),
    ]));

    let sdk = ZerobusSdk::builder()
        .endpoint(ENDPOINT)
        .unity_catalog_url(UC_URL)
        .build()?;

    let mut stream = sdk
        .stream_builder()
        .table(TABLE)
        .oauth(&client_id, &client_secret)
        .arrow(schema.clone())
        .ipc_compression(None)
        .build_arrow()
        .await?;

    let ids = Int32Array::from(vec![1, 2, 3]);
    let payload = variant_struct(&[
        r#"{"user":"alice","tags":[1,2,3]}"#,
        r#""just a string""#,
        r#"{"nested":{"a":true,"b":null,"c":3.14}}"#,
    ]);
    let batch = RecordBatch::try_new(schema.clone(), vec![Arc::new(ids) as ArrayRef, payload])?;

    let offset = stream.ingest_batch(batch).await?;
    stream.flush().await?;
    stream.close().await?;
    Ok(())
}

この例はRustで書かれています。 他の言語での同等の使用については、 Zerobus SDKリポジトリを参照してください。

IPC圧縮

既定では、Arrow IPC ペイロードは圧縮されずに送信されます。 必要に応じて、2 つのコーデックのいずれかを使用してワイヤ上で圧縮できます。

  • LZ4_FRAME:高速、低 CPU オーバーヘッド、適度な圧縮率。 これは、クライアントが CPU 制約を受けているが、ネットワーク上のバイト数を減らしたい場合に優先します。
  • ZSTD: 圧縮率が高く、バッチあたりの CPU 使用率が高くなります。 クライアントが追加の CPU コストを吸収できる場合は常に有効にします。

圧縮により、ネットワーク上のバイト数が削減されますが、クライアントの CPU コストが増えます。 ペイロードを小さくすると、ネットワークのボトルネックを回避し、ネットワーク コストを削減できます。

Python SDK で、ipc_compressionArrowStreamConfigurationOptions フィールドを設定します。

from zerobus.sdk.shared.arrow import IPCCompression, ArrowStreamConfigurationOptions

options = ArrowStreamConfigurationOptions(ipc_compression=IPCCompression.ZSTD)

Rust SDK で、ビルダーに設定します。 CompressionType列挙型はarrow-ipcクレート内に存在するため、依存関係として追加します。

cargo add arrow-ipc
use arrow_ipc::CompressionType;

let stream = sdk
    .stream_builder()
    .table(TABLE_NAME)
    .oauth(CLIENT_ID, CLIENT_SECRET)
    .arrow(schema)
    .ipc_compression(Some(CompressionType::ZSTD))
    .build_arrow()
    .await?;

ベスト プラクティス

次のガイドラインに従って、Arrow Flight のインジェストから最高のパフォーマンスと信頼性を得ます。

  • バッチごとに新しいストリームを開くのではなく、多数のバッチに対してストリームを再利用します。 ストリームの作成には、多数のバッチにわたってストリームを再利用することで償却できる大きなオーバーヘッドが伴います。
  • バッチごとに複数の行を送信します。 呼び出しごとに 1 行ではなく、自然なアプリケーション サイズのバッチから開始します。 一度に 1 行ずつ送信すると機能しますが、Arrow を使用するパフォーマンス上の利点のほとんどは否定されます。
  • 制御されたチェックポイントで flush() を呼び出します。 これにより、バッチのグループに対して明確な持続性境界が提供され、1 つ 1 つでブロックされません。
  • IPC 圧縮を有効にしてスループットを向上させます。 ZSTD クライアントに予備CPUがある場合、ほとんどのワークロードで推奨されます。 クライアントがCPU制約がある場合は LZ4_FRAME 圧縮を使うか、圧縮を行わないでください。
  • データ生成元がすでに列指向形式である場合は、Arrow Flight を使用します。 もしソースデータが自然に行単位で小規模であれば、Zerobus IngestとJSONやprotobufを使う方が簡単になることが多いです。 「 Zerobus Ingestの使用」を参照してください。

エラー処理と回復

Arrow Flight ストリームは、Zerobus Ingest の他の部分と同じ gRPC エラーカテゴリを使用します。 エラー コード、再試行ガイダンス、および完全なクライアントとサーバーの分類については、「 Zerobus 取り込みエラー処理」を参照してください。

自動復旧 (既定値) を使用して SDK を構成すると、一時的な障害時に未確認のバッチが透過的に再接続され、再生されます。 ストリームが閉じた後、サーバーが受信したが、まだ確認していないバッチを取得できます。 これは、ストリームが正常に閉じられたか、回復不能な障害が原因で閉じられたかに関係なく適用されます。 Python SDK で次の手順を実行します。

# Retry unacked_batches against a freshly created stream
if stream.is_closed:
    unacked_batches = stream.get_unacked_batches()

Rust SDK で、 stream.get_unacked_batches().await? を呼び出して、未確認のバッチを取得して再試行します。

その他のリソース

  • Zerobus Ingestを活用してください:まだZerobus Ingestを設定していない場合は、ワークスペースURLの見つけ方、ターゲットのDeltaテーブルの作成、サービスプリンシパルの設定手順をこちらからご覧ください。 これらの手順は、すべてのレコード形式で共有されます。
  • Zerobus Ingest クォータ:本番環境にデプロイする前にデフォルトの Zerobus クォータを確認しましょう。 同じスループットとレイテンシ特性がArrow Flightにも適用され、より高いワークロードに対応するためにスケール可能です。
  • Zerobus 取り込みエラー処理: gRPC エラー コードの完全な一覧と、クライアントの推奨される再試行と回復の動作については、このページを参照してください