Aller au contenu principal

Présentation de Zerobus Ingest

Zerobus Ingest est une API de streaming basée sur le push qui écrit les données directement dans les tables Delta de Unity Catalog à grande échelle, sans aucun bus de messages à exécuter. Cela supprime la couche intermédiaire que de nombreuses équipes placent entre leurs producteurs et le lakehouse. Le workflow se compose de deux étapes : créer une table, puis y envoyer des données. Un client « hello world » et une charge de travail à l’échelle du pétaoctet exécutent essentiellement le même code sans infrastructure à gérer.

L'ingestion via un bus de messages achemine les producteurs via un broker et un Job d'ingestion avant d'atteindre les tables Delta, tandis que Zerobus Ingest connecte les producteurs directement au lakehouse.

Zerobus Ingest est serverless ; il ajoute et supprime de la capacité en fonction de l'évolution de la charge. Il a ingéré plus de 1 000 milliards d'enregistrements dans une seule table en moins de 24 heures (voir le billet de blog Ingesting the Milky Way: Petabyte-Scale with Zerobus Ingest), et enregistre les données en quelques secondes.

Avantages

Zerobus Ingest simplifie l’ingestion tout en s’adaptant aux charges de travail les plus importantes :

  • Simple par conception. Créez une table, puis envoyez-y des données : il n'y a aucun broker, partition ou pipeline à gérer. Plutôt que d'acheminer les données via un bus de messages et un job d'ingestion avant qu'elles n'arrivent à destination, les producteurs écrivent directement dans la table, ce qui réduit le nombre de sauts et d'éléments mobiles à gérer.
  • Serverless et élastique. Zerobus Ingest est activé par « default » et ajoute ou supprime de la capacité à mesure que la charge évolue. Vous montez en charge en exécutant davantage de producteurs, et non en réécrivant votre application. Pour savoir comment, consultez Comment Zerobus Ingest monte en charge.
  • Workloads à haut throughput. Zerobus Ingest est conçu pour une ingestion à grande échelle, permettant de maintenir des taux d’écriture élevés dans une table unique.
  • Fraîcheur en quasi-temps réel. Les enregistrements arrivent dans Delta en quelques secondes et sont prêts à être interrogés presque dès leur arrivée.
  • high concurrency. Zerobus Ingest gère les écritures simultanées de milliers de clients dans la même table.

Lorsque votre destination est le lakehouse, Zerobus Ingest est le chemin le plus direct. D'autres outils Databricks répondent à des besoins adjacents et fonctionnent bien en complément :

  • Pour les cas d'usage où vous exécutez Kafka pour prendre en charge des consommateurs non-Lakehouse, vous pouvez également souhaiter qu'une copie des données soit effectuée dans le Lakehouse. Utilisez des connecteurs de streaming gérés pour le répliquer.
  • Pour les données déjà enregistrées sous forme de fichiers dans le stockage cloud, utilisez Auto Loader.
  • Lorsque vous avez besoin d'une latence opérationnelle inférieure à la seconde dans le chemin de traitement, utilisez le mode temps réel.

Créez une table, puis envoyez les données

L’utilisation de Zerobus Ingest est aussi simple que la création d’une table, suivie de l’envoi de données vers celle-ci. Le schéma de la table définit ce que chaque enregistrement doit contenir. Tout d'abord, créez la table cible :

SQL
CREATE TABLE main.default.air_quality (
device_name STRING,
temp INT,
humidity INT
);

Ensuite, l'ingestion d'un enregistrement ne nécessite que quelques lignes de code :

Python
from zerobus.sdk.sync import ZerobusSdk
from zerobus.sdk.shared import TableProperties

sdk = ZerobusSdk(SERVER_ENDPOINT, DATABRICKS_WORKSPACE_URL)

table_properties = TableProperties("main.default.air_quality")
stream = sdk.create_stream(CLIENT_ID, CLIENT_SECRET, table_properties)

stream.ingest_record_offset({"device_name": "sensor-1", "temp": 22, "humidity": 55})
# ingest more records...
stream.close()

Le code que vous déployez en développement peut monter en charge pour des workloads de production. Pour la procédure complète, consultez Utiliser Zerobus Ingest.

