Verwenden Sie Arrow Flight mit Zerobus Ingest

Mit Arrow Flight Ingestion kannst du Apache Arrow-DatenRecordBatch direkt an Zerobus Ingest senden, anstatt jede Zeile zuerst in JSON oder Protocol Buffers (Protobuf) umzuwandeln. Es handelt sich um eine dritte Option für das Datensatzformat in Zerobus-SDKs, die Arrow Flight unterstützen, zusammen mit JSON und Protobuf, und läuft über dieselbe gRPC-Verbindung. Er verwendet den gleichen Zerobus-Endpunkt, denselben OAuth-Fluss und dieselbe x-databricks-zerobus-table-name Headerkonvention. Das Übertragungsprotokoll ist Arrow FlightDoPut, das Arrow-IPC-Nachrichten über gRPC überträgt.

Gründe für die Verwendung von Arrow Flight

Arrow Flight eignet sich am besten in die folgenden Szenarien:

  • Ihre Anwendung erzeugt bereits Arrow-Daten, wie pyarrow.Table oder pyarrow.RecordBatch (Python), arrow_array::RecordBatch aus den arrow-rs-Crates (Rust) oder VectorSchemaRoot (Java). DataFrame-Bibliotheken, die auf Arrow basieren, wie Polars oder DataFusion, fügen sich natürlich in diesen Weg ein.
  • Sie erfassen Zeilen in Batches, anstatt jeweils einen Datensatz zu senden.
  • Ihr Schema ist breit, numerisch schwer oder analyseorientiert, wobei die Serialisierung von Zeilen nach Zeile einen spürbaren CPU-Aufwand hinzufügt.
  • Sie erstellen Sammler oder Gateways, die Daten über einen kurzen Zeitraum aggregieren und sie dann als Batch im Spaltenformat senden.

Arrow Flight ist in der Regel nicht die beste Wahl für sporadischen Datenverkehr mit jeweils nur einer Zeile. In diesen Fällen sind JSON oder Protobuf über den SDK gRPC-Pfad typischerweise einfacher. Siehe Auswählen einer Schnittstelle.

Wie das Erfassungsmodell funktioniert

Bei der Arrow-Flight-Erfassung schreibt ein Stream in genau eine Zieltabelle. Gehen Sie wie folgt vor, um Daten aufzunehmen:

  1. Definieren Sie ein Pfeilschema, das dem Ziel-Delta-Tabellenschema entspricht.
  2. Öffnen Sie einen Zerobus-Pfeilstream für diese Tabelle.
  3. Senden RecordBatch (oder Table) Nutzdaten.
  4. Warten Sie auf den letzten Offset oder rufen Sie flush() auf, um die Dauerhaftigkeit zu bestätigen.
  5. Schließen Sie den Datenstrom.

Wenn Sie ein Zerobus-SDK verwenden, übernimmt das SDK die niedrigstufigen Details des Arrow-Flight-Protokolls für Sie. Es serialisiert Ihre Arrow-Daten im IPC-Format und teilt eine große Charge automatisch in kleinere, geordnete Flugnachrichten auf.

Der Server meldet kumulativen Fortschritt, sobald die Datensätze dauerhaft werden. Ein großer logischer Batch kann daher teilweise langlebig sein, wenn während des Versands ein Fehler auftritt. ingest_batch() liefert für den übermittelten Batch weiterhin einen logischen Offset zurück, und das Warten auf diesen Offset bestätigt, dass alle seine Datensätze dauerhaft gespeichert sind. Siehe Arrow-Flight-Batches sind die Ausnahme.

Zerobus ordnet Arrow-Felder Delta-Spalten anhand ihrer Namen zu. Das Arrow-Schema muss alle erforderlichen (nicht nullablen) Delta-Spalten enthalten. Du kannst nullable Spalten weglassen, die Zerobus als NULLschreibt. Fügen Sie keine fehlenden Felder in der Zieltabelle ein. Enthaltene Felder müssen der relativen Reihenfolge des Delta-Schemas folgen, dessen Nullfähigkeit anpassen und den Typ der Zielspalte verwenden. Details finden Sie unter Schema-Abgleichsregeln.

Ein Arrow-Batch unterliegt nicht der 10 MB gRPC-Nachrichtengrößenbeschränkung. Allerdings muss jede einzelne Zeile innerhalb eines RecordBatch Rasters innerhalb des 10-MB-Limits passen. Das SDK teilt größere Chargen automatisch in mehrere Wire-Nachrichten auf, kann jedoch keine übergroße Zeile aufteilen. Siehe Zerobus Ingest-Quoten.

Entwicklung eines Clients

Die untenstehenden Beispiele verwenden die Python- und Rust-SDKs. Für andere Sprachen, die Arrow Flight unterstützen, siehe das Zerobus SDK-Repository.

Python SDK

Das Python-SDK akzeptiert bei der Stream-Erstellung ein pyarrow.Schema sowie für jeden Ingest-Aufruf ein pyarrow.RecordBatch oder pyarrow.Table.

