Usa Arrow Flight con Zerobus Ingest

L'ingestione Arrow Flight ti permette di inviare i dati di Apache ArrowRecordBatch direttamente a Zerobus Ingest invece di convertire prima ogni riga in JSON o Protocol Buffer (protobuf). È un'opzione di terzo formato di registrazione negli SDK Zerobus che supportano Arrow Flight, insieme a JSON e protobuf, e funziona sulla stessa connessione gRPC. Usa lo stesso endpoint Zerobus, lo stesso flusso OAuth e la stessa x-databricks-zerobus-table-name convenzione di intestazione. Il protocollo wire è Arrow FlightDoPut, che trasporta messaggi Arrow IPC su gRPC.

Quando usare Arrow Flight

Arrow Flight è l'ideale per gli scenari seguenti:

  • L'applicazione produce già dati Arrow, ad esempio pyarrow.Table o pyarrow.RecordBatch (Python), arrow_array::RecordBatch dai crate arrow-rs (Rust), o VectorSchemaRoot (Java). Le librerie DataFrame costruite su Arrow, come Polars o DataFusion, si inseriscono naturalmente in questo percorso.
  • Le righe vengono inserite in batch anziché inviare un record alla volta.
  • Lo schema è ampio, ricco di campi numerici o orientato all’analitica, dove la serializzazione riga per riga aggiunge un significativo sovraccarico della CPU.
  • Si stanno creando agenti di raccolta o gateway che aggregano i dati per un breve intervallo e quindi lo inviano come batch in formato colonna.

Arrow Flight di solito non è la scelta migliore per il traffico sporadico, una riga alla volta. In questi casi, JSON o protobuf sul percorso SDK gRPC sono tipicamente più semplici. Vedere Scegliere un'interfaccia.

Funzionamento del modello di inserimento

Con l'acquisizione tramite Arrow Flight, un flusso scrive su un'unica tabella di destinazione. Per inserire dati, seguire questa sequenza:

  1. Definire uno schema Arrow corrispondente allo schema della tabella Delta di destinazione.
  2. Apri uno stream Zerobus Arrow per quella tabella.
  3. Inviare RecordBatch (o Table) payload.
  4. Attendere l'ultimo offset o chiamare flush() per confermare la durabilità.
  5. Chiudere il flusso.

Se utilizzi un SDK Zerobus, lo SDK gestisce per te i dettagli a basso livello del protocollo Arrow Flight. Serializza i tuoi dati Arrow in formato IPC e divide automaticamente un grande lotto in messaggi Flight più piccoli e ordinati.

Il server riporta i progressi cumulativi man mano che i record diventano durevoli. Un grande lotto logico può quindi essere parzialmente durabile se si verifica un guasto durante l'invio. ingest_batch() restituisce comunque un offset logico per il lotto inviato, e aspettare quel offset conferma che tutti i suoi record sono durevoli. Consulta I batch di Arrow Flight costituiscono l'eccezione.

Zerobus abbina i campi Arrow alle colonne Delta per nome. Lo schema Arrow deve includere tutte le colonne Delta richieste (non nullabili). Puoi omettere colonne nullabili, che Zerobus scrive come NULL. Non includere campi assenti dalla tabella target. I campi inclusi devono seguire l'ordine relativo dello schema Delta, corrispondere alla sua nullabilità e utilizzare il tipo della colonna di destinazione. Per dettagli, vedi regole di abbinamento degli schemi.

Un batch Arrow non è soggetto al limite di dimensione gRPC di 10 MB. Tuttavia, ogni singola riga all'interno di una RecordBatch deve rientrare nel limite di 10 MB. L'SDK suddivide automaticamente batch più grandi in più messaggi filati, ma non può dividere una riga sovradimensionata. Vedi quote di acquisizione di Zerobus.

Scrivere un client

Gli esempi seguenti utilizzano gli SDK di Python e Rust. Per altri linguaggi che supportano Arrow Flight, vedi il repository Zerobus SDK.

PYTHON SDK

L'SDK di Python accetta un pyarrow.Schema in fase di creazione del flusso e un pyarrow.RecordBatch o pyarrow.Table per ogni chiamata di inserimento.

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() accetta anche un oggetto pyarrow.Table. L'SDK lo converte in un unico RecordBatch internamente prima dell'invio. Ogni chiamata restituisce un offset logico. L'esempio crea diversi lotti e chiama flush() una volta per confermare che tutti i lotti in sospeso siano durevoli. Usalo wait_for_offset() quando devi confermare un lotto specifico prima di continuare; aspettare l'ultimo offset conferma anche tutti gli offset precedenti. Per sapere quando aspettare e come funziona il riconoscimento, vedi Blocco e riconoscimento dei messaggi.

Rust SDK

Rust SDK espone Arrow Flight tramite l'API stream_builder() , dietro la arrow-flight funzionalità Cargo. Usare la stessa versione major di Arrow dell’SDK, in modo che il tipo RecordBatch e i tipi array corrispondano in fase di compilazione.

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

