Utilisez Arrow Flight avec Zerobus Ingest
Bêta
L'ingestion Arrow Flight est en bêta.
L'ingestion Arrow Flight vous permet d'envoyer des données Apache Arrow RecordBatch directement vers Zerobus Ingest au lieu de convertir chaque ligne en JSON ou en Protocol Buffers (protobuf) au préalable. Il s'agit d'une troisième option de format d'enregistrement sur les SDK Zerobus, aux côtés de JSON et protobuf, et elle s'exécute sur la même connexion gRPC. Il utilise le même Endpoint Zerobus, le même flux OAuth et la même convention d'en-tête x-databricks-zerobus-table-name. Le protocole de communication est Arrow Flight DoPut, qui transporte les messages IPC Arrow via gRPC.
Quand utiliser Arrow Flight
Arrow Flight convient le mieux dans les scénarios suivants :
- Votre application produit déjà des données Arrow, telles que
pyarrow.Tableoupyarrow.RecordBatch(Python),arrow_array::RecordBatchà partir des crates arrow-rs (Rust), ouVectorSchemaRoot(Java). Les bibliothèques DataFrame basées sur Arrow, telles que Polars ou DataFusion, s’intègrent naturellement dans ce chemin. - Vous ingérez des lignes par batchs au lieu d'envoyer un enregistrement à la fois.
- Votre schéma est large, à forte composante numérique ou orienté analytique, où la sérialisation ligne par ligne ajoute une surcharge CPU notable.
- Vous créez des collecteurs ou des passerelles qui regroupent les données pendant un court intervalle, puis les envoient sous la forme d'un batch de format de colonne.
Arrow Flight n'est généralement pas le meilleur choix pour un trafic sparse, ligne par ligne. Dans ces cas-là, JSON ou protobuf via le chemin gRPC du SDK sont généralement plus simples. Consultez Choisir une interface.
Fonctionnement du modèle d'ingestion
Avec l'ingestion Arrow Flight, un stream écrit dans une table cible. Pour ingérer des données, suivez cette séquence :
- Définissez un schéma Arrow qui correspond au schéma de la table Delta de destination.
- Ouvrir un Zerobus Arrow Stream pour cette table.
- Envoyer
RecordBatch(ouTable) charges utiles. - Attendez le dernier décalage ou appelez
flush()pour confirmer la durabilité. - Fermez le Stream.
Si vous utilisez un SDK Zerobus, celui-ci gère pour vous les détails de bas niveau de la couche de transport Arrow Flight. Il sérialise vos données Arrow au format IPC et divise un batch volumineux en messages de transport plus petits, que le serveur accuse réception individuellement.
Arrow Flight ne garantit pas une durabilité « tout ou rien » pour l'ensemble du batch logique. Un batch Arrow peut être très volumineux et, comme le SDK le divise en messages de transport distincts qui sont acquittés au fur et à mesure de leur réception, un batch volumineux peut être partiellement durable en cas de défaillance en cours de traitement. Ceci diffère des batchs JSON et protobuf, qui commit de manière atomique et sont limités par une taille de message de 10 Mo. Voir Les batchs Arrow Flight font exception.
L’abstraction de décalage logique (logical-offset) reste valable au-dessus de ce découpage en segments. ingest_batch() renvoie un décalage logique unique pour le batch que vous avez soumis, et wait_for_offset() sur ce décalage ne se termine qu’après que chaque message de transport constituant le batch a été acquitté. (Les noms des méthodes proviennent du SDK Python ; d’autres SDK exposent des méthodes équivalentes, telles que ingestBatch et waitForOffset en Java.)
Comme pour la règle Protobuf schema, le schéma que vous transmettez au stream doit correspondre à la table Delta cible : il doit contenir au minimum toutes les colonnes non nullables. Votre schéma peut omettre des colonnes nullables qui existent dans la table Delta (ceci est traité comme un changement de schéma non bloquant), mais toute autre incohérence est rejetée. Le type de chaque champ Arrow doit être compatible avec sa colonne Delta. Pour les types Delta pris en charge, consultez Supported data types.
Comme le SDK divise les gros batchs en messages de transport, un batch Arrow n'est pas soumis à la limite de taille de message de 10 Mo, contrairement à un batch JSON ou Protocol Buffers. Comme Arrow Flight s'exécute sur le même transport gRPC, les mêmes caractéristiques de throughput, de latence et de quota s'appliquent, et toutes permettent de monter en charge pour répondre à des workloads plus élevés. Voir les quotas de Zerobus Ingest.
Écrire un client
Les exemples ci-dessous ouvrent un Arrow Flight Stream sur la même table air_quality utilisée dans les exemples Utiliser Zerobus Ingest. Ils sont présentés en Python et en Rust par souci de concision, mais le même constructeur, les mêmes options de configuration et la même séquence d'appels sont disponibles dans chaque SDK Zerobus. Adaptez la syntaxe à votre langage et consultez le repository du SDK pour connaître les types Arrow spécifiques au langage.
- Python SDK
- Rust SDK
Le SDK Python accepte un pyarrow.Schema au moment de la création du Stream et un pyarrow.RecordBatch ou pyarrow.Table pour chaque appel d'ingestion.
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() accepte également un pyarrow.Table. Le SDK le convertit en un seul RecordBatch en interne avant l'envoi. Chaque appel renvoie un offset logique. Le blocage sur l'offset est facultatif. Pour savoir quand attendre et comment fonctionne la confirmation, consultez Blocage et confirmation de message.
Le SDK Rust expose Arrow Flight via l'stream_builder() API, derrière la fonctionnalité Cargo arrow-flight. Utilisez la même version majeure d'Arrow que le SDK afin que les types RecordBatch et de tableau correspondent au moment de la compilation.
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(())
}
Le générateur sélectionne le format Arrow Flight avec .arrow(schema) et finalise le Stream avec .build_arrow(), ce qui renvoie un ZerobusArrowStream. JSON et protobuf continuent d'utiliser .json() / .compiled_proto(...) et .build().
Ingestion de colonnes VARIANT
Apache Arrow ne possède pas de type VARIANT natif. Pour ingérer dans une colonne VARIANT via Arrow Flight, construisez les champs metadata et value de support de la colonne sous forme de structure de deux colonnes LargeBinary, puis incluez cette structure dans votre RecordBatch. Via les SDK gRPC et REST, vous passez plutôt une valeur Variant sous forme de chaîne encodée en JSON. Voir Types de données pris en charge.
L'exemple Rust suivant construit une colonne de structure VARIANT à partir de lignes JSON et l'ingère :
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(())
}
Cet exemple est en Rust. Pour une utilisation équivalente dans d'autres langages, consultez le repository du SDK Zerobus.
Compression IPC
By default, les charges utiles Arrow IPC sont envoyées non compressées. Vous pouvez les compresser éventuellement sur le réseau à l'aide de l'un des deux codecs.
LZ4_FRAME: Rapide, faible surcharge CPU, taux de compression modeste. Préférez cette option lorsque le client est contraint par le processeur mais souhaite tout de même réduire le nombre d'octets sur le réseau.ZSTD: Taux de compression plus élevé, plus de CPU par batch. Activez-le chaque fois que votre client peut absorber le coût CPU supplémentaire.
La compression réduit les octets sur le réseau, mais ajoute des coûts de CPU côté client. Des charges utiles plus petites peuvent éviter les goulots d'étranglement réseau et réduire les coûts réseau.
Dans le SDK Python, définissez le champ ipc_compression sur ArrowStreamConfigurationOptions:
from zerobus.sdk.shared.arrow import IPCCompression, ArrowStreamConfigurationOptions
options = ArrowStreamConfigurationOptions(ipc_compression=IPCCompression.ZSTD)
Dans le SDK Rust, définissez-le sur le générateur. L’énumération CompressionType se trouve dans le crate arrow-ipc, ajoutez-la donc comme dépendance :
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?;
Bonnes pratiques
Suivez ces directives pour obtenir les meilleures performances et la meilleure fiabilité de l'ingestion Arrow Flight.
- Réutilisez un Stream pour de nombreux batchs au lieu d'ouvrir un nouveau Stream par batch. La création de Stream entraîne une surcharge importante que vous pouvez amortir en réutilisant un Stream sur de nombreux batchs.
- Envoyez plusieurs lignes par batch. Start avec des lots de taille naturelle adaptés à l'application, pas une ligne par appel. L'envoi d'une ligne à la fois fonctionne, mais annule la plupart de l'avantage de performance lié à l'utilisation d'Arrow.
- Appelez
flush()à des points de contrôle contrôlés. Cela vous offre une limite de durabilité claire pour un groupe de batchs sans bloquer chaque batch individuel. - Activez la compression IPC pour améliorer le throughput.
ZSTDest recommandé pour la plupart des workloads lorsque le client dispose de CPU disponible. UtilisezLZ4_FRAMEou aucune compression si le client est limité en CPU. - Utilisez Arrow Flight lorsque votre producteur est déjà colonnaire. Si vos données sources sont naturellement orientées ligne et de petite taille, l'utilisation de Zerobus Ingest avec JSON ou protobuf est souvent plus simple. Consultez Utiliser Zerobus Ingest.
Gestion des erreurs et récupération
Les streams Arrow Flight utilisent les mêmes catégories d'erreurs gRPC que le reste de Zerobus Ingest. Pour les codes d'erreur, les conseils de nouvelle tentative et la taxonomie complète client-vs-serveur, consultez la gestion des erreurs Zerobus Ingest.
Lorsque vous configurez le SDK avec la récupération automatique (le default), il se reconnecte de manière transparente et rejoue les batches non accusés de réception en cas de défaillances transitoires. Une fois le stream fermé, vous pouvez récupérer tous les batches que le serveur a reçus mais n'a pas encore accusés de réception. Cela s'applique que le stream se soit fermé correctement ou en raison d'une défaillance irrécupérable. Dans le SDK Python :
# Retry unacked_batches against a freshly created stream
if stream.is_closed:
unacked_batches = stream.get_unacked_batches()
Dans le SDK Rust, appelez stream.get_unacked_batches().await? pour récupérer les batches non acquittés pour un nouvel essai.
Ressources supplémentaires
- Utiliser Zerobus Ingest: si vous n’avez pas encore configuré Zerobus Ingest, start ici pour obtenir des instructions sur la recherche de l’URL de votre Workspace, la création de la table Delta cible et la configuration d’un Service Principal. Ces étapes sont partagées entre tous les formats d'enregistrement.
- Quotas Zerobus Ingest: examinez les quotas Zerobus « default » avant le déploiement en production. Les mêmes caractéristiques de throughput et de latence s’appliquent à Arrow Flight, et montent en charge pour répondre à des workloads plus élevés.
- Gestion des erreurs Zerobus Ingest: Consultez cette page pour obtenir une liste complète des codes d'erreur gRPC et du comportement de nouvelle tentative et de récupération recommandé pour votre client.