pip install "databricks-zerobus-ingest-sdk[arrow]"
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)

try:
    for start in range(0, 10_000, 1_000):
        end = start + 1_000
        batch = pa.record_batch(
            {
                "device_name": [f"sensor-{i}" for i in range(start, end)],
                "temp": [20 + (i % 5) for i in range(start, end)],
                "humidity": [55 + (i % 10) for i in range(start, end)],
            },
            schema=schema,
        )
        stream.ingest_batch(batch)
    stream.flush()
finally:
    stream.close()

stream.ingest_batch() akzeptiert auch ein pyarrow.Table. Das SDK wandelt es intern in ein einzelnes RecordBatch um, bevor es gesendet wird. Jeder Aufruf gibt einen logischen Offset zurück. Das Beispiel erstellt mehrere Batches und ruft flush() einmal auf, um zu bestätigen, dass alle ausstehenden Batches dauerhaft gespeichert sind. Verwenden Sie wait_for_offset(), wenn Sie einen bestimmten Batch bestätigen müssen, bevor Sie fortfahren; das Warten auf den letzten Offset bestätigt auch alle vorherigen Offsets. Für wann man warten sollte und wie Bestätigung funktioniert, siehe Nachrichtenblockierung und Bestätigung.

Rust SDK

Das Rust SDK macht Arrow Flight über die stream_builder() API hinter dem arrow-flight Cargo-Feature verfügbar. Verwenden Sie dieselbe Hauptversion von Arrow wie für das SDK, damit die Typen RecordBatch und Array beim Kompilieren übereinstimmen.

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, true),
        Field::new("temp", DataType::Int32, true),
        Field::new("humidity", DataType::Int64, true),
    ]));

    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?;

    for start in (0_i32..10_000).step_by(1_000) {
        let end = start + 1_000;
        let batch = RecordBatch::try_new(
            Arc::clone(&schema),
            vec![
                Arc::new(LargeStringArray::from(
                    (start..end)
                        .map(|i| format!("sensor-{i}"))
                        .collect::<Vec<_>>(),
                )),
                Arc::new(Int32Array::from(
                    (start..end).map(|i| 20 + (i % 5)).collect::<Vec<_>>(),
                )),
                Arc::new(Int64Array::from(
                    (start..end)
                        .map(|i| 55 + (i % 10) as i64)
                        .collect::<Vec<_>>(),
                )),
            ],
        )?;
        stream.ingest_batch(batch).await?;
    }
    stream.flush().await?;
    stream.close().await?;

    Ok(())
}

Der Builder wählt das Arrow-Flight-Format mit .arrow(schema) und schließt den Stream mit .build_arrow() ab, wodurch ein ZerobusArrowStream zurückgegeben wird. JSON und protobuf verwenden weiterhin .json() / .compiled_proto(...) und .build().

Erfassung von VARIANT-Spalten

Apache Arrow hat keinen eigenen VARIANT Typ. Um Daten über Arrow Flight in eine VARIANT-Spalte einzuspeisen, erstellen Sie die zugrunde liegenden metadata- und value-Felder der Spalte als eine Struktur aus zwei LargeBinary-Spalten und fügen diese Struktur dann in Ihr RecordBatch ein. Über die gRPC-SDKs und REST übergeben Sie stattdessen einen Variant-Wert als JSON-codierte Zeichenkette. Weitere Informationen finden Sie unter Unterstützte Datentypen.

Das folgende Rust-Beispiel erstellt eine VARIANT Struct-Spalte aus JSON-Zeilen und nimmt sie ein:

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(())
}

Dieses Beispiel ist in Rust. Für eine gleichwertige Verwendung in anderen Sprachen siehe das Zerobus SDK-Repository.

IPC-Komprimierung

Standardmäßig werden Arrow-IPC-Payloads unkomprimiert gesendet. Sie können sie optional mit einem von zwei Codecs auf dem Draht komprimieren.

  • LZ4_FRAME: Schnell, geringer CPU-Overhead, bescheidenes Komprimierungsverhältnis. Bevorzugen Sie dies, wenn der Client CPU-eingeschränkt ist, aber dennoch Bytes auf dem Draht reduzieren möchte.
  • ZSTD: Höhere Komprimierungsquote, mehr CPU pro Batch. Aktivieren Sie sie, wenn Ihr Client die zusätzlichen CPU-Kosten aufnehmen kann.

Die Komprimierung reduziert Bytes auf dem Draht, fügt jedoch CPU-Kosten auf dem Client hinzu. Kleinere Nutzlasten können Netzwerkengpässe vermeiden und Netzwerkkosten reduzieren.

Python SDK

Stelle das Feld ipc_compression auf :ArrowStreamConfigurationOptions

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

