Aller au contenu principal

Utilisez Arrow Flight avec Zerobus Ingest

info

Bêta

Cette fonctionnalité 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 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 de Protocol Buffers, et elle fonctionne 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 connexion est Arrow Flight DoPut, qui transporte des messages Arrow IPC sur 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.Table ou pyarrow.RecordBatch (Python), arrow_array::RecordBatch des crâtes arrow-rs (Rust), ou VectorSchemaRoot (Java). Les bibliothèques DataFrame basées sur Arrow — telles que Polars ou DataFusion — s'inscrivent naturellement dans cette approche.
  • 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 le trafic clairsemé, ligne par ligne. Dans ces cas, JSON ou les buffers de protocole via le chemin gRPC du SDK sont généralement plus simples. Voir 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 :

  1. Définissez un schéma Arrow qui correspond au schéma de la table Delta de destination.
  2. Ouvrir un Zerobus Arrow Stream pour cette table.
  3. Envoyer RecordBatch (ou Table) charges utiles.
  4. Attendez le dernier décalage ou appelez flush() pour confirmer la durabilité.
  5. Fermez le Stream.

Si vous utilisez un SDK Zerobus, le SDK gère pour vous les détails de la liaison filaire Arrow Flight de bas niveau. Il sérialise vos données Arrow au format IPC et divise automatiquement les batches surdimensionnés sur plusieurs messages Flight. Le serveur confirme chaque batch durablement, et le SDK présente ces confirmations sous forme de décalages de batch logiques.

Comme pour la règle de schéma Protobuf, le schéma que vous transmettez au stream doit correspondre à la table Delta cible 1:1. Votre schéma peut omettre les colonnes pouvant être nulles qui existent dans la table Delta (ceci est traité comme une modification de schéma non bloquante), mais tout autre écart est refusé.

Chaque ligne individuelle à l'intérieur d'un RecordBatch doit tenir dans la limite de taille de message gRPC de 10 Mo. Arrow Flight divise automatiquement les batchs surdimensionnés en plusieurs messages filaires, de sorte que les grands batchs ne posent pas de problème — mais une seule ligne dont la taille sérialisée dépasse la limite ne peut pas être divisée et est rejetée. Le throughput, la latence, le quota et les limites des tables partitionnées s'appliquent également à Arrow Flight, car il s'exécute sur le même transport gRPC. Voir les limitations du connecteur Zerobus Ingest.

Écrire un client

Les exemples ci-dessous ouvrent un Arrow Stream Flight sur la même table air_quality utilisée dans les exemples de Utiliser le connecteur Zerobus Ingest. Ils sont présentés en Python et Rust pour faire court, mais le même constructeur, les mêmes options de configuration et la même séquence d'appels sont disponibles dans chaque Zerobus SDK. Adaptez la syntaxe à votre langage et consultez le repository SDK pour les types Arrow spécifiques au langage.

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.

Bash
pip install "databricks-zerobus-ingest-sdk[arrow]" pyarrow
Python
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)
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 retourne un offset logique. stream.wait_for_offset(offset) bloque jusqu'à ce que le serveur ait persisté durablement ce batch.

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:

Python
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 :

Bash
cargo add arrow-ipc
Rust
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.
  • Activer la compression IPC pour améliorer le throughput. ZSTD est recommandé pour la plupart des workloads lorsque le client dispose de CPU de rechange ; utilisez LZ4_FRAME ou aucune compression si le client est limité en CPU.
  • Utilisez Arrow Flight lorsque votre producteur est déjà en colonnes. Si vos données source sont naturellement orientées lignes et petites, l'utilisation du connecteur Zerobus Ingest avec JSON ou Protocol Buffers est souvent plus simple. Consultez Utilisez le connecteur 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 :

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 le connecteur Zerobus Ingest: si vous n'avez pas encore configuré Zerobus Ingest, start ici pour obtenir des instructions sur la façon de trouver l'URL de votre Workspace, de créer la table Delta cible et de configurer un Service Principal. Ces étapes sont partagées entre tous les formats d'enregistrement.
  • Limitations du connecteur Zerobus Ingest: examinez les quotas Zerobus avant de déployer en production. Les limites de throughput, de latence et de tables partitionnées s'appliquent toutes à Arrow Flight.
  • 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.