Aller au contenu principal

Utiliser Zerobus Ingest

Cette page décrit comment ingérer des données à l'aide de Zerobus Ingest dans Lakeflow Connect.

Premiers pas avec Zerobus Ingest

remarque

Si vous disposez d'un pare-feu côté client, ajoutez l'adresse IP utilisée par Zerobus Ingest à votre liste d'autorisation. Pour afficher les adresses IP par région, voir Adresses IP et domaines pour les services et assets Databricks.

Avant de start, confirmez que Zerobus Ingest est disponible dans la région de votre Workspace. Voir Disponibilité de l'ingestion.

  1. Obtenez une URL Zerobus Ingest.
  2. Créez ou identifiez la table dans laquelle vous souhaitez ingérer des données.
  3. Créez un Service Principal et accordez des privilèges sur la table.
  4. Connectez un client ou un exportateur pour start à envoyer des données.

Choisissez le guide correspondant à votre cas d'utilisation :

  • Ingérez vos propres données : utilisez les SDK ou l'API REST de Zerobus Ingest avec un schéma que vous définissez. Veuillez suivre les instructions figurant sur cette page.

  • Ingérer des données OpenTelemetry : utilisez des SDK ou des collecteurs OpenTelemetry standard pour envoyer des traces, des logs et des métriques dans des schémas de table prédéfinis. Pour obtenir des instructions complètes, consultez Ingest OpenTelemetry data with Zerobus Ingest.

Choisir une interface

Zerobus Ingest prend en charge plusieurs interfaces, écrivant toutes directement dans des tables Delta de Unity Catalog. En bref :

Pour une comparaison complète et savoir comment choisir, consultez les protocoles API. Via les SDK, vous pouvez également choisir un format d'enregistrement (JSON, Protocol Buffers (protobuf) ou Apache Arrow). Voir Types de messages. Le reste de cette page utilise les SDK et l'API REST.

Obtenez l'URL de votre Workspace et l'Endpoint Zerobus Ingest

Votre URL de workspace apparaît dans le navigateur lorsque vous vous connectez. Bien que l'URL complète suive le format https://<databricks-instance>.com/o=XXXXX, l'URL du workspace comprend tout ce qui précède le /o=XXXXX. Par exemple, étant donné l'URL complète suivante, vous pouvez déterminer l'URL du workspace et l'ID du workspace.

  • URL complète : https://abcd-teste2-test-spcse2.cloud.databricks.com/?o=2281745829657864#
  • URL du workspace : https://abcd-teste2-test-spcse2.cloud.databricks.com
  • ID du Workspace : 2281745829657864

L'endpoint de serveur dépend du workspace et de la région :

  • Endpoint du serveur : <workspace-id>.zerobus.<region>.cloud.databricks.com

Pour trouver la région de votre workspace, ouvrez le sélecteur de workspace dans la barre de navigation supérieure de l'interface utilisateur Databricks. La région est affichée sous chaque nom de workspace (par exemple, us-west-2). Vous pouvez également la trouver dans la console du compte sous Workspaces .

Pour connaître la disponibilité par région, consultez les quotas de Zerobus Ingest.

Créez ou identifiez la table cible

Identifiez la table cible dans laquelle vous souhaitez ingérer des données. Pour créer une nouvelle table cible, exécutez la commande CREATE TABLE SQL. Par exemple, créez une nouvelle table nommée unity.default.air_quality.

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

Zerobus Ingest peut écrire à la fois dans des tables Delta gérées et dans des tables de streaming, qui fonctionnent de la même manière, avec les mêmes limites et quotas.

remarque

Pour l'ingestion OpenTelemetry, les tables doivent utiliser des schémas prédéfinis pour chaque type de signal (traces, logs, métriques). Consultez Créer des tables cibles dans Unity Catalog.

Le schéma de votre table constitue le contrat de ce que Zerobus Ingest accepte, et Zerobus Ingest ne le fait jamais évoluer automatiquement. Planifiez les changements de schéma de manière proactive : faites d'abord évoluer la table, puis mettez à jour les producteurs. Zerobus Ingest écrit les enregistrements qui ne correspondent plus après un changement de table critique vers un emplacement de fallback durable au lieu de les supprimer. Voir Gestion de schémas et Récupération des données depuis l'emplacement de fallback durable.

Par default, Zerobus Ingest rejette les enregistrements dont les champs ne correspondent pas au schéma de la table cible. Pour capturer ces champs au lieu de les perdre, configurez une colonne de données sauvées. Voir colonne de données sauvées Zerobus.

Créer un Service Principal et accorder des autorisations

Un service principal est une identité spécialisée qui offre une meilleure sécurité que les comptes personnalisés. Pour plus d’information sur les Service Principal et sur la manière de les utiliser pour l’authentification, consultez Autoriser l’accès d’un Service Principal à Databricks avec OAuth.