options = ArrowStreamConfigurationOptions(ipc_compression=IPCCompression.ZSTD)
stream = sdk.create_arrow_stream(
    TABLE_NAME, schema, CLIENT_ID, CLIENT_SECRET, options=options
)

Rust SDK

Stellen Sie den Kompressionstyp im Stream-Builder ein. Das CompressionType-Enum befindet sich im arrow-ipc-Crate, fügen Sie es daher als Abhängigkeit hinzu:

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?;

Bewährte Methoden

Befolgen Sie diese Richtlinien, um bei der Datenaufnahme mit Arrow Flight die beste Leistung und Zuverlässigkeit zu erzielen.

  • Verwenden Sie einen Datenstrom für viele Batches, anstatt einen neuen Datenstrom pro Batch zu öffnen. Die Streamerstellung führt zu einem erheblichen Aufwand, den Sie amortisieren können, indem Sie einen Datenstrom über viele Batches hinweg wiederverwenden.
  • Mehrere Zeilen pro Batch senden. Beginnen Sie mit Stapeln in natürlicher Anwendungsgröße, nicht mit einer Zeile pro Aufruf. Das Senden jeweils einzelner Zeilen funktioniert, macht aber den Großteil des Performancevorteils von Arrow zunichte.
  • Rufen Sie flush() an kontrollierten Prüfpunkten an. Dadurch erhalten Sie eine klare Haltbarkeitsgrenze für eine Gruppe von Batches, ohne jede einzelne zu blockieren.
  • Aktivieren Sie die IPC-Komprimierung, um den Durchsatz zu verbessern. ZSTD wird für die meisten Workloads empfohlen, wenn der Client eine CPU übrig hat. Verwenden Sie LZ4_FRAME oder gar keine Komprimierung, wenn der Client nur über begrenzte CPU-Ressourcen verfügt.
  • Verwenden Sie Arrow Flight, wenn Ihre Datenquelle bereits spaltenbasiert ist. Wenn deine Quelldaten von Natur aus zeilenorientiert und klein sind, ist die Verwendung von Zerobus Ingest mit JSON oder Protobuf oft einfacher. Siehe Use Zerobus Ingest.

Fehlerbehandlung und Wiederherstellung

Arrow-Flight-Streams verwenden dieselben gRPC-Fehlerkategorien wie der Rest von Zerobus Ingest. Fehlercodes, Anleitungen zum Wiederholen und die vollständige Client-vs-Server-Taxonomie finden Sie unter Zerobus Ingest-Fehlerbehandlung.

Wenn Sie das SDK mit der automatischen Wiederherstellung (standardmäßig) konfigurieren, stellt es bei vorübergehenden Fehlern die Verbindung nahtlos wieder her und sendet nicht bestätigte Batches erneut. Nachdem ein Stream mit nicht bestätigter Arbeit geschlossen wird, behält das SDK Batches, die der Client akzeptiert, der Server aber nicht bestätigt hat, einschließlich Chargen, die möglicherweise noch nicht gesendet wurden.

Nachdem die automatische Wiederherstellung ausgeschöpft ist, beheben Sie die Ursache des Fehlers und rufen Sie close() auf, um die nicht bestätigten Batches des Streams abzuschließen. Da der Stream bereits fehlgeschlagen ist, könnte close() derselbe Terminalfehler zurückgegeben werden, obwohl der Abschluss erfolgreich ist.

Du kannst erst anrufen get_unacked_batches() , nachdem der Stream geschlossen ist. Es gibt die aufbewahrten Chargen für Persistenz oder eine von der Anwendung verwaltete Wiedergabe zurück. Wie Sie einen Ersatzstream erstellen, die Batches speichern und erneut versuchen, hängt von der Wiederherstellungsrichtlinie Ihrer Anwendung ab.

Python SDK

from zerobus.sdk.shared import ZerobusException

try:
    stream.close()
except ZerobusException:
    # The terminal error can be returned after closure is finalized.
    pass

unacked_batches = stream.get_unacked_batches()

Rust SDK

// The terminal error can be returned after closure is finalized.
let _ = stream.close().await;
let unacked_batches = stream.get_unacked_batches().await?;

Weitere Ressourcen

  • Zerobus Ingest verwenden: Falls Sie Zerobus Ingest noch nicht eingerichtet haben, beginnen Sie hier mit Anweisungen zum Finden Ihrer Workspace-URL, zur Erstellung der Ziel-Delta-Tabelle und zur Konfiguration eines Service Principal. Diese Schritte gelten für alle Datensatzformate.
  • Zerobus-Ingest-Quoten: Überprüfen Sie die Zerobus-Quoten, bevor Sie in die Produktion gehen. Durchsatz-, Latenz- und Partitionierungstabellengrenzwerte gelten für Arrow Flight.
  • Zerobus Ingest Fehlerbehandlung: Konsultieren Sie diese Seite für eine vollständige Liste der gRPC-Fehlercodes sowie empfohlenes Wiederholungs- und Wiederherstellungsverhalten für Ihren Kunden.