Aller au contenu principal

Utilisez MQTT avec Zerobus Ingest

info

Bêta

Cette fonctionnalité est en version bêta. Pour l'utiliser, un administrateur de workspace doit activer Zerobus Ingest MQTT Endpoint depuis la page Previews . Consultez Gérer les aperçus Databricks.

Zerobus Ingest fournit un endpoint MQTT v5 pour les applications qui publient déjà des messages MQTT. Chaque message contient un objet JSON et écrit directement dans une table Delta du Unity Catalog existante. Cette interface convient parfaitement aux applications de l’Internet des objets (IoT), de périphérie (« edge ») et de télémétrie qui utilisent un client MQTT v5 et ne nécessitent pas de fonctionnalités d’abonnement de courtier.

Prérequis​

Avant de vous connecter, assurez-vous de disposer des ressources et des accès suivants :

Confirmez que Zerobus Ingest est disponible dans la région de votre workspace. Voir la disponibilité régionale de Zerobus Ingest. La disponibilité de MQTT peut différer de la disponibilité générale de Zerobus Ingest.

Modèle de connexion​

Chaque connexion MQTT cible une table. Définissez la table lorsque le client envoie CONNECT, et utilisez le même nom de table à trois niveaux comme rubrique pour chaque paquet PUBLISH :

Text
catalog.schema.table

Chaque charge utile PUBLISH doit être un objet JSON encodé en UTF-8 qui correspond au schéma de la table cible. Une connexion ne peut pas changer de table. Ouvrez une autre connexion pour publier vers une autre table.

L’endpoint MQTT utilise TLS sur le port 8883 avec le format de hostname suivant :

<workspace-id>.zerobus.<region>.cloud.databricks.com

Il s'agit du même Hostname que l'Endpoint du serveur Zerobus Ingest. Configurez votre client MQTT avec ce Hostname et le port 8883, sans schéma https://.

Authentification​

S’authentifier pendant CONNECT avec ces propriétés utilisateur MQTT v5 :

Propriété utilisateur

Valeur

Authorization

Bearer <token>, où <token> est un jeton OAuth Databricks étendu à la table cible.

x-databricks-zerobus-table-name

Le nom complet de la table cible, tel que main.default.air_quality.

Propriété utilisateur

Valeur

Authorization

Bearer <token>, où <token> est un jeton OAuth Databricks étendu à la table cible.

x-databricks-zerobus-table-name

Le nom complet de la table cible, tel que main.default.air_quality.

Obtenez le jeton OAuth avec le flux d’identifiants client du service principal. Incluez authorization_details pour le catalogue, le schéma et la table, et utilisez la ressource zerobusDirectWriteApi pour votre workspace. L’ exemple Python démontre cet échange.

Les identifiants OAuth s'appliquent uniquement lorsqu'une connexion start. Les connexions ont une durée de vie limitée côté serveur, après quoi Zerobus Ingest déconnecte le client. Lorsque vous vous reconnectez, obtenez un nouveau jeton et envoyez-le dans les propriétés du nouveau CONNECT. L'authentification renforcée MQTT et la réauthentification au sein de la connexion ne sont pas prises en charge.

Durabilité et accusés de réception​

Zerobus Ingest prend en charge les niveaux de qualité de service (QoS) MQTT suivants :

  • QoS 0 est fourni dans la mesure du possible. Zerobus Ingest tente d'envoyer des enregistrements valides via le pipeline d'ingestion normal, mais n'envoie ni accusé de réception de réussite ni réponse d'échec par enregistrement. Un enregistrement peut être rejeté ou perdu sans notification du client ; par conséquent, une connexion ouverte ne confirme pas la livraison. QoS 0 ne garantit pas une diffusion au moins unique. Utilisez-le uniquement lorsque votre application peut tolérer la perte d'enregistrements.
  • QoS 1 renvoie un PUBACK une fois que Zerobus Ingest a rendu l’enregistrement durable. Considérez l’enregistrement comme accusé de réception uniquement lorsque le PUBACK correspond à l’ID du paquet PUBLISH et possède un code de motif de réussite.

Un PUBACK réussi confirme la durabilité, et non une visibilité immédiate dans la table Delta. Zerobus Ingest matérialise les enregistrements durables dans la table de manière asynchrone. Voir latence.

