Aller au contenu principal

Utilisez le connecteur Zerobus Ingest

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

Choisir une interface

Zerobus Ingest prend en charge les interfaces gRPC, REST et OpenTelemetry (OTLP). Les SDK fournissent des clients personnalisés basés sur gRPC avec une interface conviviale pour les développeurs afin de créer des applications à haut throughput. L'interface REST gère les contraintes architecturales lorsqu'il s'agit de flottes massives de dispositifs "volumineux". L'interface OTLP accepte les données OpenTelemetry standard sans nécessiter de bibliothèques personnalisées.

  • Les SDK avec la « taxe de connexion » gRPC : gRPC se spécialise dans les performances à haut throughput grâce à des connexions persistantes. Chaque Stream ouvert est comptabilisé dans vos quotas de simultanéité. Les SDK prennent en charge trois formats d'enregistrement sur gRPC :

    • JSON : le plus simple, aucune définition de schéma requise.
    • Protocol Buffers : Recommandé pour les charges de travail de production orientées lignes.
    • Apache Arrow Flight : recommandé pour le format de colonne ou les charges de travail par batch. Voir Utiliser Arrow Flight avec Zerobus Ingest.
  • La « taxe de throughput » REST : REST nécessite une négociation complète pour chaque mise à jour, ce qui le rend sans état. REST convient bien aux cas d'utilisation des appareils en périphérie où l'état est rarement rapporté.

  • OpenTelemetry (OTLP) : Si vous utilisez déjà les SDK ou collecteurs OpenTelemetry, l'endpoint OTLP ingère les traces, les logs et les métriques dans les tables Delta de Unity Catalog sans aucune intégration personnalisée requise. Pour plus d'informations, consultez Ingérer des données OpenTelemetry avec Zerobus Ingest.

Utilisez les SDK avec proto pour les flux à volume élevé orientés lignes, ou les SDK avec Arrow Flight pour le format de colonne ou les charges de travail en batch. Utilisez REST pour les flottes d'appareils massives à basse fréquence, et OTLP pour les environnements déjà instrumentés avec OpenTelemetry.

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 de Databricks — la région est affichée sous chaque nom de Workspace (par exemple, us-west-2). Vous pouvez également le trouver dans la console de compte sous **Workspaces**.

Pour la disponibilité régionale, consultez les limites du connecteur 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);
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.

By default, les enregistrements dont les champs ne correspondent pas au schéma de la table cible sont rejetés. Pour capturer ces champs au lieu de les perdre, configurez une colonne de récupération. Consultez Capturez les champs non conformes avec la colonne de récupération Zerobus.

Ingérer dans une table de streaming

info

Bêta

L'ingestion dans les tables en streaming à l'aide du connecteur Zerobus Ingest est en bêta. Les administrateurs du Workspace peuvent contrôler l'accès à cette fonctionnalité depuis la page Aperçus . Consultez Gérer les aperçus Databricks.

Pour créer une nouvelle table de streaming, exécutez la commande SQL CREATE STREAMING TABLE. Par exemple :

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

Une fois la table de streaming créée, ingérez-y des données en utilisant l’une des interfaces de Écrire un client, exactement comme vous le feriez pour une table Delta standard.

Créer un Service Principal et accorder des autorisations

Un Service Principal est une identité spécialisée qui offre plus de sécurité que les comptes personnalisés. Pour plus d'information concernant les service principals et comment les utiliser pour l'authentification, consultez Autoriser l'accès des Service Principal à Databricks avec OAuth.

  1. Pour créer un Service Principal, accédez à **Paramètres** > **Identité et accès**.

  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. Dans la page **Service principal**, accédez à l'**tab** Configurations.
    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.

Python 3.9 ou une version ultérieure est requis. Le SDK utilise les liaisons PyO3 vers le SDK Rust haute performance, offrant un throughput jusqu'à 40 fois supérieur à celui de Python pur et des E/S réseau efficaces grâce au runtime asynchrone de Rust. Il prend en charge le JSON (le plus simple) et les Protocol Buffers (recommandé pour la production). Le SDK prend également en charge les implémentations synchrones et asynchrones, ainsi que 3 méthodes d'ingestion différentes (basées sur le futur, basées sur le décalage et « fire-and-forget »).

Bash
pip install databricks-zerobus-ingest-sdk

Exemple JSON :

Python
import logging
from zerobus.sdk.sync import ZerobusSdk
from zerobus.sdk.shared import RecordType, StreamConfigurationOptions, 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)
options = StreamConfigurationOptions(record_type=RecordType.JSON)
stream = sdk.create_stream(CLIENT_ID, CLIENT_SECRET, table_properties, options)

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

# Optional: Wait for durability confirmation
stream.wait_for_offset(offset)
finally:
stream.close()

**Rappel d'accusé de réception :** Pour suivre la progression de l'ingestion de manière asynchrone, utilisez ack_callback l'option. Passez une sous-classe de AckCallback avec les méthodes on_ack(offset: int) et on_error(offset: int, error_message: str), qui sont appelées lorsque les enregistrements sont soit accusés de réception, soit en échec.

Python
from zerobus.sdk.shared import AckCallback, StreamConfigurationOptions, RecordType

class MyAckCallback(AckCallback):
def on_ack(self, offset: int) -> None:
print(f"Record acknowledged at offset: {offset}")

def on_error(self, offset: int, error_message: str) -> None:
print(f"Error at offset {offset}: {error_message}")

options = StreamConfigurationOptions(
record_type = RecordType.JSON,
ack_callback = MyAckCallback()
)

Protocol Buffers : Pour une ingestion sécurisée, utilisez Protocol Buffers avec RecordType.PROTO (default) et fournissez un descriptorProto dans les propriétés de la table.

Arrow Flight (Beta) : Pour une ingestion orientée colonne ou par batch des données Apache Arrow RecordBatch sur la même connexion gRPC, voir Utiliser Arrow Flight avec Zerobus Ingest. Nécessite le [arrow] supplémentaire : 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.

Étapes suivantes