Il generatore seleziona il formato Arrow Flight con .arrow(schema) e finalizza il flusso con .build_arrow(), che restituisce un oggetto ZerobusArrowStream. JSON e protobuf continuano a usare .json() / .compiled_proto(...) e ..build()

Ingestione delle colonne VARIANT

Apache Arrow non ha un tipo nativo VARIANT . Per importare in una colonna VARIANT tramite Arrow Flight, crea i campi metadata e value sottostanti della colonna come una struct composta da due colonne LargeBinary, quindi includi tale struct nel tuo RecordBatch. Sui gRPC SDK e REST, invece si passa un valore Variant come stringa codificata in JSON. Vedere Tipi di dati supportati.

Il seguente esempio di Rust costruisce una VARIANT colonna struct a partire dalle righe JSON e la ingerisce:

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

Questo esempio si trova in Rust. Per un uso equivalente in altri linguaggi, vedi il repository Zerobus SDK.

Compressione IPC

Per impostazione predefinita, i payload IPC arrow vengono inviati non compressi. Facoltativamente, è possibile comprimerli durante la trasmissione con uno dei due codec.

  • LZ4_FRAME: sovraccarico rapido e basso della CPU, rapporto di compressione modesto. Preferisce questo quando il client è vincolato dalla CPU, ma vuole comunque ridurre i byte in transito.
  • ZSTD: maggiore rapporto di compressione, più CPU per batch. Abilitarla ogni volta che il client può assorbire il costo aggiuntivo della CPU.

La compressione riduce i byte in transito, ma aggiunge costi di CPU nel client. I payload più piccoli possono evitare colli di bottiglia di rete e ridurre i costi di rete.

PYTHON SDK

Imposta il ipc_compression campo su 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

Imposta il tipo di compressione nel generatore di flussi. L'enum CompressionType si trova nel crate arrow-ipc, quindi aggiungila come dipendenza:

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

Procedure consigliate

Segui queste linee guida per ottenere prestazioni e affidabilità ottimali durante l'acquisizione con Arrow Flight.

  • Riutilizzare un flusso per molti batch invece di aprire un nuovo flusso per batch. La creazione del flusso comporta un sovraccarico significativo che è possibile ammortizzare riutilizzando un flusso in molti batch.
  • Inviare più righe per batch. Inizia con batch di dimensioni naturali per l'applicazione, non una riga a chiamata. L'invio di una riga alla volta funziona, ma nega la maggior parte del vantaggio delle prestazioni dell'uso di Arrow.
  • Chiama flush() ai posti di blocco controllati. In questo modo è possibile definire un limite di durabilità chiaro per un gruppo di batch senza bloccarne uno.
  • Abilitare la compressione IPC per migliorare la velocità effettiva. ZSTD è consigliato per la maggior parte dei carichi di lavoro quando il client ha una CPU di riserva. Uso LZ4_FRAME o nessuna compressione se il client è vincolato dalla CPU.
  • Usa Arrow Flight quando il produttore è già in formato colonnare. Se i tuoi dati sorgente sono naturalmente orientati alle righe e piccoli, usare Zerobus Ingest con JSON o protobuf è spesso più semplice. Vedi Usa Zerobus Ingest.

Gestione e ripristino degli errori

I flussi Arrow Flight usano le stesse categorie di errore gRPC del resto di Zerobus Ingest. Per i codici di errore, linee guida per i tentativi e la tassonomia completa client-vs-server, vedere Gestione degli errori di inserimento zerobus.

Quando si configura l'SDK con il ripristino automatico (impostazione predefinita), si riconnette in modo trasparente e riproduce batch non riconosciuti in caso di errori temporanei. Dopo che uno stream si chiude con lavori non riconosciuti, l'SDK conserva i lotti che il client ha accettato ma che il server non ha confermato, inclusi i lotti che potrebbero non essere ancora stati inviati.

Dopo che il recupero automatico è esaurito, correggi la causa del guasto e chiama close() per finalizzare i lotti non confermati del stream. Poiché il flusso è già fallito, close() potrebbe restituire lo stesso errore terminale anche se la finalizzazione ha successo.

Puoi chiamare get_unacked_batches() solo dopo che lo stream è stato chiuso. Restituisce i lotti conservati per persistenza o riproduzione gestita dall'applicazione. Come crei un flusso sostitutivo, persisti i lotti e li riprovi dipende dalla politica di recupero della tua applicazione.

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

Risorse aggiuntive

  • Usa Zerobus Ingest: Se non hai ancora configurato Zerobus Ingest, inizia da qui per le istruzioni su come trovare l'URL del tuo workspace, creare la tabella Delta target e configurare un service principal. Questi passaggi vengono condivisi in tutti i formati di record.
  • Quote di inserimento di Zerobus: Verifica le quote di Zerobus prima della distribuzione in produzione. Tutti i limiti di velocità effettiva, latenza e tabella partizionata si applicano a Arrow Flight.
  • Gestione degli errori di Zerobus Ingest: Consulta questa pagina per l'elenco completo dei codici di errore gRPC e del comportamento consigliato per i nuovi tentativi e il ripristino del client.