La livraison QoS 1 peut contenir des doublons. Si une connexion se ferme avant que votre client ne reçoive PUBACK, l'enregistrement est peut-être déjà persistant. Le fait de relancer cet enregistrement peut créer un doublon. Les applications qui réessaient des enregistrements non acquittés doivent tolérer les doublons ou les dédupliquer dans le lakehouse.

Un enregistrement QoS 1 n'est pas acquitté si son PUBACK comporte un code de motif d'échec ou si la connexion se ferme avant l'arrivée de PUBACK. Appliquez la règle de nouvelles tentatives de votre application aux enregistrements non acquittés. Un paquet PUBLISH supérieur à la limite de taille des paquets ferme la connexion et est toujours rejeté ; par conséquent, ne réessayez pas.

Exemple : publier un enregistrement​

L'exemple suivant utilise Paho MQTT 2.x et requests.

remarque

Cet exemple d'introduction désactive le comportement de reconnexion automatique de Paho. Dans une application longue durée, reconnectez-vous avec un jeton OAuth nouvellement récupéré. Suivez les enregistrements sans PUBACK réussi et appliquez votre propre règle de nouvelle tentative, en tenant compte du risque de duplication décrit dans la page Durability and acknowledgments.

L'exemple s'attend à ce qu'une table possède ce schéma :

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

Installez les dépendances :

Bash
pip install "paho-mqtt>=2,<3" requests

Définissez les variables d’environnement requises. Définissez ZEROBUS_MQTT_HOST sur le hostname MQTT décrit dans la section Modèle de connexion.

Bash
export DATABRICKS_WORKSPACE_URL="https://<databricks-instance>"
export DATABRICKS_WORKSPACE_ID="<workspace-id>"
export DATABRICKS_CLIENT_ID="<service-principal-client-id>"
export DATABRICKS_CLIENT_SECRET="<service-principal-client-secret>"
export DATABRICKS_TABLE_NAME="main.default.air_quality"
export ZEROBUS_MQTT_HOST="<cloud-specific-zerobus-hostname>"

Exécutez le script suivant :

Python
import json
import os
import ssl
import threading
import uuid

import paho.mqtt.client as mqtt
import requests

CONNECT_TIMEOUT_SECONDS = 10
CONNACK_TIMEOUT_SECONDS = 30
PUBACK_TIMEOUT_SECONDS = 30


def required_env(name):
value = os.environ.get(name)
if not value:
raise RuntimeError(f"Set the {name} environment variable")
return value


workspace_url = required_env("DATABRICKS_WORKSPACE_URL").rstrip("/")
workspace_id = required_env("DATABRICKS_WORKSPACE_ID")
client_id = required_env("DATABRICKS_CLIENT_ID")
client_secret = required_env("DATABRICKS_CLIENT_SECRET")
table_name = required_env("DATABRICKS_TABLE_NAME")
mqtt_host = required_env("ZEROBUS_MQTT_HOST")


def fetch_zerobus_token():
catalog, schema, _ = 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={
"grant_type": "client_credentials",
"scope": "all-apis",
"resource": f"api://databricks/workspaces/{workspace_id}/zerobusDirectWriteApi",
"authorization_details": json.dumps(authorization_details),
},
timeout=30,
)
response.raise_for_status()
return response.json()["access_token"]


connack_received = threading.Event()
puback_received = threading.Event()
connack_result = {}
puback_result = {}
disconnect_result = {}


def on_connect(client, userdata, flags, reason_code, properties):
connack_result["reason_code"] = reason_code
connack_received.set()


def on_publish(client, userdata, mid, reason_code, properties):
puback_result["mid"] = mid
puback_result["reason_code"] = reason_code
puback_received.set()


def on_disconnect(client, userdata, disconnect_flags, reason_code, properties):
disconnect_result["reason_code"] = reason_code
connack_received.set()
puback_received.set()


connect_properties = mqtt.Properties(mqtt.PacketTypes.CONNECT)
connect_properties.UserProperty = [
("Authorization", f"Bearer {fetch_zerobus_token()}"),
("x-databricks-zerobus-table-name", table_name),
]