Vous pouvez créer et gérer des Service Principals par programmation avec l'API REST ou les SDK Databricks, ou via l'interface utilisateur du workspace comme décrit ci-dessous. Les attributions de permissions à la fin de cette section sont du SQL que vous pouvez exécuter depuis n'importe quel client.

  1. Pour créer un Service Principal, accédez à Settings > Identity and Access .

  2. Sous Service Principal , sélectionnez Gérer .

  3. Cliquez sur Ajouter un Service Principal .

  4. Dans la fenêtre Ajouter un Service Principal , créez un nouveau Service Principal en cliquant sur Ajouter .

  5. Générez et enregistrez l'ID client et le secret client pour le Service Principal.

  6. Accordez au service principal les autorisations requises pour le catalogue, le schéma et la table.

    1. Sur la page Service Principal , accédez à l'onglet tab .
    2. Copiez l' ID d'application (UUID).
    3. Utilisez le SQL suivant pour accorder des autorisations, en remplaçant l'UUID d'exemple et le catalogue, le nom du schéma et les noms de table si nécessaire.
    SQL
    GRANT USE CATALOG ON CATALOG <catalog> TO `<UUID>`;
    GRANT USE SCHEMA ON SCHEMA <catalog.schema> TO `<UUID>`;
    GRANT MODIFY, SELECT ON TABLE <catalog.schema.table_name> TO `<UUID>`;

Écrire un client

Utilisez un SDK Zerobus dans votre langage de programmation préféré ou l'API REST pour ingérer des données dans votre table cible. Les SDK sont open source. Pour la bibliothèque complète, la documentation spécifique au langage et des exemples supplémentaires, consultez le repository du SDK Zerobus.

Les exemples ci-dessous utilisent ingest_record_offset, qui préserve l’ordre dans lequel vous envoyez les enregistrements.

Python 3.9 ou une version ultérieure est requis. Le SDK offre un throughput élevé et des E/S réseau efficaces grâce à un runtime asynchrone. Il prend en charge JSON (le plus simple) et Protocol Buffers (recommandé pour la production). Le SDK prend également en charge les implémentations synchrones et asynchrones, ainsi que les méthodes d'ingestion basées sur les offsets et sur les futures.

Bash
pip install databricks-zerobus-ingest-sdk

Exemple JSON :

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

# See "Get your workspace URL and Zerobus Ingest endpoint" for information on obtaining these values.
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"

sdk = ZerobusSdk(
SERVER_ENDPOINT,
DATABRICKS_WORKSPACE_URL
)

table_properties = TableProperties(TABLE_NAME)
stream = sdk.create_stream(CLIENT_ID, CLIENT_SECRET, table_properties)

try:
for i in range(1000):
record_dict = {
"device_name": f"sensor-{i}",
"temp": 20 + i % 15,
"humidity": 50 + i % 40
}
stream.ingest_record_offset(record_dict)
finally:
stream.close()

Les exemples ci-dessus utilisent la méthode ingest_record_offset basée sur le décalage sans attendre le décalage renvoyé. Pour en savoir plus sur les méthodes d'ingestion disponibles, sur le moment où attendre une confirmation de durabilité sur un décalage et sur la façon de suivre la progression avec un rappel d'accusé de réception, consultez Message blocking and acknowledgment.

Protocol Buffers : pour une ingestion typée en toute sécurité, transmettez un descripteur protobuf à TableProperties (le format est sélectionné automatiquement). Générez un schéma à partir de votre table en utilisant l’outil generate_proto, compilez-le avec protoc, puis transmettez le descripteur compilé pour créer le Stream.

Arrow Flight : pour une ingestion en colonnes ou par batch de données Apache Arrow RecordBatch via la même connexion gRPC, consultez Utiliser Arrow Flight avec Zerobus Ingest. Nécessite le supplément [arrow] : pip install "databricks-zerobus-ingest-sdk[arrow]" pyarrow.

Pour une documentation complète, les options de configuration, l'ingestion par lots et les exemples de tampon de protocole, consultez le repository Python du SDK.

Gérer les erreurs

Les exemples ci-dessus illustrent le chemin nominal. En production, encapsulez l'ingestion dans une gestion des erreurs. Le SDK réessaie automatiquement les erreurs transitoires, telles que les problèmes réseau, grâce à sa récupération intégrée. Les échecs dont il ne peut pas se remettre, tels que des identifiants non valides ou une table manquante, apparaissent sous la forme ZerobusException:

Python
from zerobus.sdk.shared import ZerobusException

try:
stream.ingest_record_offset(record)
except ZerobusException as e:
# Handle the failure: log it, fix the cause, recover on a new stream, or stop.
...

Les SDK se remettent également automatiquement des défaillances transitoires et vous permettent de récupérer les enregistrements non acquittés lorsqu'un Stream échoue de manière permanente. Pour les modèles de client résilient et la référence complète des erreurs, consultez Modèles de récupération et de nouvelle tentative et Gestion des erreurs Zerobus Ingest.

Étapes suivantes