Utilisez Arrow Flight avec Zerobus Ingest
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 dans les SDK Zerobus qui prennent en charge Arrow Flight, 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 automatiquement un large batch en messages Flight plus petits et ordonnés.
Le serveur signale la progression cumulée à mesure que les enregistrements deviennent durables. Un batch logique volumineux peut donc être partiellement durable si une défaillance survient pendant son envoi. ingest_batch() renvoie toujours un décalage logique pour le batch soumis, et l'attente de ce décalage confirme que tous ses enregistrements sont durables. Consultez Les batchs Arrow Flight sont l'exception.
Zerobus fait correspondre les champs Arrow aux colonnes Delta par leur nom. Le schéma Arrow doit inclure toutes les colonnes Delta requises (non nullables). Vous pouvez omettre les colonnes nullables, que Zerobus écrit sous la forme NULL. N’incluez pas les champs absents de la table cible. Les champs inclus doivent respecter l’ordre relatif du schéma Delta, correspondre à sa nullabilité et utiliser le type de la colonne cible. Pour plus de détails, consultez les règles de correspondance de schéma.
Un batch Arrow n’est pas soumis à la limite de taille de message gRPC de 10 Mo. Cependant, chaque ligne individuelle à l’intérieur d’un RecordBatch doit respecter la limite de 10 Mo. Le SDK divise automatiquement les batches plus volumineux en plusieurs messages réseau, mais il ne peut pas diviser une ligne trop volumineuse. Voir les quotas de Zerobus Ingest.
Écrire un client
Les exemples ci-dessous utilisent les SDK Python et Rust. Pour les autres langages prenant en charge Arrow Flight, consultez le repository du SDK Zerobus.
- 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]"
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() accepte également un pyarrow.Table. Le SDK le convertit en un seul RecordBatch en interne avant l'envoi. Chaque appel renvoie un offset logique. L’exemple crée plusieurs batchs et appelle flush() une fois pour confirmer que tous les batchs en attente sont durables. Utilisez wait_for_offset() lorsque vous devez confirmer un batch spécifique avant de continuer ; l’attente du dernier décalage confirme également tous les décalages précédents. Pour savoir quand attendre et comment fonctionne l’accusé de réception, consultez Blocage et accusé de réception des messages.
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, 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(())
}
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.
- Python SDK
- Rust SDK
Définissez le champ ipc_compression sur 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
)
Définissez le type de compression sur le générateur de Stream. L’énumération CompressionType se trouve dans le crate arrow-ipc, ajoutez-la donc en tant que 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 (la valeur par « default »), il se reconnecte de manière transparente et relit les batches non acquittés en cas de défaillances transitoires. Après la fermeture d’un Stream avec des travaux non acquittés, le SDK conserve les batches que le client a acceptés mais que le serveur n’a pas acquittés, y compris les batches qui pourraient ne pas encore avoir été envoyés.
Une fois la récupération automatique épuisée, corrigez la cause de l’échec et appelez close() pour finaliser les batchs non acquittés du Stream. Comme le Stream a déjà échoué, close() pourrait renvoyer la même erreur terminale, même si la finalisation réussit.
Vous ne pouvez appeler get_unacked_batches() qu'une fois le Stream fermé. Il renvoie les batchs conservés pour la persistance ou la relecture gérée par l'application. La manière dont vous créez un Stream de remplacement, persistez les batchs et les relancez dépend de la politique de récupération de votre application.
- Python SDK
- Rust 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()
// The terminal error can be returned after closure is finalized.
let _ = stream.close().await;
let unacked_batches = stream.get_unacked_batches().await?;
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 de Zerobus Ingest: examinez les quotas Zerobus avant le déploiement en production. Les limites de throughput, de latence et de table partitionnée s'appliquent toutes à Arrow Flight.
- Gestion des erreurs de Zerobus Ingest: consultez cette page pour obtenir la liste complète des codes d'erreur gRPC ainsi que les comportements de nouvelle tentative et de récupération recommandés pour votre client.