client = mqtt.Client(
mqtt.CallbackAPIVersion.VERSION2,
client_id=f"zerobus-mqtt-{uuid.uuid4()}",
protocol=mqtt.MQTTv5,
reconnect_on_failure=False,
)
client.on_connect = on_connect
client.on_publish = on_publish
client.on_disconnect = on_disconnect
client.connect_timeout = CONNECT_TIMEOUT_SECONDS
client.tls_set(cert_reqs=ssl.CERT_REQUIRED, tls_version=ssl.PROTOCOL_TLS_CLIENT)

loop_started = False
try:
result = client.connect(
mqtt_host,
port=8883,
keepalive=60,
clean_start=True,
properties=connect_properties,
)
if result != mqtt.MQTT_ERR_SUCCESS:
raise RuntimeError(f"MQTT connect failed: {mqtt.error_string(result)}")

result = client.loop_start()
if result != mqtt.MQTT_ERR_SUCCESS:
raise RuntimeError(f"MQTT network loop failed: {mqtt.error_string(result)}")
loop_started = True

if not connack_received.wait(CONNACK_TIMEOUT_SECONDS):
raise TimeoutError("Timed out waiting for CONNACK")
if "reason_code" not in connack_result:
raise RuntimeError(
f"Disconnected before CONNACK: {disconnect_result['reason_code']}"
)
connack_reason = connack_result["reason_code"]
if connack_reason.is_failure:
raise RuntimeError(f"CONNECT rejected: {connack_reason}")

record = {"device_name": "sensor-1", "temp": 22, "humidity": 55}
message = client.publish(table_name, json.dumps(record), qos=1)
if message.rc != mqtt.MQTT_ERR_SUCCESS:
raise RuntimeError(f"PUBLISH failed: {mqtt.error_string(message.rc)}")

if not puback_received.wait(PUBACK_TIMEOUT_SECONDS):
raise TimeoutError("Timed out waiting for PUBACK")
if "reason_code" not in puback_result:
raise RuntimeError(
f"Disconnected before PUBACK: {disconnect_result['reason_code']}"
)
if puback_result["mid"] != message.mid:
raise RuntimeError("Received PUBACK for a different PUBLISH packet")
puback_reason = puback_result["reason_code"]
if puback_reason.value != 0:
raise RuntimeError(f"PUBLISH rejected: {puback_reason}")

print(f"Record {message.mid} is durable")
finally:
client.disconnect()
if loop_started:
client.loop_stop()

Limitations​

L’interface MQTT a les limites de protocole suivantes :

Élément

Aide

Version du protocole

Réservé à MQTT v5.

Qualité de service

QoS 0 et QoS 1. QoS 2 n'est pas pris en charge.

Messages QoS 1 en cours de transmission

Le maximum de réception CONNACK est de 4 096 messages par connexion. Les clients doivent respecter la limite annoncée.

CONNECT taille des paquets

16 Kio pour le paquet initial complet CONNECT, y compris les propriétés.

Taille du paquet après CONNECT

64 Kio pour le paquet MQTT complet, y compris l'en-tête fixe, l'en-tête variable, les propriétés, la rubrique, l'ID de paquet et la charge utile. CONNACK annonce cette limite en tant que taille maximale de paquet (Maximum Packet Size). Un paquet plus volumineux ferme la connexion avec le code de motif Packet too large.

Identifiant client

Obligatoire, non vide et ne dépassant pas 128 octets. Utilisez un identifiant client unique pour chaque connexion active.

État de session

Le serveur ne conserve pas et ne reprend pas l'état de la session et définit l'intervalle d'expiration de la session à 0.

Intervalle Keep-alive

Entre 5 et 300 secondes. Si la valeur CONNECT est 0, le serveur utilise 60 secondes ; sinon, le serveur limite la valeur à cet intervalle. Respectez la valeur Server Keep Alive renvoyée dans CONNACK. Le serveur ferme la connexion s'il ne reçoit aucun paquet complet pendant 1,5 fois l'intervalle effectif. Un paquet est considéré comme reçu uniquement lorsque son dernier octet arrive. Par conséquent, sur les Link lents, choisissez un intervalle de persistance (« keep-alive ») suffisamment long pour envoyer votre enregistrement le plus volumineux.

Front-end Link

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

PUBLISH Propriétés