Cas d'utilisation courants

  • IoT et télémétrie des appareils : Stream les données de capteurs, de véhicules et d'appareils intelligents provenant de grands parcs distribués directement dans des tables Delta gouvernées.
  • On-premise vers cloud : connectez les systèmes on-premise et hybrides au lakehouse sans mettre en place d'infrastructure de courtage intermédiaire. Pour la connectivité privée et la configuration du pare-feu, consultez Considérations relatives au réseau.
  • Événements d'application et de flux de clics : envoyez des événements depuis des applications cloud et edge pour une analytique en quasi temps réel.
  • Change data capture (CDC) : enregistrez les changements de lignes des systèmes opérationnels dans Delta.
  • Données d'observabilité : envoyez des traces, des logs et des métriques OpenTelemetry dans les tables Delta que vous possédez. Voir Ingérer des données OpenTelemetry avec Zerobus Ingest.

Comment ça marche

Un producteur ouvre un Stream vers Zerobus Ingest et envoie des enregistrements vers une table Delta cible. Le service valide chaque enregistrement par rapport au schéma de la table et le rend durable. Une fois qu'un enregistrement est durable, Zerobus Ingest l'acquitte rapidement, afin que votre producteur puisse continuer à envoyer des enregistrements sans attendre pour chacun d'eux. Les données sont matérialisées dans la table peu après, généralement en quelques secondes. La conception dynamique et sans partition de Zerobus Ingest rend l'ingestion élastique, de sorte que son compute serverless monte en charge avec vos workloads.

Comment fonctionne Zerobus Ingest : les producteurs envoient des enregistrements vers l'endpoint Zerobus Ingest, qui les valide, les rend durables, en accuse réception et les matérialise dans des tables Delta Unity Catalog

Pour une explication plus approfondie des flux et de la manière dont Zerobus Ingest monte en charge, consultez les concepts de Zerobus Ingest. Pour le modèle de communication asynchrone entre client et serveur, consultez Communication asynchrone.

Méthodes d'envoi de données

Zerobus Ingest est un endpoint qui prend en charge plusieurs interfaces, afin que vous puissiez choisir la solution la plus adaptée pour chaque producteur :

  • SDK via gRPC : clients de streaming à haut throughput en Python, Java, Rust, Go, TypeScript et (en version bêta) C++ et C# / .NET. Idéal pour l'ingestion ordonnée de gros volumes. Voir Écrire un client.
  • API REST : une interface sans état pour les clients légers ou « bavards », tels que les larges flottes d'appareils périphériques. Voir Écrire un client.
  • OpenTelemetry (OTLP) : pointez les collecteurs OpenTelemetry existants vers Zerobus Ingest pour enregistrer les traces, les logs et les métriques sans intégration personnalisée. Voir Ingérer des données OpenTelemetry avec Zerobus Ingest.
  • APIs compatibles avec Kafka (Bêta) : pointez un producteur Apache Kafka existant vers Zerobus Ingest, sans SDK Databricks. Consultez Utiliser les APIs compatibles avec Kafka avec Zerobus Ingest.

Architecture de mise à l'échelle de Zerobus Ingest : les sources envoient des enregistrements Protocol Buffers (protobuf), JSON et Arrow via les APIs gRPC, REST, OpenTelemetry et compatibles Kafka, qui transitent par un système de mise à l'échelle automatique et d'équilibrage de charge vers un pool de nœuds Zerobus sans état évolutif horizontalement, chacun doté d'un journal de pré-écriture (write-ahead log) et d'un writer Lakehouse qui effectue des commit par batch des enregistrements dans une table Delta gérée par Unity Catalog

Ils écrivent tous directement dans des tables Delta. Pour une comparaison complète et savoir comment choisir, consultez Protocoles API. Pour écrire votre premier client, consultez Utiliser Zerobus Ingest.

Coût

Les frais pour Zerobus Ingest sont facturés sous la « SKU Jobs Serverless ». Les tarifs sont disponibles sur la page des tarifs de Lakeflow Connect.

monitoring de votre utilisation

Vous pouvez surveiller vos dépenses via la table système d’utilisation facturable. Consultez la référence de la table système d’utilisation facturable. Filtrer l’utilisation de Zerobus Ingest avec :

  • billing_origin_product = 'LAKEFLOW_CONNECT'
  • product_features.lakeflow_connect.zerobus_request_type identifie la manière dont les données ont été ingérées : 'GRPC' (streaming SDK), 'HTTP' (REST), 'OTEL_GRPC' et 'OTEL_HTTP' (OpenTelemetry/OTLP), ou 'KAFKA' (APIs compatibles Kafka).

Ressources supplémentaires