Aller au contenu principal

Utilisez des API compatibles Kafka avec Zerobus Ingest

info

Bêta

Cette fonctionnalité est en bêta et est disponible uniquement sur AWS et Azure.

Zerobus Ingest propose des APIs de producteur compatibles avec Kafka qui vous permettent d'ingérer des données en utilisant n'importe quel client producteur Apache Kafka sans SDK Databricks. [[ ## completed ##]] Vous pointez un producteur Kafka existant vers l'Endpoint Zerobus et produisez vers un sujet nommé d'après votre table cible, et les enregistrements sont stockés dans une table Delta Unity Catalog. Les APIs compatibles avec Kafka sont adaptées si vous disposez déjà d'un producteur Kafka, d'un collecteur utilisant le protocole Kafka ou d'outils émettant vers Kafka, et que vous souhaitez acheminer ces données vers Delta avec un minimum de modifications de code. [[ ## completed ##]]

Les APIs compatibles Kafka s'exécutent sur TCP avec SASL_SSL et le mécanisme OAUTHBEARER et implémentent le sous-ensemble côté producteur du protocole Kafka, qui couvre Produce, Metadata, ApiVersions et les APIs de handshake SASL. [[ ## completed ##]] Les APIs de consommateur, d'administration et transactionnelles ne sont pas disponibles. Les API sont en écriture seule.

Quand utiliser les APIs compatibles avec Kafka

Les APIs compatibles avec Kafka sont les plus adaptées dans les scénarios suivants :

  • Vous souhaitez envoyer des données vers Delta sans adopter de SDK Zerobus, et vous exécutez déjà un producteur Kafka ou une application, un agent ou un collecteur qui émet vers Kafka.
  • Vous souhaitez réutiliser votre configuration de producteur Kafka, votre traitement par batch et vos outils opérationnels existants.
  • Vous envoyez des enregistrements JSON et n'avez pas besoin de Protocol Buffers ou d'Apache Arrow.

Si vous créez un nouveau client à partir de zéro et que vous souhaitez obtenir le throughput le plus élevé, des accusés de réception par enregistrement et une récupération automatique, utilisez un SDK Zerobus via gRPC plutôt que les APIs compatibles avec Kafka. [[ ## completed ##]] Consultez Choisir une interface. Pour les workloads en colonnes ou par batch, consultez Utiliser Arrow Flight avec Zerobus Ingest.

Comment fonctionne le modèle d'ingestion

Les APIs compatibles Kafka mappent les concepts Kafka sur Zerobus Ingest comme suit :

  • Sujets. Un nom de sujet Kafka est le nom complet de la table Unity Catalog à trois niveaux (catalog.schema.table). La table cible doit déjà exister, car Zerobus ne crée jamais de sujets.
  • Enregistrements. Zerobus ingère uniquement la valeur de l'enregistrement, qui doit être un objet JSON encodé en UTF-8 correspondant au schéma de la table Delta cible. Il ignore les clés d'enregistrement, les en-têtes, ainsi que la partition et le Timestamp fournis par le client, et ne les conserve pas.
  • Authentification. Chaque connexion s’authentifie avec un jeton OAuth Databricks limité à la table cible, présenté via SASL/OAUTHBEARER. Voir Authentification.
  • Remerciements. Zerobus renvoie une réponse Produce uniquement après avoir persisté durablement les enregistrements. Configurez votre producteur avec acks=all.

Zerobus est sans partition par conception. L’endpoint Zerobus est un broker logique unique et une partition unique, de sorte que chaque enregistrement pour un sujet est accusé de réception sur la partition 0. Les producteurs n’ont pas besoin d’en tenir compte. Zerobus Ingest monte en charge horizontalement pour gérer la charge entrante.

De plus, Zerobus Ingest assure une livraison au moins une fois. Une connexion unique prend en charge environ 50 000 messages par seconde, ce qui est inférieur au chemin gRPC du SDK. Pour un throughput maximal, utilisez un SDK Zerobus plutôt que les APIs compatibles Kafka. Les limites de latence, de quota, de taille d’enregistrement et de table partitionnée sont partagées avec le reste de Zerobus Ingest. Voir les quotas du connecteur Zerobus Ingest.

Authentification

Les APIs compatibles avec Kafka utilisent SASL/OAUTHBEARER. Le jeton porteur est un jeton OAuth Databricks que vous obtenez avec les identifiants client d’un Service Principal ayant accès à la table cible. Le jeton est limité à cette table via OAuth authorization_details et utilise la ressource zerobusDirectWriteApi, le même flux que l’ API REST Zerobus.

Les jetons OAuth expirent après une heure ; fournissez donc le jeton via le rappel du fournisseur de jetons de votre client Kafka plutôt que sous forme de chaîne statique. Le client récupère ensuite un nouveau jeton à chaque reconnexion. Les connexions ont également une durée de vie limitée côté serveur. Lorsqu'une connexion atteint cette limite, Zerobus la ferme, et le producteur se reconnecte et se réauthentifie de lui-même. La réauthentification sur une connexion active n'est pas prise en charge.

Accordez au Service Principal les privilèges Unity Catalog requis sur la table cible avant de vous connecter. Consultez Créer un Service Principal et accorder des autorisations.

Écrire un client

L'exemple ci-dessous produit vers la même table air_quality utilisée dans les exemples Utiliser le connecteur Zerobus Ingest. Il utilise kafka-python, mais tout client producteur Kafka prenant en charge SASL_SSL avec le mécanisme OAUTHBEARER fonctionne. Adaptez le modèle de fournisseur de jetons à votre bibliothèque cliente.

Le producteur se connecte au serveur bootstrap Zerobus sur le port 9092. Trouvez votre ID de workspace et votre région comme décrit dans Obtenir votre URL de workspace et votre endpoint Zerobus Ingest.

  • Serveur bootstrap : <workspace-id>.zerobus.<region>.cloud.databricks.com:9092
Bash
pip install kafka-python requests

Étape 1 : Créer un fournisseur de jetons

Zerobus authentifie chaque connexion avec un jeton OAuth Databricks éphémère et limité à la table. Comme le jeton expire, transmettez un rappel qui génère un nouveau jeton à la demande plutôt qu'un jeton statique.

La fonction fetch_zerobus_token() échange vos identifiants de Service Principal contre un jeton limité à la table cible, et ZerobusTokenProvider l’encapsule dans l’interface de rappel attendue par kafka-python.

Python
import json

import requests
from kafka.sasl.oauth import AbstractTokenProvider

# See "Get your workspace URL and Zerobus Ingest endpoint" in zerobus-ingest.md.
WORKSPACE_ID = "1234567890123456"
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"


def fetch_zerobus_token():
catalog, schema, table = TABLE_NAME.split(".")
authorization_details = [
{
"type": "unity_catalog_privileges",
"privileges": ["USE CATALOG"],
"object_type": "CATALOG",
"object_full_path": catalog,
},
{
"type": "unity_catalog_privileges",
"privileges": ["USE SCHEMA"],
"object_type": "SCHEMA",
"object_full_path": f"{catalog}.{schema}",
},
{
"type": "unity_catalog_privileges",
"privileges": ["SELECT", "MODIFY"],
"object_type": "TABLE",
"object_full_path": TABLE_NAME,
},
]

response = requests.post(
f"{WORKSPACE_URL}/oidc/v1/token",
auth=(CLIENT_ID, CLIENT_SECRET),
data={
&quot;grant_type&quot;: &quot;client_credentials&quot;,
&quot;scope&quot;: &quot;all-apis&quot;,
&quot;resource&quot;: f&quot;api://databricks/workspaces/{WORKSPACE_ID}/zerobusDirectWriteApi&quot;,
&quot;authorization_details&quot;: json.dumps(authorization_details),
},
timeout=30,
)
response.raise_for_status()
return response.json()["access_token"]


# kafka-python calls token() whenever it needs a fresh OAuth token.
class ZerobusTokenProvider(AbstractTokenProvider):
def token(self):
return fetch_zerobus_token()

Étape 2 : Configurer le producteur et envoyer les enregistrements

Pointez le producteur vers le serveur bootstrap, configurez SASL_SSL avec le mécanisme OAUTHBEARER et transmettez le fournisseur de jetons de l’étape 1. Utilisez acks="all" afin que chaque batch soit confirmé uniquement après avoir été durablement persisté, et envoyez les enregistrements sans compression.

Python
from kafka import KafkaProducer

BOOTSTRAP_SERVERS = "1234567890123456.zerobus.us-west-2.cloud.databricks.com:9092"

producer = KafkaProducer(
bootstrap_servers=BOOTSTRAP_SERVERS,
security_protocol="SASL_SSL",
sasl_mechanism="OAUTHBEARER",
sasl_oauth_token_provider=ZerobusTokenProvider(),
# Wait for durable acknowledgement before treating a record as ingested.
acks="all",
# Compression is not supported by the endpoint; send records uncompressed.
compression_type=None,
)

# Each send() returns a future immediately. The topic is the full table name.
futures = [
producer.send(
topic=TABLE_NAME,
value=json.dumps(
{"device_name": f"sensor-{i}", "temp": 20 + i % 15, "humidity": 50 + i % 40}
).encode("utf-8"),
)
for i in range(1000)
]

producer.flush()

# Block on each future to confirm every record was durably acknowledged.
for future in futures:
future.get(timeout=30)

producer.close()
print("All records ingested successfully")

Options de configuration

Les APIs compatibles Kafka implémentent le sous-ensemble producteur du protocole Kafka. Configurez votre producteur selon les options suivantes.

Option

Détails

Format d’enregistrement

JSON uniquement. Chaque valeur d'enregistrement doit être un objet JSON encodé en UTF-8 qui correspond au schéma de la table cible. Pour envoyer des Protocol Buffers ou des Avro, utilisez un SDK Zerobus.

Compression

Non pris en charge. Envoyez des batchs non compressés, par exemple compression.type=none. Zerobus rejette les batchs gzip, snappy, lz4 et zstd avec le code d'erreur UNSUPPORTED_COMPRESSION_TYPE.

Champs d'enregistrement

Valeur uniquement. Zerobus ingère la valeur de l'enregistrement et ne conserve pas les clés, les en-têtes, les affectations de partition ou les Timestamp.

Prise en charge des API

Écriture seule. Zerobus accepte les requêtes Produce ainsi que les métadonnées et les requêtes SASL nécessaires pour établir une session. Les APIs de consommateur, d'administration et transactionnelles ne sont pas prises en charge.

Private Link front-end

Non pris en charge. Connectez-vous plutôt via l’endpoint public.

Schéma

Appliqué. Zerobus rejette les enregistrements dont les champs ne correspondent pas au schéma de la table cible et traite les colonnes Delta supplémentaires autorisant les valeurs NULL comme une modification non bloquante. Pour capturer les champs non correspondants au lieu de les rejeter, configurez une colonne de secours.

Option

Détails

Format d’enregistrement

JSON uniquement. Chaque valeur d'enregistrement doit être un objet JSON encodé en UTF-8 qui correspond au schéma de la table cible. Pour envoyer des Protocol Buffers ou des Avro, utilisez un SDK Zerobus.

Compression

Non pris en charge. Envoyez des batchs non compressés, par exemple compression.type=none. Zerobus rejette les batchs gzip, snappy, lz4 et zstd avec le code d'erreur UNSUPPORTED_COMPRESSION_TYPE.

Champs d'enregistrement

Valeur uniquement. Zerobus ingère la valeur de l'enregistrement et ne conserve pas les clés, les en-têtes, les affectations de partition ou les Timestamp.

Prise en charge des API

Écriture seule. Zerobus accepte les requêtes Produce ainsi que les métadonnées et les requêtes SASL nécessaires pour établir une session. Les APIs de consommateur, d'administration et transactionnelles ne sont pas prises en charge.

Private Link front-end

Non pris en charge. Connectez-vous plutôt via l’endpoint public.

Schéma

Appliqué. Zerobus rejette les enregistrements dont les champs ne correspondent pas au schéma de la table cible et traite les colonnes Delta supplémentaires autorisant les valeurs NULL comme une modification non bloquante. Pour capturer les champs non correspondants au lieu de les rejeter, configurez une colonne de secours.

Pour plus d'informations sur le Private Link front-end, consultez Concepts de Private Link.

Pour les limites de latence, de quota, de taille d'enregistrement et de table partitionnée, consultez les quotas du connecteur Zerobus Ingest.

Bonnes pratiques

Suivez ces directives pour obtenir les meilleures performances et la meilleure fiabilité des APIs compatibles avec Kafka. Une seule connexion peut supporter environ 50 000 messages par seconde.

  • Réutilisez un producteur à longue durée de vie pour de nombreux enregistrements plutôt que d’en créer un par batch, car la création du producteur et la négociation SASL entraînent des coûts de configuration.
  • Laissez le producteur accumuler les enregistrements dans des batchs, par exemple en ajustant linger.ms et batch.size, au lieu d'effectuer un vidage après chaque enregistrement. Le traitement par batch est le levier le plus important pour le throughput.
  • Utilisez acks=all pour obtenir un accusé de réception durable pour chaque batch, en respectant la sémantique « au moins une fois » de Zerobus.
  • Récupérez les jetons OAuth via le rappel du fournisseur de jetons de votre client afin qu'ils refresh automatiquement lors de la reconnexion, plutôt que de transmettre un jeton statique qui expire.
  • Exécutez le producteur dans la même région cloud que l'endpoint Zerobus pour un throughput maximal.

Gestion des erreurs

Zerobus signale les échecs en utilisant les codes d'erreur Kafka standard sur le sujet et la partition concernés. Les codes courants incluent :

Erreur Kafka

Signification

SASL_AUTHENTICATION_FAILED

Le jeton OAuth est manquant, invalide ou ne dispose pas des privilèges Unity Catalog requis sur la table.

UNKNOWN_TOPIC_OR_PARTITION

La table cible n'existe pas, a été supprimée ou le jeton n'est pas autorisé à y écrire.

INVALID_RECORD

Un enregistrement a échoué à la validation du schéma ou n'a pas pu être décodé en tant que JSON UTF-8.

MESSAGE_TOO_LARGE

Un seul enregistrement dépasse la limite de taille d’enregistrement de 10 Mo. Voir Taille de l’enregistrement.

UNSUPPORTED_COMPRESSION_TYPE

Le batch a été compressé. Envoyez les enregistrements non compressés.

Erreur Kafka

Signification

SASL_AUTHENTICATION_FAILED

Le jeton OAuth est manquant, invalide ou ne dispose pas des privilèges Unity Catalog requis sur la table.

UNKNOWN_TOPIC_OR_PARTITION

La table cible n'existe pas, a été supprimée ou le jeton n'est pas autorisé à y écrire.

INVALID_RECORD

Un enregistrement a échoué à la validation du schéma ou n'a pas pu être décodé en tant que JSON UTF-8.

MESSAGE_TOO_LARGE

Un seul enregistrement dépasse la limite de taille d’enregistrement de 10 Mo. Voir Taille de l’enregistrement.

UNSUPPORTED_COMPRESSION_TYPE

Le batch a été compressé. Envoyez les enregistrements non compressés.

Après l'échec d'une requête Produce, Zerobus renvoie les codes d'erreur par partition et ferme la connexion. Les producteurs Kafka se reconnectent automatiquement, mais concevez votre client de manière à faire remonter les échecs d'envoi, par exemple en inspectant le résultat de chaque envoi, afin que les enregistrements ne soient pas perdus silencieusement.

Ressources supplémentaires