Seuls le sujet, la QoS et la charge utile affectent l'importation. Le serveur ne conserve ni n'applique l'indicateur de format de charge utile (Payload Format Indicator), le type de contenu (Content Type), les données de corrélation (Correlation Data), l'intervalle d'expiration des messages (Message Expiry Interval), le sujet de réponse (Response Topic) ou les propriétés utilisateur (User Properties). L'identifiant d'abonnement n'est pas pris en charge.

Fonctionnalités de courtier

Les abonnements, les messages conservés, les alias de rubriques, les messages Last Will et l'authentification renforcée ne sont pas pris en charge.

Élément

Aide

Version du protocole

Réservé à MQTT v5.

Qualité de service

QoS 0 et QoS 1. QoS 2 n'est pas pris en charge.

Messages QoS 1 en cours de transmission

Le maximum de réception CONNACK est de 4 096 messages par connexion. Les clients doivent respecter la limite annoncée.

CONNECT taille des paquets

16 Kio pour le paquet initial complet CONNECT, y compris les propriétés.

Taille du paquet après CONNECT

64 Kio pour le paquet MQTT complet, y compris l'en-tête fixe, l'en-tête variable, les propriétés, la rubrique, l'ID de paquet et la charge utile. CONNACK annonce cette limite en tant que taille maximale de paquet (Maximum Packet Size). Un paquet plus volumineux ferme la connexion avec le code de motif Packet too large.

Identifiant client

Obligatoire, non vide et ne dépassant pas 128 octets. Utilisez un identifiant client unique pour chaque connexion active.

État de session

Le serveur ne conserve pas et ne reprend pas l'état de la session et définit l'intervalle d'expiration de la session à 0.

Intervalle Keep-alive

Entre 5 et 300 secondes. Si la valeur CONNECT est 0, le serveur utilise 60 secondes ; sinon, le serveur limite la valeur à cet intervalle. Respectez la valeur Server Keep Alive renvoyée dans CONNACK. Le serveur ferme la connexion s'il ne reçoit aucun paquet complet pendant 1,5 fois l'intervalle effectif. Un paquet est considéré comme reçu uniquement lorsque son dernier octet arrive. Par conséquent, sur les Link lents, choisissez un intervalle de persistance (« keep-alive ») suffisamment long pour envoyer votre enregistrement le plus volumineux.

Front-end Link

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

PUBLISH Propriétés

Seuls le sujet, la QoS et la charge utile affectent l'importation. Le serveur ne conserve ni n'applique l'indicateur de format de charge utile (Payload Format Indicator), le type de contenu (Content Type), les données de corrélation (Correlation Data), l'intervalle d'expiration des messages (Message Expiry Interval), le sujet de réponse (Response Topic) ou les propriétés utilisateur (User Properties). L'identifiant d'abonnement n'est pas pris en charge.

Fonctionnalités de courtier

Les abonnements, les messages conservés, les alias de rubriques, les messages Last Will et l'authentification renforcée ne sont pas pris en charge.

Maintenez la connexion ouverte lors de la publication de plusieurs enregistrements dans la même table.

Dépannage​

Utilisez les vérifications suivantes pour diagnostiquer les défaillances courantes :

  • Échecs de connexion : vérifiez le hostname spécifique au cloud, le port 8883, la résolution DNS, les règles de pare-feu sortantes et la vérification du certificat TLS. Confirmez également que MQTT est activé pour le workspace.
  • Rejeté CONNECT: obtenez un nouveau jeton OAuth limité à la table. Vérifiez à la fois les noms de propriété utilisateur, le nom de la table à trois niveaux, la table cible et les attributions du service principal.
  • Élément PUBACK erroné ou manquant : ne signalez pas l'enregistrement comme durable. Confirmez que le sujet correspond au nom de la table et que la charge utile est un objet JSON en UTF-8 conforme au schéma. Un élément PUBACK manquant ne prouve pas que l'enregistrement n'a pas pu devenir durable.
  • Déconnexions : vérifiez la présence d'un sujet non concordant, d'un paquet supérieur à 64 KiB, d'une fonctionnalité MQTT non prise en charge, d'un trafic de conservation de connexion manqué ou de la durée de vie limitée de la connexion côté serveur. Obtenez un nouveau jeton avant de vous reconnecter, et appliquez la politique de nouvelles tentatives et de déduplication de votre application aux enregistrements non reconnus.